Compare commits

..

1 Commits

Author SHA1 Message Date
gitea 114b7eebdf feat(ui): okno postepu z logiem podczas pisania horoskopu
Testy / Testy warstwy logicznej (silnik) (push) Successful in 11m57s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 9m52s
Testy / Build obrazu silnika B (swisseph) (push) Successful in 38s
Testy / Kontrola składni wszystkich warstw (push) Successful in 24s
Testy / Testy warstwy logicznej (silnik) (pull_request) Successful in 11m21s
Testy / Testy warstwy prezentacji (dostęp do baz) (pull_request) Successful in 9m50s
Testy / Build obrazu silnika B (swisseph) (pull_request) Successful in 34s
Testy / Kontrola składni wszystkich warstw (pull_request) Successful in 21s
Generowanie trwa minutami, a zwykly POST nie dawal zadnego sygnalu — aplikacja
wygladala na zawieszona. Teraz w trakcie pracy pojawia sie okno z logiem,
zegarem i spinnerem.

Log pokazuje RZECZYWISTE zdarzenia z serwera, nie udawany pasek postepu:
- app/progress.py — strumien NDJSON; praca leci w watku roboczym, generator
  odpompowuje kolejke, wiec zdarzenia docz w TRAKCIE pracy, nie na koncu;
  heartbeat co 10s, zeby proxy nie uznalo polaczenia za martwe,
- providers.generate(..., on_event) — raportuje kazda ture (start, czas trwania,
  liczba znakow, czy urwana), bo to tura trwa,
- POST /chart/horoscope/stream w logice + proxy /horoscope/stream w prezentacji.

Wynik: ostatnie zdarzenie niesie GOTOWY HTML wyrenderowany z tego samego
szablonu, ktory renderuje przeladowanie strony (_prompt_result.html wydzielony
z _prompt_block.html). Jedno zrodlo prawdy dla wygladu wyniku — okno wstawia go
bez przeladowania.

Degradacja: bez strumieniowania w przegladarce formularz idzie klasycznie
i wszystko dziala jak wczesniej, tylko bez okna. Blad polaczenia konczy sie
komunikatem w logu, nie cisza.

BLAD ZNALEZIONY PRZY TESCIE NA ZYWO: petla kontynuacji odejmowala od budzetu
ZAMOWIONY limit tury zamiast tokenow faktycznie wyprodukowanych — pierwsza tura
zjadala caly budzet, wiec urwana odpowiedz nigdy nie doczekala sie dokonczenia
i wracala do uzytkownika jako calosc. Naprawione i pokryte testem regresyjnym.

