From 91644a44e32038281c34e23070f01ebf42b6ad47 Mon Sep 17 00:00:00 2001 From: migatu Date: Wed, 22 Jul 2026 23:26:57 +0200 Subject: [PATCH 1/2] fix(prezentacja): limit zadan po adresie klienta, nie proxy (PRE-16) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Po wlaczeniu TLS aplikacja stanie za Ingressem, a wtedy `request.client.host` to adres POD-a Traefika — jednakowy dla wszystkich. Limiter wrzucalby caly ruch do jednego wiadra 120/min i pierwsza osoba, ktora go wyklika, odcielaby pozostalych. Cicha regresja, ktora ujawnilaby sie dopiero na produkcji. Nowe `client_ip()` czyta adres z naglowka, ale WYLACZNIE przy TRUST_PROXY — bo inaczej wystarczyloby dopisywac wlasny X-Forwarded-For, zeby przy kazdym zadaniu wygladac na kogos innego i ominac limit calkowicie. Z tego samego powodu bierzemy OSTATNI wpis listy: to jedyny, ktory dopisal nasz proxy; wczesniejsze mogl podstawic klient, wiec nie znacza nic. Szesc testow, w tym dwa istotne: - podszycie sie pod X-Forwarded-For NIE resetuje wiadra przy wylaczonym TRUST_PROXY (inaczej baze dalo by sie pompowac bez ograniczen), - za proxy dwa rozne adresy dostaja osobne wiadra i nie odcinaja sie nawzajem. Oba sprawdzone celowym zepsuciem implementacji (zawsze ufaj naglowkowi + bierz pierwszy wpis) — testy wtedy czerwienieja. 23 passed. Co-Authored-By: Claude Opus 4.8 --- services/presentation/app/security.py | 37 ++++++++++- services/presentation/tests/test_security.py | 68 ++++++++++++++++++++ 2 files changed, 102 insertions(+), 3 deletions(-) diff --git a/services/presentation/app/security.py b/services/presentation/app/security.py index de8577a..0eddb25 100644 --- a/services/presentation/app/security.py +++ b/services/presentation/app/security.py @@ -7,7 +7,9 @@ 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ń). + * **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`). Świadomie NIE logujemy treści żądań ani promptów — logi to kolejny nośnik wycieku. @@ -46,6 +48,11 @@ 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/",) @@ -74,6 +81,31 @@ 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: @@ -105,8 +137,7 @@ def install(app) -> None: if _is_public(request.url.path): return await call_next(request) - client = request.client.host if request.client else "?" - if _rate_limited(client): + if _rate_limited(client_ip(request)): return JSONResponse( {"detail": "Zbyt wiele żądań — spróbuj za chwilę."}, status_code=429, headers={"Retry-After": "60"}, diff --git a/services/presentation/tests/test_security.py b/services/presentation/tests/test_security.py index 9d2a54e..92ec096 100644 --- a/services/presentation/tests/test_security.py +++ b/services/presentation/tests/test_security.py @@ -131,3 +131,71 @@ 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" -- 2.52.0 From 3ce3911f55cbad489b535d472ee88d35c554a7b9 Mon Sep 17 00:00:00 2001 From: migatu Date: Wed, 22 Jul 2026 23:52:18 +0200 Subject: [PATCH 2/2] feat(bezpieczenstwo): szyfrowanie lacz miedzy warstwami AES-256-GCM (PRE-16) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Warstwy rozmawialy ze soba jawnym tekstem wewnatrz klastra. Token miedzywarstwowy (LOG-32) mowil KTO pyta, ale nie ukrywal CZEGO dotyczy odpowiedz — a plyna nia surowe wiersze oryginalnych baz interpretacyjnych, czyli rdzen produktu. Kto podsluchal ruch wewnatrz sieci (drugi pod, mirror portu na switchu, zrzut z wezla), mial je w calosci. Nowy modul link_crypto (kopia w kazdej z trzech uslug — nie maja wspolnej biblioteki; test pilnuje, ze kopie sa identyczne): - AES-256-GCM na ciele kazdego zadania i odpowiedzi. GCM daje poufnosc I uwierzytelnienie naraz, wiec nie ma wariantu „zaszyfrowane, ale podatne na modyfikacje". - DWA niezalezne klucze, po jednym na pare rozmowcow (prezentacja-logika, logika-dane). Przejecie klucza prezentacji nie otwiera warstwy danych, gdzie leza cale bazy. Z kazdego klucza lacza HKDF wyprowadza osobne podklucze na kierunek, wiec zadanie i odpowiedz nigdy nie szyfruja sie tym samym kluczem. - Do materialu uwierzytelnianego (AAD) wchodza kierunek, sciezka, znacznik czasu i numer ramki — wiec ramki nie da sie przekleic na inny endpoint, odtworzyc po czasie (okno MAX_SKEW) ani przestawic w strumieniu. - Strona serwerowa to czyste ASGI: podmienia cialo zanim zobaczy je FastAPI i przepuszcza odpowiedz strumieniowa kawalek po kawalku (okno postepu dziala dalej). Fail-closed: przy ustawionym kluczu jawne zadanie dostaje odmowe. Najgrozniejszy blad wyszedl dopiero z PODSLUCHU prawdziwego gniazda, nie z testow: klient BEZ klucza wysylal pytanie jawnym tekstem, ZANIM serwer zdazyl odmowic. Odpowiedz byla chroniona, zapytanie juz nie — a to wlasnie ono niesie sygnifikatory. Stad LINK_ENCRYPTION_REQUIRED: klient nie wysyla niczego, a usluga nie wstaje, jesli klucza brak. Ta sama zasada co przy sekrecie logowania — wolimy pod w CrashLoop niz usluge, ktora wstala i po cichu nie chroni niczego. Klient prezentacji przepuszczony przez jeden punkt `_post()`: dopoki kazda metoda skladala zadanie sama, dolozenie nowej znaczylo, ze latwo zapomniec o tokenie albo kluczu (401 wyszedl juz raz dopiero na produkcji). Test strukturalny rozszerzony: kazde wyjscie w dol musi miec i token, i klucz lacza, a surowe httpx wolno tylko na sciezkach wyjetych spod szyfrowania (/health). Weryfikacja: - 23 testy link_crypto, w tym dowod, ze tajny opis NIE wystepuje w bajtach lecacych po sieci, oraz odrzucenie: obcego klucza, przestawionego bitu, przekleconej sciezki, przestawionej ramki, przeterminowanej koperty, urwanego strumienia i absurdalnej dlugosci ramki. - E2E na prawdziwym uvicornie z proxy zrzucajacym gniazdo do pliku: przy wlaczonym wymaganiu tresci baz NIE MA na kablu (grep = 0) w obie strony; bez klucza / ze zlym kluczem — odmowa; klucz jednej pary nie otwiera drugiej. - Calosc: logika 200 passed/1 skipped, prezentacja 25 passed. docs/wdrozenie-pre16.md: instrukcja krok po kroku (klucze -> cert-manager -> DNS -> merge aplikacji -> merge manifestow -> zaufanie CA -> weryfikacja), z uzasadnieniem kolejnosci i tabela diagnostyki. Co-Authored-By: Claude Opus 4.8 --- docs/wdrozenie-pre16.md | 296 +++++++++++ services/data/app/link_crypto.py | 466 ++++++++++++++++++ services/data/app/main.py | 5 +- services/data/requirements.txt | 2 + services/logic/app/clients/data_client.py | 16 +- services/logic/app/link_crypto.py | 466 ++++++++++++++++++ services/logic/app/main.py | 5 +- services/logic/requirements.txt | 2 + services/logic/tests/test_link_crypto.py | 351 +++++++++++++ .../presentation/app/clients/logic_client.py | 53 +- services/presentation/app/link_crypto.py | 466 ++++++++++++++++++ services/presentation/requirements.txt | 2 + .../presentation/tests/test_client_auth.py | 64 ++- 13 files changed, 2152 insertions(+), 42 deletions(-) create mode 100644 docs/wdrozenie-pre16.md create mode 100644 services/data/app/link_crypto.py create mode 100644 services/logic/app/link_crypto.py create mode 100644 services/logic/tests/test_link_crypto.py create mode 100644 services/presentation/app/link_crypto.py diff --git a/docs/wdrozenie-pre16.md b/docs/wdrozenie-pre16.md new file mode 100644 index 0000000..89f7c8b --- /dev/null +++ b/docs/wdrozenie-pre16.md @@ -0,0 +1,296 @@ +# 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:'' https://astrololo.czernobog.pl/ +curl -sI -u astrololo:'' 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. diff --git a/services/data/app/link_crypto.py b/services/data/app/link_crypto.py new file mode 100644 index 0000000..5eef78b --- /dev/null +++ b/services/data/app/link_crypto.py @@ -0,0 +1,466 @@ +"""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") diff --git a/services/data/app/main.py b/services/data/app/main.py index 6184dd8..29d3338 100644 --- a/services/data/app/main.py +++ b/services/data/app/main.py @@ -10,7 +10,7 @@ from contextlib import asynccontextmanager from fastapi import FastAPI -from app import security +from app import link_crypto, security from app.config import settings from app.models import HealthInfo, SearchQuery, SearchResult from app.providers.factory import build_provider @@ -26,6 +26,9 @@ 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) diff --git a/services/data/requirements.txt b/services/data/requirements.txt index 7b91692..352a2b8 100644 --- a/services/data/requirements.txt +++ b/services/data/requirements.txt @@ -7,3 +7,5 @@ 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 diff --git a/services/logic/app/clients/data_client.py b/services/logic/app/clients/data_client.py index 9895ba6..61f67ec 100644 --- a/services/logic/app/clients/data_client.py +++ b/services/logic/app/clients/data_client.py @@ -10,6 +10,7 @@ from typing import Any import httpx +from app import link_crypto from app.config import settings @@ -19,6 +20,13 @@ 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("/") @@ -33,11 +41,13 @@ 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: - r = client.post(f"{self.base_url}/search", json=payload, headers=_auth_headers()) - r.raise_for_status() - return r.json() + return link_crypto.call_json(client, "POST", f"{self.base_url}/search", + payload=payload, headers=_auth_headers(), + link=_link()) 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() diff --git a/services/logic/app/link_crypto.py b/services/logic/app/link_crypto.py new file mode 100644 index 0000000..5eef78b --- /dev/null +++ b/services/logic/app/link_crypto.py @@ -0,0 +1,466 @@ +"""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") diff --git a/services/logic/app/main.py b/services/logic/app/main.py index 9f88a5a..d83a52e 100644 --- a/services/logic/app/main.py +++ b/services/logic/app/main.py @@ -12,7 +12,7 @@ import httpx from fastapi import FastAPI, HTTPException from pydantic import BaseModel -from app import security +from app import link_crypto, security from app.clients.data_client import DataClient from app.models import QueryRequest, QueryResponse from app.service import QueryService @@ -20,6 +20,9 @@ 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 diff --git a/services/logic/requirements.txt b/services/logic/requirements.txt index 4eadce8..115e762 100644 --- a/services/logic/requirements.txt +++ b/services/logic/requirements.txt @@ -4,3 +4,5 @@ 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 diff --git a/services/logic/tests/test_link_crypto.py b/services/logic/tests/test_link_crypto.py new file mode 100644 index 0000000..e75f28f --- /dev/null +++ b/services/logic/tests/test_link_crypto.py @@ -0,0 +1,351 @@ +"""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") diff --git a/services/presentation/app/clients/logic_client.py b/services/presentation/app/clients/logic_client.py index 7df6731..a6a6938 100644 --- a/services/presentation/app/clients/logic_client.py +++ b/services/presentation/app/clients/logic_client.py @@ -10,6 +10,7 @@ from typing import Any import httpx +from app import link_crypto from app.config import settings @@ -19,16 +20,29 @@ 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} - 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() + return self._post("/api/query", payload, settings.http_timeout) def positions( self, @@ -51,20 +65,15 @@ class LogicClient: "zodiac": zodiac, } # stacje wymagają root-findów — dłuższy 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() + timeout = max(settings.http_timeout, 60.0) if stations else settings.http_timeout + return self._post("/chart/positions", payload, timeout) 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} - 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() + return self._post("/chart/report", payload, max(settings.http_timeout, 30.0)) def prompt( self, profile: str, when_utc_iso: str, lat: float, lon: float, @@ -82,10 +91,7 @@ class LogicClient: } if from_date and to_date: payload["from_date"], payload["to_date"] = from_date, to_date - 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() + return self._post("/chart/prompt", payload, max(settings.http_timeout, 60.0)) def horoscope( self, profile: str, when_utc_iso: str, lat: float, lon: float, @@ -106,17 +112,13 @@ class LogicClient: payload["model"] = model if from_date and to_date: payload["from_date"], payload["to_date"] = from_date, to_date - 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() + return self._post("/chart/horoscope", payload, max(settings.http_timeout, 300.0)) 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: - r = client.get(f"{self.base_url}/llm/models", headers=_auth_headers()) - r.raise_for_status() - return r.json() + return link_crypto.call_json(client, "GET", f"{self.base_url}/llm/models", + headers=_auth_headers(), link=_link()) def timeline( self, when_utc_iso: str, lat: float, lon: float, @@ -127,7 +129,4 @@ class LogicClient: "when_utc": when_utc_iso, "lat": lat, "lon": lon, "from_date": from_date, "to_date": to_date, "interpret": interpret, } - 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() + return self._post("/chart/timeline", payload, max(settings.http_timeout, 60.0)) diff --git a/services/presentation/app/link_crypto.py b/services/presentation/app/link_crypto.py new file mode 100644 index 0000000..5eef78b --- /dev/null +++ b/services/presentation/app/link_crypto.py @@ -0,0 +1,466 @@ +"""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") diff --git a/services/presentation/requirements.txt b/services/presentation/requirements.txt index 71982e7..faca6fc 100644 --- a/services/presentation/requirements.txt +++ b/services/presentation/requirements.txt @@ -3,3 +3,5 @@ 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 diff --git a/services/presentation/tests/test_client_auth.py b/services/presentation/tests/test_client_auth.py index e29d6cc..bea5e80 100644 --- a/services/presentation/tests/test_client_auth.py +++ b/services/presentation/tests/test_client_auth.py @@ -1,30 +1,44 @@ -"""Niezmiennik: KAŻDE wyjście HTTP w dół niesie token międzywarstwowy (LOG-32). +"""Niezmiennik: KAŻDE wyjście HTTP w dół niesie token międzywarstwowy (LOG-32) +oraz klucz szyfrujący łącze (PRE-16). 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: nie ma wywołania bez `headers=`. +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. """ 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]]: - """(nazwa_metody_http, linia, czy_ma_headers) dla każdego client.post/get.""" + +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ół.""" 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 - 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 + target = node.func.value has_headers = any(kw.arg == "headers" for kw in node.keywords) - out.append((node.func.attr, node.lineno, has_headers)) + 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)) return out @@ -35,13 +49,43 @@ 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} client.{verb}()" for verb, line, ok in calls if not ok] + missing = [f"{CLIENT.name}:{line} {what}" for what, 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.""" -- 2.52.0