"""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")