Testy: 176 passed / 1 skipped (logika) + 17 (prezentacja). Zweryfikowane na zywo
z wolna atrapa modelu: zdarzenia z poprawnymi czasami, okno z 11 liniami logu,
wynik wstawiony bez przeladowania strony.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-22 22:58:56 +02:00
25 changed files with 553 additions and 2334 deletions
-296
View File
@@ -1,296 +0,0 @@
# Wdrożenie PRE-16 — HTTPS na wejściu i szyfrowanie łączy między warstwami
Instrukcja krok po kroku. **Kolejność ma znaczenie** — punkt „Dlaczego taka
kolejność" niżej tłumaczy, co się stanie, jeśli ją zamienić.
Dotyczy dwóch pull requestów:
| Repo | PR | Co wnosi |
|---|---|---|
| `gitea/astrololo` | [#21](https://gitea.czernobog.pl/gitea/astrololo/pulls/21) | kod: szyfrowanie łączy, limit żądań za proxy |
| `gitea/deploy` | [#4](https://gitea.czernobog.pl/gitea/deploy/pulls/4) | manifesty: Ingress, certyfikat, klucze łączy |
---
## Co się właściwie zmienia
**Na wejściu do aplikacji.** Dotąd logowanie szło przez HTTP Basic po zwykłym
http — czyli hasło leciało siecią w postaci trywialnej do podsłuchania (base64 to
nie szyfrowanie). Po zmianie wejście jest po https, a http odsyła na https.
Przy okazji **odblokowują się dwie funkcje zepsute dziś z tego samego powodu**:
geolokalizacja („Tu i teraz") i kopiowanie promptu do schowka działają wyłącznie
w tzw. secure context i po http po prostu odmawiały.
**Między warstwami.** Prezentacja, logika i dane rozmawiały ze sobą otwartym
tekstem wewnątrz klastra. Token międzywarstwowy mówił *kto* pyta, ale nie ukrywał
*czego dotyczy odpowiedź* — a płyną nią surowe wiersze oryginalnych baz. Teraz
każde ciało żądania i odpowiedzi jest szyfrowane **AES-256-GCM**, osobnym kluczem
na każdą parę rozmówców.
**Wejście na świat pozostaje jedno: prompt do modelu.** Ta zmiana niczego tu nie
rusza — dotyczy wyłącznie ruchu wewnątrz sieci i wejścia z przeglądarki.
---
## Zanim zaczniesz — stan wyjściowy
```bash
kubectl -n astrololo get deploy,svc
kubectl -n astrololo get secret # powinny być: astrololo-auth, gitea-registry
kubectl -n kube-system get svc traefik -o jsonpath='{.status.loadBalancer.ingress[*].ip}'; echo
```
Zanotuj adres Traefika — będzie potrzebny w kroku 3. Sprawdź też, czy działa
aplikacja w obecnej postaci (przez NodePort), żeby mieć punkt odniesienia.
---
## Krok 1 — sekret z kluczami łączy
**Przed czymkolwiek innym.** Klucze muszą istnieć, zanim pody spróbują wstać
z nową konfiguracją, bo bez nich celowo **nie wystartują**.
```bash
kubectl -n astrololo create secret generic astrololo-link \
--from-literal=LINK_KEY_PRESENTATION_LOGIC="$(openssl rand -hex 32)" \
--from-literal=LINK_KEY_LOGIC_DATA="$(openssl rand -hex 32)"
```
Kluczy nikt nigdy nie musi oglądać — służą tylko usługom. Nie ma ich w repo
GitOps i **nie ma ich tam wkładać**: cokolwiek trafi do gita, zostaje w historii
na zawsze.
Dwa osobne klucze to nie ozdobnik. Przejęcie klucza prezentacji nie daje dostępu
do warstwy danych, gdzie leżą całe bazy. Logika dostaje oba, bo rozmawia w obie
strony; prezentacja i dane dostają wyłącznie swój.
Sprawdź:
```bash
kubectl -n astrololo get secret astrololo-link -o jsonpath='{.data}' | tr ',' '\n'
# oczekiwane: dwa klucze, każdy 64 znaki po odkodowaniu (32 bajty)
```
---
## Krok 2 — cert-manager
Jednorazowo, na cały klaster:
```bash
kubectl apply -f https://github.com/cert-manager/cert-manager/releases/download/v1.21.0/cert-manager.yaml
kubectl -n cert-manager rollout status deploy/cert-manager deploy/cert-manager-webhook --timeout=180s
```
Poczekaj, aż **webhook** będzie gotowy — dopóki nie wstanie, tworzenie obiektów
`Certificate` kończy się błędem połączenia i wygląda jak zepsuty manifest.
Sprawdź:
```bash
kubectl get crd | grep cert-manager | head -3 # muszą się pojawić
```
> **Dlaczego własne CA, a nie Let's Encrypt.** Klaster stoi w LAN (Traefik trzyma
> LoadBalancera na adresach 192.168.1.x), więc walidacja HTTP-01 nie ma jak dojść
> z internetu, a DNS-01 wymagałby trzymania w klastrze tokena API do domeny.
> Własne CA nie potrzebuje niczego z zewnątrz i odnawia certyfikaty samo. Cena:
> raz na urządzenie importujesz korzeń (krok 6).
---
## Krok 3 — DNS
Wpis `astrololo.czernobog.pl` → adres Traefika z kroku „stan wyjściowy”.
W routerze, lokalnym DNS-ie albo doraźnie w `/etc/hosts`:
```bash
echo "192.168.1.73 astrololo.czernobog.pl" | sudo tee -a /etc/hosts
```
**To nie jest krok opcjonalny.** Service `presentation` przestaje być NodePortem
(był drugą, nieszyfrowaną drogą do aplikacji — czyli obejściem całego PRE-16),
więc po wdrożeniu manifestów nazwa jest jedynym wejściem. Awaryjnie zawsze zostaje:
```bash
kubectl -n astrololo port-forward svc/presentation 8000:8000 # http://localhost:8000
```
---
## Krok 4 — merge PR-a aplikacji (astrololo #21)
Teraz, **przed** manifestami.
```bash
tea pr merge --login gitea --repo gitea/astrololo 21
```
Po merge'u CI zbuduje obrazy, a image-updater sam podbije tagi w repo `deploy`,
skąd ArgoCD wymieni pody. Poczekaj, aż to się przetoczy:
```bash
kubectl -n astrololo rollout status deploy/presentation deploy/logic deploy/data
kubectl -n astrololo get pods -o jsonpath='{range .items[*]}{.spec.containers[0].image}{"\n"}{end}'
```
Na tym etapie **nic się jeszcze nie szyfruje** — nowy kod to potrafi, ale zmienne
z kluczami dokłada dopiero PR do `deploy`. Aplikacja działa dokładnie jak dotąd.
To celowe: chcemy, żeby *cała* obsada podów umiała szyfrować, zanim ktokolwiek
tego zażąda.
---
## Krok 5 — merge PR-a manifestów (deploy #4)
```bash
tea pr merge --login gitea --repo gitea/deploy 4
```
ArgoCD zsynchronizuje się sam (`automated`, `selfHeal`). Wjeżdża naraz: Ingress,
certyfikat, zmienne z kluczami, `TRUST_PROXY` i zdjęcie NodePortu.
```bash
kubectl -n argocd get application astrololo
kubectl -n astrololo rollout status deploy/presentation deploy/logic deploy/data
kubectl -n astrololo get certificate # astrololo-ca i astrololo-tls: READY=True
```
> **Spodziewaj się kilkudziesięciu sekund błędów w trakcie.** Pody wymieniają się
> po kolei, więc przez chwilę stara prezentacja (jeszcze bez klucza) rozmawia
> z nową logiką (już z kluczem) i dostaje odmowę. To zamierzone: alternatywą byłby
> tryb „przyjmuj i szyfrowane, i jawne”, który zwykle zostaje włączony na zawsze.
Merge nie cofnie tagów obrazów — PR dotyka w `kustomization.yaml` wyłącznie listy
`resources`, nie bloku `images`, więc git złoży to z nowszymi tagami z mastera.
---
## Krok 6 — zaufanie do własnego CA (raz na urządzenie)
Bez tego przeglądarka pokaże ostrzeżenie o certyfikacie. Korzeń jest ważny 10 lat,
więc robisz to raz:
```bash
kubectl -n astrololo get secret astrololo-ca -o jsonpath='{.data.ca\.crt}' \
| base64 -d > astrololo-ca.crt
# macOS — do systemowego zaufania (poprosi o hasło administratora)
sudo security add-trusted-cert -d -r trustRoot \
-k /Library/Keychains/System.keychain astrololo-ca.crt
# Linux (Debian/Ubuntu)
sudo cp astrololo-ca.crt /usr/local/share/ca-certificates/ && sudo update-ca-certificates
```
Firefox ma **własny** magazyn certyfikatów — import przez *Ustawienia →
Prywatność i bezpieczeństwo → Wyświetl certyfikaty → Organy certyfikacji*.
---
## Krok 7 — sprawdzenie, że działa to, co miało zadziałać
### Wejście po https
```bash
curl -sI http://astrololo.czernobog.pl/ | head -2 # 301 → https
curl -s -o /dev/null -w "bez hasła: %{http_code}\n" https://astrololo.czernobog.pl/
curl -s -o /dev/null -w "z hasłem: %{http_code}\n" -u astrololo:'<hasło>' https://astrololo.czernobog.pl/
curl -sI -u astrololo:'<hasło>' https://astrololo.czernobog.pl/ | grep -i strict-transport
```
Oczekiwane: **301**, **401**, **200**, nagłówek HSTS obecny. Brak ostrzeżenia
o certyfikacie w przeglądarce oznacza, że krok 6 się udał.
### W przeglądarce
Kliknij **„Tu i teraz"** — powinno pobrać lokalizację (po http odmawiało).
Wygeneruj prompt i kliknij **kopiuj** — schowek powinien zadziałać bez obejść.
### Szyfrowanie łączy — sprawdzenie wprost
Najmocniejszy test to próba obejścia. Z wnętrza klastra, **bez klucza**:
```bash
kubectl -n astrololo exec deploy/presentation -- \
python -c "
import httpx, os
r = httpx.post('http://logic:8001/chart/report',
json={'when_utc':'1984-04-30T09:20:00+00:00','lat':50.06,'lon':19.94},
headers={'X-Astrololo-Token': os.environ['INTERNAL_TOKEN']})
print(r.status_code, r.text[:120])
"
```
Oczekiwane: **400** i `Łącze międzywarstwowe wymaga szyfrowania.` Zwróć uwagę, że
żądanie miało **prawidłowy token** — sam token już nie wystarcza, i o to chodziło.
To samo w dół, do warstwy danych:
```bash
kubectl -n astrololo exec deploy/logic -- \
python -c "
import httpx, os
r = httpx.post('http://data:8002/search',
json={'key':'significator','value':'[Sat','exact':False,'limit':5},
headers={'X-Astrololo-Token': os.environ['INTERNAL_TOKEN']})
print(r.status_code, r.text[:120])
"
```
### Logi startowe
```bash
kubectl -n astrololo logs deploy/logic | grep -i "łącze\|UWAGA"
```
Powinno być `łącze szyfrowane (AES-256-GCM…)`. Jeśli widzisz ostrzeżenie
o rozmowie **jawnym tekstem** — klucz nie doszedł do poda.
---
## Dlaczego taka kolejność
| Kolejność | Skutek zamiany |
|---|---|
| Sekret **przed** manifestami | `LINK_ENCRYPTION_REQUIRED=true` bez klucza celowo wywraca start. Pody wpadną w CrashLoop i będą tak siedzieć do czasu utworzenia sekretu. |
| cert-manager **przed** manifestami | API odrzuci `Certificate`/`Issuer` jako nieznane rodzaje zasobów, ArgoCD pokaże aplikację jako niezsynchronizowaną i sam tego nie naprawi. |
| DNS **przed** manifestami | NodePort znika razem z nimi. Bez wpisu DNS zostaje tylko `port-forward`. |
| Aplikacja **przed** manifestami | Odwrotnie: manifesty włączyłyby szyfrowanie na obrazach, które go nie znają — wszystkie żądania kończyłyby się odmową do czasu przebudowy obrazów. |
Fail-closed w obie strony jest zamierzony. Usługa, która wstała i **po cichu nie
szyfruje**, jest gorsza niż pod w CrashLoop — awarii nie widać, a bazy jadą
otwartym tekstem.
---
## Wycofanie
Manifestów: `git revert` merge'a w `deploy` — ArgoCD samo wróci do NodePortu
i ruchu bez szyfrowania. Kod aplikacji **nie wymaga wycofania**: bez zmiennych
`LINK_KEY_*` moduł przepuszcza ruch jak dotąd (i głośno o tym mówi w logach).
Certyfikat i CA zostają w namespace; usunięcie: `kubectl -n astrololo delete
certificate astrololo-ca astrololo-tls`. cert-managera można zostawić — nie
przeszkadza.
---
## Gdy coś nie gra
| Objaw | Przyczyna | Co zrobić |
|---|---|---|
| Pody w `CrashLoopBackOff`, w logach `LINK_ENCRYPTION_REQUIRED … nie ustawiony` | brak sekretu `astrololo-link` | krok 1, potem `rollout restart` |
| `400 Łącze międzywarstwowe wymaga szyfrowania` przy normalnym korzystaniu | jedna warstwa ma klucz, druga nie (albo trwa rollout) | `rollout status`; sprawdź, czy wszystkie trzy pody mają zmienną |
| `400 Nie udało się odczytać zaszyfrowanego żądania` | klucze po obu stronach łącza są **różne** | wymień sekret i zrestartuj **wszystkie trzy** naraz |
| `Certificate` stoi w `READY=False` | webhook cert-managera jeszcze nie wstał | `kubectl -n cert-manager get pods`, poczekaj i sprawdź `kubectl -n astrololo describe certificate astrololo-tls` |
| Przeglądarka: „połączenie nie jest prywatne” | korzeń CA nieimportowany na tym urządzeniu | krok 6 (pamiętaj, że Firefox ma osobny magazyn) |
| `404` z Traefika pod adresem aplikacji | DNS wskazuje gdzie indziej niż LoadBalancer Traefika | porównaj `dig +short astrololo.czernobog.pl` z adresem z kroku „stan wyjściowy” |
| Limit żądań odcina wszystkich naraz | brak `TRUST_PROXY=true` — cały ruch liczony jako jeden klient | sprawdź zmienną w `deploy/presentation` |
---
## Czego to nie załatwia
- **Szyfrowane są ciała żądań, nie nagłówki.** Ścieżka (`/search`) i token
międzywarstwowy jadą czytelnie. Sam token nikomu nic nie daje — bez klucza łącza
każde żądanie kończy się odmową — ale metadanych to nie ukrywa. Pełne ukrycie
wymagałoby mTLS.
- **Własne CA to nie publiczne zaufanie.** Każde nowe urządzenie wymaga importu
korzenia. Gdyby aplikacja miała kiedyś wyjść na świat, właściwą drogą jest
Let's Encrypt przez DNS-01.
- **NFS z plikami baz** stoi obok aplikacji — kto ma dostęp do share'u, bierze
pliki z pominięciem wszystkich powyższych zabezpieczeń. Do zamknięcia po stronie
infrastruktury (eksport tylko dla IP węzłów, `root_squash`, najlepiej read-only).
- **Sekrety w etcd** są tylko zakodowane base64. Docelowo: szyfrowanie etcd
at-rest albo Sealed Secrets / SOPS.
-466
View File
@@ -1,466 +0,0 @@
"""Szyfrowanie łączy między warstwami (PRE-16 / LOG-33).
Do tej pory warstwy rozmawiały ze sobą zwykłym HTTP-em wewnątrz klastra. Token
międzywarstwowy (LOG-32) mówił KTO pyta, ale nie ukrywał CZEGO dotyczy odpowiedź
— a płyną nią surowe wiersze oryginalnych baz interpretacyjnych, czyli rdzeń
produktu. Kto podsłuchał ruch wewnątrz sieci (drugi pod, port mirror na switchu,
zrzut z węzła), miał je w całości.
Ten moduł zamyka tę drogę: **AES-256-GCM** na ciele każdego żądania i odpowiedzi.
GCM daje jednocześnie poufność i uwierzytelnienie — cudzy albo podmieniony bajt
nie odszyfruje się w ogóle, więc nie ma osobnego problemu „zaszyfrowane, ale
podatne na modyfikację".
**Dwa niezależne klucze**, po jednym na parę rozmówców:
* ``LINK_KEY_PRESENTATION_LOGIC`` — prezentacja ↔ logika,
* ``LINK_KEY_LOGIC_DATA`` — logika ↔ dane.
Dzięki temu przejęcie klucza prezentacji nie daje dostępu do warstwy danych,
gdzie leżą całe bazy. Logika trzyma oba, bo rozmawia w obie strony.
Z każdego klucza łącza wyprowadzamy **osobne podklucze na kierunek** (HKDF).
Żądanie i odpowiedź nigdy nie szyfrują się tym samym kluczem, więc powtórzenie
losowej jednorazówki w jedną stronę nie osłabia drugiej.
Format ramki (bo strumień odpowiedzi może iść kawałkami — patrz okno postępu):
[4 bajty długości][magia "AL1"][12 bajtów jednorazówki][szyfrogram + znacznik]
Do materiału uwierzytelnianego (AAD) wchodzą kierunek, ścieżka, znacznik czasu
i numer ramki. Skutek: ramki nie da się przekleić do innego endpointu, odtworzyć
po czasie (dopuszczalny poślizg ``MAX_SKEW``) ani przestawić w strumieniu.
Bez ustawionego klucza moduł **przepuszcza ruch otwartym tekstem** (dev, zgodność
wstecz) i krzyczy o tym przy starcie. Gdy klucz JEST ustawiony, warstwa serwerowa
działa fail-closed: nieszyfrowane żądanie dostaje odmowę, żeby przypadkowa
regresja po stronie klienta nie oznaczała cichego powrotu do jawnego ruchu.
"""
from __future__ import annotations
import base64
import binascii
import logging
import os
import struct
import time
from typing import Iterable, Iterator
from cryptography.exceptions import InvalidTag
from cryptography.hazmat.primitives import hashes
from cryptography.hazmat.primitives.ciphers.aead import AESGCM
from cryptography.hazmat.primitives.kdf.hkdf import HKDF
log = logging.getLogger("astrololo.link")
MAGIC = b"AL1"
VERSION = "v1"
NONCE_BYTES = 12
KEY_BYTES = 32 # AES-256
LENGTH_PREFIX = 4
MAX_FRAME = 64 * 1024 * 1024 # zapora przed alokacją z podanej długości
MAX_SKEW_SECONDS = 300.0
HEADER_ENC = "X-Astrololo-Enc"
HEADER_TS = "X-Astrololo-Enc-Ts"
CONTENT_TYPE = "application/vnd.astrololo.enc"
ENV_PRESENTATION_LOGIC = "LINK_KEY_PRESENTATION_LOGIC"
ENV_LOGIC_DATA = "LINK_KEY_LOGIC_DATA"
ENV_REQUIRED = "LINK_ENCRYPTION_REQUIRED"
REQUEST, RESPONSE = b"req", b"res"
# Sondy k8s pukają tu bez klucza i tak ma zostać — inaczej pierwsza literówka
# w sekrecie kładłaby pody zamiast pokazać błąd w aplikacji.
PUBLIC_PATHS = frozenset({"/health"})
class LinkError(Exception):
"""Cokolwiek poszło nie tak z kopertą — celowo bez szczegółów na zewnątrz."""
# --------------------------------------------------------------------- klucze
def parse_key(raw: str) -> bytes:
"""Klucz z konfiguracji: hex (64 znaki) albo base64. Zawsze 32 bajty."""
text = raw.strip()
if not text:
raise LinkError("pusty klucz łącza")
try:
key = bytes.fromhex(text)
except ValueError:
try:
key = base64.b64decode(text, validate=True)
except (binascii.Error, ValueError) as exc:
raise LinkError("klucz łącza nie jest ani hexem, ani base64") from exc
if len(key) != KEY_BYTES:
raise LinkError(
f"klucz łącza ma {len(key)} B zamiast {KEY_BYTES} — wygeneruj przez "
f"`openssl rand -hex 32`"
)
return key
def key_from_env(env_name: str) -> bytes | None:
"""Klucz albo None. Zły klucz to wyjątek OD RAZU — nie przy pierwszym żądaniu."""
raw = os.getenv(env_name, "")
return parse_key(raw) if raw.strip() else None
def encryption_required() -> bool:
"""Czy brak klucza ma być błędem, a nie cichym powrotem do jawnego ruchu.
Serwer sam z siebie broni się fail-closed, ale to za mało: klient BEZ klucza
wysyła pytanie otwartym tekstem i dopiero potem dostaje odmowę — czyli treść
zapytania zdążyła już przelecieć przez sieć. Ta flaga zatrzymuje go, zanim
cokolwiek opuści proces. Ustawiana razem z kluczami we wdrożeniu.
"""
return os.getenv(ENV_REQUIRED, "").strip().lower() in {"1", "true", "yes", "on"}
def _subkey(link_key: bytes, direction: bytes) -> bytes:
return HKDF(
algorithm=hashes.SHA256(), length=KEY_BYTES, salt=None,
info=b"astrololo/link/" + direction,
).derive(link_key)
class Link:
"""Jedna para rozmówców: klucz plus wyprowadzone z niego podklucze."""
def __init__(self, link_key: bytes) -> None:
self._by_direction = {
REQUEST: AESGCM(_subkey(link_key, REQUEST)),
RESPONSE: AESGCM(_subkey(link_key, RESPONSE)),
}
# ---------------------------------------------------------- pojedyncza ramka
def _aad(self, direction: bytes, path: str, stamp: str, seq: int) -> bytes:
return b"|".join([MAGIC, direction, path.encode("utf-8"),
stamp.encode("ascii"), str(seq).encode("ascii")])
def seal(self, direction: bytes, path: str, stamp: str, seq: int,
plaintext: bytes) -> bytes:
nonce = os.urandom(NONCE_BYTES)
sealed = self._by_direction[direction].encrypt(
nonce, plaintext, self._aad(direction, path, stamp, seq))
return MAGIC + nonce + sealed
def open(self, direction: bytes, path: str, stamp: str, seq: int,
frame: bytes) -> bytes:
if not frame.startswith(MAGIC):
raise LinkError("ramka bez znacznika astrololo")
body = frame[len(MAGIC):]
if len(body) <= NONCE_BYTES:
raise LinkError("ramka za krótka")
nonce, sealed = body[:NONCE_BYTES], body[NONCE_BYTES:]
try:
return self._by_direction[direction].decrypt(
nonce, sealed, self._aad(direction, path, stamp, seq))
except InvalidTag as exc:
# Jeden komunikat na wszystkie przypadki: zły klucz, podmieniony bajt,
# przeklejenie z innej ścieżki, przestawiona ramka. Rozróżnianie ich
# na zewnątrz podpowiadałoby atakującemu, w co trafił.
raise LinkError("nie udało się odszyfrować — zły klucz albo naruszone dane") from exc
# ------------------------------------------------------------ strumień ramek
def seal_stream(self, direction: bytes, path: str, stamp: str,
chunks: Iterable[bytes]) -> Iterator[bytes]:
for seq, chunk in enumerate(chunks):
yield frame_out(self.seal(direction, path, stamp, seq, chunk))
def open_stream(self, direction: bytes, path: str, stamp: str,
raw: bytes) -> Iterator[bytes]:
for seq, frame in enumerate(frames_in(raw)):
yield self.open(direction, path, stamp, seq, frame)
def open_all(self, direction: bytes, path: str, stamp: str, raw: bytes) -> bytes:
return b"".join(self.open_stream(direction, path, stamp, raw))
# ---------------------------------------------------------------- ramkowanie
def frame_out(payload: bytes) -> bytes:
return struct.pack(">I", len(payload)) + payload
def frames_in(raw: bytes) -> Iterator[bytes]:
"""Rozbiera bufor na ramki. Ucięty strumień to błąd, nie cicha strata danych."""
offset = 0
while offset < len(raw):
if offset + LENGTH_PREFIX > len(raw):
raise LinkError("urwana ramka (brak nagłówka długości)")
(size,) = struct.unpack(">I", raw[offset:offset + LENGTH_PREFIX])
if size > MAX_FRAME:
raise LinkError("ramka ponad dopuszczalny rozmiar")
offset += LENGTH_PREFIX
if offset + size > len(raw):
raise LinkError("urwana ramka (za mało danych)")
yield raw[offset:offset + size]
offset += size
def unframe_incremental(buffer: bytearray) -> Iterator[bytes]:
"""Wyjmuje z bufora KOMPLETNE ramki i zjada je; resztę zostawia na później.
Dla odbioru na żywo: kawałki przychodzą podzielone dowolnie i ramka potrafi
rozjechać się między dwa odczyty.
"""
while True:
if len(buffer) < LENGTH_PREFIX:
return
(size,) = struct.unpack(">I", buffer[:LENGTH_PREFIX])
if size > MAX_FRAME:
raise LinkError("ramka ponad dopuszczalny rozmiar")
if len(buffer) < LENGTH_PREFIX + size:
return
frame = bytes(buffer[LENGTH_PREFIX:LENGTH_PREFIX + size])
del buffer[:LENGTH_PREFIX + size]
yield frame
# ------------------------------------------------------------- świeżość ruchu
def stamp_now() -> str:
return f"{time.time():.3f}"
def check_stamp(stamp: str) -> None:
"""Odrzuca ramki spoza okna czasowego — inaczej podsłuchane żądanie dałoby się
odtworzyć w dowolnym momencie w przyszłości."""
try:
sent = float(stamp)
except (TypeError, ValueError) as exc:
raise LinkError("brak albo błędny znacznik czasu") from exc
if abs(time.time() - sent) > MAX_SKEW_SECONDS:
raise LinkError("znacznik czasu poza dopuszczalnym oknem")
# =========================================================== strona serwerowa
class LinkCryptoMiddleware:
"""Rozszyfrowuje wchodzące żądania i zaszyfrowuje wychodzące odpowiedzi.
Napisane jako czyste ASGI, nie ``@app.middleware("http")``, bo trzeba
podmienić CIAŁO żądania jeszcze zanim zobaczy je FastAPI, oraz przepuścić
odpowiedź strumieniową kawałek po kawałku, bez zbierania jej w pamięci.
"""
def __init__(self, app, link: Link | None, layer: str) -> None:
self.app = app
self.link = link
self.layer = layer
async def __call__(self, scope, receive, send):
if scope["type"] != "http" or self.link is None or scope["path"] in PUBLIC_PATHS:
return await self.app(scope, receive, send)
path = scope["path"]
headers = {k.decode("latin-1").lower(): v.decode("latin-1") for k, v in scope["headers"]}
if headers.get(HEADER_ENC.lower()) != VERSION:
# Fail-closed. Klucz jest ustawiony, więc jawne żądanie oznacza albo
# pomyłkę w konfiguracji, albo kogoś obcego — w obu wypadkach nie
# chcemy po cichu wrócić do jawnego ruchu.
log.warning("warstwa %s: odrzucone żądanie bez szyfrowania łącza (%s)",
self.layer, path)
return await _refuse(send, "Łącze międzywarstwowe wymaga szyfrowania.")
stamp = headers.get(HEADER_TS.lower(), "")
try:
check_stamp(stamp)
plaintext = self.link.open_all(REQUEST, path, stamp, await _read_body(receive))
except LinkError as exc:
log.warning("warstwa %s: %s (%s)", self.layer, exc, path)
return await _refuse(send, "Nie udało się odczytać zaszyfrowanego żądania.")
scope = dict(scope)
scope["headers"] = _rewritten_headers(scope["headers"], len(plaintext))
await self.app(scope, _replay(plaintext, receive), self._sealing_send(send, path))
def _sealing_send(self, send, path: str):
state: dict = {"stamp": "", "seq": 0}
async def sealing(message):
if message["type"] == "http.response.start":
state["stamp"] = stamp_now()
keep = [(k, v) for k, v in message.get("headers", [])
if k.lower() not in (b"content-length", b"content-type")]
message = dict(message)
message["headers"] = keep + [
(b"content-type", CONTENT_TYPE.encode()),
(HEADER_ENC.lower().encode(), VERSION.encode()),
(HEADER_TS.lower().encode(), state["stamp"].encode()),
]
return await send(message)
if message["type"] == "http.response.body":
chunk = message.get("body", b"")
sealed = b""
if chunk:
sealed = frame_out(self.link.seal(
RESPONSE, path, state["stamp"], state["seq"], chunk))
state["seq"] += 1
return await send({"type": "http.response.body", "body": sealed,
"more_body": message.get("more_body", False)})
return await send(message)
return sealing
def _rewritten_headers(raw: Iterable[tuple[bytes, bytes]], length: int):
"""Po odszyfrowaniu ciało ma inną długość i zwykły typ — inaczej FastAPI
próbowałby sparsować JSON o cudzej deklarowanej wielkości."""
kept = [(k, v) for k, v in raw if k.lower() not in (b"content-length", b"content-type")]
kept.append((b"content-length", str(length).encode()))
if length:
kept.append((b"content-type", b"application/json"))
return kept
async def _read_body(receive) -> bytes:
body = bytearray()
while True:
message = await receive()
if message["type"] == "http.disconnect":
raise LinkError("rozłączenie w trakcie odbioru żądania")
body += message.get("body", b"")
if not message.get("more_body", False):
return bytes(body)
def _replay(body: bytes, original):
"""Podstawia odszyfrowane ciało jako jedyną porcję wejścia dla aplikacji.
Po oddaniu ciała oddajemy głos ORYGINALNEMU `receive`, zamiast od razu
zgłaszać rozłączenie. Odpowiedź strumieniowa nasłuchuje bowiem rozłączenia
równolegle do wysyłania i przerywa się, gdy je zobaczy — na skróconej wersji
okno postępu dostawało pustą odpowiedź, choć zwykłe żądania działały.
"""
delivered = False
async def receive():
nonlocal delivered
if delivered:
return await original()
delivered = True
return {"type": "http.request", "body": body, "more_body": False}
return receive
async def _refuse(send, detail: str) -> None:
"""Odmowa leci JAWNIE — rozmówca właśnie pokazał, że nie umie odszyfrować,
więc zaszyfrowany komunikat o błędzie byłby dla niego nieczytelny."""
payload = f'{{"detail":"{detail}"}}'.encode("utf-8")
await send({"type": "http.response.start", "status": 400, "headers": [
(b"content-type", b"application/json"),
(b"content-length", str(len(payload)).encode()),
]})
await send({"type": "http.response.body", "body": payload})
def install(app, env_name: str, layer: str):
"""Podpina szyfrowanie łącza. Wołać PO `security.install`, żeby także odmowa
tokenowa (401) wracała zaszyfrowana — inaczej klient by jej nie odczytał."""
link_key = key_from_env(env_name)
if link_key is None and encryption_required():
# Celowo wywracamy start. Ta sama zasada co przy sekrecie logowania:
# wolimy widoczną awarię niż usługę, która wstała i po cichu nie chroni
# niczego. Pod w CrashLoop widać od razu, jawny ruch — nie.
raise LinkError(
f"{ENV_REQUIRED} jest włączone, ale {env_name} nie ustawiony — "
f"warstwa {layer} nie wystartuje bez klucza łącza"
)
if link_key is None:
log.warning(
"UWAGA: %s nie ustawiony — warstwa %s rozmawia z sąsiadem JAWNYM tekstem, "
"więc treść baz interpretacyjnych jest widoczna dla każdego, kto podsłucha "
"ruch wewnątrz sieci.", env_name, layer,
)
return None
link = Link(link_key)
app.add_middleware(LinkCryptoMiddleware, link=link, layer=layer)
log.info("warstwa %s: łącze szyfrowane (AES-256-GCM, klucz z %s)", layer, env_name)
return link
# ============================================================ strona kliencka
def call(client, method: str, url: str, *, payload=None,
headers: dict[str, str] | None = None, link: Link | None) -> bytes:
"""Żądanie do sąsiedniej warstwy; zwraca odszyfrowane ciało odpowiedzi.
Ścieżkę do materiału uwierzytelnianego bierzemy Z URL-a, a nie z osobnego
argumentu — gdyby klient i serwer liczyły ją inaczej, każde żądanie kończyłoby
się niejasnym błędem odszyfrowania.
"""
import json as _json
import httpx
request_headers = dict(headers or {})
if link is None:
if encryption_required():
# Zatrzymujemy się PRZED wysłaniem. Gdyby polecieć jawnie i dopiero
# zebrać odmowę, pytanie byłoby już na kablu — a to właśnie ono niesie
# sygnifikatory, o które pytamy bazę.
raise LinkError(
f"{ENV_REQUIRED} jest włączone, ale brak klucza łącza — żądanie "
f"NIE zostało wysłane, żeby jego treść nie poszła jawnym tekstem"
)
response = client.request(method, url, json=payload, headers=request_headers)
response.raise_for_status()
return response.content
path = httpx.URL(url).path
stamp = stamp_now()
plaintext = b"" if payload is None else _json.dumps(payload).encode("utf-8")
body = frame_out(link.seal(REQUEST, path, stamp, 0, plaintext))
request_headers.update({HEADER_ENC: VERSION, HEADER_TS: stamp,
"Content-Type": CONTENT_TYPE})
response = client.request(method, url, content=body, headers=request_headers)
if response.status_code >= 400 and response.headers.get(HEADER_ENC) != VERSION:
log.error("łącze %s odmówiło: %s", path, response.text[:200])
response.raise_for_status()
if response.headers.get(HEADER_ENC) != VERSION:
raise LinkError("odpowiedź przyszła nieszyfrowana, choć klucz łącza jest ustawiony")
reply_stamp = response.headers.get(HEADER_TS, "")
check_stamp(reply_stamp)
return link.open_all(RESPONSE, path, reply_stamp, response.content)
def call_json(client, method: str, url: str, *, payload=None,
headers: dict[str, str] | None = None, link: Link | None):
import json as _json
return _json.loads(call(client, method, url, payload=payload,
headers=headers, link=link))
def open_response_stream(response, link: Link | None) -> Iterator[bytes]:
"""Odbiór odpowiedzi płynącej kawałkami (okno postępu).
Ramka potrafi rozjechać się między dwa odczyty z gniazda, więc składamy ją
w buforze zamiast zakładać, że każdy kawałek to komplet.
"""
if link is None:
yield from response.iter_bytes()
return
if response.headers.get(HEADER_ENC) != VERSION:
raise LinkError("strumień przyszedł nieszyfrowany, choć klucz łącza jest ustawiony")
stamp = response.headers.get(HEADER_TS, "")
check_stamp(stamp)
path = response.request.url.path
buffer = bytearray()
seq = 0
for chunk in response.iter_bytes():
buffer += chunk
for frame in unframe_incremental(buffer):
yield link.open(RESPONSE, path, stamp, seq, frame)
seq += 1
if buffer:
raise LinkError("strumień urwał się w połowie ramki")
+1 -4
View File
@@ -10,7 +10,7 @@ from contextlib import asynccontextmanager
from fastapi import FastAPI
from app import link_crypto, security
from app import security
from app.config import settings
from app.models import HealthInfo, SearchQuery, SearchResult
from app.providers.factory import build_provider
@@ -26,9 +26,6 @@ async def lifespan(app: FastAPI):
app = FastAPI(title="astrololo · warstwa bazodanowa", lifespan=lifespan)
security.install(app, "danych") # token międzywarstwowy (LOG-32)
# Szyfrowanie łącza od logiki. PO `security.install`, żeby także odmowa
# tokenowa wracała zaszyfrowana — inaczej klient nie umiałby jej odczytać.
link_crypto.install(app, link_crypto.ENV_LOGIC_DATA, "danych")
@app.post("/search", response_model=SearchResult)
-2
View File
@@ -7,5 +7,3 @@ openpyxl>=3.1
pyarrow>=18.0
SQLAlchemy>=2.0
pydantic>=2.10
# Szyfrowanie łącza między warstwami (PRE-16): AES-256-GCM + HKDF
cryptography>=44.0
+3 -13
View File
@@ -10,7 +10,6 @@ from typing import Any
import httpx
from app import link_crypto
from app.config import settings
@@ -20,13 +19,6 @@ def _auth_headers() -> dict[str, str]:
return {"X-Astrololo-Token": token} if token else {}
def _link() -> link_crypto.Link | None:
"""Klucz łącza logika↔dane. Czytany przy każdym wywołaniu, bo konfiguracja
może się zmienić bez restartu procesu (testy, podmiana sekretu)."""
key = link_crypto.key_from_env(link_crypto.ENV_LOGIC_DATA)
return link_crypto.Link(key) if key else None
class DataClient:
def __init__(self, base_url: str | None = None) -> None:
self.base_url = (base_url or settings.data_url).rstrip("/")
@@ -41,13 +33,11 @@ class DataClient:
) -> dict[str, Any]:
payload = {"key": key, "value": value, "exact": exact, "limit": limit, "fields": fields}
with httpx.Client(timeout=max(settings.http_timeout, 30.0)) as client:
return link_crypto.call_json(client, "POST", f"{self.base_url}/search",
payload=payload, headers=_auth_headers(),
link=_link())
r = client.post(f"{self.base_url}/search", json=payload, headers=_auth_headers())
r.raise_for_status()
return r.json()
def health(self) -> dict[str, Any]:
# /health celowo poza szyfrowaniem — pukają tu sondy k8s, które klucza
# nie mają, a nie przechodzi tędy nic z baz.
with httpx.Client(timeout=settings.http_timeout) as client:
r = client.get(f"{self.base_url}/health", headers=_auth_headers())
r.raise_for_status()
-466
View File
@@ -1,466 +0,0 @@
"""Szyfrowanie łączy między warstwami (PRE-16 / LOG-33).
Do tej pory warstwy rozmawiały ze sobą zwykłym HTTP-em wewnątrz klastra. Token
międzywarstwowy (LOG-32) mówił KTO pyta, ale nie ukrywał CZEGO dotyczy odpowiedź
— a płyną nią surowe wiersze oryginalnych baz interpretacyjnych, czyli rdzeń
produktu. Kto podsłuchał ruch wewnątrz sieci (drugi pod, port mirror na switchu,
zrzut z węzła), miał je w całości.
Ten moduł zamyka tę drogę: **AES-256-GCM** na ciele każdego żądania i odpowiedzi.
GCM daje jednocześnie poufność i uwierzytelnienie — cudzy albo podmieniony bajt
nie odszyfruje się w ogóle, więc nie ma osobnego problemu „zaszyfrowane, ale
podatne na modyfikację".
**Dwa niezależne klucze**, po jednym na parę rozmówców:
* ``LINK_KEY_PRESENTATION_LOGIC`` — prezentacja ↔ logika,
* ``LINK_KEY_LOGIC_DATA`` — logika ↔ dane.
Dzięki temu przejęcie klucza prezentacji nie daje dostępu do warstwy danych,
gdzie leżą całe bazy. Logika trzyma oba, bo rozmawia w obie strony.
Z każdego klucza łącza wyprowadzamy **osobne podklucze na kierunek** (HKDF).
Żądanie i odpowiedź nigdy nie szyfrują się tym samym kluczem, więc powtórzenie
losowej jednorazówki w jedną stronę nie osłabia drugiej.
Format ramki (bo strumień odpowiedzi może iść kawałkami — patrz okno postępu):
[4 bajty długości][magia "AL1"][12 bajtów jednorazówki][szyfrogram + znacznik]
Do materiału uwierzytelnianego (AAD) wchodzą kierunek, ścieżka, znacznik czasu
i numer ramki. Skutek: ramki nie da się przekleić do innego endpointu, odtworzyć
po czasie (dopuszczalny poślizg ``MAX_SKEW``) ani przestawić w strumieniu.
Bez ustawionego klucza moduł **przepuszcza ruch otwartym tekstem** (dev, zgodność
wstecz) i krzyczy o tym przy starcie. Gdy klucz JEST ustawiony, warstwa serwerowa
działa fail-closed: nieszyfrowane żądanie dostaje odmowę, żeby przypadkowa
regresja po stronie klienta nie oznaczała cichego powrotu do jawnego ruchu.
"""
from __future__ import annotations
import base64
import binascii
import logging
import os
import struct
import time
from typing import Iterable, Iterator
from cryptography.exceptions import InvalidTag
from cryptography.hazmat.primitives import hashes
from cryptography.hazmat.primitives.ciphers.aead import AESGCM
from cryptography.hazmat.primitives.kdf.hkdf import HKDF
log = logging.getLogger("astrololo.link")
MAGIC = b"AL1"
VERSION = "v1"
NONCE_BYTES = 12
KEY_BYTES = 32 # AES-256
LENGTH_PREFIX = 4
MAX_FRAME = 64 * 1024 * 1024 # zapora przed alokacją z podanej długości
MAX_SKEW_SECONDS = 300.0
HEADER_ENC = "X-Astrololo-Enc"
HEADER_TS = "X-Astrololo-Enc-Ts"
CONTENT_TYPE = "application/vnd.astrololo.enc"
ENV_PRESENTATION_LOGIC = "LINK_KEY_PRESENTATION_LOGIC"
ENV_LOGIC_DATA = "LINK_KEY_LOGIC_DATA"
ENV_REQUIRED = "LINK_ENCRYPTION_REQUIRED"
REQUEST, RESPONSE = b"req", b"res"
# Sondy k8s pukają tu bez klucza i tak ma zostać — inaczej pierwsza literówka
# w sekrecie kładłaby pody zamiast pokazać błąd w aplikacji.
PUBLIC_PATHS = frozenset({"/health"})
class LinkError(Exception):
"""Cokolwiek poszło nie tak z kopertą — celowo bez szczegółów na zewnątrz."""
# --------------------------------------------------------------------- klucze
def parse_key(raw: str) -> bytes:
"""Klucz z konfiguracji: hex (64 znaki) albo base64. Zawsze 32 bajty."""
text = raw.strip()
if not text:
raise LinkError("pusty klucz łącza")
try:
key = bytes.fromhex(text)
except ValueError:
try:
key = base64.b64decode(text, validate=True)
except (binascii.Error, ValueError) as exc:
raise LinkError("klucz łącza nie jest ani hexem, ani base64") from exc
if len(key) != KEY_BYTES:
raise LinkError(
f"klucz łącza ma {len(key)} B zamiast {KEY_BYTES} — wygeneruj przez "
f"`openssl rand -hex 32`"
)
return key
def key_from_env(env_name: str) -> bytes | None:
"""Klucz albo None. Zły klucz to wyjątek OD RAZU — nie przy pierwszym żądaniu."""
raw = os.getenv(env_name, "")
return parse_key(raw) if raw.strip() else None
def encryption_required() -> bool:
"""Czy brak klucza ma być błędem, a nie cichym powrotem do jawnego ruchu.
Serwer sam z siebie broni się fail-closed, ale to za mało: klient BEZ klucza
wysyła pytanie otwartym tekstem i dopiero potem dostaje odmowę — czyli treść
zapytania zdążyła już przelecieć przez sieć. Ta flaga zatrzymuje go, zanim
cokolwiek opuści proces. Ustawiana razem z kluczami we wdrożeniu.
"""
return os.getenv(ENV_REQUIRED, "").strip().lower() in {"1", "true", "yes", "on"}
def _subkey(link_key: bytes, direction: bytes) -> bytes:
return HKDF(
algorithm=hashes.SHA256(), length=KEY_BYTES, salt=None,
info=b"astrololo/link/" + direction,
).derive(link_key)
class Link:
"""Jedna para rozmówców: klucz plus wyprowadzone z niego podklucze."""
def __init__(self, link_key: bytes) -> None:
self._by_direction = {
REQUEST: AESGCM(_subkey(link_key, REQUEST)),
RESPONSE: AESGCM(_subkey(link_key, RESPONSE)),
}
# ---------------------------------------------------------- pojedyncza ramka
def _aad(self, direction: bytes, path: str, stamp: str, seq: int) -> bytes:
return b"|".join([MAGIC, direction, path.encode("utf-8"),
stamp.encode("ascii"), str(seq).encode("ascii")])
def seal(self, direction: bytes, path: str, stamp: str, seq: int,
plaintext: bytes) -> bytes:
nonce = os.urandom(NONCE_BYTES)
sealed = self._by_direction[direction].encrypt(
nonce, plaintext, self._aad(direction, path, stamp, seq))
return MAGIC + nonce + sealed
def open(self, direction: bytes, path: str, stamp: str, seq: int,
frame: bytes) -> bytes:
if not frame.startswith(MAGIC):
raise LinkError("ramka bez znacznika astrololo")
body = frame[len(MAGIC):]
if len(body) <= NONCE_BYTES:
raise LinkError("ramka za krótka")
nonce, sealed = body[:NONCE_BYTES], body[NONCE_BYTES:]
try:
return self._by_direction[direction].decrypt(
nonce, sealed, self._aad(direction, path, stamp, seq))
except InvalidTag as exc:
# Jeden komunikat na wszystkie przypadki: zły klucz, podmieniony bajt,
# przeklejenie z innej ścieżki, przestawiona ramka. Rozróżnianie ich
# na zewnątrz podpowiadałoby atakującemu, w co trafił.
raise LinkError("nie udało się odszyfrować — zły klucz albo naruszone dane") from exc
# ------------------------------------------------------------ strumień ramek
def seal_stream(self, direction: bytes, path: str, stamp: str,
chunks: Iterable[bytes]) -> Iterator[bytes]:
for seq, chunk in enumerate(chunks):
yield frame_out(self.seal(direction, path, stamp, seq, chunk))
def open_stream(self, direction: bytes, path: str, stamp: str,
raw: bytes) -> Iterator[bytes]:
for seq, frame in enumerate(frames_in(raw)):
yield self.open(direction, path, stamp, seq, frame)
def open_all(self, direction: bytes, path: str, stamp: str, raw: bytes) -> bytes:
return b"".join(self.open_stream(direction, path, stamp, raw))
# ---------------------------------------------------------------- ramkowanie
def frame_out(payload: bytes) -> bytes:
return struct.pack(">I", len(payload)) + payload
def frames_in(raw: bytes) -> Iterator[bytes]:
"""Rozbiera bufor na ramki. Ucięty strumień to błąd, nie cicha strata danych."""
offset = 0
while offset < len(raw):
if offset + LENGTH_PREFIX > len(raw):
raise LinkError("urwana ramka (brak nagłówka długości)")
(size,) = struct.unpack(">I", raw[offset:offset + LENGTH_PREFIX])
if size > MAX_FRAME:
raise LinkError("ramka ponad dopuszczalny rozmiar")
offset += LENGTH_PREFIX
if offset + size > len(raw):
raise LinkError("urwana ramka (za mało danych)")
yield raw[offset:offset + size]
offset += size
def unframe_incremental(buffer: bytearray) -> Iterator[bytes]:
"""Wyjmuje z bufora KOMPLETNE ramki i zjada je; resztę zostawia na później.
Dla odbioru na żywo: kawałki przychodzą podzielone dowolnie i ramka potrafi
rozjechać się między dwa odczyty.
"""
while True:
if len(buffer) < LENGTH_PREFIX:
return
(size,) = struct.unpack(">I", buffer[:LENGTH_PREFIX])
if size > MAX_FRAME:
raise LinkError("ramka ponad dopuszczalny rozmiar")
if len(buffer) < LENGTH_PREFIX + size:
return
frame = bytes(buffer[LENGTH_PREFIX:LENGTH_PREFIX + size])
del buffer[:LENGTH_PREFIX + size]
yield frame
# ------------------------------------------------------------- świeżość ruchu
def stamp_now() -> str:
return f"{time.time():.3f}"
def check_stamp(stamp: str) -> None:
"""Odrzuca ramki spoza okna czasowego — inaczej podsłuchane żądanie dałoby się
odtworzyć w dowolnym momencie w przyszłości."""
try:
sent = float(stamp)
except (TypeError, ValueError) as exc:
raise LinkError("brak albo błędny znacznik czasu") from exc
if abs(time.time() - sent) > MAX_SKEW_SECONDS:
raise LinkError("znacznik czasu poza dopuszczalnym oknem")
# =========================================================== strona serwerowa
class LinkCryptoMiddleware:
"""Rozszyfrowuje wchodzące żądania i zaszyfrowuje wychodzące odpowiedzi.
Napisane jako czyste ASGI, nie ``@app.middleware("http")``, bo trzeba
podmienić CIAŁO żądania jeszcze zanim zobaczy je FastAPI, oraz przepuścić
odpowiedź strumieniową kawałek po kawałku, bez zbierania jej w pamięci.
"""
def __init__(self, app, link: Link | None, layer: str) -> None:
self.app = app
self.link = link
self.layer = layer
async def __call__(self, scope, receive, send):
if scope["type"] != "http" or self.link is None or scope["path"] in PUBLIC_PATHS:
return await self.app(scope, receive, send)
path = scope["path"]
headers = {k.decode("latin-1").lower(): v.decode("latin-1") for k, v in scope["headers"]}
if headers.get(HEADER_ENC.lower()) != VERSION:
# Fail-closed. Klucz jest ustawiony, więc jawne żądanie oznacza albo
# pomyłkę w konfiguracji, albo kogoś obcego — w obu wypadkach nie
# chcemy po cichu wrócić do jawnego ruchu.
log.warning("warstwa %s: odrzucone żądanie bez szyfrowania łącza (%s)",
self.layer, path)
return await _refuse(send, "Łącze międzywarstwowe wymaga szyfrowania.")
stamp = headers.get(HEADER_TS.lower(), "")
try:
check_stamp(stamp)
plaintext = self.link.open_all(REQUEST, path, stamp, await _read_body(receive))
except LinkError as exc:
log.warning("warstwa %s: %s (%s)", self.layer, exc, path)
return await _refuse(send, "Nie udało się odczytać zaszyfrowanego żądania.")
scope = dict(scope)
scope["headers"] = _rewritten_headers(scope["headers"], len(plaintext))
await self.app(scope, _replay(plaintext, receive), self._sealing_send(send, path))
def _sealing_send(self, send, path: str):
state: dict = {"stamp": "", "seq": 0}
async def sealing(message):
if message["type"] == "http.response.start":
state["stamp"] = stamp_now()
keep = [(k, v) for k, v in message.get("headers", [])
if k.lower() not in (b"content-length", b"content-type")]
message = dict(message)
message["headers"] = keep + [
(b"content-type", CONTENT_TYPE.encode()),
(HEADER_ENC.lower().encode(), VERSION.encode()),
(HEADER_TS.lower().encode(), state["stamp"].encode()),
]
return await send(message)
if message["type"] == "http.response.body":
chunk = message.get("body", b"")
sealed = b""
if chunk:
sealed = frame_out(self.link.seal(
RESPONSE, path, state["stamp"], state["seq"], chunk))
state["seq"] += 1
return await send({"type": "http.response.body", "body": sealed,
"more_body": message.get("more_body", False)})
return await send(message)
return sealing
def _rewritten_headers(raw: Iterable[tuple[bytes, bytes]], length: int):
"""Po odszyfrowaniu ciało ma inną długość i zwykły typ — inaczej FastAPI
próbowałby sparsować JSON o cudzej deklarowanej wielkości."""
kept = [(k, v) for k, v in raw if k.lower() not in (b"content-length", b"content-type")]
kept.append((b"content-length", str(length).encode()))
if length:
kept.append((b"content-type", b"application/json"))
return kept
async def _read_body(receive) -> bytes:
body = bytearray()
while True:
message = await receive()
if message["type"] == "http.disconnect":
raise LinkError("rozłączenie w trakcie odbioru żądania")
body += message.get("body", b"")
if not message.get("more_body", False):
return bytes(body)
def _replay(body: bytes, original):
"""Podstawia odszyfrowane ciało jako jedyną porcję wejścia dla aplikacji.
Po oddaniu ciała oddajemy głos ORYGINALNEMU `receive`, zamiast od razu
zgłaszać rozłączenie. Odpowiedź strumieniowa nasłuchuje bowiem rozłączenia
równolegle do wysyłania i przerywa się, gdy je zobaczy — na skróconej wersji
okno postępu dostawało pustą odpowiedź, choć zwykłe żądania działały.
"""
delivered = False
async def receive():
nonlocal delivered
if delivered:
return await original()
delivered = True
return {"type": "http.request", "body": body, "more_body": False}
return receive
async def _refuse(send, detail: str) -> None:
"""Odmowa leci JAWNIE — rozmówca właśnie pokazał, że nie umie odszyfrować,
więc zaszyfrowany komunikat o błędzie byłby dla niego nieczytelny."""
payload = f'{{"detail":"{detail}"}}'.encode("utf-8")
await send({"type": "http.response.start", "status": 400, "headers": [
(b"content-type", b"application/json"),
(b"content-length", str(len(payload)).encode()),
]})
await send({"type": "http.response.body", "body": payload})
def install(app, env_name: str, layer: str):
"""Podpina szyfrowanie łącza. Wołać PO `security.install`, żeby także odmowa
tokenowa (401) wracała zaszyfrowana — inaczej klient by jej nie odczytał."""
link_key = key_from_env(env_name)
if link_key is None and encryption_required():
# Celowo wywracamy start. Ta sama zasada co przy sekrecie logowania:
# wolimy widoczną awarię niż usługę, która wstała i po cichu nie chroni
# niczego. Pod w CrashLoop widać od razu, jawny ruch — nie.
raise LinkError(
f"{ENV_REQUIRED} jest włączone, ale {env_name} nie ustawiony — "
f"warstwa {layer} nie wystartuje bez klucza łącza"
)
if link_key is None:
log.warning(
"UWAGA: %s nie ustawiony — warstwa %s rozmawia z sąsiadem JAWNYM tekstem, "
"więc treść baz interpretacyjnych jest widoczna dla każdego, kto podsłucha "
"ruch wewnątrz sieci.", env_name, layer,
)
return None
link = Link(link_key)
app.add_middleware(LinkCryptoMiddleware, link=link, layer=layer)
log.info("warstwa %s: łącze szyfrowane (AES-256-GCM, klucz z %s)", layer, env_name)
return link
# ============================================================ strona kliencka
def call(client, method: str, url: str, *, payload=None,
headers: dict[str, str] | None = None, link: Link | None) -> bytes:
"""Żądanie do sąsiedniej warstwy; zwraca odszyfrowane ciało odpowiedzi.
Ścieżkę do materiału uwierzytelnianego bierzemy Z URL-a, a nie z osobnego
argumentu — gdyby klient i serwer liczyły ją inaczej, każde żądanie kończyłoby
się niejasnym błędem odszyfrowania.
"""
import json as _json
import httpx
request_headers = dict(headers or {})
if link is None:
if encryption_required():
# Zatrzymujemy się PRZED wysłaniem. Gdyby polecieć jawnie i dopiero
# zebrać odmowę, pytanie byłoby już na kablu — a to właśnie ono niesie
# sygnifikatory, o które pytamy bazę.
raise LinkError(
f"{ENV_REQUIRED} jest włączone, ale brak klucza łącza — żądanie "
f"NIE zostało wysłane, żeby jego treść nie poszła jawnym tekstem"
)
response = client.request(method, url, json=payload, headers=request_headers)
response.raise_for_status()
return response.content
path = httpx.URL(url).path
stamp = stamp_now()
plaintext = b"" if payload is None else _json.dumps(payload).encode("utf-8")
body = frame_out(link.seal(REQUEST, path, stamp, 0, plaintext))
request_headers.update({HEADER_ENC: VERSION, HEADER_TS: stamp,
"Content-Type": CONTENT_TYPE})
response = client.request(method, url, content=body, headers=request_headers)
if response.status_code >= 400 and response.headers.get(HEADER_ENC) != VERSION:
log.error("łącze %s odmówiło: %s", path, response.text[:200])
response.raise_for_status()
if response.headers.get(HEADER_ENC) != VERSION:
raise LinkError("odpowiedź przyszła nieszyfrowana, choć klucz łącza jest ustawiony")
reply_stamp = response.headers.get(HEADER_TS, "")
check_stamp(reply_stamp)
return link.open_all(RESPONSE, path, reply_stamp, response.content)
def call_json(client, method: str, url: str, *, payload=None,
headers: dict[str, str] | None = None, link: Link | None):
import json as _json
return _json.loads(call(client, method, url, payload=payload,
headers=headers, link=link))
def open_response_stream(response, link: Link | None) -> Iterator[bytes]:
"""Odbiór odpowiedzi płynącej kawałkami (okno postępu).
Ramka potrafi rozjechać się między dwa odczyty z gniazda, więc składamy ją
w buforze zamiast zakładać, że każdy kawałek to komplet.
"""
if link is None:
yield from response.iter_bytes()
return
if response.headers.get(HEADER_ENC) != VERSION:
raise LinkError("strumień przyszedł nieszyfrowany, choć klucz łącza jest ustawiony")
stamp = response.headers.get(HEADER_TS, "")
check_stamp(stamp)
path = response.request.url.path
buffer = bytearray()
seq = 0
for chunk in response.iter_bytes():
buffer += chunk
for frame in unframe_incremental(buffer):
yield link.open(RESPONSE, path, stamp, seq, frame)
seq += 1
if buffer:
raise LinkError("strumień urwał się w połowie ramki")
+31 -2
View File
@@ -124,7 +124,14 @@ class _Driver:
"""(tekst, czy_ucięta, zużycie, nazwa_modelu, powód_zakończenia)."""
raise NotImplementedError
def generate(self, prompt: str, max_tokens: int) -> Completion:
def generate(self, prompt: str, max_tokens: int, on_event=None) -> Completion:
"""`on_event(dict)` dostaje zdarzenia postępu — UI pokazuje z nich log.
Raportujemy KAŻDĄ turę, bo to ona trwa; bez tego pasek postępu byłby
ozdobnikiem, a nie informacją."""
def emit(kind: str, message: str, **extra):
if on_event:
on_event({"type": kind, "message": message, **extra})
messages: list[dict] = [{"role": "user", "content": prompt}]
parts: list[str] = []
usage: dict = {}
@@ -137,9 +144,29 @@ class _Driver:
while turns <= MAX_CONTINUATIONS:
turns += 1
budget = max(256, min(remaining, TURN_TOKENS_CAP))
emit("turn_start",
f"Tura {turns}: wysyłam do modelu {self.model} (limit {budget} tokenów)…",
turn=turns)
started = time.monotonic()
text, truncated, turn_usage, model_name, stop = self._turn(messages, budget)
took = time.monotonic() - started
_merge_usage(usage, turn_usage)
remaining -= budget
# Odejmujemy tokeny FAKTYCZNIE wyprodukowane, nie zamówiony limit tury.
# Inaczej pierwsza tura zjadałaby cały budżet i urwana odpowiedź nigdy
# nie doczekałaby się kontynuacji — wracałby do użytkownika fragment
# udający całość.
produced = (turn_usage or {}).get("completion_tokens")
if produced is None:
produced = (turn_usage or {}).get("output_tokens")
if produced is None:
produced = max(1, int(len(text) / 3.6))
remaining -= max(1, int(produced))
emit("turn_end",
f"Tura {turns}: odebrano {len(text.strip())} znaków w {took:.1f}s"
+ (" — odpowiedź urwana, poproszę o dokończenie" if truncated else ""),
turn=turns, chars=len(text.strip()), truncated=truncated)
chunk = text.strip()
if chunk:
@@ -172,6 +199,8 @@ class _Driver:
final = _join(parts)
if not final:
raise LLMError(_explain_empty(turns, usage, stop))
emit("generated", f"Gotowe: {len(final)} znaków w {turns} turach.",
chars=len(final), turns=turns)
usage["turns"] = turns
return Completion(text=final, model=model_name, provider=self.name,
+67 -4
View File
@@ -12,7 +12,7 @@ import httpx
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
from app import link_crypto, security
from app import security
from app.clients.data_client import DataClient
from app.models import QueryRequest, QueryResponse
from app.service import QueryService
@@ -20,9 +20,6 @@ from app.service import QueryService
app = FastAPI(title="astrololo · warstwa logiczna")
service = QueryService()
security.install(app, "logiczna") # token międzywarstwowy (LOG-32)
# Szyfrowanie łącza od prezentacji. PO `security.install`, żeby także odmowa
# tokenowa wracała zaszyfrowana — inaczej klient nie umiałby jej odczytać.
link_crypto.install(app, link_crypto.ENV_PRESENTATION_LOGIC, "logiczna")
# --- silnik efemeryd (LOG-24): budowany leniwie, by nie wymagać Skyfielda do startu ---
_engine = None
@@ -268,6 +265,72 @@ def chart_horoscope(req: HoroscopeRequest) -> dict:
return out
@app.post("/chart/horoscope/stream")
def chart_horoscope_stream(req: HoroscopeRequest):
"""To samo co /chart/horoscope, ale strumieniuje POSTĘP w trakcie pracy.
Pisanie horoskopu trwa minutami — bez sygnału aplikacja wygląda na zawieszoną.
Strumień (NDJSON, jedna linia = jedno zdarzenie) niesie RZECZYWISTE etapy:
budowę promptu, limity modelu i każdą turę generowania. Ostatnie zdarzenie
(`result`) ma identyczny kształt co odpowiedź zwykłego endpointu.
"""
from fastapi.responses import StreamingResponse
from app.progress import stream
def work(emit) -> dict:
from app.llm.base import LLMError
from app.llm.factory import build_provider
from app.llm.limits import plan
emit({"type": "stage", "message": "Liczę horoskop i szukam wskazań w bazach…"})
out = chart_prompt(req)
st = out.get("stats", {})
emit({"type": "stage", "message":
f"Prompt gotowy: {st.get('chars', 0)} znaków, "
f"wskazań {st.get('included', 0)}"
+ (f", pominięto {st['omitted']}" if st.get("omitted") else "")})
if out.get("data_error"):
emit({"type": "warn", "message": out["data_error"]})
try:
provider = build_provider(req.provider, req.model)
emit({"type": "stage", "message":
f"Dostawca: {provider.name}, model: {provider.model}"
+ ("" if not provider.leaves_lan else " — dane opuszczają sieć")})
emit({"type": "stage", "message": "Liczę tokeny promptu…"})
prompt_tokens = provider.count_tokens(out["prompt"])
budget = plan(provider.name, provider.model, prompt_tokens, req.max_tokens)
out["token_plan"] = budget
emit({"type": "stage", "message":
f"Prompt {prompt_tokens} tok. · okno modelu {budget['context_window']} · "
f"na odpowiedź {budget['max_output']}"})
for warning in budget["warnings"]:
emit({"type": "warn", "message": warning})
if budget["warnings"]:
out["warnings"] = budget["warnings"]
if not budget["fits"]:
out["llm_error"] = " ".join(budget["warnings"])
return out
result = provider.generate(out["prompt"], budget["max_output"], on_event=emit)
except LLMError as e:
emit({"type": "warn", "message": f"Model zawiódł: {e}"})
out["llm_error"] = str(e)
return out
out.update(horoscope=result.text, provider=result.provider, model=result.model,
leaves_lan=result.leaves_lan, usage=result.usage)
return out
return StreamingResponse(
stream(work),
media_type="application/x-ndjson",
headers={"Cache-Control": "no-store", "X-Accel-Buffering": "no"},
)
@app.get("/llm/models")
def llm_models() -> dict:
"""Podpowiedzi modeli per dostawca — UI buduje z tego listę wyboru.
+70
View File
@@ -0,0 +1,70 @@
"""Strumień postępu długiej operacji (NDJSON).
Po co: pisanie horoskopu trwa czasem minuty. Bez sygnału aplikacja wygląda na
zawieszoną. Zamiast udawanego paska postępu strumieniujemy **rzeczywiste**
zdarzenia z kolejnych etapów, żeby log pokazywał to, co faktycznie się dzieje.
Dlaczego NDJSON, a nie SSE: `EventSource` w przeglądarce obsługuje wyłącznie GET,
a to jest POST z ciałem. Strumień jedna linia = jeden obiekt JSON" czyta się
zwykłym `fetch()` i jest trywialny do sparsowania.
Dlaczego wątek: właściwa praca (silnik, baza, model) jest synchroniczna. Puszczamy
w wątku roboczym, a generator odpompowuje kolejkę zdarzeń dzięki temu
zdarzenia docierają w trakcie pracy, a nie dopiero na końcu.
"""
from __future__ import annotations
import json
import queue
import threading
import traceback
from collections.abc import Iterator
from typing import Any, Callable
_HEARTBEAT_SECONDS = 10.0
_DONE = object()
def line(kind: str, message: str, **extra: Any) -> str:
return json.dumps({"type": kind, "message": message, **extra}, ensure_ascii=False) + "\n"
def stream(work: Callable[[Callable[[dict], None]], dict]) -> Iterator[str]:
"""Uruchamia `work(emit)` w wątku i strumieniuje zdarzenia w czasie rzeczywistym.
`work` dostaje funkcję `emit(zdarzenie)` i zwraca końcowy wynik, który leci
jako ostatnie zdarzenie typu `result`. Wyjątek zamienia się w zdarzenie `error`
połączenie nigdy nie urywa się bez wyjaśnienia.
"""
events: queue.Queue = queue.Queue()
def emit(event: dict) -> None:
events.put(event)
def run() -> None:
try:
result = work(emit)
events.put({"type": "result", "message": "Gotowe.", "result": result})
except Exception as e: # noqa: BLE001 — zgłaszamy KAŻDY błąd
events.put({
"type": "error",
"message": f"{type(e).__name__}: {e}",
"detail": traceback.format_exc(limit=3),
})
finally:
events.put(_DONE)
worker = threading.Thread(target=run, daemon=True)
worker.start()
while True:
try:
event = events.get(timeout=_HEARTBEAT_SECONDS)
except queue.Empty:
# cisza dłuższa niż heartbeat: dajemy znak życia, żeby pośredniki
# (proxy, load balancer) nie uznały połączenia za martwe
yield line("ping", "")
continue
if event is _DONE:
break
yield json.dumps(event, ensure_ascii=False) + "\n"
-2
View File
@@ -4,5 +4,3 @@ httpx>=0.28
pydantic>=2.10
# Silnik własny (ścieżka A, permisywny): Skyfield (MIT) + dane JPL (public domain)
skyfield>=1.49
# Szyfrowanie łącza między warstwami (PRE-16): AES-256-GCM + HKDF
cryptography>=44.0
-351
View File
@@ -1,351 +0,0 @@
"""Szyfrowanie łącza między warstwami (PRE-16).
Sedno: przez to łącze płyną surowe wiersze oryginalnych baz interpretacyjnych.
Testy nie sprawdzają więc tylko, czy coś się zaszyfrowało i odszyfrowało"
sprawdzają, czy tajnego tekstu FAKTYCZNIE NIE MA w bajtach lecących po sieci
oraz czy każda znana droga na skróty (brak szyfrowania, obcy klucz, podmieniony
bajt, przeklejenie na inny endpoint, odtworzenie po czasie) kończy się odmową.
"""
import json
import pathlib
import time
import pytest
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from pydantic import BaseModel
from starlette.testclient import TestClient
from app import link_crypto
from app.link_crypto import ENV_LOGIC_DATA, ENV_PRESENTATION_LOGIC, Link, LinkError
SECRET = "Saturn w VII domu — opis z oryginalnej bazy interpretacyjnej"
KEY_A = "11" * 32 # hex, 32 bajty
KEY_B = "22" * 32
class Payload(BaseModel):
question: str
def _app(link: Link | None) -> FastAPI:
app = FastAPI()
if link is not None:
app.add_middleware(link_crypto.LinkCryptoMiddleware, link=link, layer="testowa")
@app.get("/health")
def health():
return {"status": "ok"}
@app.post("/search")
def search(payload: Payload):
return {"echo": payload.question, "interpretation": SECRET}
@app.get("/catalog")
def catalog():
return {"models": ["a", "b"]}
@app.post("/stream")
def stream():
def lines():
for i in range(4):
yield json.dumps({"step": i, "note": SECRET}).encode() + b"\n"
return StreamingResponse(lines(), media_type="application/x-ndjson")
return app
@pytest.fixture
def link():
return Link(link_crypto.parse_key(KEY_A))
# ------------------------------------------------------------------- klucze
def test_key_accepts_hex_and_base64():
import base64
raw = bytes(range(32))
assert link_crypto.parse_key(raw.hex()) == raw
assert link_crypto.parse_key(base64.b64encode(raw).decode()) == raw
def test_key_of_wrong_length_is_rejected_loudly():
"""Krótki klucz to nie „słabsze szyfrowanie", tylko błąd konfiguracji."""
with pytest.raises(LinkError, match="32"):
link_crypto.parse_key("aabb")
def test_key_that_is_neither_hex_nor_base64_is_rejected():
with pytest.raises(LinkError, match="hex"):
link_crypto.parse_key("!!! to nie jest klucz !!!")
def test_key_from_env_is_lazy(monkeypatch):
monkeypatch.delenv(ENV_PRESENTATION_LOGIC, raising=False)
assert link_crypto.key_from_env(ENV_PRESENTATION_LOGIC) is None
monkeypatch.setenv(ENV_PRESENTATION_LOGIC, KEY_A)
assert link_crypto.key_from_env(ENV_PRESENTATION_LOGIC) == bytes.fromhex(KEY_A)
def test_directions_use_different_subkeys(link):
"""Żądanie i odpowiedź nie dzielą klucza — powtórzenie jednorazówki w jedną
stronę nie osłabia drugiej."""
stamp = link_crypto.stamp_now()
sealed = link.seal(link_crypto.REQUEST, "/search", stamp, 0, b"tajne")
with pytest.raises(LinkError):
link.open(link_crypto.RESPONSE, "/search", stamp, 0, sealed)
# --------------------------------------------------------- podstawowy obieg
def test_round_trip_delivers_plaintext_to_the_app(link):
with TestClient(_app(link)) as client:
got = link_crypto.call_json(client, "POST", "http://testserver/search",
payload={"question": "Saturn"}, link=link)
assert got["echo"] == "Saturn"
assert got["interpretation"] == SECRET
def test_get_without_body_also_works(link):
"""GET nie ma ciała, ale i tak pieczętujemy pustą kopertę — to ona dowodzi,
że pytający ma klucz, i ona wymusza zaszyfrowanie odpowiedzi."""
with TestClient(_app(link)) as client:
got = link_crypto.call_json(client, "GET", "http://testserver/catalog", link=link)
assert got == {"models": ["a", "b"]}
def test_without_key_traffic_stays_plaintext(link):
"""Dev bez sekretów ma działać jak dotąd — inaczej nikt nie odpali projektu lokalnie."""
with TestClient(_app(None)) as client:
got = link_crypto.call_json(client, "POST", "http://testserver/search",
payload={"question": "Saturn"}, link=None)
assert got["interpretation"] == SECRET
def test_health_stays_open_for_kubernetes_probes(link):
"""Sondy k8s klucza nie mają. Gdyby /health wymagał szyfrowania, literówka
w sekrecie kładłaby pody zamiast pokazać błąd w aplikacji."""
with TestClient(_app(link)) as client:
assert client.get("http://testserver/health").json() == {"status": "ok"}
# ------------------------------------------- czy na kablu naprawdę nic nie widać
def _raw_exchange(client, link, payload):
"""Wysyła zapieczętowane żądanie i zwraca SUROWE bajty obu stron."""
stamp = link_crypto.stamp_now()
body = link_crypto.frame_out(
link.seal(link_crypto.REQUEST, "/search", stamp, 0, json.dumps(payload).encode()))
response = client.request(
"POST", "http://testserver/search", content=body,
headers={link_crypto.HEADER_ENC: link_crypto.VERSION,
link_crypto.HEADER_TS: stamp,
"Content-Type": link_crypto.CONTENT_TYPE})
return body, response
def test_request_bytes_do_not_contain_the_question(link):
with TestClient(_app(link)) as client:
body, _ = _raw_exchange(client, link, {"question": "Saturn w VII"})
assert b"Saturn" not in body
assert b"question" not in body
def test_response_bytes_do_not_contain_the_interpretation(link):
"""To jest właściwy powód istnienia całego modułu."""
with TestClient(_app(link)) as client:
_, response = _raw_exchange(client, link, {"question": "Saturn"})
assert response.status_code == 200
assert SECRET.encode() not in response.content
assert b"interpretation" not in response.content
assert response.headers[link_crypto.HEADER_ENC] == link_crypto.VERSION
assert response.headers["content-type"] == link_crypto.CONTENT_TYPE
# ------------------------------------------------------- drogi na skróty i ataki
def test_plaintext_request_is_refused_when_key_is_set(link):
"""Fail-closed: regresja po stronie klienta nie może oznaczać cichego
powrotu do jawnego ruchu."""
with TestClient(_app(link)) as client:
response = client.post("http://testserver/search", json={"question": "Saturn"})
assert response.status_code == 400
assert SECRET.encode() not in response.content
def test_foreign_key_cannot_read_the_link(link):
"""Klucze są osobne dla każdej pary warstw — przejęcie jednego nie otwiera drugiej."""
intruder = Link(link_crypto.parse_key(KEY_B))
with TestClient(_app(link)) as client:
stamp = link_crypto.stamp_now()
body = link_crypto.frame_out(
intruder.seal(link_crypto.REQUEST, "/search", stamp, 0, b'{"question":"x"}'))
response = client.request(
"POST", "http://testserver/search", content=body,
headers={link_crypto.HEADER_ENC: link_crypto.VERSION,
link_crypto.HEADER_TS: stamp,
"Content-Type": link_crypto.CONTENT_TYPE})
assert response.status_code == 400
assert SECRET.encode() not in response.content
def test_single_flipped_bit_is_rejected(link):
"""GCM uwierzytelnia, więc nie ma wariantu „odszyfrowało się, ale zmienione"."""
with TestClient(_app(link)) as client:
stamp = link_crypto.stamp_now()
sealed = bytearray(link.seal(link_crypto.REQUEST, "/search", stamp, 0,
b'{"question":"x"}'))
sealed[-1] ^= 0x01
response = client.request(
"POST", "http://testserver/search",
content=link_crypto.frame_out(bytes(sealed)),
headers={link_crypto.HEADER_ENC: link_crypto.VERSION,
link_crypto.HEADER_TS: stamp,
"Content-Type": link_crypto.CONTENT_TYPE})
assert response.status_code == 400
def test_frame_cannot_be_replayed_against_another_endpoint(link):
"""Ścieżka wchodzi do materiału uwierzytelnianego, więc podsłuchanej koperty
nie da się przekleić tam, gdzie odpowiedź byłaby ciekawsza."""
stamp = link_crypto.stamp_now()
sealed = link.seal(link_crypto.REQUEST, "/catalog", stamp, 0, b"")
with pytest.raises(LinkError):
link.open(link_crypto.REQUEST, "/search", stamp, 0, sealed)
def test_frames_cannot_be_reordered(link):
"""Numer ramki jest uwierzytelniony — przestawienie kolejności w strumieniu
to błąd, a nie po cichu pomieszany horoskop."""
stamp = link_crypto.stamp_now()
second = link.seal(link_crypto.RESPONSE, "/stream", stamp, 1, b"druga")
with pytest.raises(LinkError):
link.open(link_crypto.RESPONSE, "/stream", stamp, 0, second)
def test_stale_frame_is_refused(link, monkeypatch):
"""Bez okna czasowego podsłuchane żądanie dałoby się odtworzyć kiedykolwiek."""
old = f"{time.time() - link_crypto.MAX_SKEW_SECONDS - 60:.3f}"
with TestClient(_app(link)) as client:
body = link_crypto.frame_out(
link.seal(link_crypto.REQUEST, "/search", old, 0, b'{"question":"x"}'))
response = client.request(
"POST", "http://testserver/search", content=body,
headers={link_crypto.HEADER_ENC: link_crypto.VERSION,
link_crypto.HEADER_TS: old,
"Content-Type": link_crypto.CONTENT_TYPE})
assert response.status_code == 400
def test_truncated_stream_is_an_error_not_silent_loss(link):
stamp = link_crypto.stamp_now()
full = link_crypto.frame_out(link.seal(link_crypto.RESPONSE, "/x", stamp, 0, b"abc"))
with pytest.raises(LinkError, match="urwana"):
link.open_all(link_crypto.RESPONSE, "/x", stamp, full[:-2])
def test_absurd_frame_length_does_not_allocate(link):
"""Zadeklarowana długość pochodzi z sieci — nie wolno jej wierzyć na słowo."""
import struct
with pytest.raises(LinkError, match="rozmiar"):
list(link_crypto.frames_in(struct.pack(">I", 2 ** 31) + b"nic"))
# ---------------------------------------------------------- odpowiedź strumieniowa
def test_streaming_response_survives_encryption(link):
"""Okno postępu dostaje kolejne linie na żywo — muszą dojść po kolei
i w komplecie, mimo że każda jedzie w osobnej kopercie."""
with TestClient(_app(link)) as client:
stamp = link_crypto.stamp_now()
body = link_crypto.frame_out(link.seal(link_crypto.REQUEST, "/stream", stamp, 0, b""))
with client.stream("POST", "http://testserver/stream", content=body,
headers={link_crypto.HEADER_ENC: link_crypto.VERSION,
link_crypto.HEADER_TS: stamp,
"Content-Type": link_crypto.CONTENT_TYPE}) as response:
chunks = list(link_crypto.open_response_stream(response, link))
steps = [json.loads(line) for line in b"".join(chunks).splitlines()]
assert [s["step"] for s in steps] == [0, 1, 2, 3]
assert all(s["note"] == SECRET for s in steps)
def test_incremental_unframing_handles_split_frames(link):
"""Ramka potrafi rozjechać się między dwa odczyty z gniazda — składamy ją
w buforze, zamiast zakładać, że każdy kawałek to komplet."""
stamp = link_crypto.stamp_now()
stream = b"".join(link.seal_stream(link_crypto.RESPONSE, "/x", stamp,
[b"raz", b"dwa", b"trzy"]))
buffer = bytearray()
opened, seq = [], 0
for i in range(0, len(stream), 5): # ciachamy w poprzek ramek
buffer += stream[i:i + 5]
for frame in link_crypto.unframe_incremental(buffer):
opened.append(link.open(link_crypto.RESPONSE, "/x", stamp, seq, frame))
seq += 1
assert opened == [b"raz", b"dwa", b"trzy"]
assert not buffer, "bufor musi zostać pusty — inaczej gdzieś zgubiliśmy ramkę"
# ----------------------------------------------- trzy kopie muszą być identyczne
def test_all_three_services_share_the_same_module():
"""Moduł jest skopiowany do trzech niezależnych usług (nie mają wspólnej
biblioteki). Rozjazd między kopiami objawiłby się dopiero na produkcji jako
nie da się odszyfrować" — więc pilnujemy tego testem."""
root = pathlib.Path(__file__).resolve().parents[3]
copies = {svc: (root / "services" / svc / "app" / "link_crypto.py")
for svc in ("presentation", "logic", "data")}
missing = [svc for svc, path in copies.items() if not path.is_file()]
assert not missing, f"brak modułu w warstwach: {missing}"
contents = {svc: path.read_bytes() for svc, path in copies.items()}
assert len(set(contents.values())) == 1, (
"kopie link_crypto.py rozjechały się między warstwami: "
+ ", ".join(f"{svc}={len(body)}B" for svc, body in contents.items())
)
def test_env_names_are_two_distinct_keys():
"""Wymóg wprost: osobny klucz dla pary prezentacja-logika i logika-dane."""
assert ENV_PRESENTATION_LOGIC != ENV_LOGIC_DATA
# ------------------------------------------------- klient też musi być fail-closed
#
# To wyszło dopiero z podsłuchu prawdziwego gniazda, nie z testów: przy kliencie
# BEZ klucza serwer owszem odmawiał, ale pytanie leciało po drodze otwartym
# tekstem. Odpowiedź była chroniona — zapytanie już nie.
def test_client_without_key_sends_nothing_when_encryption_required(monkeypatch):
monkeypatch.setenv(link_crypto.ENV_REQUIRED, "true")
sent = []
class Tripwire:
def request(self, *args, **kwargs):
sent.append(args)
raise AssertionError("żądanie NIE powinno opuścić procesu")
with pytest.raises(LinkError, match=link_crypto.ENV_REQUIRED):
link_crypto.call(Tripwire(), "POST", "http://logic/search",
payload={"value": "[Sat"}, link=None)
assert not sent, "treść zapytania wyszłaby jawnym tekstem"
def test_plaintext_still_allowed_in_dev(monkeypatch, link):
"""Bez tej flagi lokalne uruchomienie bez sekretów ma dalej działać."""
monkeypatch.delenv(link_crypto.ENV_REQUIRED, raising=False)
with TestClient(_app(None)) as client:
got = link_crypto.call_json(client, "POST", "http://testserver/search",
payload={"question": "Saturn"}, link=None)
assert got["interpretation"] == SECRET
def test_service_refuses_to_start_without_key_when_required(monkeypatch):
"""Pod w CrashLoop widać od razu; usługę, która wstała i nie szyfruje — nie."""
monkeypatch.setenv(link_crypto.ENV_REQUIRED, "true")
monkeypatch.delenv(ENV_LOGIC_DATA, raising=False)
with pytest.raises(LinkError, match=ENV_LOGIC_DATA):
link_crypto.install(FastAPI(), ENV_LOGIC_DATA, "danych")
+30
View File
@@ -405,3 +405,33 @@ def test_max_budget_differs_between_models(monkeypatch):
haiku = prompt_token_budget(*factory.resolve_model("anthropic", "claude-haiku-4-5"))
local = prompt_token_budget(*factory.resolve_model("local", "llama3.1:8b"))
assert opus > haiku > local > 0
def test_turn_budget_counts_produced_not_requested(monkeypatch):
"""Regresja: odejmowanie ZAMOWIONEGO limitu tury zamiast wyprodukowanych
tokenow konczylo petle po jednej turze urwany fragment wracal jako calosc."""
seq = [_chat("Fragment 1. ", "length"), _chat("Fragment 2. ", "length"),
_chat("Zakonczenie. KONIEC", "stop")]
def handler(request):
return httpx.Response(200, json=seq.pop(0) if seq else seq[-1])
monkeypatch.setattr(httpx, "Client", _mock_client(handler))
# budzet 8000 < TURN_TOKENS_CAP: przy starej logice byla dokladnie jedna tura
out = ChatCompletionsProvider("local", "http://x/v1", "m").generate("p", 8000)
assert out.usage["turns"] == 3, "urwana odpowiedz musi byc kontynuowana"
assert "Fragment 1." in out.text and "Zakonczenie." in out.text
def test_generate_reports_progress_events(monkeypatch):
"""Log w UI ma pokazywac RZECZYWISTE tury, nie udawany pasek postepu."""
seq = [_chat("Czesc. ", "length"), _chat("Reszta. KONIEC", "stop")]
monkeypatch.setattr(httpx, "Client", _mock_client(
lambda r: httpx.Response(200, json=seq.pop(0) if seq else seq[-1])))
events = []
ChatCompletionsProvider("local", "http://x/v1", "m").generate(
"p", 8000, on_event=events.append)
kinds = [e["type"] for e in events]
assert kinds.count("turn_start") == 2 and kinds.count("turn_end") == 2
assert kinds[-1] == "generated"
assert any("urwana" in e["message"] for e in events if e["type"] == "turn_end")
@@ -10,7 +10,6 @@ from typing import Any
import httpx
from app import link_crypto
from app.config import settings
@@ -20,29 +19,16 @@ def _auth_headers() -> dict[str, str]:
return {"X-Astrololo-Token": token} if token else {}
def _link() -> link_crypto.Link | None:
"""Klucz łącza prezentacja↔logika. Czytany przy każdym wywołaniu, bo
konfiguracja może się zmienić bez restartu procesu (testy, podmiana sekretu)."""
key = link_crypto.key_from_env(link_crypto.ENV_PRESENTATION_LOGIC)
return link_crypto.Link(key) if key else None
class LogicClient:
def __init__(self, base_url: str | None = None) -> None:
self.base_url = (base_url or settings.logic_url).rstrip("/")
def _post(self, path: str, payload: dict[str, Any], timeout: float) -> dict[str, Any]:
"""Jedyna droga w dół. Celowo JEDNA: dopóki każda metoda składała żądanie
sama, dołożenie nowej znaczyło, że łatwo zapomnieć o tokenie albo kluczu
łącza i tak się już raz stało (401 wyszedł dopiero na produkcji)."""
with httpx.Client(timeout=timeout) as client:
return link_crypto.call_json(client, "POST", f"{self.base_url}{path}",
payload=payload, headers=_auth_headers(),
link=_link())
def query(self, query: str, field: str, exact: bool, limit: int) -> dict[str, Any]:
payload = {"query": query, "field": field, "exact": exact, "limit": limit}
return self._post("/api/query", payload, settings.http_timeout)
with httpx.Client(timeout=settings.http_timeout) as client:
r = client.post(f"{self.base_url}/api/query", json=payload, headers=_auth_headers())
r.raise_for_status()
return r.json()
def positions(
self,
@@ -65,15 +51,20 @@ class LogicClient:
"zodiac": zodiac,
}
# stacje wymagają root-findów — dłuższy timeout
timeout = max(settings.http_timeout, 60.0) if stations else settings.http_timeout
return self._post("/chart/positions", payload, timeout)
with httpx.Client(timeout=max(settings.http_timeout, 60.0) if stations else settings.http_timeout) as client:
r = client.post(f"{self.base_url}/chart/positions", json=payload, headers=_auth_headers())
r.raise_for_status()
return r.json()
def report(
self, when_utc_iso: str, lat: float, lon: float, limit: int = 5000, group: bool = False
) -> dict[str, Any]:
"""Sygnifikatory z obliczeń szukane w bazie — woła logic /chart/report."""
payload = {"when_utc": when_utc_iso, "lat": lat, "lon": lon, "limit": limit, "group": group}
return self._post("/chart/report", payload, max(settings.http_timeout, 30.0))
with httpx.Client(timeout=max(settings.http_timeout, 30.0)) as client:
r = client.post(f"{self.base_url}/chart/report", json=payload, headers=_auth_headers())
r.raise_for_status()
return r.json()
def prompt(
self, profile: str, when_utc_iso: str, lat: float, lon: float,
@@ -91,7 +82,10 @@ class LogicClient:
}
if from_date and to_date:
payload["from_date"], payload["to_date"] = from_date, to_date
return self._post("/chart/prompt", payload, max(settings.http_timeout, 60.0))
with httpx.Client(timeout=max(settings.http_timeout, 60.0)) as client:
r = client.post(f"{self.base_url}/chart/prompt", json=payload, headers=_auth_headers())
r.raise_for_status()
return r.json()
def horoscope(
self, profile: str, when_utc_iso: str, lat: float, lon: float,
@@ -112,13 +106,31 @@ class LogicClient:
payload["model"] = model
if from_date and to_date:
payload["from_date"], payload["to_date"] = from_date, to_date
return self._post("/chart/horoscope", payload, max(settings.http_timeout, 300.0))
with httpx.Client(timeout=max(settings.http_timeout, 300.0)) as client:
r = client.post(f"{self.base_url}/chart/horoscope", json=payload, headers=_auth_headers())
r.raise_for_status()
return r.json()
def horoscope_stream(self, payload: dict[str, Any]):
"""Strumień postępu pisania horoskopu (NDJSON) — przekazywany do przeglądarki.
Timeout jest długi, bo generowanie trwa; strumień i tak niesie heartbeat,
więc cisza na łączu nie zostanie wzięta za zerwanie.
"""
with httpx.Client(timeout=httpx.Timeout(None, connect=15.0)) as client:
with client.stream("POST", f"{self.base_url}/chart/horoscope/stream",
json=payload, headers=_auth_headers()) as r:
r.raise_for_status()
for chunk in r.iter_lines():
if chunk:
yield chunk
def llm_models(self) -> dict[str, Any]:
"""Katalog modeli per dostawca (podpowiedzi do pola wyboru w UI)."""
with httpx.Client(timeout=settings.http_timeout) as client:
return link_crypto.call_json(client, "GET", f"{self.base_url}/llm/models",
headers=_auth_headers(), link=_link())
r = client.get(f"{self.base_url}/llm/models", headers=_auth_headers())
r.raise_for_status()
return r.json()
def timeline(
self, when_utc_iso: str, lat: float, lon: float,
@@ -129,4 +141,7 @@ class LogicClient:
"when_utc": when_utc_iso, "lat": lat, "lon": lon,
"from_date": from_date, "to_date": to_date, "interpret": interpret,
}
return self._post("/chart/timeline", payload, max(settings.http_timeout, 60.0))
with httpx.Client(timeout=max(settings.http_timeout, 60.0)) as client:
r = client.post(f"{self.base_url}/chart/timeline", json=payload, headers=_auth_headers())
r.raise_for_status()
return r.json()
-466
View File
@@ -1,466 +0,0 @@
"""Szyfrowanie łączy między warstwami (PRE-16 / LOG-33).
Do tej pory warstwy rozmawiały ze sobą zwykłym HTTP-em wewnątrz klastra. Token
międzywarstwowy (LOG-32) mówił KTO pyta, ale nie ukrywał CZEGO dotyczy odpowiedź
a płyną nią surowe wiersze oryginalnych baz interpretacyjnych, czyli rdzeń
produktu. Kto podsłuchał ruch wewnątrz sieci (drugi pod, port mirror na switchu,
zrzut z węzła), miał je w całości.
Ten moduł zamyka drogę: **AES-256-GCM** na ciele każdego żądania i odpowiedzi.
GCM daje jednocześnie poufność i uwierzytelnienie cudzy albo podmieniony bajt
nie odszyfruje się w ogóle, więc nie ma osobnego problemu zaszyfrowane, ale
podatne na modyfikację".
**Dwa niezależne klucze**, po jednym na parę rozmówców:
* ``LINK_KEY_PRESENTATION_LOGIC`` prezentacja logika,
* ``LINK_KEY_LOGIC_DATA`` logika dane.
Dzięki temu przejęcie klucza prezentacji nie daje dostępu do warstwy danych,
gdzie leżą całe bazy. Logika trzyma oba, bo rozmawia w obie strony.
Z każdego klucza łącza wyprowadzamy **osobne podklucze na kierunek** (HKDF).
Żądanie i odpowiedź nigdy nie szyfrują się tym samym kluczem, więc powtórzenie
losowej jednorazówki w jedną stronę nie osłabia drugiej.
Format ramki (bo strumień odpowiedzi może iść kawałkami patrz okno postępu):
[4 bajty długości][magia "AL1"][12 bajtów jednorazówki][szyfrogram + znacznik]
Do materiału uwierzytelnianego (AAD) wchodzą kierunek, ścieżka, znacznik czasu
i numer ramki. Skutek: ramki nie da się przekleić do innego endpointu, odtworzyć
po czasie (dopuszczalny poślizg ``MAX_SKEW``) ani przestawić w strumieniu.
Bez ustawionego klucza moduł **przepuszcza ruch otwartym tekstem** (dev, zgodność
wstecz) i krzyczy o tym przy starcie. Gdy klucz JEST ustawiony, warstwa serwerowa
działa fail-closed: nieszyfrowane żądanie dostaje odmowę, żeby przypadkowa
regresja po stronie klienta nie oznaczała cichego powrotu do jawnego ruchu.
"""
from __future__ import annotations
import base64
import binascii
import logging
import os
import struct
import time
from typing import Iterable, Iterator
from cryptography.exceptions import InvalidTag
from cryptography.hazmat.primitives import hashes
from cryptography.hazmat.primitives.ciphers.aead import AESGCM
from cryptography.hazmat.primitives.kdf.hkdf import HKDF
log = logging.getLogger("astrololo.link")
MAGIC = b"AL1"
VERSION = "v1"
NONCE_BYTES = 12
KEY_BYTES = 32 # AES-256
LENGTH_PREFIX = 4
MAX_FRAME = 64 * 1024 * 1024 # zapora przed alokacją z podanej długości
MAX_SKEW_SECONDS = 300.0
HEADER_ENC = "X-Astrololo-Enc"
HEADER_TS = "X-Astrololo-Enc-Ts"
CONTENT_TYPE = "application/vnd.astrololo.enc"
ENV_PRESENTATION_LOGIC = "LINK_KEY_PRESENTATION_LOGIC"
ENV_LOGIC_DATA = "LINK_KEY_LOGIC_DATA"
ENV_REQUIRED = "LINK_ENCRYPTION_REQUIRED"
REQUEST, RESPONSE = b"req", b"res"
# Sondy k8s pukają tu bez klucza i tak ma zostać — inaczej pierwsza literówka
# w sekrecie kładłaby pody zamiast pokazać błąd w aplikacji.
PUBLIC_PATHS = frozenset({"/health"})
class LinkError(Exception):
"""Cokolwiek poszło nie tak z kopertą — celowo bez szczegółów na zewnątrz."""
# --------------------------------------------------------------------- klucze
def parse_key(raw: str) -> bytes:
"""Klucz z konfiguracji: hex (64 znaki) albo base64. Zawsze 32 bajty."""
text = raw.strip()
if not text:
raise LinkError("pusty klucz łącza")
try:
key = bytes.fromhex(text)
except ValueError:
try:
key = base64.b64decode(text, validate=True)
except (binascii.Error, ValueError) as exc:
raise LinkError("klucz łącza nie jest ani hexem, ani base64") from exc
if len(key) != KEY_BYTES:
raise LinkError(
f"klucz łącza ma {len(key)} B zamiast {KEY_BYTES} — wygeneruj przez "
f"`openssl rand -hex 32`"
)
return key
def key_from_env(env_name: str) -> bytes | None:
"""Klucz albo None. Zły klucz to wyjątek OD RAZU — nie przy pierwszym żądaniu."""
raw = os.getenv(env_name, "")
return parse_key(raw) if raw.strip() else None
def encryption_required() -> bool:
"""Czy brak klucza ma być błędem, a nie cichym powrotem do jawnego ruchu.
Serwer sam z siebie broni się fail-closed, ale to za mało: klient BEZ klucza
wysyła pytanie otwartym tekstem i dopiero potem dostaje odmowę czyli treść
zapytania zdążyła już przelecieć przez sieć. Ta flaga zatrzymuje go, zanim
cokolwiek opuści proces. Ustawiana razem z kluczami we wdrożeniu.
"""
return os.getenv(ENV_REQUIRED, "").strip().lower() in {"1", "true", "yes", "on"}
def _subkey(link_key: bytes, direction: bytes) -> bytes:
return HKDF(
algorithm=hashes.SHA256(), length=KEY_BYTES, salt=None,
info=b"astrololo/link/" + direction,
).derive(link_key)
class Link:
"""Jedna para rozmówców: klucz plus wyprowadzone z niego podklucze."""
def __init__(self, link_key: bytes) -> None:
self._by_direction = {
REQUEST: AESGCM(_subkey(link_key, REQUEST)),
RESPONSE: AESGCM(_subkey(link_key, RESPONSE)),
}
# ---------------------------------------------------------- pojedyncza ramka
def _aad(self, direction: bytes, path: str, stamp: str, seq: int) -> bytes:
return b"|".join([MAGIC, direction, path.encode("utf-8"),
stamp.encode("ascii"), str(seq).encode("ascii")])
def seal(self, direction: bytes, path: str, stamp: str, seq: int,
plaintext: bytes) -> bytes:
nonce = os.urandom(NONCE_BYTES)
sealed = self._by_direction[direction].encrypt(
nonce, plaintext, self._aad(direction, path, stamp, seq))
return MAGIC + nonce + sealed
def open(self, direction: bytes, path: str, stamp: str, seq: int,
frame: bytes) -> bytes:
if not frame.startswith(MAGIC):
raise LinkError("ramka bez znacznika astrololo")
body = frame[len(MAGIC):]
if len(body) <= NONCE_BYTES:
raise LinkError("ramka za krótka")
nonce, sealed = body[:NONCE_BYTES], body[NONCE_BYTES:]
try:
return self._by_direction[direction].decrypt(
nonce, sealed, self._aad(direction, path, stamp, seq))
except InvalidTag as exc:
# Jeden komunikat na wszystkie przypadki: zły klucz, podmieniony bajt,
# przeklejenie z innej ścieżki, przestawiona ramka. Rozróżnianie ich
# na zewnątrz podpowiadałoby atakującemu, w co trafił.
raise LinkError("nie udało się odszyfrować — zły klucz albo naruszone dane") from exc
# ------------------------------------------------------------ strumień ramek
def seal_stream(self, direction: bytes, path: str, stamp: str,
chunks: Iterable[bytes]) -> Iterator[bytes]:
for seq, chunk in enumerate(chunks):
yield frame_out(self.seal(direction, path, stamp, seq, chunk))
def open_stream(self, direction: bytes, path: str, stamp: str,
raw: bytes) -> Iterator[bytes]:
for seq, frame in enumerate(frames_in(raw)):
yield self.open(direction, path, stamp, seq, frame)
def open_all(self, direction: bytes, path: str, stamp: str, raw: bytes) -> bytes:
return b"".join(self.open_stream(direction, path, stamp, raw))
# ---------------------------------------------------------------- ramkowanie
def frame_out(payload: bytes) -> bytes:
return struct.pack(">I", len(payload)) + payload
def frames_in(raw: bytes) -> Iterator[bytes]:
"""Rozbiera bufor na ramki. Ucięty strumień to błąd, nie cicha strata danych."""
offset = 0
while offset < len(raw):
if offset + LENGTH_PREFIX > len(raw):
raise LinkError("urwana ramka (brak nagłówka długości)")
(size,) = struct.unpack(">I", raw[offset:offset + LENGTH_PREFIX])
if size > MAX_FRAME:
raise LinkError("ramka ponad dopuszczalny rozmiar")
offset += LENGTH_PREFIX
if offset + size > len(raw):
raise LinkError("urwana ramka (za mało danych)")
yield raw[offset:offset + size]
offset += size
def unframe_incremental(buffer: bytearray) -> Iterator[bytes]:
"""Wyjmuje z bufora KOMPLETNE ramki i zjada je; resztę zostawia na później.
Dla odbioru na żywo: kawałki przychodzą podzielone dowolnie i ramka potrafi
rozjechać się między dwa odczyty.
"""
while True:
if len(buffer) < LENGTH_PREFIX:
return
(size,) = struct.unpack(">I", buffer[:LENGTH_PREFIX])
if size > MAX_FRAME:
raise LinkError("ramka ponad dopuszczalny rozmiar")
if len(buffer) < LENGTH_PREFIX + size:
return
frame = bytes(buffer[LENGTH_PREFIX:LENGTH_PREFIX + size])
del buffer[:LENGTH_PREFIX + size]
yield frame
# ------------------------------------------------------------- świeżość ruchu
def stamp_now() -> str:
return f"{time.time():.3f}"
def check_stamp(stamp: str) -> None:
"""Odrzuca ramki spoza okna czasowego — inaczej podsłuchane żądanie dałoby się
odtworzyć w dowolnym momencie w przyszłości."""
try:
sent = float(stamp)
except (TypeError, ValueError) as exc:
raise LinkError("brak albo błędny znacznik czasu") from exc
if abs(time.time() - sent) > MAX_SKEW_SECONDS:
raise LinkError("znacznik czasu poza dopuszczalnym oknem")
# =========================================================== strona serwerowa
class LinkCryptoMiddleware:
"""Rozszyfrowuje wchodzące żądania i zaszyfrowuje wychodzące odpowiedzi.
Napisane jako czyste ASGI, nie ``@app.middleware("http")``, bo trzeba
podmienić CIAŁO żądania jeszcze zanim zobaczy je FastAPI, oraz przepuścić
odpowiedź strumieniową kawałek po kawałku, bez zbierania jej w pamięci.
"""
def __init__(self, app, link: Link | None, layer: str) -> None:
self.app = app
self.link = link
self.layer = layer
async def __call__(self, scope, receive, send):
if scope["type"] != "http" or self.link is None or scope["path"] in PUBLIC_PATHS:
return await self.app(scope, receive, send)
path = scope["path"]
headers = {k.decode("latin-1").lower(): v.decode("latin-1") for k, v in scope["headers"]}
if headers.get(HEADER_ENC.lower()) != VERSION:
# Fail-closed. Klucz jest ustawiony, więc jawne żądanie oznacza albo
# pomyłkę w konfiguracji, albo kogoś obcego — w obu wypadkach nie
# chcemy po cichu wrócić do jawnego ruchu.
log.warning("warstwa %s: odrzucone żądanie bez szyfrowania łącza (%s)",
self.layer, path)
return await _refuse(send, "Łącze międzywarstwowe wymaga szyfrowania.")
stamp = headers.get(HEADER_TS.lower(), "")
try:
check_stamp(stamp)
plaintext = self.link.open_all(REQUEST, path, stamp, await _read_body(receive))
except LinkError as exc:
log.warning("warstwa %s: %s (%s)", self.layer, exc, path)
return await _refuse(send, "Nie udało się odczytać zaszyfrowanego żądania.")
scope = dict(scope)
scope["headers"] = _rewritten_headers(scope["headers"], len(plaintext))
await self.app(scope, _replay(plaintext, receive), self._sealing_send(send, path))
def _sealing_send(self, send, path: str):
state: dict = {"stamp": "", "seq": 0}
async def sealing(message):
if message["type"] == "http.response.start":
state["stamp"] = stamp_now()
keep = [(k, v) for k, v in message.get("headers", [])
if k.lower() not in (b"content-length", b"content-type")]
message = dict(message)
message["headers"] = keep + [
(b"content-type", CONTENT_TYPE.encode()),
(HEADER_ENC.lower().encode(), VERSION.encode()),
(HEADER_TS.lower().encode(), state["stamp"].encode()),
]
return await send(message)
if message["type"] == "http.response.body":
chunk = message.get("body", b"")
sealed = b""
if chunk:
sealed = frame_out(self.link.seal(
RESPONSE, path, state["stamp"], state["seq"], chunk))
state["seq"] += 1
return await send({"type": "http.response.body", "body": sealed,
"more_body": message.get("more_body", False)})
return await send(message)
return sealing
def _rewritten_headers(raw: Iterable[tuple[bytes, bytes]], length: int):
"""Po odszyfrowaniu ciało ma inną długość i zwykły typ — inaczej FastAPI
próbowałby sparsować JSON o cudzej deklarowanej wielkości."""
kept = [(k, v) for k, v in raw if k.lower() not in (b"content-length", b"content-type")]
kept.append((b"content-length", str(length).encode()))
if length:
kept.append((b"content-type", b"application/json"))
return kept
async def _read_body(receive) -> bytes:
body = bytearray()
while True:
message = await receive()
if message["type"] == "http.disconnect":
raise LinkError("rozłączenie w trakcie odbioru żądania")
body += message.get("body", b"")
if not message.get("more_body", False):
return bytes(body)
def _replay(body: bytes, original):
"""Podstawia odszyfrowane ciało jako jedyną porcję wejścia dla aplikacji.
Po oddaniu ciała oddajemy głos ORYGINALNEMU `receive`, zamiast od razu
zgłaszać rozłączenie. Odpowiedź strumieniowa nasłuchuje bowiem rozłączenia
równolegle do wysyłania i przerywa się, gdy je zobaczy na skróconej wersji
okno postępu dostawało pustą odpowiedź, choć zwykłe żądania działały.
"""
delivered = False
async def receive():
nonlocal delivered
if delivered:
return await original()
delivered = True
return {"type": "http.request", "body": body, "more_body": False}
return receive
async def _refuse(send, detail: str) -> None:
"""Odmowa leci JAWNIE — rozmówca właśnie pokazał, że nie umie odszyfrować,
więc zaszyfrowany komunikat o błędzie byłby dla niego nieczytelny."""
payload = f'{{"detail":"{detail}"}}'.encode("utf-8")
await send({"type": "http.response.start", "status": 400, "headers": [
(b"content-type", b"application/json"),
(b"content-length", str(len(payload)).encode()),
]})
await send({"type": "http.response.body", "body": payload})
def install(app, env_name: str, layer: str):
"""Podpina szyfrowanie łącza. Wołać PO `security.install`, żeby także odmowa
tokenowa (401) wracała zaszyfrowana inaczej klient by jej nie odczytał."""
link_key = key_from_env(env_name)
if link_key is None and encryption_required():
# Celowo wywracamy start. Ta sama zasada co przy sekrecie logowania:
# wolimy widoczną awarię niż usługę, która wstała i po cichu nie chroni
# niczego. Pod w CrashLoop widać od razu, jawny ruch — nie.
raise LinkError(
f"{ENV_REQUIRED} jest włączone, ale {env_name} nie ustawiony — "
f"warstwa {layer} nie wystartuje bez klucza łącza"
)
if link_key is None:
log.warning(
"UWAGA: %s nie ustawiony — warstwa %s rozmawia z sąsiadem JAWNYM tekstem, "
"więc treść baz interpretacyjnych jest widoczna dla każdego, kto podsłucha "
"ruch wewnątrz sieci.", env_name, layer,
)
return None
link = Link(link_key)
app.add_middleware(LinkCryptoMiddleware, link=link, layer=layer)
log.info("warstwa %s: łącze szyfrowane (AES-256-GCM, klucz z %s)", layer, env_name)
return link
# ============================================================ strona kliencka
def call(client, method: str, url: str, *, payload=None,
headers: dict[str, str] | None = None, link: Link | None) -> bytes:
"""Żądanie do sąsiedniej warstwy; zwraca odszyfrowane ciało odpowiedzi.
Ścieżkę do materiału uwierzytelnianego bierzemy Z URL-a, a nie z osobnego
argumentu gdyby klient i serwer liczyły inaczej, każde żądanie kończyłoby
się niejasnym błędem odszyfrowania.
"""
import json as _json
import httpx
request_headers = dict(headers or {})
if link is None:
if encryption_required():
# Zatrzymujemy się PRZED wysłaniem. Gdyby polecieć jawnie i dopiero
# zebrać odmowę, pytanie byłoby już na kablu — a to właśnie ono niesie
# sygnifikatory, o które pytamy bazę.
raise LinkError(
f"{ENV_REQUIRED} jest włączone, ale brak klucza łącza — żądanie "
f"NIE zostało wysłane, żeby jego treść nie poszła jawnym tekstem"
)
response = client.request(method, url, json=payload, headers=request_headers)
response.raise_for_status()
return response.content
path = httpx.URL(url).path
stamp = stamp_now()
plaintext = b"" if payload is None else _json.dumps(payload).encode("utf-8")
body = frame_out(link.seal(REQUEST, path, stamp, 0, plaintext))
request_headers.update({HEADER_ENC: VERSION, HEADER_TS: stamp,
"Content-Type": CONTENT_TYPE})
response = client.request(method, url, content=body, headers=request_headers)
if response.status_code >= 400 and response.headers.get(HEADER_ENC) != VERSION:
log.error("łącze %s odmówiło: %s", path, response.text[:200])
response.raise_for_status()
if response.headers.get(HEADER_ENC) != VERSION:
raise LinkError("odpowiedź przyszła nieszyfrowana, choć klucz łącza jest ustawiony")
reply_stamp = response.headers.get(HEADER_TS, "")
check_stamp(reply_stamp)
return link.open_all(RESPONSE, path, reply_stamp, response.content)
def call_json(client, method: str, url: str, *, payload=None,
headers: dict[str, str] | None = None, link: Link | None):
import json as _json
return _json.loads(call(client, method, url, payload=payload,
headers=headers, link=link))
def open_response_stream(response, link: Link | None) -> Iterator[bytes]:
"""Odbiór odpowiedzi płynącej kawałkami (okno postępu).
Ramka potrafi rozjechać się między dwa odczyty z gniazda, więc składamy
w buforze zamiast zakładać, że każdy kawałek to komplet.
"""
if link is None:
yield from response.iter_bytes()
return
if response.headers.get(HEADER_ENC) != VERSION:
raise LinkError("strumień przyszedł nieszyfrowany, choć klucz łącza jest ustawiony")
stamp = response.headers.get(HEADER_TS, "")
check_stamp(stamp)
path = response.request.url.path
buffer = bytearray()
seq = 0
for chunk in response.iter_bytes():
buffer += chunk
for frame in unframe_incremental(buffer):
yield link.open(RESPONSE, path, stamp, seq, frame)
seq += 1
if buffer:
raise LinkError("strumień urwał się w połowie ramki")
+57 -1
View File
@@ -9,11 +9,12 @@ Strona główna „/" = wprowadzenie danych horoskopu i podgląd policzonych poz
"""
from __future__ import annotations
import json
from datetime import datetime, timedelta, timezone
import httpx
from fastapi import FastAPI, Form, HTTPException, Query, Request
from fastapi.responses import HTMLResponse
from fastapi.responses import HTMLResponse, JSONResponse
from fastapi.staticfiles import StaticFiles
from fastapi.templating import Jinja2Templates
@@ -228,6 +229,61 @@ def timeline_run(
return templates.TemplateResponse(request, "timeline.html", ctx)
# ---------------- Postęp pisania horoskopu (strumień do okna z logiem) ----------------
@app.post("/horoscope/stream")
def horoscope_stream(
profile: str = Form("natal"),
date: str = Form(...),
time: str = Form(...),
tz_offset: float = Form(0.0),
lat: float = Form(0.0),
lon: float = Form(0.0),
prompt_budget: str = Form("medium"),
llm_provider: str = Form("local"),
llm_model: str = Form(""),
from_date: str = Form(""),
to_date: str = Form(""),
):
"""Przekazuje strumień postępu z logiki i DOKLEJA gotowy HTML wyniku.
Dzięki temu okno postępu wstawia dokładnie ten sam widok, który wyrenderowałoby
przeładowanie strony jedno źródło prawdy dla wyglądu wyniku.
"""
from fastapi.responses import StreamingResponse
try:
iso_utc, _ = _build_utc(date, time, tz_offset)
except ValueError as e:
return JSONResponse({"detail": f"Niepoprawne dane wejściowe: {e}"}, status_code=422)
payload: dict = {
"profile": profile, "when_utc": iso_utc, "lat": lat, "lon": lon,
"budget": prompt_budget, "provider": llm_provider, "model": llm_model,
}
if profile == "period" and from_date and to_date:
payload["from_date"], payload["to_date"] = from_date, to_date
def relay():
try:
for raw in logic.horoscope_stream(payload):
try:
event = json.loads(raw)
except ValueError:
continue
if event.get("type") == "result":
html = templates.get_template("_prompt_result.html").render(
prompt_result=event.get("result") or {}
)
event["html"] = html
yield json.dumps(event, ensure_ascii=False) + "\n"
except httpx.HTTPError as e:
yield json.dumps({"type": "error", "message": _logic_error(e)},
ensure_ascii=False) + "\n"
return StreamingResponse(relay(), media_type="application/x-ndjson",
headers={"Cache-Control": "no-store", "X-Accel-Buffering": "no"})
# ---------------- Geokoder (proxy OSM/Nominatim dla wyszukiwarki lokalizacji) ----------------
@app.get("/geocode")
def geocode_search(q: str = Query("", description="Nazwa / adres / POI do wyszukania")):
+3 -34
View File
@@ -7,9 +7,7 @@ jakiegokolwiek modelu językowego. Ten moduł zamyka tę drogę.
Dwa mechanizmy:
* **HTTP Basic** wejście do aplikacji; włącza się, gdy ustawiono APP_PASSWORD.
* **limit żądań** hamuje masowe odpytywanie (eksfiltrację przez pętlę zapytań);
rozliczany per adres klienta, a za odwrotnym proxy po TRUST_PROXY=true
per adres z nagłówka, nie per adres proxy (patrz `client_ip`).
* **limit żądań** hamuje masowe odpytywanie (eksfiltrację przez pętlę zapytań).
Świadomie NIE logujemy treści żądań ani promptów logi to kolejny nośnik wycieku.
@@ -48,11 +46,6 @@ def app_password() -> str:
def rate_limit_per_min() -> int:
return int(os.getenv("RATE_LIMIT_PER_MIN", "120"))
def trust_proxy() -> bool:
return os.getenv("TRUST_PROXY", "").strip().lower() in {"1", "true", "yes", "on"}
PUBLIC_PATHS = frozenset({"/health"})
PUBLIC_PREFIXES = ("/static/",)
@@ -81,31 +74,6 @@ def _authorized(header: str | None) -> bool:
return ok_user and ok_pass
def client_ip(request: Request) -> str:
"""Adres, po którym rozliczamy limit żądań.
Za odwrotnym proxy (u nas: Ingress/Traefik po włączeniu TLS PRE-16)
`request.client.host` to adres POD-a proxy, jednakowy dla wszystkich. Bez
poprawki cały ruch trafiałby do jednego wiadra i pierwsza osoba, która
wyklika limit, odcięłaby pozostałe.
Nagłówkom wierzymy WYŁĄCZNIE przy TRUST_PROXY bo inaczej wystarczyłoby
dopisać własny `X-Forwarded-For`, żeby przy każdym żądaniu wyglądać na kogoś
innego i ominąć limit całkowicie. Z tego samego powodu bierzemy OSTATNI wpis
listy: to jedyny, który dopisał nasz proxy. Wcześniejsze mógł podstawić
klient, więc nie znaczą nic.
"""
peer = request.client.host if request.client else "?"
if not trust_proxy():
return peer
forwarded = request.headers.get("x-forwarded-for", "")
if forwarded:
last = forwarded.rsplit(",", 1)[-1].strip()
if last:
return last
return request.headers.get("x-real-ip", "").strip() or peer
def _rate_limited(client: str) -> bool:
cap = rate_limit_per_min()
if cap <= 0:
@@ -137,7 +105,8 @@ def install(app) -> None:
if _is_public(request.url.path):
return await call_next(request)
if _rate_limited(client_ip(request)):
client = request.client.host if request.client else "?"
if _rate_limited(client):
return JSONResponse(
{"detail": "Zbyt wiele żądań — spróbuj za chwilę."},
status_code=429, headers={"Retry-After": "60"},
@@ -0,0 +1,133 @@
// Okno postępu przy pisaniu horoskopu.
//
// Problem: generowanie trwa minutami, a zwykły POST formularza nie daje żadnego
// sygnału — aplikacja wygląda na zawieszoną. Zamiast udawanego paska postępu
// czytamy strumień RZECZYWISTYCH zdarzeń z serwera (NDJSON) i wypisujemy je
// jako log: budowa promptu, limity modelu, każda tura generowania.
//
// Degradacja: jeśli przeglądarka nie umie strumieniować `fetch`, nie przechwytujemy
// wysyłki — formularz idzie klasycznie i wszystko działa jak wcześniej, tylko bez okna.
document.addEventListener('DOMContentLoaded', function () {
const canStream = typeof fetch === 'function' && typeof ReadableStream === 'function' &&
typeof TextDecoder === 'function';
const form = document.querySelector('form[action="/interpret"], form[action="/timeline"]');
if (!canStream || !form) return;
const profile = form.getAttribute('action') === '/timeline' ? 'period' : 'natal';
const btn = form.querySelector('button[value="horoscope"]');
if (!btn) return;
// --- okno ---------------------------------------------------------------
const overlay = document.createElement('div');
overlay.className = 'progress-overlay';
overlay.hidden = true;
overlay.innerHTML =
'<div class="progress-box" role="dialog" aria-modal="true" aria-label="Postęp generowania">' +
'<div class="progress-head">' +
'<span class="progress-spinner" aria-hidden="true"></span>' +
'<strong id="progressTitle">Piszę horoskop…</strong>' +
'<span class="progress-clock" id="progressClock">0:00</span>' +
'</div>' +
'<ol class="progress-log" id="progressLog"></ol>' +
'<p class="muted small">Nie zamykaj tej karty — generowanie trwa na serwerze.</p>' +
'<div class="actions"><button type="button" class="ghost" id="progressClose" hidden>Zamknij</button></div>' +
'</div>';
document.body.appendChild(overlay);
const logEl = overlay.querySelector('#progressLog');
const clockEl = overlay.querySelector('#progressClock');
const titleEl = overlay.querySelector('#progressTitle');
const closeEl = overlay.querySelector('#progressClose');
let timer = null;
function addLine(text, kind) {
const li = document.createElement('li');
if (kind) li.className = 'log-' + kind;
const now = new Date();
li.textContent = String(now.getHours()).padStart(2, '0') + ':' +
String(now.getMinutes()).padStart(2, '0') + ':' +
String(now.getSeconds()).padStart(2, '0') + ' ' + text;
logEl.appendChild(li);
logEl.scrollTop = logEl.scrollHeight;
}
function startClock() {
const t0 = Date.now();
timer = setInterval(function () {
const s = Math.floor((Date.now() - t0) / 1000);
clockEl.textContent = Math.floor(s / 60) + ':' + String(s % 60).padStart(2, '0');
}, 1000);
}
function finish(title, allowClose) {
if (timer) { clearInterval(timer); timer = null; }
titleEl.textContent = title;
overlay.querySelector('.progress-spinner').style.visibility = 'hidden';
if (allowClose) closeEl.hidden = false;
}
closeEl.addEventListener('click', function () { overlay.hidden = true; });
// --- przechwycenie wysyłki ----------------------------------------------
btn.addEventListener('click', function (event) {
event.preventDefault();
if (!form.reportValidity()) return; // te same reguły co przy zwykłej wysyłce
const data = new FormData(form);
data.set('profile', profile);
logEl.innerHTML = '';
closeEl.hidden = true;
overlay.querySelector('.progress-spinner').style.visibility = '';
titleEl.textContent = 'Piszę horoskop…';
overlay.hidden = false;
startClock();
addLine('Wysyłam żądanie…');
fetch('/horoscope/stream', { method: 'POST', body: data })
.then(function (response) {
if (!response.ok || !response.body) throw new Error('HTTP ' + response.status);
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = '';
function pump() {
return reader.read().then(function (chunk) {
if (chunk.done) return;
buffer += decoder.decode(chunk.value, { stream: true });
const lines = buffer.split('\n');
buffer = lines.pop(); // ostatni może być niepełny
lines.forEach(function (raw) {
if (!raw.trim()) return;
let ev;
try { ev = JSON.parse(raw); } catch (e) { return; }
if (ev.type === 'ping') return; // sam heartbeat, nie logujemy
if (ev.type === 'result') {
addLine(ev.message || 'Gotowe.', 'ok');
if (ev.html) {
const host = document.getElementById('promptResult');
if (host) host.innerHTML = ev.html;
}
finish('Gotowe', true);
overlay.hidden = true;
const host = document.getElementById('promptResult');
if (host) host.scrollIntoView({ behavior: 'smooth', block: 'start' });
return;
}
addLine(ev.message || ev.type, ev.type === 'error' ? 'err'
: ev.type === 'warn' ? 'warn' : null);
if (ev.type === 'error') finish('Nie udało się', true);
});
return pump();
});
}
return pump();
})
.catch(function (e) {
addLine('Połączenie przerwane: ' + e.message, 'err');
addLine('Możesz spróbować ponownie — nic nie zostało utracone.');
finish('Nie udało się', true);
});
});
});
@@ -88,3 +88,27 @@ textarea.prompt { width: 100%; margin-top: .5rem; padding: .7rem .8rem; box-sizi
border-radius: 10px; font-family: ui-monospace, SFMono-Regular, Menlo, monospace;
font-size: .82rem; line-height: 1.45; resize: vertical; white-space: pre; }
textarea.prompt:focus { outline: 2px solid var(--accent); outline-offset: 1px; }
/* Okno postępu przy generowaniu horoskopu */
.progress-overlay { position: fixed; inset: 0; background: rgba(8,9,20,.72);
display: flex; align-items: center; justify-content: center;
z-index: 50; padding: 1rem; }
.progress-overlay[hidden] { display: none; }
.progress-box { background: var(--panel); border: 1px solid var(--line); border-radius: 14px;
padding: 1rem 1.1rem; width: min(38rem, 100%); box-shadow: 0 18px 48px rgba(0,0,0,.45); }
.progress-head { display: flex; align-items: center; gap: .6rem; margin-bottom: .6rem; }
.progress-clock { margin-left: auto; font-family: ui-monospace, Menlo, monospace;
color: var(--muted); font-size: .9rem; }
.progress-spinner { width: 14px; height: 14px; border-radius: 50%; flex: none;
border: 2px solid var(--line); border-top-color: var(--accent);
animation: spin .8s linear infinite; }
@keyframes spin { to { transform: rotate(360deg); } }
@media (prefers-reduced-motion: reduce) { .progress-spinner { animation: none; } }
.progress-log { list-style: none; margin: 0 0 .6rem; padding: .6rem .7rem;
background: #12132a; border: 1px solid var(--line); border-radius: 10px;
max-height: 15rem; overflow-y: auto;
font-family: ui-monospace, Menlo, monospace; font-size: .78rem; line-height: 1.6; }
.progress-log li { color: var(--ink); white-space: pre-wrap; }
.progress-log li.log-ok { color: #7fd18b; }
.progress-log li.log-warn { color: #e0c060; }
.progress-log li.log-err { color: #ef6b6b; }
@@ -36,80 +36,4 @@
sieć</strong>. Możesz najpierw obejrzeć prompt, a dopiero potem wysłać.
</p>
{% if prompt_result %}
{% set st = prompt_result.stats %}
{# --- wynik: napisany horoskop (LOG-31) + transparentność (PRE-15) --- #}
{% if prompt_result.horoscope %}
<div class="meta">
Horoskop napisany przez: <strong>{{ prompt_result.provider }}</strong> ·
model: {{ prompt_result.model }}
{% if prompt_result.usage and prompt_result.usage.completion_tokens %}
· tokeny odpowiedzi: {{ prompt_result.usage.completion_tokens }}
{% endif %}
{% if prompt_result.leaves_lan %}
· <strong class="retro">dane opuściły sieć</strong>
{% else %}
· dane nie opuściły sieci
{% endif %}
</div>
<div class="actions">
<button type="button" class="ghost" data-copy="#horoscopeText">Kopiuj horoskop</button>
</div>
<textarea id="horoscopeText" class="prompt" rows="20" readonly>{{ prompt_result.horoscope }}</textarea>
<p class="muted small">
Treść wygenerował model językowy na podstawie {{ st.included }} wskazań z baz.
<strong>Nie stanowi porady medycznej, prawnej ani finansowej.</strong>
</p>
{% endif %}
{% if prompt_result.warnings %}
{% for w in prompt_result.warnings %}
<p class="muted small"><strong>Uwaga:</strong> {{ w }}</p>
{% endfor %}
{% endif %}
{% if prompt_result.token_plan %}
{% set tp = prompt_result.token_plan %}
<p class="muted small">
Tokeny: prompt {{ tp.prompt_tokens }} · okno modelu {{ tp.context_window }} ·
zarezerwowane na odpowiedź {{ tp.max_output }}
{% if prompt_result.usage and prompt_result.usage.turns and prompt_result.usage.turns > 1 %}
· odpowiedź złożona z {{ prompt_result.usage.turns }} tur (model dokańczał urwany tekst)
{% endif %}
</p>
{% endif %}
{% if prompt_result.llm_error %}
<div class="error">
Nie udało się napisać horoskopu: {{ prompt_result.llm_error }}<br>
Prompt poniżej jest gotowy — możesz go skopiować i użyć ręcznie.
</div>
{% endif %}
{# --- prompt: zawsze dostępny do podglądu i skopiowania --- #}
<div class="meta">
Prompt · {{ st.chars }} znaków (~{{ st.est_tokens }} tokenów) ·
budżet: {{ st.budget }} ·
wskazań: <strong>{{ st.included }}</strong>
{% if st.omitted %}· pominięto: <strong>{{ st.omitted }}</strong>{% endif %}
{% if st.deduplicated %}· scalono powtórek: {{ st.deduplicated }}{% endif %}
</div>
{% if st.omitted %}
<p class="muted small">
Pominięto {{ st.omitted }} najsłabszych wskazań (próg wagi {{ st.min_score_included }}).
Chcesz komplet — wybierz obszerniejszy budżet i wygeneruj ponownie.
</p>
{% endif %}
{% if prompt_result.data_error %}
<div class="error">{{ prompt_result.data_error }} — prompt złożony z samych wyliczeń.</div>
{% endif %}
<div class="actions">
<button type="button" class="ghost" data-copy="#promptText">Kopiuj prompt</button>
<span class="muted small">Możesz też wkleić go samodzielnie do ChatGPT lub Claude.</span>
</div>
<textarea id="promptText" class="prompt" rows="14" readonly>{{ prompt_result.prompt }}</textarea>
{% endif %}
<div id="promptResult">{% include "_prompt_result.html" %}</div>
@@ -0,0 +1,80 @@
{# Blok WYNIKU generowania (prompt/horoskop). Wydzielony, bo wstawia go także
okno postępu po zakończeniu strumienia — dzięki temu jest jedno źródło
prawdy dla wyglądu wyniku, niezależnie od drogi, którą przyszedł. #}
{% if prompt_result %}
{% set st = prompt_result.stats %}
{# --- wynik: napisany horoskop (LOG-31) + transparentność (PRE-15) --- #}
{% if prompt_result.horoscope %}
<div class="meta">
Horoskop napisany przez: <strong>{{ prompt_result.provider }}</strong> ·
model: {{ prompt_result.model }}
{% if prompt_result.usage and prompt_result.usage.completion_tokens %}
· tokeny odpowiedzi: {{ prompt_result.usage.completion_tokens }}
{% endif %}
{% if prompt_result.leaves_lan %}
· <strong class="retro">dane opuściły sieć</strong>
{% else %}
· dane nie opuściły sieci
{% endif %}
</div>
<div class="actions">
<button type="button" class="ghost" data-copy="#horoscopeText">Kopiuj horoskop</button>
</div>
<textarea id="horoscopeText" class="prompt" rows="20" readonly>{{ prompt_result.horoscope }}</textarea>
<p class="muted small">
Treść wygenerował model językowy na podstawie {{ st.included }} wskazań z baz.
<strong>Nie stanowi porady medycznej, prawnej ani finansowej.</strong>
</p>
{% endif %}
{% if prompt_result.warnings %}
{% for w in prompt_result.warnings %}
<p class="muted small"><strong>Uwaga:</strong> {{ w }}</p>
{% endfor %}
{% endif %}
{% if prompt_result.token_plan %}
{% set tp = prompt_result.token_plan %}
<p class="muted small">
Tokeny: prompt {{ tp.prompt_tokens }} · okno modelu {{ tp.context_window }} ·
zarezerwowane na odpowiedź {{ tp.max_output }}
{% if prompt_result.usage and prompt_result.usage.turns and prompt_result.usage.turns > 1 %}
· odpowiedź złożona z {{ prompt_result.usage.turns }} tur (model dokańczał urwany tekst)
{% endif %}
</p>
{% endif %}
{% if prompt_result.llm_error %}
<div class="error">
Nie udało się napisać horoskopu: {{ prompt_result.llm_error }}<br>
Prompt poniżej jest gotowy — możesz go skopiować i użyć ręcznie.
</div>
{% endif %}
{# --- prompt: zawsze dostępny do podglądu i skopiowania --- #}
<div class="meta">
Prompt · {{ st.chars }} znaków (~{{ st.est_tokens }} tokenów) ·
budżet: {{ st.budget }} ·
wskazań: <strong>{{ st.included }}</strong>
{% if st.omitted %}· pominięto: <strong>{{ st.omitted }}</strong>{% endif %}
{% if st.deduplicated %}· scalono powtórek: {{ st.deduplicated }}{% endif %}
</div>
{% if st.omitted %}
<p class="muted small">
Pominięto {{ st.omitted }} najsłabszych wskazań (próg wagi {{ st.min_score_included }}).
Chcesz komplet — wybierz obszerniejszy budżet i wygeneruj ponownie.
</p>
{% endif %}
{% if prompt_result.data_error %}
<div class="error">{{ prompt_result.data_error }} — prompt złożony z samych wyliczeń.</div>
{% endif %}
<div class="actions">
<button type="button" class="ghost" data-copy="#promptText">Kopiuj prompt</button>
<span class="muted small">Możesz też wkleić go samodzielnie do ChatGPT lub Claude.</span>
</div>
<textarea id="promptText" class="prompt" rows="14" readonly>{{ prompt_result.prompt }}</textarea>
{% endif %}
@@ -91,4 +91,5 @@
<script src="/static/now.js"></script>
<script src="/static/copy.js"></script>
<script src="/static/models.js"></script>
<script src="/static/progress.js"></script>
{% endblock %}
@@ -81,4 +81,5 @@
<script src="/static/now.js"></script>
<script src="/static/copy.js"></script>
<script src="/static/models.js"></script>
<script src="/static/progress.js"></script>
{% endblock %}
-2
View File
@@ -3,5 +3,3 @@ uvicorn[standard]>=0.34
httpx>=0.28
jinja2>=3.1
python-multipart>=0.0.20
# Szyfrowanie łącza między warstwami (PRE-16): AES-256-GCM + HKDF
cryptography>=44.0
+10 -54
View File
@@ -1,44 +1,30 @@
"""Niezmiennik: KAŻDE wyjście HTTP w dół niesie token międzywarstwowy (LOG-32)
oraz klucz szyfrujący łącze (PRE-16).
"""Niezmiennik: KAŻDE wyjście HTTP w dół niesie token międzywarstwowy (LOG-32).
Powód istnienia tego testu: token dodano do klienta na gałęzi, która odbiła się od
mastera zanim powstały metody `prompt()` i `horoscope()`. Git zmergował obie zmiany
czysto (różne linie), ale nowe metody wyszły BEZ tokenu i dostawały 401 dopiero na
produkcji. Zwykły test jednej metody by tego nie złapał, więc sprawdzamy regułę
strukturalnie.
Po dołożeniu szyfrowania ta sama klasa błędu ma gorszy objaw: wywołanie bez `link=`
nie wywala się widocznie, tylko po cichu wysyła treść JAWNYM tekstem. Dlatego dla
wywołań przez `link_crypto` wymagamy obu argumentów naraz.
strukturalnie: nie ma wywołania bez `headers=`.
"""
import ast
import pathlib
CLIENT = pathlib.Path(__file__).resolve().parents[1] / "app" / "clients" / "logic_client.py"
HTTP_VERBS = ("post", "get", "put", "patch", "delete")
LINK_CALLS = ("call", "call_json")
def _http_calls(path: pathlib.Path) -> list[tuple[str, int, bool, bool]]:
"""(opis, linia, czy_ma_headers, czy_wymaga_i_ma_link) dla każdego wyjścia w dół."""
def _http_calls(path: pathlib.Path) -> list[tuple[str, int, bool]]:
"""(nazwa_metody_http, linia, czy_ma_headers) dla każdego client.post/get."""
tree = ast.parse(path.read_text(encoding="utf-8"))
out = []
for node in ast.walk(tree):
if not isinstance(node, ast.Call) or not isinstance(node.func, ast.Attribute):
continue
target = node.func.value
if node.func.attr not in ("post", "get", "put", "patch", "delete"):
continue
if not (isinstance(node.func.value, ast.Name) and node.func.value.id == "client"):
continue
has_headers = any(kw.arg == "headers" for kw in node.keywords)
has_link = any(kw.arg == "link" for kw in node.keywords)
if (node.func.attr in HTTP_VERBS
and isinstance(target, ast.Name) and target.id == "client"):
# Surowe wywołanie httpx — omija szyfrowanie, więc dopuszczalne tylko
# dla ścieżek wyjętych spod niego (patrz link_crypto.PUBLIC_PATHS).
out.append((f"client.{node.func.attr}()", node.lineno, has_headers, True))
elif (node.func.attr in LINK_CALLS
and isinstance(target, ast.Name) and target.id == "link_crypto"):
out.append((f"link_crypto.{node.func.attr}()", node.lineno, has_headers, has_link))
out.append((node.func.attr, node.lineno, has_headers))
return out
@@ -49,43 +35,13 @@ def test_client_module_exists():
def test_every_outbound_call_sends_auth_header():
calls = _http_calls(CLIENT)
assert calls, "nie znaleziono żadnego wywołania HTTP — test przestał cokolwiek pilnować"
missing = [f"{CLIENT.name}:{line} {what}" for what, line, ok, _ in calls if not ok]
missing = [f"{CLIENT.name}:{line} client.{verb}()" for verb, line, ok in calls if not ok]
assert not missing, (
"Wywołania w dół bez tokenu międzywarstwowego (dostaną 401 przy włączonej "
"ochronie): " + ", ".join(missing)
)
def test_every_outbound_call_passes_link_key():
"""Brak `link=` nie boli od razu — po prostu treść leci jawnym tekstem."""
calls = _http_calls(CLIENT)
missing = [f"{CLIENT.name}:{line} {what}" for what, line, _, ok in calls if not ok]
assert not missing, (
"Wywołania w dół bez klucza łącza — poszłyby NIEZASZYFROWANE: " + ", ".join(missing)
)
def test_raw_http_calls_only_on_paths_exempt_from_encryption():
"""Surowe `client.get/post` wolno wołać wyłącznie tam, gdzie szyfrowania nie ma
z założenia (`/health` dla sond k8s). Każde inne to obejście PRE-16."""
import re
from app import link_crypto
source = CLIENT.read_text(encoding="utf-8").splitlines()
offenders = []
for what, line, _, _ in _http_calls(CLIENT):
if not what.startswith("client."):
continue
url = re.search(r'f"\{self\.base_url\}([^"]*)"', source[line - 1])
if url is None or url.group(1) not in link_crypto.PUBLIC_PATHS:
offenders.append(f"{CLIENT.name}:{line} {what}")
assert not offenders, (
"Surowe wywołania HTTP poza ścieżkami wyjętymi spod szyfrowania: "
+ ", ".join(offenders)
)
def test_auth_headers_helper_is_lazy():
"""Token czytany przy wywołaniu, nie przy imporcie — inaczej pod wystartowałby
z pustym tokenem, gdyby zmienna pojawiła się później."""
@@ -131,71 +131,3 @@ def test_rate_limit_disabled_when_zero(monkeypatch):
monkeypatch.delenv("APP_PASSWORD", raising=False)
client = TestClient(_app())
assert all(client.get("/significators").status_code == 200 for _ in range(30))
# ------------------------------------------ adres klienta za odwrotnym proxy
#
# Po włączeniu TLS (PRE-16) aplikacja stoi za Ingressem, więc bezpośredni peer
# to zawsze POD proxy. Te testy pilnują obu stron kompromisu: żeby limit dalej
# rozróżniał ludzi, a jednocześnie żeby nagłówek nie stał się furtką do jego
# ominięcia.
class _Req:
"""Minimalny zamiennik Request — `client_ip` czyta tylko te dwa pola."""
def __init__(self, peer: str | None, **headers: str):
self.client = type("C", (), {"host": peer})() if peer else None
self.headers = {k.replace("_", "-"): v for k, v in headers.items()}
def test_client_ip_ignores_headers_without_trust_proxy(monkeypatch):
"""Bez TRUST_PROXY nagłówek jest bezwartościowy — każdy może go dopisać."""
monkeypatch.delenv("TRUST_PROXY", raising=False)
req = _Req("10.42.0.7", x_forwarded_for="1.2.3.4", x_real_ip="5.6.7.8")
assert security.client_ip(req) == "10.42.0.7"
def test_client_ip_takes_last_forwarded_entry(monkeypatch):
"""Ostatni wpis dopisał NASZ proxy; wcześniejsze mógł podstawić klient."""
monkeypatch.setenv("TRUST_PROXY", "true")
req = _Req("10.42.0.7", x_forwarded_for="1.1.1.1, 2.2.2.2, 192.168.1.50")
assert security.client_ip(req) == "192.168.1.50"
def test_client_ip_falls_back_to_real_ip(monkeypatch):
monkeypatch.setenv("TRUST_PROXY", "true")
req = _Req("10.42.0.7", x_real_ip="192.168.1.50")
assert security.client_ip(req) == "192.168.1.50"
def test_client_ip_falls_back_to_peer_when_headers_missing(monkeypatch):
monkeypatch.setenv("TRUST_PROXY", "true")
assert security.client_ip(_Req("10.42.0.7")) == "10.42.0.7"
def test_spoofed_forwarded_header_cannot_dodge_the_limit(monkeypatch):
"""Sedno sprawy: bez zaufania do proxy podszywanie się NIE resetuje wiadra.
Gdyby limiter brał pierwszy lepszy `X-Forwarded-For`, wystarczyłoby zmieniać
go co żądanie, żeby pompować bazę bez ograniczeń.
"""
monkeypatch.delenv("TRUST_PROXY", raising=False)
monkeypatch.setenv("RATE_LIMIT_PER_MIN", "3")
monkeypatch.delenv("APP_PASSWORD", raising=False)
client = TestClient(_app())
codes = [client.get("/significators", headers={"X-Forwarded-For": f"9.9.9.{i}"}).status_code
for i in range(6)]
assert 429 in codes
def test_proxied_clients_get_separate_buckets(monkeypatch):
"""Za proxy dwie różne osoby nie mogą się nawzajem odcinać."""
monkeypatch.setenv("TRUST_PROXY", "true")
monkeypatch.setenv("RATE_LIMIT_PER_MIN", "3")
monkeypatch.delenv("APP_PASSWORD", raising=False)
client = TestClient(_app())
first = [client.get("/significators", headers={"X-Forwarded-For": "192.168.1.50"}).status_code
for _ in range(5)]
second = client.get("/significators", headers={"X-Forwarded-For": "192.168.1.51"})
assert 429 in first, "limit musi zadziałać dla pierwszego adresu"
assert second.status_code == 200, "drugi adres ma własne wiadro"