56fdf01d7d
Testy / Testy warstwy logicznej (silnik) (pull_request) Successful in 10m45s
Testy / Testy warstwy prezentacji (dostęp do baz) (pull_request) Successful in 9m31s
Testy / Build obrazu silnika B (swisseph) (pull_request) Successful in 31s
Testy / Kontrola składni wszystkich warstw (pull_request) Successful in 15s
build / build (push) Successful in 4m15s
Testy / Testy warstwy logicznej (silnik) (push) Successful in 11m41s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 9m36s
Testy / Build obrazu silnika B (swisseph) (push) Successful in 33s
Testy / Kontrola składni wszystkich warstw (push) Successful in 15s
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. Strumien postepu (okno pisania horoskopu) tez idzie przez szyfrowane lacze: link_crypto.stream_lines() pieczetuje zadanie i odszyfrowuje odpowiedz ramka po ramce (granice ramek != granice linii NDJSON), zachowujac dostarczanie na zywo. Bez tego przy wlaczonym LINK_ENCRYPTION_REQUIRED serwer odrzucalby strumien (400) i okno postepu przestaloby dzialac. Fail-closed obejmuje takze strumien: klient bez klucza nie wysyla nic, zamiast puscic dane urodzenia jawnym tekstem, zanim serwer zdazy odmowic. Najgrozniejszy blad wyszedl z PODSLUCHU prawdziwego gniazda, nie z testow: klient bez klucza wysylal pytanie jawnym tekstem, ZANIM serwer zdazyl odmowic. Stad LINK_ENCRYPTION_REQUIRED: klient nie wysyla niczego, a usluga nie wstaje, jesli klucza brak. Ta sama zasada co przy sekrecie logowania. Klient prezentacji przepuszczony przez jeden punkt _post()/stream_lines: 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: kazde wyjscie w dol musi miec i token, i klucz lacza (takze strumien), a surowe httpx wolno tylko na sciezkach wyjetych spod szyfrowania. Weryfikacja: - testy link_crypto (round-trip, brak tresci baz w bajtach na sieci, odrzucenie obcego klucza / przestawionego bitu / przekleconej sciezki / przestawionej ramki / przeterminowanej koperty / urwanego strumienia; round-trip strumienia i fail-closed klienta i serwera dla strumienia), - e2e na prawdziwym uvicornie z proxy zrzucajacym gniazdo: tresci baz brak na kablu w obie strony (grep=0), takze dla strumienia horoskopu; klucz jednej pary nie otwiera drugiej, - calosc: logika 234 passed / 1 skipped, prezentacja 25 passed. docs/wdrozenie-pre16.md: instrukcja krok po kroku z uzasadnieniem kolejnosci. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
522 lines
22 KiB
Python
522 lines
22 KiB
Python
"""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")
|
|
|
|
|
|
def stream_lines(client, url: str, *, payload, headers: dict[str, str] | None = None,
|
|
link: Link | None) -> Iterator[str]:
|
|
"""Strumieniowe POST zwracające kolejne NIEPUSTE linie NDJSON — na żywo.
|
|
|
|
Dla okna postępu: linie muszą docierać w trakcie pracy, nie na końcu, więc
|
|
czytamy strumień, a nie całe ciało. Gdy łącze ma klucz, żądanie jest
|
|
pieczętowane, a odpowiedź odszyfrowywana ramka po ramce; granice ramek NIE
|
|
pokrywają się z granicami linii, więc sklejamy bajty w buforze i tniemy je
|
|
dopiero na znakach nowej linii.
|
|
|
|
Bez klucza zachowuje się jak dotąd (surowy strumień), żeby dev bez sekretów
|
|
działał bez zmian.
|
|
"""
|
|
import json as _json
|
|
|
|
import httpx as _httpx
|
|
|
|
request_headers = dict(headers or {})
|
|
if link is None:
|
|
if encryption_required():
|
|
# Ten sam kontrakt co w `call`: nie wypuszczamy jawnego żądania, gdy
|
|
# szyfrowanie jest wymagane. Bez tego serwer owszem odrzuca (400), ale
|
|
# ciało żądania — tu dane urodzenia — zdążyłoby już pójść w eter.
|
|
raise LinkError(
|
|
f"{ENV_REQUIRED} jest włączone, ale brak klucza łącza — strumień "
|
|
f"NIE został wysłany, żeby jego treść nie poszła jawnym tekstem"
|
|
)
|
|
with client.stream("POST", url, json=payload, headers=request_headers) as response:
|
|
response.raise_for_status()
|
|
for text_line in response.iter_lines():
|
|
if text_line:
|
|
yield text_line
|
|
return
|
|
|
|
path = _httpx.URL(url).path
|
|
stamp = stamp_now()
|
|
body = frame_out(link.seal(REQUEST, path, stamp, 0, _json.dumps(payload).encode("utf-8")))
|
|
request_headers.update({HEADER_ENC: VERSION, HEADER_TS: stamp, "Content-Type": CONTENT_TYPE})
|
|
with client.stream("POST", url, content=body, headers=request_headers) as response:
|
|
response.raise_for_status()
|
|
buffer = bytearray()
|
|
for plain in open_response_stream(response, link):
|
|
buffer += plain
|
|
while True:
|
|
nl = buffer.find(b"\n")
|
|
if nl < 0:
|
|
break
|
|
text_line = bytes(buffer[:nl])
|
|
del buffer[:nl + 1]
|
|
if text_line:
|
|
yield text_line.decode("utf-8")
|
|
if buffer:
|
|
yield bytes(buffer).decode("utf-8")
|