Files
astrololo/services/logic/app/link_crypto.py
T
gitea 2e9d3706ec
Testy / Testy warstwy logicznej (silnik) (push) Failing after 4m50s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 9m28s
Testy / Testy warstwy bazodanowej (ochrona baz) (push) Successful in 9m25s
Testy / Testy astrodemo (push) Successful in 9m25s
Testy / Build obrazu silnika B (swisseph) (push) Successful in 8s
Testy / Kontrola składni wszystkich warstw (push) Successful in 5s
Testy / Testy warstwy logicznej (silnik) (pull_request) Failing after 4m43s
Testy / Testy warstwy prezentacji (dostęp do baz) (pull_request) Successful in 9m28s
Testy / Testy warstwy bazodanowej (ochrona baz) (pull_request) Successful in 9m25s
Testy / Testy astrodemo (pull_request) Successful in 9m25s
Testy / Build obrazu silnika B (swisseph) (pull_request) Successful in 6s
Testy / Kontrola składni wszystkich warstw (pull_request) Successful in 4s
astrodemo: zmiana nazwy, rebase na mastera i domknięcie wycieków
Pierwszy z pięciu kroków budowy trzech produktów: astrodemo (dwie funkcje) →
astroklient (pełne astro bez AI) → astrololo (wszystko). Nazwa „astroklient"
zostaje zwolniona dla warstwy pośredniej, więc dotychczasowe astroklient-demo
nazywa się teraz astrodemo.

REBASE. Cztery commity demo przeniesione na aktualnego mastera. Konflikt był
jeden — rejestr wymagań (xlsx, binarny, git go nie scali). Master dodał LOG-34,
gałąź demo PRE-28 i PRE-29; żaden wspólny wiersz się nie różnił, więc scalone
ręcznie: 148 pozycji, wszystkie trzy obecne.

ZMIANA NAZWY. Katalog, ciasteczko sesji (astrodemo_sesja), zmienne
ASTRODEMO_USERS/USER/PASSWORD, CI, README, docstringi. Ponieważ PR z demo nigdy
nie został zmergowany, usługa nie jest nigdzie wdrożona — zmiana nazw niczego
nie migruje i nikogo nie wylogowuje.

WYCIEKI. astrodemo powstało przed audytem z PRE-27, więc miało komplet tych
samych dziur:

- /static omijało bramkę (PUBLIC_PREFIXES), a pierwszy komentarz w styles.css
  brzmiał „Nie kopiujemy stylów pełnej aplikacji" — czyli anonimowy curl
  dowiadywał się, że istnieje pełna aplikacja. Zasoby idą teraz trasą z jawną
  listą, komentarze są zdejmowane przy serwowaniu.
- Komunikat awarii wypisywał na ekran treść wyjątku httpx, z nazwą usługi
  i portem. Teraz jedno neutralne zdanie, szczegóły do dziennika.
- /health oddawał nazwę warstwy. Teraz samo „ok".
- Dziesięć komentarzy i docstringów tłumaczyło decyzje przez porównanie
  z „pełną aplikacją". Obraz tej usługi się KOMUŚ ODDAJE, więc kto go dostanie,
  przeczyta też komentarze. Przepisane tak, żeby opisywały tę usługę samą
  w sobie.
- Nagłówek main.py twierdził, że demo dzieli pulę plików z produkcją. To
  nieprawda od PRE-29 (pule per konto) — opis poprawiony.

ZAPORA SŁOWNIKOWA, dwupoziomowa. Poziom „wszędzie" (także w kodzie serwera, bo
obraz się oddaje) obejmuje wzmianki o większym rodzeństwie, o modelu językowym
i o funkcjach, których tu nie ma. Poziom „do przeglądarki" dokłada słownictwo
mechanizmów. Test sprawdza odpowiedzi ORAZ drzewo plików.

Jeden wyjątek jest jawny i opisany: stałe protokołu łącza (X-Astrololo-Token,
X-Astrololo-Enc, typ treści, etykieta HKDF) niosą nazwę rodziny produktów. Są
wspólne z warstwą logiczną, więc zmiana wymaga jednoczesnej podmiany we
wszystkich usługach i rotacji — osobna decyzja. Osobny test pilnuje warunku, pod
jakim to zostaje: że nie docierają do przeglądarki. Wcześniej przechodziły tylko
dlatego, że regex nie dopasowywał po myślniku — przypadek, nie decyzja.

Przy okazji: komunikat „ramka bez znacznika astrololo" zmieniony na neutralny we
WSZYSTKICH PIĘCIU kopiach link_crypto.py (presentation, astrodemo, logic, data,
render), żeby nie rozjechały się przed scaleniem w rdzeń. Te kopie to 2625 linii
tego samego kodu.

Testy: astrodemo 27, presentation 358, logic 342, data 42, render 41.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-26 12:05:40 +02:00

526 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"
# Trzecia para: prezentacja ↔ render (PRE-24). Osobny klucz, jak przy pozostałych —
# usługa render dostaje CAŁY raport (dane urodzeniowe + opisy z baz), więc przejęcie
# jej klucza nie może otwierać łącza do logiki ani do danych.
ENV_PRESENTATION_RENDER = "LINK_KEY_PRESENTATION_RENDER"
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 protokołu")
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")