Compare commits

..

9 Commits

Author SHA1 Message Date
gitea 3ce3911f55 feat(bezpieczenstwo): szyfrowanie lacz miedzy warstwami AES-256-GCM (PRE-16)
Testy / Testy warstwy logicznej (silnik) (push) Successful in 10m50s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 9m50s
Testy / Build obrazu silnika B (swisseph) (push) Successful in 30s
Testy / Kontrola składni wszystkich warstw (push) Successful in 21s
Testy / Testy warstwy logicznej (silnik) (pull_request) Successful in 10m47s
Testy / Testy warstwy prezentacji (dostęp do baz) (pull_request) Successful in 9m52s
Testy / Build obrazu silnika B (swisseph) (pull_request) Successful in 35s
Testy / Kontrola składni wszystkich warstw (pull_request) Successful in 17s
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 <noreply@anthropic.com>
2026-07-22 23:52:18 +02:00
gitea 91644a44e3 fix(prezentacja): limit zadan po adresie klienta, nie proxy (PRE-16)
Testy / Testy warstwy logicznej (silnik) (push) Successful in 10m47s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 9m33s
Testy / Build obrazu silnika B (swisseph) (push) Successful in 31s
Testy / Kontrola składni wszystkich warstw (push) Successful in 21s
Testy / Testy warstwy logicznej (silnik) (pull_request) Successful in 10m46s
Testy / Testy warstwy prezentacji (dostęp do baz) (pull_request) Successful in 9m53s
Testy / Build obrazu silnika B (swisseph) (pull_request) Successful in 32s
Testy / Kontrola składni wszystkich warstw (pull_request) Successful in 20s
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 <noreply@anthropic.com>
2026-07-22 23:26:57 +02:00
gitea 163ace4283 feat(ui): interaktywny wybor modelu u dostawcy
Testy / Testy warstwy logicznej (silnik) (pull_request) Successful in 11m14s
Testy / Testy warstwy prezentacji (dostęp do baz) (pull_request) Successful in 10m19s
Testy / Build obrazu silnika B (swisseph) (pull_request) Successful in 46s
Testy / Kontrola składni wszystkich warstw (pull_request) Successful in 27s
build / build (push) Successful in 1m18s
Testy / Testy warstwy logicznej (silnik) (push) Successful in 11m45s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 10m2s
Testy / Build obrazu silnika B (swisseph) (push) Successful in 36s
Testy / Kontrola składni wszystkich warstw (push) Successful in 23s
Uzytkownik wybiera nie tylko dostawce, ale konkretny model — Fable czy Opus
u Anthropica, gpt-4o-mini czy gpt-5 u OpenAI, cokolwiek ma pobrane lokalnie.

- app/llm/catalog.py: podpowiedzi modeli per dostawca wraz z oknem kontekstu,
  nadpisywalne przez <DOSTAWCA>_MODELS; GET /llm/models wystawia je dla UI.
- Pole modelu w UI jest TEKSTOWE z datalista, nie zamknietym <select> — konto
  moze miec dostep do modeli, o ktorych kod nie wie, a nowe wychodza szybciej,
  niz aktualizuje sie katalog. Puste pole = model domyslny dostawcy.
- static/models.js: zmiana dostawcy przelacza podpowiedzi, podmienia placeholder
  na model domyslny i pokazuje okno kontekstu wybranego modelu.
- build_provider(name, model) — model z zadania wygrywa nad konfiguracja.

WAZNE (znalezione przy tescie e2e): liczenie budzetu „maksymalny kontekst modelu"
szlo przez build_provider(), ktory WYMAGA klucza API — bez klucza budzet cicho
spadal do wartosci zapasowej i byl identyczny dla wszystkich modeli Anthropic.
Budzet zalezy wylacznie od okna kontekstu, wiec doszlo resolve_model(), ktore
rozwiazuje nazwe modelu bez budowania dostawcy. Teraz budzet realnie sie rozni:
Opus/Fable 3,48 mln znakow, Haiku 536 tys., gpt-4o-mini 438 tys., llama3.1 8 tys.

Pewnosc danych w katalogu: modele Anthropic pochodza z oficjalnej dokumentacji
API (okna i limity zgodne z limits.py); modele OpenAI to podpowiedzi, ktorych
nie weryfikowalem; lokalne zaleza od tego, co masz pobrane.

Testy: 174 passed / 1 skipped (logika) + 17 (prezentacja). Nowy test strukturalny
pilnuje, ze KAZDE wywolanie w dol niesie wybrany model i dostawce — dokladnie ta
klasa bledu zlapala brakujacy parametr przy horoskopie okresowym.
Zweryfikowane w przegladarce: przelaczanie dostawcy podmienia liste modeli.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-22 21:40:19 +02:00
gitea f33cdc88f4 feat(llm): horoskop powstaje zawsze — kontynuacja, okna kontekstu, budzet max
PRZYCZYNA PUSTYCH ODPOWIEDZI NA ANTHROPICU (potwierdzona w dokumentacji API):
domyslnym modelem byl `claude-sonnet-5`, ktory przy POMINIETYM parametrze
`thinking` wlacza myslenie adaptacyjne, a `thinking.display` domyslnie jest
"omitted". Tokeny myslenia licza sie do max_tokens, wiec przy LLM_MAX_TOKENS=2000
cala tura wychodzila jako bloki `thinking` z pustym tekstem — parser filtrowal
type=="text" i zwracal pusty string. Opus 4.8 bez `thinking` nie mysli, wiec tam
objaw by nie wystapil.

Gwarancja niepustej odpowiedzi (wszyscy trzej dostawcy):
- generate() to teraz PETLA, nie pojedynczy strzal: tura -> jesli urwana na
  limicie, dopisz ture „kontynuuj" w tej samej rozmowie i sklej tekst,
- tura zlozona z samego myslenia traktowana jak urwana (nie jak pustka),
- pusta i NIE urwana -> jedna proba z podpowiedzia, dopiero potem blad,
- kontynuacja konczy sie tura UZYTKOWNIKA — Claude odrzuca prefill asystenta (400),
- `thinking` konfigurowany JAWNIE (adaptive + effort=high; ANTHROPIC_THINKING=off).

Okna kontekstu i rezerwa na odpowiedz (app/llm/limits.py):
- tabela okien/limitow wyjscia per model + nadpisanie z ENV,
- plan() liczy okno odpowiedzi jako okno - prompt - margines i NIGDY nie oddaje
  calego kontekstu promptowi,
- Anthropic liczy tokeny DOKLADNIE (/v1/messages/count_tokens), reszta szacuje,
- >90 tys. tokenow promptu -> ostrzezenie, ale wyslanie NADAL mozliwe i z pelnym
  oknem odpowiedzi.

UI: suwak budzetu rozszerzony o „bardzo obszerny" i „maksymalny kontekst modelu"
(liczony z okna wybranego modelu po odjeciu rezerwy); przy wyniku widac plan
tokenow, liczbe tur i ostrzezenia.

Domyslny model Anthropic: claude-opus-4-8.

Testy: 170 passed / 1 skipped (logika) + 15 (prezentacja). Nowe testy pokrywaja
sklejanie kontynuacji, brak prefillu asystenta, ture z samego myslenia, rezerwe
na odpowiedz i prog ostrzezenia. Zweryfikowane e2e na atrapie Anthropica
odtwarzajacej zgloszony objaw: 3 tury, obie czesci tekstu obecne.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-22 21:40:19 +02:00
gitea 877ec91ff0 fix(llm): pusta odpowiedz modelu to blad, nie pusta strona
build / build (push) Successful in 57s
Testy / Testy warstwy logicznej (silnik) (push) Successful in 13m7s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 9m57s
Testy / Build obrazu silnika B (swisseph) (push) Successful in 40s
Testy / Kontrola składni wszystkich warstw (push) Successful in 25s
Objaw zgloszony przez uzytkownika: w Kalendarzu horoskop wyswietla sie
poprawnie, w Interpretacjach zapytanie wychodzi, wraca — i NIC sie nie
pokazuje. Bez zadnego komunikatu.

Przyczyna: model potrafi oddac pusta tresc (finish_reason=length,
completion_tokens=0), a generate() zwracalo wtedy pusty tekst BEZ bledu.
Widok sprawdza {% if prompt_result.horoscope %} -> falsz -> nie renderuje nic,
a llm_error nie jest ustawiony -> zero wyjasnienia. Cicha awaria.

Asymetria miedzy ekranami wynika z rozmiaru promptu: natalny (13 obiektow x
fasety x opisy z bazy) wypelnia okno kontekstu modelu lokalnego i na odpowiedz
nie zostaje miejsca; okresowy jest mniejszy i sie miesci.

- wspolny straznik _require_text() dla obu dostawcow: pusta lub bialoznakowa
  odpowiedz podnosi LLMError,
- komunikat PROWADZI DO PRZYCZYNY: podaje finish_reason i zuzycie tokenow oraz
  radzi zmniejszyc budzet promptu / zwiekszyc num_ctx / LLM_MAX_TOKENS,
- Anthropic sprowadzony do wspolnego ksztaltu diagnostyki (input/output_tokens).

Dzieki temu uzytkownik widzi powod ORAZ gotowy prompt do recznego uzycia.

Zweryfikowane na zywym stosie z atrapa modelu oddajaca pusta tresc: zamiast
pustej strony pojawia sie pelny komunikat z diagnostyka. Testy: 148 passed.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-22 18:05:20 +00:00
gitea 487dbb8fd3 fix(ci): smoke test silnika B bez kontenera w tle i bez sieci
Testy / Testy warstwy logicznej (silnik) (pull_request) Successful in 10m54s
Testy / Testy warstwy prezentacji (dostęp do baz) (pull_request) Successful in 9m54s
Testy / Build obrazu silnika B (swisseph) (pull_request) Successful in 39s
Testy / Kontrola składni wszystkich warstw (pull_request) Successful in 21s
build / build (push) Successful in 56s
Testy / Testy warstwy logicznej (silnik) (push) Successful in 10m53s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 9m51s
Testy / Build obrazu silnika B (swisseph) (push) Successful in 39s
Testy / Kontrola składni wszystkich warstw (push) Successful in 20s
Job „Build obrazu silnika B" padal na kazdym przebiegu po pierwszym:

    Conflict. The container name "/swe" is already in use

Dwie wady starego kroku, obie moje:
1. startowal kontener w tle (docker run -d --name swe) i NIGDY go nie usuwal,
   wiec nazwa zostawala zajeta na runnerze i kolejne przebiegi sie wywalaly;
2. pukal curl-em w localhost:8003, podczas gdy job Gitea Actions sam dziala
   w kontenerze, a -p publikuje port na HOSCIE — to nie ten sam localhost,
   wiec health-check i tak nie mial prawa dojsc.

Teraz test biegnie WEWNATRZ obrazu (docker run --rm ... python -), wolajac
funkcje endpointow wprost. Omija oba problemy, nie zostawia niczego po sobie,
a sprawdza to samo i wiecej: obraz sie zbudowal, pyswisseph liczy, komplet 13
obiektow, Slonce w oczekiwanym zakresie, SN = NN + 180.

Dodany krok sprzatajacy osierocony kontener „swe" ze starych przebiegow.

Zweryfikowane lokalnie na realnym pyswisseph: Sun=40.2102, NN=68.1530,
13 obiektow — zgodnie z wyrocznia.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-22 17:29:27 +02:00
gitea 5203ba9e76 fix(llm): konfiguracja per dostawca — przelacznik w UI byl iluzja
Testy / Testy warstwy logicznej (silnik) (pull_request) Successful in 10m49s
Testy / Testy warstwy prezentacji (dostęp do baz) (pull_request) Successful in 9m55s
Testy / Build obrazu silnika B (swisseph) (pull_request) Failing after 28s
Testy / Kontrola składni wszystkich warstw (pull_request) Successful in 23s
build / build (push) Successful in 1m14s
Testy / Testy warstwy logicznej (silnik) (push) Successful in 11m0s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 10m0s
Testy / Build obrazu silnika B (swisseph) (push) Failing after 50s
Testy / Kontrola składni wszystkich warstw (push) Successful in 22s
UI pozwala wybrac dostawce przy KAZDYM zadaniu, ale factory czytalo jedna
wspolna trojke LLM_MODEL / LLM_BASE_URL / LLM_API_KEY dla wszystkich. Na
klastrze LLM_BASE_URL trzeba ustawic na lokalny model (localhost:11434 w podzie
nie istnieje) — i wtedy:
  - wybor „OpenAI" wysylal zadanie do Ollamy,
  - LLM_MODEL=llama3.1:8b kazal Anthropic uzyc modelu llama,
  - jeden LLM_API_KEY nie moze byc kluczem OpenAI i Anthropic naraz.
Czyli nie bylo miejsca, w ktore dalo sie sensownie wpisac klucze do chmury.

- konfiguracja per dostawca: <DOSTAWCA>_MODEL / _BASE_URL / _API_KEY
  (LOCAL_*, OPENAI_*, ANTHROPIC_*),
- zgodnosc wstecz: wspolne LLM_* dziala nadal, ale stosuje sie WYLACZNIE do
  dostawcy domyslnego (LLM_PROVIDER) — instalacja jednodostawcowa bez zmian,
- Anthropic dostal brakujaca walidacje klucza (mial ja tylko OpenAI),
- komunikat bledu wskazuje konkretna zmienna do ustawienia.

Testy regresyjne pilnuja, ze ustawienia jednego dostawcy NIE przeciekaja na
pozostalych. Calosc: 143 passed / 1 skipped.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-21 22:32:04 +02:00
gitea 64d1afc76d fix(presentation): brakujacy token przy /chart/prompt i /chart/horoscope
build / build (push) Successful in 59s
Testy / Testy warstwy logicznej (silnik) (push) Successful in 10m57s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 9m52s
Testy / Build obrazu silnika B (swisseph) (push) Failing after 30s
Testy / Kontrola składni wszystkich warstw (push) Successful in 18s
Blad z mergu: _auth_headers() dodano na galezi hardeningu, ktora odbila sie
od mastera ZANIM powstaly metody prompt() i horoscope() (LOG-29/30, LOG-31).
Git zmergowal obie zmiany czysto — byly w roznych liniach — ale semantycznie
nowe metody wyszly bez tokenu i dostawaly 401 przy wlaczonej ochronie.

Skutek dla uzytkownika: przyciski „Generuj prompt (AI)" i „Napisz horoskop"
nie dzialaly po wdrozeniu INTERNAL_TOKEN, mimo ze reszta aplikacji dzialala.

- naprawione oba wywolania,
- nowy test strukturalny (AST): KAZDE wyjscie HTTP w dol musi niesc headers=.
  Test jednej metody by tego nie zlapal — regula musi byc pilnowana calosciowo.

Straznik zweryfikowany sabotazem: po usunieciu naglowka test pada ze
wskazaniem konkretnej linii; po przywroceniu 15/15 przechodzi.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-21 19:02:40 +00:00
gitea 4da5f5fe7e docs: braki bezpieczenstwa jako wymagania (PRE-16/17, DAN-25/26, LOG-33)
Testy / Testy warstwy logicznej (silnik) (pull_request) Successful in 10m53s
Testy / Testy warstwy prezentacji (dostęp do baz) (pull_request) Successful in 9m56s
Testy / Build obrazu silnika B (swisseph) (pull_request) Failing after 44s
Testy / Kontrola składni wszystkich warstw (pull_request) Successful in 34s
build / build (push) Successful in 52s
Testy / Testy warstwy logicznej (silnik) (push) Successful in 10m51s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 9m56s
Testy / Build obrazu silnika B (swisseph) (push) Failing after 30s
Testy / Kontrola składni wszystkich warstw (push) Successful in 20s
Obecny poziom ochrony wystarcza do developmentu, ale luki musza byc zapisane,
zeby nie wyparowaly przed produkcja.

- PRE-16 (Must) HTTPS/TLS — dzis Basic Auth leci po http, czyli haslo da sie
  podsluchac. Odblokowuje przy okazji DWIE funkcje zepsute z tego samego
  powodu: geolokalizacje (Tu i teraz) i kopiowanie do schowka — oba wymagaja
  secure context.
- PRE-17 (Should) konta imienne + slad audytowy zamiast jednego wspolnego
  hasla; bez tego nie wiadomo, kto pobieral dane, ani jak odciac jedna osobe.
- DAN-25 (Must) ograniczenie udzialu NFS — kto ma do niego dostep, bierze
  komplet baz z pominieciem aplikacji. Dzis najkrotsza droga do wycieku.
- DAN-26 (Should) znakowanie baz rekordami-pulapkami — zabezpieczenie
  detekcyjne: pozwala udowodnic zrodlo wycieku.
- LOG-33 (Should) sekrety w spoczynku (etcd to tylko base64) + rotacja.

Q-12 odnotowane jako rozstrzygniete: bazy zostaly KUPIONE, wiec zgoda jest —
ale to nie zwalnia z ochrony. LOG-32 przestawione na "W trakcie" z wykazem,
co juz wdrozone, a co zostaje.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-21 20:47:21 +02:00
29 changed files with 3416 additions and 126 deletions
+41 -15
View File
@@ -82,25 +82,51 @@ jobs:
- name: docker build
run: docker build -t astrololo/engine-swisseph:ci services/engine-swisseph
- name: Smoke test (health + pozycje)
# Test biegnie WEWNĄTRZ obrazu, bez sieci i bez kontenera w tle. Poprzednia
# wersja startowała kontener w tle (--name swe) i pukała curl-em w
# localhost:8003 — co miało dwie wady:
# 1. nie sprzątała kontenera, więc każdy kolejny przebieg padał na
# konflikcie nazwy (Conflict. The container name "/swe" is already in use),
# 2. job Gitea Actions sam działa w kontenerze, a -p publikuje port na
# HOŚCIE — więc localhost joba to nie ten sam localhost.
# Wywołanie funkcji endpointów wprost omija oba problemy, a sprawdza to samo:
# obraz się zbudował, pyswisseph liczy, kontrakt /positions się zgadza.
# --rm gwarantuje, że nic nie zostaje po przebiegu.
- name: Smoke test (health + pozycje) wewnątrz obrazu
run: |
docker run -d --name swe -p 8003:8003 astrololo/engine-swisseph:ci
for i in $(seq 1 30); do
curl -fsS http://localhost:8003/health >/dev/null 2>&1 && break
sleep 1
done
curl -fsS http://localhost:8003/health
echo
docker run --rm astrololo/engine-swisseph:ci python - <<'PY'
from datetime import datetime, timezone
from app.main import DEFAULT_OBJECTS, PositionsRequest, health, positions
h = health()
assert h["status"] == "ok", h
print("health:", h)
# Horoskop referencyjny (30.04.1984) — ten sam, na którym opieramy testy
# silnika własnego; sprawdzamy, że silnik B faktycznie liczy.
curl -fsS -X POST http://localhost:8003/positions \
-H 'Content-Type: application/json' \
-d '{"when_utc":"1984-04-30T09:20:00Z","lat":50.0647,"lon":19.9450}'
echo
req = PositionsRequest(when_utc=datetime(1984, 4, 30, 9, 20, tzinfo=timezone.utc),
lat=50.0647, lon=19.9450)
out = positions(req)
by = {p["name"]: p for p in out["positions"]}
- name: Logi kontenera (gdy coś padło)
if: failure()
run: docker logs swe || true
assert out["engine"] == "swisseph", out["engine"]
assert len(by) == len(DEFAULT_OBJECTS), sorted(by)
sun = by["Sun"]["longitude"]
assert 39.5 < sun < 41.0, f"Slonce poza oczekiwanym zakresem: {sun}"
nn, sn = by["North Node"]["longitude"], by["South Node"]["longitude"]
assert abs(((sn - nn) % 360.0) - 180.0) < 1e-6, (nn, sn)
print(f"Sun={sun:.4f} NN={nn:.4f} obiektow={len(by)}")
print("SMOKE OK")
PY
# Sprzątanie po POPRZEDNICH przebiegach starej wersji workflow, która
# zostawiała kontener „swe" na runnerze i blokowała nazwę. Nowa wersja
# kontenera w tle nie tworzy, więc to tylko jednorazowe uprzątnięcie.
- name: Usuń osierocony kontener ze starych przebiegów
if: always()
run: docker rm -f swe 2>/dev/null || true
compile-all:
name: Kontrola składni wszystkich warstw
Binary file not shown.
+296
View File
@@ -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:'<hasło>' https://astrololo.czernobog.pl/
curl -sI -u astrololo:'<hasło>' https://astrololo.czernobog.pl/ | grep -i strict-transport
```
Oczekiwane: **301**, **401**, **200**, nagłówek HSTS obecny. Brak ostrzeżenia
o certyfikacie w przeglądarce oznacza, że krok 6 się udał.
### W przeglądarce
Kliknij **„Tu i teraz"** — powinno pobrać lokalizację (po http odmawiało).
Wygeneruj prompt i kliknij **kopiuj** — schowek powinien zadziałać bez obejść.
### Szyfrowanie łączy — sprawdzenie wprost
Najmocniejszy test to próba obejścia. Z wnętrza klastra, **bez klucza**:
```bash
kubectl -n astrololo exec deploy/presentation -- \
python -c "
import httpx, os
r = httpx.post('http://logic:8001/chart/report',
json={'when_utc':'1984-04-30T09:20:00+00:00','lat':50.06,'lon':19.94},
headers={'X-Astrololo-Token': os.environ['INTERNAL_TOKEN']})
print(r.status_code, r.text[:120])
"
```
Oczekiwane: **400** i `Łącze międzywarstwowe wymaga szyfrowania.` Zwróć uwagę, że
żądanie miało **prawidłowy token** — sam token już nie wystarcza, i o to chodziło.
To samo w dół, do warstwy danych:
```bash
kubectl -n astrololo exec deploy/logic -- \
python -c "
import httpx, os
r = httpx.post('http://data:8002/search',
json={'key':'significator','value':'[Sat','exact':False,'limit':5},
headers={'X-Astrololo-Token': os.environ['INTERNAL_TOKEN']})
print(r.status_code, r.text[:120])
"
```
### Logi startowe
```bash
kubectl -n astrololo logs deploy/logic | grep -i "łącze\|UWAGA"
```
Powinno być `łącze szyfrowane (AES-256-GCM…)`. Jeśli widzisz ostrzeżenie
o rozmowie **jawnym tekstem** — klucz nie doszedł do poda.
---
## Dlaczego taka kolejność
| Kolejność | Skutek zamiany |
|---|---|
| Sekret **przed** manifestami | `LINK_ENCRYPTION_REQUIRED=true` bez klucza celowo wywraca start. Pody wpadną w CrashLoop i będą tak siedzieć do czasu utworzenia sekretu. |
| cert-manager **przed** manifestami | API odrzuci `Certificate`/`Issuer` jako nieznane rodzaje zasobów, ArgoCD pokaże aplikację jako niezsynchronizowaną i sam tego nie naprawi. |
| DNS **przed** manifestami | NodePort znika razem z nimi. Bez wpisu DNS zostaje tylko `port-forward`. |
| Aplikacja **przed** manifestami | Odwrotnie: manifesty włączyłyby szyfrowanie na obrazach, które go nie znają — wszystkie żądania kończyłyby się odmową do czasu przebudowy obrazów. |
Fail-closed w obie strony jest zamierzony. Usługa, która wstała i **po cichu nie
szyfruje**, jest gorsza niż pod w CrashLoop — awarii nie widać, a bazy jadą
otwartym tekstem.
---
## Wycofanie
Manifestów: `git revert` merge'a w `deploy` — ArgoCD samo wróci do NodePortu
i ruchu bez szyfrowania. Kod aplikacji **nie wymaga wycofania**: bez zmiennych
`LINK_KEY_*` moduł przepuszcza ruch jak dotąd (i głośno o tym mówi w logach).
Certyfikat i CA zostają w namespace; usunięcie: `kubectl -n astrololo delete
certificate astrololo-ca astrololo-tls`. cert-managera można zostawić — nie
przeszkadza.
---
## Gdy coś nie gra
| Objaw | Przyczyna | Co zrobić |
|---|---|---|
| Pody w `CrashLoopBackOff`, w logach `LINK_ENCRYPTION_REQUIRED … nie ustawiony` | brak sekretu `astrololo-link` | krok 1, potem `rollout restart` |
| `400 Łącze międzywarstwowe wymaga szyfrowania` przy normalnym korzystaniu | jedna warstwa ma klucz, druga nie (albo trwa rollout) | `rollout status`; sprawdź, czy wszystkie trzy pody mają zmienną |
| `400 Nie udało się odczytać zaszyfrowanego żądania` | klucze po obu stronach łącza są **różne** | wymień sekret i zrestartuj **wszystkie trzy** naraz |
| `Certificate` stoi w `READY=False` | webhook cert-managera jeszcze nie wstał | `kubectl -n cert-manager get pods`, poczekaj i sprawdź `kubectl -n astrololo describe certificate astrololo-tls` |
| Przeglądarka: „połączenie nie jest prywatne” | korzeń CA nieimportowany na tym urządzeniu | krok 6 (pamiętaj, że Firefox ma osobny magazyn) |
| `404` z Traefika pod adresem aplikacji | DNS wskazuje gdzie indziej niż LoadBalancer Traefika | porównaj `dig +short astrololo.czernobog.pl` z adresem z kroku „stan wyjściowy” |
| Limit żądań odcina wszystkich naraz | brak `TRUST_PROXY=true` — cały ruch liczony jako jeden klient | sprawdź zmienną w `deploy/presentation` |
---
## Czego to nie załatwia
- **Szyfrowane są ciała żądań, nie nagłówki.** Ścieżka (`/search`) i token
międzywarstwowy jadą czytelnie. Sam token nikomu nic nie daje — bez klucza łącza
każde żądanie kończy się odmową — ale metadanych to nie ukrywa. Pełne ukrycie
wymagałoby mTLS.
- **Własne CA to nie publiczne zaufanie.** Każde nowe urządzenie wymaga importu
korzenia. Gdyby aplikacja miała kiedyś wyjść na świat, właściwą drogą jest
Let's Encrypt przez DNS-01.
- **NFS z plikami baz** stoi obok aplikacji — kto ma dostęp do share'u, bierze
pliki z pominięciem wszystkich powyższych zabezpieczeń. Do zamknięcia po stronie
infrastruktury (eksport tylko dla IP węzłów, `root_squash`, najlepiej read-only).
- **Sekrety w etcd** są tylko zakodowane base64. Docelowo: szyfrowanie etcd
at-rest albo Sealed Secrets / SOPS.
+466
View File
@@ -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")
+4 -1
View File
@@ -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)
+2
View File
@@ -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
+13 -3
View File
@@ -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()
+466
View File
@@ -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")
+75
View File
@@ -0,0 +1,75 @@
"""Katalog modeli do wyboru w UI (LOG-31).
To są **podpowiedzi**, nie zamknięta lista. Pole modelu w UI jest tekstowe z
datalistą, więc można wpisać dowolny identyfikator — konto może mieć dostęp do
modeli, których tu nie ma, a nowe wychodzą szybciej, niż aktualizuje się kod.
Puste pole = model domyślny dostawcy.
Uwaga o pewności danych:
* modele **Anthropic** pochodzą z oficjalnej dokumentacji API (okna kontekstu
i limity wyjścia zgadzają się z `app/llm/limits.py`);
* modele **OpenAI** to podpowiedzi — nie weryfikowałem ich katalogu, więc
traktuj je jako wygodę, a nie źródło prawdy;
* modele **lokalne** zależą wyłącznie od tego, co masz pobrane w Ollamie/vLLM.
Katalog można nadpisać/rozszerzyć zmienną `<DOSTAWCA>_MODELS` (lista po przecinku),
np. `OPENAI_MODELS="gpt-5,gpt-4o"`.
"""
from __future__ import annotations
import os
from app.llm.limits import limits_for
# dostawca -> [(id modelu, krótki opis dla człowieka)]
_CATALOG: dict[str, list[tuple[str, str]]] = {
"anthropic": [
("claude-opus-4-8", "Opus 4.8 — domyślny, bardzo zdolny, 1M kontekstu"),
("claude-fable-5", "Fable 5 — najbardziej zdolny, do najtrudniejszych zadań"),
("claude-sonnet-5", "Sonnet 5 — szybszy i tańszy, jakość blisko Opusa"),
("claude-opus-4-7", "Opus 4.7 — poprzednia generacja Opusa"),
("claude-haiku-4-5", "Haiku 4.5 — najszybszy i najtańszy, mniejsze okno"),
],
"openai": [
("gpt-4o-mini", "GPT-4o mini — tani i szybki"),
("gpt-4o", "GPT-4o"),
("gpt-5", "GPT-5 — jeśli Twoje konto ma dostęp"),
("gpt-4.1", "GPT-4.1"),
("gpt-4.1-mini", "GPT-4.1 mini"),
],
"local": [
("llama3.1:8b", "Llama 3.1 8B"),
("llama3.2", "Llama 3.2"),
("qwen2.5", "Qwen 2.5 — większe okno kontekstu"),
("mistral", "Mistral"),
],
}
def models_for(provider: str) -> list[dict]:
"""Podpowiedzi modeli dla dostawcy, wraz z oknem kontekstu.
Okno kontekstu podajemy, bo wprost przekłada się na opcję „maksymalny
kontekst modelu" — użytkownik widzi, na ile budżetu promptu może liczyć.
"""
override = os.getenv(f"{provider.upper()}_MODELS", "").strip()
if override:
entries = [(m.strip(), "") for m in override.split(",") if m.strip()]
else:
entries = _CATALOG.get(provider, [])
out = []
for model_id, label in entries:
context_window, max_output = limits_for(provider, model_id)
out.append({
"id": model_id,
"label": label or model_id,
"context_window": context_window,
"max_output": max_output,
})
return out
def catalog() -> dict[str, list[dict]]:
"""Pełny katalog dla UI — jedno żądanie zamiast trzech."""
return {provider: models_for(provider) for provider in ("local", "anthropic", "openai")}
+64 -13
View File
@@ -4,13 +4,26 @@ Domyślny jest **model lokalny**: prompt niesie oryginalne opisy z baz, więc
domyślnie nic nie opuszcza naszej sieci (LOG-32). Chmurę włącza się świadomie —
przez konfigurację albo pojedyncze żądanie.
Konfiguracja jest **per dostawca**, bo UI pozwala przełączać go przy każdym żądaniu.
Wspólne `LLM_*` nie wystarczy: ustawienie `LLM_BASE_URL` na lokalny model kierowałoby
tam także żądania do OpenAI, a `LLM_MODEL=llama3.1:8b` kazałoby Anthropic użyć modelu
llama. Dlatego każdy dostawca ma własny komplet zmiennych.
Zmienne środowiskowe:
LLM_PROVIDER local (domyślnie) | openai | anthropic
LLM_MODEL nazwa modelu (domyślna zależy od dostawcy)
LLM_BASE_URL adres API (domyślnie: lokalny serwer zgodny z OpenAI)
LLM_API_KEY klucz — WYŁĄCZNIE z sekretu; niepotrzebny dla modelu lokalnego
LLM_TIMEOUT sekundy (domyślnie 120)
LLM_MAX_TOKENS limit długości odpowiedzi (domyślnie 2000)
LLM_PROVIDER local (domyślnie) | openai | anthropic — dostawca domyślny
LLM_TIMEOUT sekundy (domyślnie 120)
LLM_MAX_TOKENS limit długości odpowiedzi (domyślnie 2000)
<DOSTAWCA>_MODEL / _BASE_URL / _API_KEY — konfiguracja konkretnego dostawcy:
LOCAL_MODEL, LOCAL_BASE_URL (klucz zwykle zbędny)
OPENAI_MODEL, OPENAI_BASE_URL, OPENAI_API_KEY
ANTHROPIC_MODEL, ANTHROPIC_BASE_URL, ANTHROPIC_API_KEY
Klucze WYŁĄCZNIE z sekretu — nigdy w repo, w UI ani w logach.
Zgodność wstecz: wspólne `LLM_MODEL` / `LLM_BASE_URL` / `LLM_API_KEY` nadal działają,
ale stosują się TYLKO do dostawcy domyślnego (LLM_PROVIDER) — czyli konfiguracja
instalacji jednodostawcowej zostaje nietknięta, a pozostali dostawcy jej nie dziedziczą.
"""
from __future__ import annotations
@@ -27,7 +40,11 @@ PROVIDERS = (LOCAL, OPENAI, ANTHROPIC)
_DEFAULT_MODEL = {
LOCAL: "llama3.1:8b",
OPENAI: "gpt-4o-mini",
ANTHROPIC: "claude-sonnet-5",
# Opus 4.8 świadomie zamiast Sonnet 5: Sonnet uruchamia myślenie adaptacyjne,
# gdy pominąć parametr `thinking`, a jego tokeny liczą się do max_tokens —
# przy ciasnym limicie cała tura wychodziła jako samo myślenie z pustym
# tekstem. To była przyczyna pustych odpowiedzi na Anthropicu.
ANTHROPIC: "claude-opus-4-8",
}
_DEFAULT_URL = {
# Ollama i vLLM wystawiają zgodne API pod /v1
@@ -49,20 +66,54 @@ def timeout() -> float:
return float(os.getenv("LLM_TIMEOUT", "120"))
def build_provider(name: str | None = None) -> LLMProvider:
def setting(provider: str, suffix: str, fallback: str = "") -> str:
"""Ustawienie dostawcy: <DOSTAWCA>_<SUFIKS> → LLM_<SUFIKS> → wbudowana domyślna.
Wspólne `LLM_*` stosuje się WYŁĄCZNIE do dostawcy domyślnego — inaczej adres
lokalnego modelu przejąłby żądania do chmury (i odwrotnie).
"""
specific = os.getenv(f"{provider.upper()}_{suffix}")
if specific:
return specific
if provider == default_provider_name():
generic = os.getenv(f"LLM_{suffix}")
if generic:
return generic
return fallback
def resolve_model(name: str | None = None, model: str | None = None) -> tuple[str, str]:
"""(dostawca, model) BEZ budowania dostawcy — czyli bez wymogu klucza API.
Rozmiar budżetu promptu zależy tylko od okna kontekstu modelu, więc nie może
zależeć od tego, czy klucz jest już skonfigurowany.
"""
provider = (name or default_provider_name()).lower()
if provider not in PROVIDERS:
provider = default_provider_name()
chosen = (model or "").strip() or setting(provider, "MODEL", _DEFAULT_MODEL[provider])
return provider, chosen
def build_provider(name: str | None = None, model: str | None = None) -> LLMProvider:
"""Dostawca modelu. `model` z żądania wygrywa nad konfiguracją — użytkownik
wybiera model w UI, a konfiguracja podaje tylko wartość domyślną."""
name = (name or default_provider_name()).lower()
if name not in PROVIDERS:
raise LLMError(f"Nieznany dostawca LLM: {name!r} (dostępne: {', '.join(PROVIDERS)})")
model = os.getenv("LLM_MODEL") or _DEFAULT_MODEL[name]
base_url = os.getenv("LLM_BASE_URL") or _DEFAULT_URL[name]
api_key = os.getenv("LLM_API_KEY", "")
model = (model or "").strip() or setting(name, "MODEL", _DEFAULT_MODEL[name])
base_url = setting(name, "BASE_URL", _DEFAULT_URL[name])
api_key = setting(name, "API_KEY")
if name in (OPENAI, ANTHROPIC) and not api_key:
raise LLMError(
f"Brak klucza dla dostawcy {name} — ustaw {name.upper()}_API_KEY "
f"(z sekretu). Model lokalny klucza nie wymaga."
)
if name == ANTHROPIC:
return AnthropicProvider(base_url, model, api_key, timeout())
if name == OPENAI:
if not api_key:
raise LLMError("Brak LLM_API_KEY — dostawca openai wymaga klucza.")
return ChatCompletionsProvider(OPENAI, base_url, model, api_key, timeout(),
leaves_lan=True)
# lokalny — klucz zwykle zbędny; treść NIE opuszcza sieci
+133
View File
@@ -0,0 +1,133 @@
"""Okna kontekstu modeli i planowanie budżetu tokenów (LOG-30/31).
Po co to istnieje: horoskop MA powstać niezależnie od objętości promptu. Żeby to
zagwarantować, trzeba wiedzieć dwie rzeczy o każdym modelu — ile zmieści na
wejściu (okno kontekstu) i ile maksymalnie wypisze na wyjściu. Bez tego łatwo
wysłać prompt, który wypełnia całe okno i **nie zostawia miejsca na odpowiedź** —
model kończy wtedy na `max_tokens` z pustą albo uciętą treścią.
Zasada naczelna: **zawsze rezerwuj miejsce na odpowiedź.** Budżet promptu liczy
się jako `okno_kontekstu zarezerwowane_wyjście margines`, nigdy odwrotnie.
Wartości są zaszyte jako rozsądne domyślne i nadpisywalne środowiskiem
(`<DOSTAWCA>_CONTEXT_WINDOW`, `<DOSTAWCA>_MAX_OUTPUT`) — modele wychodzą szybciej,
niż aktualizuje się ten plik.
"""
from __future__ import annotations
import os
# model -> (okno kontekstu, maksymalne wyjście) w tokenach
_MODEL_LIMITS: dict[str, tuple[int, int]] = {
# Anthropic
"claude-opus-4-8": (1_000_000, 128_000),
"claude-opus-4-7": (1_000_000, 128_000),
"claude-opus-4-6": (1_000_000, 128_000),
"claude-sonnet-5": (1_000_000, 128_000),
"claude-sonnet-4-6": (1_000_000, 128_000),
"claude-fable-5": (1_000_000, 128_000),
"claude-haiku-4-5": (200_000, 64_000),
# OpenAI
"gpt-4o": (128_000, 16_384),
"gpt-4o-mini": (128_000, 16_384),
"gpt-4.1": (1_000_000, 32_768),
"gpt-4.1-mini": (1_000_000, 32_768),
# lokalne (Ollama/vLLM) — zwykle małe okno, dlatego ostrożna domyślna
"llama3.1": (8_192, 4_096),
"llama3.2": (8_192, 4_096),
"qwen2.5": (32_768, 8_192),
"mistral": (32_768, 8_192),
}
# gdy modelu nie ma w tabeli — zachowawczo, żeby nie obiecywać nieistniejącego okna
_FALLBACK: dict[str, tuple[int, int]] = {
"anthropic": (200_000, 32_000),
"openai": (128_000, 16_384),
"local": (8_192, 4_096),
}
# ile tokenów zostawiamy jako bufor na narzut protokołu i niedokładność liczenia
SAFETY_MARGIN = 2_000
# poniżej tylu tokenów wyjścia nie ma sensu wołać modelu — nie zmieści horoskopu
MIN_OUTPUT = 1_500
# powyżej tylu tokenów promptu ostrzegamy użytkownika (nadal pozwalając wysłać)
WARN_PROMPT_TOKENS = 90_000
def _env_int(provider: str, suffix: str) -> int | None:
raw = os.getenv(f"{provider.upper()}_{suffix}")
if not raw:
return None
try:
value = int(raw)
except ValueError:
return None
return value if value > 0 else None
def limits_for(provider: str, model: str) -> tuple[int, int]:
"""(okno kontekstu, maksymalne wyjście) dla modelu — z nadpisaniem z ENV.
Dopasowanie po prefiksie, bo nazwy modeli lokalnych niosą tag (`llama3.1:8b`).
"""
env_ctx = _env_int(provider, "CONTEXT_WINDOW")
env_out = _env_int(provider, "MAX_OUTPUT")
key = (model or "").strip().lower()
known: tuple[int, int] | None = _MODEL_LIMITS.get(key)
if known is None:
for name, pair in _MODEL_LIMITS.items():
if key.startswith(name):
known = pair
break
if known is None:
known = _FALLBACK.get(provider, _FALLBACK["local"])
return (env_ctx or known[0], env_out or known[1])
def plan(provider: str, model: str, prompt_tokens: int,
want_output: int | None = None) -> dict:
"""Ile tokenów wyjścia zamówić dla promptu tej wielkości.
Zwraca plan z jawną diagnostyką — UI ma z czego zbudować ostrzeżenie, a błąd
ma czym wytłumaczyć, dlaczego się nie udało.
"""
context_window, model_max_output = limits_for(provider, model)
room = context_window - prompt_tokens - SAFETY_MARGIN
target = want_output or model_max_output
max_output = max(0, min(model_max_output, target, room))
warnings: list[str] = []
if prompt_tokens > WARN_PROMPT_TOKENS:
warnings.append(
f"Prompt ma ~{prompt_tokens} tokenów — to dużo. Zapytanie zostanie wysłane, "
f"ale potrwa dłużej i będzie odpowiednio kosztowne."
)
if max_output < MIN_OUTPUT:
warnings.append(
f"Po zmieszczeniu promptu zostaje tylko {max_output} tokenów na odpowiedź "
f"(minimum {MIN_OUTPUT}). Zmniejsz budżet promptu albo wybierz model "
f"z większym oknem kontekstu."
)
return {
"provider": provider,
"model": model,
"context_window": context_window,
"model_max_output": model_max_output,
"prompt_tokens": prompt_tokens,
"max_output": max_output,
"fits": max_output >= MIN_OUTPUT,
"warnings": warnings,
}
def prompt_token_budget(provider: str, model: str, reserve_output: int | None = None) -> int:
"""Ile tokenów promptu wolno wysłać, ZAWSZE zostawiając miejsce na odpowiedź.
To jest podstawa opcji „maksymalny kontekst modelu" w UI.
"""
context_window, model_max_output = limits_for(provider, model)
reserve = reserve_output or model_max_output
return max(0, context_window - reserve - SAFETY_MARGIN)
+199 -36
View File
@@ -5,9 +5,24 @@ vLLM, llama.cpp) i OpenAI mówią **tym samym** protokołem `/chat/completions`,
więc jedna implementacja obsługuje oba — różni je tylko adres i klucz. Anthropic
ma własny kształt `/v1/messages`, stąd druga klasa. Mniej zależności, mniej
powierzchni ataku, pełna kontrola nad tym, co wychodzi z sieci.
**Gwarancja niepustej odpowiedzi.** Horoskop ma powstać niezależnie od objętości
promptu, więc `generate()` nie jest pojedynczym strzałem, tylko pętlą:
1. wyślij turę z policzonym limitem wyjścia,
2. jeśli model urwał na limicie — dopisz turę „kontynuuj" i sklej tekst,
3. jeśli tura nie dała ani znaku tekstu — ponów z podpowiedzią,
4. dopiero brak tekstu po wszystkich próbach jest błędem (z diagnostyką).
Kontynuacja jest pewniejsza niż jedno wielkie żądanie: każda tura mieści się
w timeoucie HTTP, a długość odpowiedzi przestaje być ograniczona jedną turą.
**Anthropic i myślenie.** Modele Claude potrafią mieć włączone myślenie, którego
tokeny liczą się do `max_tokens`. Przy ciasnym limicie cała tura potrafi wyjść
jako same bloki `thinking` z pustym tekstem — dokładnie ten objaw, który
zgłoszono. Traktujemy taką turę jak ucięcie i kontynuujemy, zamiast zwracać pustkę.
"""
from __future__ import annotations
import os
import time
import httpx
@@ -17,6 +32,23 @@ from app.llm.base import Completion, LLMError, LLMProvider
RETRY_STATUSES = {429, 500, 502, 503, 504}
MAX_ATTEMPTS = 3
# ile razy wolno poprosić model o dokończenie urwanej odpowiedzi
MAX_CONTINUATIONS = 12
# ile tokenów zamawiać na jedną turę — mieści się w timeoucie, a pętla i tak
# dociągnie resztę; zbyt duża wartość ryzykuje zerwanie połączenia w trakcie
TURN_TOKENS_CAP = 16_000
_CONTINUE = (
"Kontynuuj dokładnie od miejsca, w którym przerwałeś — nie powtarzaj tego, "
"co już napisałeś, i nie zaczynaj od nowa. Jeśli skończyłeś całą odpowiedź, "
"napisz wyłącznie: KONIEC"
)
_NUDGE = (
"Nie otrzymałem żadnej treści. Napisz odpowiedź zgodnie z powyższym poleceniem, "
"zaczynając od razu od treści horoskopu."
)
_DONE_MARKER = "KONIEC"
def _post_with_retry(url: str, headers: dict, payload: dict, timeout: float) -> dict:
"""POST z ponawianiem i backoffem — chroni przed chwilowym 429/5xx."""
@@ -37,8 +69,8 @@ def _post_with_retry(url: str, headers: dict, payload: dict, timeout: float) ->
time.sleep(2 ** attempt)
continue
raise LLMError(
f"Model nie odpowiedział w czasie {timeout:.0f}s. Dłuższe horoskopy "
f"wymagają większego LLM_TIMEOUT albo mniejszego budżetu promptu."
f"Model nie odpowiedział w czasie {timeout:.0f}s. Zwiększ LLM_TIMEOUT "
f"albo zmniejsz budżet promptu."
) from e
except httpx.HTTPError as e:
last = e
@@ -49,7 +81,108 @@ def _post_with_retry(url: str, headers: dict, payload: dict, timeout: float) ->
raise LLMError(f"Nie udało się wywołać modelu: {last}")
class ChatCompletionsProvider(LLMProvider):
def _merge_usage(total: dict, turn: dict) -> dict:
"""Sumuje zużycie tokenów przez wszystkie tury jednej odpowiedzi."""
for key, value in (turn or {}).items():
if isinstance(value, int):
total[key] = total.get(key, 0) + value
return total
def _join(parts: list[str]) -> str:
return "".join(parts).strip()
def _explain_empty(turns: int, usage: dict, stop: str | None) -> str:
detail = []
if stop:
detail.append(f"powód zakończenia: {stop}")
for key in ("completion_tokens", "output_tokens"):
if usage.get(key) is not None:
detail.append(f"tokeny odpowiedzi: {usage[key]}")
break
suffix = f" ({', '.join(detail)})" if detail else ""
return (
f"Model nie zwrócił żadnej treści po {turns} próbach{suffix}. "
f"Najczęstsza przyczyna: prompt wypełnił okno kontekstu i nie zostało miejsca "
f"na odpowiedź. Zmniejsz budżet promptu albo wybierz model z większym oknem."
)
class _Driver:
"""Wspólna pętla: tura → ewentualna kontynuacja → sklejony tekst.
Podklasy dostarczają tylko `_turn()` — reszta (kontynuacje, ponawianie pustej
tury, sumowanie zużycia) jest identyczna dla obu protokołów.
"""
name: str
model: str
leaves_lan: bool
def _turn(self, messages: list[dict], max_tokens: int):
"""(tekst, czy_ucięta, zużycie, nazwa_modelu, powód_zakończenia)."""
raise NotImplementedError
def generate(self, prompt: str, max_tokens: int) -> Completion:
messages: list[dict] = [{"role": "user", "content": prompt}]
parts: list[str] = []
usage: dict = {}
model_name = self.model
remaining = max(max_tokens, 256)
stop: str | None = None
turns = 0
nudged = False
while turns <= MAX_CONTINUATIONS:
turns += 1
budget = max(256, min(remaining, TURN_TOKENS_CAP))
text, truncated, turn_usage, model_name, stop = self._turn(messages, budget)
_merge_usage(usage, turn_usage)
remaining -= budget
chunk = text.strip()
if chunk:
if chunk.endswith(_DONE_MARKER): # model zgłasza koniec
parts.append(("\n" if parts else "") + chunk[: -len(_DONE_MARKER)].rstrip())
break
parts.append(("\n" if parts else "") + chunk)
if not truncated:
break
elif not truncated:
# pusta i NIE ucięta: jedna próba z podpowiedzią, potem koniec
if nudged or parts:
break
nudged = True
messages = messages + [
{"role": "assistant", "content": ""},
{"role": "user", "content": _NUDGE},
]
continue
# ucięta (także tura złożona z samego myślenia) — poproś o dokończenie
if remaining < 256:
break
messages = [
{"role": "user", "content": prompt},
{"role": "assistant", "content": _join(parts) or ""},
{"role": "user", "content": _CONTINUE},
]
final = _join(parts)
if not final:
raise LLMError(_explain_empty(turns, usage, stop))
usage["turns"] = turns
return Completion(text=final, model=model_name, provider=self.name,
leaves_lan=self.leaves_lan, usage=usage)
def count_tokens(self, prompt: str) -> int:
"""Szacunek tokenów promptu. Dostawcy z własnym licznikiem nadpisują."""
return int(len(prompt) / 3.6)
class ChatCompletionsProvider(_Driver, LLMProvider):
"""Protokół OpenAI `/chat/completions` — lokalny serwer modelu ORAZ OpenAI."""
def __init__(self, name: str, base_url: str, model: str, api_key: str = "",
@@ -67,24 +200,20 @@ class ChatCompletionsProvider(LLMProvider):
h["Authorization"] = f"Bearer {self.api_key}"
return h
def generate(self, prompt: str, max_tokens: int) -> Completion:
def _turn(self, messages: list[dict], max_tokens: int):
data = _post_with_retry(
f"{self.base_url}/chat/completions", self._headers(),
{
"model": self.model,
"max_tokens": max_tokens,
"messages": [{"role": "user", "content": prompt}],
},
{"model": self.model, "max_tokens": max_tokens, "messages": messages},
self.timeout,
)
try:
text = data["choices"][0]["message"]["content"]
choice = data["choices"][0]
text = choice["message"].get("content") or ""
except (KeyError, IndexError, TypeError) as e:
raise LLMError(f"Nieoczekiwany kształt odpowiedzi modelu: {str(data)[:300]}") from e
return Completion(
text=text, model=data.get("model", self.model), provider=self.name,
leaves_lan=self.leaves_lan, usage=data.get("usage") or {},
)
stop = choice.get("finish_reason")
return (text, stop == "length", data.get("usage") or {},
data.get("model", self.model), stop)
def health(self) -> dict:
info = {"provider": self.name, "model": self.model, "leaves_lan": self.leaves_lan}
@@ -97,7 +226,7 @@ class ChatCompletionsProvider(LLMProvider):
return info
class AnthropicProvider(LLMProvider):
class AnthropicProvider(_Driver, LLMProvider):
"""Protokół Anthropic `/v1/messages`."""
leaves_lan = True
@@ -110,34 +239,68 @@ class AnthropicProvider(LLMProvider):
self.api_key = api_key
self.timeout = timeout
def generate(self, prompt: str, max_tokens: int) -> Completion:
def _headers(self) -> dict:
return {
"Content-Type": "application/json",
"x-api-key": self.api_key,
"anthropic-version": "2023-06-01",
}
def _thinking(self) -> dict:
"""Konfiguracja myślenia. Domyślnie adaptacyjne — podnosi jakość tekstu.
UWAGA: tokeny myślenia liczą się do `max_tokens`, więc przy ciasnym limicie
cała tura potrafi wyjść jako samo myślenie z pustym tekstem. Pętla
kontynuacji to obsługuje, ale ANTHROPIC_THINKING=off wyłącza myślenie,
gdy zależy nam na przewidywalnym zużyciu tokenów.
"""
mode = os.getenv("ANTHROPIC_THINKING", "adaptive").lower()
if mode in ("off", "disabled", "0", "false"):
return {"thinking": {"type": "disabled"}}
return {
"thinking": {"type": "adaptive"},
"output_config": {"effort": os.getenv("ANTHROPIC_EFFORT", "high")},
}
def _turn(self, messages: list[dict], max_tokens: int):
if not self.api_key:
raise LLMError("Brak LLM_API_KEY — dostawca anthropic wymaga klucza.")
data = _post_with_retry(
f"{self.base_url}/v1/messages",
{
"Content-Type": "application/json",
"x-api-key": self.api_key,
"anthropic-version": "2023-06-01",
},
{
"model": self.model,
"max_tokens": max_tokens,
"messages": [{"role": "user", "content": prompt}],
},
self.timeout,
)
raise LLMError("Brak ANTHROPIC_API_KEY — dostawca anthropic wymaga klucza.")
payload = {"model": self.model, "max_tokens": max_tokens, "messages": messages}
payload.update(self._thinking())
data = _post_with_retry(f"{self.base_url}/v1/messages", self._headers(),
payload, self.timeout)
try:
text = "".join(b.get("text", "") for b in data["content"] if b.get("type") == "text")
blocks = data["content"]
text = "".join(b.get("text", "") for b in blocks if b.get("type") == "text")
except (KeyError, TypeError) as e:
raise LLMError(f"Nieoczekiwany kształt odpowiedzi modelu: {str(data)[:300]}") from e
return Completion(
text=text, model=data.get("model", self.model), provider=self.name,
leaves_lan=True, usage=data.get("usage") or {},
stop = data.get("stop_reason")
# tura złożona z samego myślenia = budżet poszedł na rozumowanie; traktujemy
# jak ucięcie, żeby pętla poprosiła o treść zamiast zwrócić pustkę
thinking_only = not text.strip() and any(
b.get("type") in ("thinking", "redacted_thinking") for b in blocks
)
return (text, stop == "max_tokens" or thinking_only, data.get("usage") or {},
data.get("model", self.model), stop)
def count_tokens(self, prompt: str) -> int:
"""Dokładny licznik Anthropic — nie szacunek. Od tego zależy, czy po
zmieszczeniu promptu zostanie miejsce na odpowiedź."""
if not self.api_key:
return super().count_tokens(prompt)
try:
data = _post_with_retry(
f"{self.base_url}/v1/messages/count_tokens", self._headers(),
{"model": self.model, "messages": [{"role": "user", "content": prompt}]},
min(self.timeout, 30.0),
)
return int(data.get("input_tokens") or super().count_tokens(prompt))
except LLMError:
return super().count_tokens(prompt)
def health(self) -> dict:
return {
"provider": self.name, "model": self.model, "leaves_lan": True,
"status": "ok (klucz ustawiony)" if self.api_key else "brak LLM_API_KEY",
"status": "ok (klucz ustawiony)" if self.api_key else "brak ANTHROPIC_API_KEY",
}
+52 -10
View File
@@ -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
@@ -129,8 +132,10 @@ class PromptRequest(BaseModel):
when_utc: datetime
lat: float = 0.0
lon: float = 0.0
budget: str = "medium" # concise | medium | extensive
budget: str = "medium" # concise | medium | extensive | huge | max
limit: int = 5000
provider: str | None = None # do wyliczenia budżetu „max" wg okna modelu
model: str | None = None
# tylko dla profilu period:
from_date: str | None = None
to_date: str | None = None
@@ -149,7 +154,7 @@ def chart_prompt(req: PromptRequest) -> dict:
"""
from app.engine.chart import build_chart
from app.engine.models import ChartMoment
from app.prompt import build_natal_prompt, build_period_prompt
from app.prompt import CHARS_PER_TOKEN, MAX_BUDGET, build_natal_prompt, build_period_prompt
engine = get_engine()
moment = ChartMoment(when_utc=req.when_utc, lat=req.lat, lon=req.lon)
@@ -157,6 +162,17 @@ def chart_prompt(req: PromptRequest) -> dict:
label = req.when_utc.strftime("%Y-%m-%d %H:%M UTC")
data_error = None
# Budżet „maksymalny kontekst modelu": limit znaków liczymy z okna kontekstu
# WYBRANEGO modelu, zawsze po odjęciu miejsca zarezerwowanego na odpowiedź.
budget_chars = None
if req.budget == MAX_BUDGET:
from app.llm.factory import resolve_model
from app.llm.limits import prompt_token_budget
# celowo bez build_provider(): budżet zależy TYLKO od okna kontekstu modelu,
# więc nie może wymagać skonfigurowanego klucza API
provider_name, model_name = resolve_model(req.provider, req.model)
budget_chars = int(prompt_token_budget(provider_name, model_name) * CHARS_PER_TOKEN)
# Warstwa danych dokłada wyłącznie WSKAZANIA. Wyliczenia (horoskop, oś czasu) są
# od niej niezależne — gdy padnie, prompt musi zachować wszystko, co policzyliśmy.
try:
@@ -171,7 +187,7 @@ def chart_prompt(req: PromptRequest) -> dict:
)
except httpx.HTTPError as e:
data_error = f"Warstwa danych niedostępna: {e}"
out = build_natal_prompt(chart, report, req.budget, label)
out = build_natal_prompt(chart, report, req.budget, label, budget_chars)
elif req.profile == "period":
if not (req.from_date and req.to_date):
@@ -191,7 +207,7 @@ def chart_prompt(req: PromptRequest) -> dict:
except httpx.HTTPError as e:
data_error = f"Warstwa danych niedostępna: {e}" # oś czasu zostaje
out = build_period_prompt(chart, events, req.from_date, req.to_date,
req.budget, label)
req.budget, label, budget_chars)
else:
raise HTTPException(422, f"Nieznany profil: {req.profile!r} (natal | period)")
except ValueError as e: # nieznany budżet
@@ -204,8 +220,7 @@ def chart_prompt(req: PromptRequest) -> dict:
class HoroscopeRequest(PromptRequest):
"""Jak PromptRequest + wybór dostawcy modelu (LOG-31)."""
provider: str | None = None # local (dom.) | openai | anthropic
"""Jak PromptRequest (niesie już provider i model) + limit wyjścia (LOG-31)."""
max_tokens: int | None = None
@@ -218,12 +233,26 @@ def chart_horoscope(req: HoroscopeRequest) -> dict:
(LOG-32); prezentacja ma na tej podstawie ostrzegać.
"""
from app.llm.base import LLMError
from app.llm.factory import build_provider, max_tokens
from app.llm.factory import build_provider
from app.llm.limits import plan
out = chart_prompt(req) # ten sam prompt co w podglądzie
try:
provider = build_provider(req.provider)
result = provider.generate(out["prompt"], req.max_tokens or max_tokens())
provider = build_provider(req.provider, req.model)
# Ile tokenów ma naprawdę ten prompt i ile zostaje na odpowiedź. Anthropic
# liczy dokładnie (własny endpoint), reszta szacuje — od tego zależy, czy
# w oknie kontekstu w ogóle zmieści się miejsce na horoskop.
prompt_tokens = provider.count_tokens(out["prompt"])
budget = plan(provider.name, provider.model, prompt_tokens, req.max_tokens)
out["token_plan"] = budget
if budget["warnings"]:
out["warnings"] = budget["warnings"]
if not budget["fits"]:
out["llm_error"] = " ".join(budget["warnings"])
return out
result = provider.generate(out["prompt"], budget["max_output"])
except LLMError as e:
# prompt zostaje — użytkownik moze go skopiowac i uzyc recznie
out["llm_error"] = str(e)
@@ -239,6 +268,19 @@ def chart_horoscope(req: HoroscopeRequest) -> dict:
return out
@app.get("/llm/models")
def llm_models() -> dict:
"""Podpowiedzi modeli per dostawca — UI buduje z tego listę wyboru.
To nie jest lista zamknięta: pole modelu jest tekstowe, więc można wpisać
dowolny identyfikator, do którego konto ma dostęp.
"""
from app.llm.catalog import catalog
from app.llm.factory import _DEFAULT_MODEL
return {"providers": catalog(), "defaults": dict(_DEFAULT_MODEL)}
@app.get("/llm/health")
def llm_health(provider: str | None = None) -> dict:
"""Czy model jest osiągalny i skonfigurowany (bez generowania czegokolwiek)."""
+18 -5
View File
@@ -28,8 +28,14 @@ BUDGETS: dict[str, int] = {
"concise": 4_000,
"medium": 12_000,
"extensive": 30_000,
"huge": 120_000,
# „maksymalny kontekst modelu" — wyliczany dynamicznie z okna kontekstu
# wybranego modelu, ZAWSZE po odjęciu miejsca zarezerwowanego na odpowiedź.
# Wartość poniżej jest tylko zapasem, gdy limity modelu są nieznane.
"max": 400_000,
}
DEFAULT_BUDGET = "medium"
MAX_BUDGET = "max"
MAX_EFFECT_CHARS = 320 # dłuższe opisy skracamy (krok 5 redukcji)
CHARS_PER_TOKEN = 4.0 # zgrubny szacunek tokenów do podglądu w UI
@@ -284,11 +290,16 @@ def _assemble(task: str, chart_sec: str, ind_sec: str, output: str) -> str:
return "\n\n".join([task, chart_sec, ind_sec, output, f"# ZASTRZEŻENIE\n{DISCLAIMER}"])
def _budget_chars(budget: str) -> int:
def _budget_chars(budget: str, budget_chars: int | None = None) -> int:
"""Limit znaków promptu. `budget_chars` nadpisuje tabelę — używane dla opcji
„maksymalny kontekst modelu", gdzie limit zależy od wybranego modelu i musi
być policzony po odjęciu miejsca zarezerwowanego na odpowiedź."""
if budget not in BUDGETS:
raise ValueError(
f"Nieznany budżet: {budget!r} (dostępne: {', '.join(BUDGETS)})"
)
if budget_chars and budget_chars > 0:
return budget_chars
return BUDGETS[budget]
@@ -305,9 +316,10 @@ def _finish(prompt: str, budget: str, limit: int, stats: dict, profile: str) ->
def build_natal_prompt(chart: dict, report: dict, budget: str = DEFAULT_BUDGET,
moment_label: str | None = None) -> dict:
moment_label: str | None = None,
budget_chars: int | None = None) -> dict:
"""Prompt na horoskop urodzeniowy (ekran „Interpretacje")."""
limit = _budget_chars(budget)
limit = _budget_chars(budget, budget_chars)
chart_sec = _chart_section(chart, moment_label)
fixed = len(_NATAL_TASK) + len(chart_sec) + len(_NATAL_OUTPUT) + len(DISCLAIMER) + 200
chosen, stats = reduce_indications(natal_indications(report), max(limit - fixed, 500))
@@ -316,9 +328,10 @@ def build_natal_prompt(chart: dict, report: dict, budget: str = DEFAULT_BUDGET,
def build_period_prompt(chart: dict, events: list[dict], from_date: str, to_date: str,
budget: str = DEFAULT_BUDGET, moment_label: str | None = None) -> dict:
budget: str = DEFAULT_BUDGET, moment_label: str | None = None,
budget_chars: int | None = None) -> dict:
"""Prompt na horoskop okresowy (ekran „Kalendarz")."""
limit = _budget_chars(budget)
limit = _budget_chars(budget, budget_chars)
chart_sec = _chart_section(chart, moment_label)
task = f"{_PERIOD_TASK}\nZakres prognozy: **{from_date}{to_date}**."
+2
View File
@@ -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
+91
View File
@@ -0,0 +1,91 @@
"""Okna kontekstu i planowanie budżetu tokenów (LOG-30/31).
Naczelna zasada, której pilnują te testy: **zawsze zostaje miejsce na odpowiedź**.
Prompt nigdy nie może wypełnić całego okna kontekstu, bo wtedy model kończy na
`max_tokens` z pustą albo uciętą treścią — to był zgłoszony błąd.
"""
import pytest
from app.llm import limits
def test_known_models_have_documented_limits():
ctx, out = limits.limits_for("anthropic", "claude-opus-4-8")
assert ctx == 1_000_000 and out == 128_000
ctx, out = limits.limits_for("anthropic", "claude-haiku-4-5")
assert ctx == 200_000 and out == 64_000
def test_local_model_tag_is_matched_by_prefix():
"""Modele lokalne niosą tag (`llama3.1:8b`) — dopasowanie musi to znieść."""
assert limits.limits_for("local", "llama3.1:8b") == limits.limits_for("local", "llama3.1")
def test_unknown_model_falls_back_conservatively():
ctx, out = limits.limits_for("local", "jakis-egzotyczny-model")
assert ctx == 8_192 and out == 4_096
def test_env_overrides_win(monkeypatch):
"""Modele wychodzą szybciej, niż aktualizuje się tabela."""
monkeypatch.setenv("LOCAL_CONTEXT_WINDOW", "131072")
monkeypatch.setenv("LOCAL_MAX_OUTPUT", "8192")
assert limits.limits_for("local", "llama3.1:8b") == (131_072, 8_192)
def test_env_override_ignores_garbage(monkeypatch):
monkeypatch.setenv("LOCAL_CONTEXT_WINDOW", "nie-liczba")
assert limits.limits_for("local", "llama3.1:8b")[0] == 8_192
# ------------------------------------------------- rezerwa miejsca na odpowiedź
def test_prompt_budget_always_reserves_room_for_answer():
ctx, out = limits.limits_for("anthropic", "claude-opus-4-8")
budget = limits.prompt_token_budget("anthropic", "claude-opus-4-8")
assert budget + out + limits.SAFETY_MARGIN <= ctx
assert budget > 0
def test_small_context_model_still_leaves_room():
budget = limits.prompt_token_budget("local", "llama3.1:8b")
ctx, out = limits.limits_for("local", "llama3.1:8b")
assert budget + out + limits.SAFETY_MARGIN <= ctx
def test_plan_shrinks_output_when_prompt_is_huge():
"""Duży prompt nie może dostać pełnego okna wyjścia — musi się zmieścić."""
ctx, model_out = limits.limits_for("local", "llama3.1:8b")
p = limits.plan("local", "llama3.1:8b", prompt_tokens=6_000)
assert p["max_output"] <= ctx - 6_000 - limits.SAFETY_MARGIN
assert p["max_output"] < model_out
def test_plan_reports_not_fitting_instead_of_failing_silently():
p = limits.plan("local", "llama3.1:8b", prompt_tokens=8_000)
assert p["fits"] is False
assert any("odpowied" in w for w in p["warnings"])
def test_plan_warns_above_90k_but_still_fits():
"""Powyżej 90 tys. tokenów ostrzegamy — ale wysłanie MA być nadal możliwe."""
p = limits.plan("anthropic", "claude-opus-4-8", prompt_tokens=120_000)
assert p["fits"] is True, "duży prompt nadal musi dać się wysłać"
assert p["max_output"] >= limits.MIN_OUTPUT
assert any("dużo" in w for w in p["warnings"])
def test_no_warning_below_threshold():
p = limits.plan("anthropic", "claude-opus-4-8", prompt_tokens=10_000)
assert p["warnings"] == [] and p["fits"] is True
def test_requested_output_is_capped_by_model_maximum():
p = limits.plan("anthropic", "claude-haiku-4-5", prompt_tokens=1_000, want_output=999_999)
assert p["max_output"] == 64_000
@pytest.mark.parametrize("tokens", [0, 1_000, 50_000, 200_000, 900_000])
def test_plan_never_returns_negative_output(tokens):
p = limits.plan("anthropic", "claude-opus-4-8", prompt_tokens=tokens)
assert p["max_output"] >= 0
+351
View File
@@ -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")
+260 -2
View File
@@ -3,6 +3,8 @@
Transport podstawiamy przez httpx.MockTransport, więc testy są szybkie,
deterministyczne i nic nie wychodzi na zewnątrz.
"""
import json
import httpx
import pytest
@@ -115,7 +117,7 @@ def test_anthropic_generates(monkeypatch):
def test_anthropic_requires_key():
with pytest.raises(LLMError, match="LLM_API_KEY"):
with pytest.raises(LLMError, match="ANTHROPIC_API_KEY"):
AnthropicProvider("https://api.anthropic.com", "m", "").generate("p", 10)
@@ -133,7 +135,9 @@ def test_default_provider_is_local(monkeypatch):
def test_openai_requires_key(monkeypatch):
monkeypatch.delenv("LLM_API_KEY", raising=False)
with pytest.raises(LLMError, match="LLM_API_KEY"):
monkeypatch.delenv("OPENAI_API_KEY", raising=False)
# komunikat wskazuje ZMIENNĄ DO USTAWIENIA dla tego dostawcy, nie ogólne LLM_API_KEY
with pytest.raises(LLMError, match="OPENAI_API_KEY"):
factory.build_provider("openai")
@@ -147,3 +151,257 @@ def test_env_overrides_model_and_url(monkeypatch):
monkeypatch.setenv("LLM_BASE_URL", "http://serwer:8000/v1")
p = factory.build_provider("local")
assert p.model == "moj-model" and p.base_url == "http://serwer:8000/v1"
# ---------------------------------------- konfiguracja per dostawca (regresja LOG-31)
# UI pozwala przelaczac dostawce przy kazdym zadaniu, wiec ustawienia JEDNEGO nie moga
# przeciekac na pozostalych. Wczesniej wspolne LLM_BASE_URL/LLM_MODEL kierowaly zadania
# do OpenAI na adres lokalnej Ollamy i prosily Anthropic o model llama.
def _clear(monkeypatch):
for v in ("LLM_PROVIDER", "LLM_MODEL", "LLM_BASE_URL", "LLM_API_KEY",
"LOCAL_MODEL", "LOCAL_BASE_URL", "LOCAL_API_KEY",
"OPENAI_MODEL", "OPENAI_BASE_URL", "OPENAI_API_KEY",
"ANTHROPIC_MODEL", "ANTHROPIC_BASE_URL", "ANTHROPIC_API_KEY"):
monkeypatch.delenv(v, raising=False)
def test_local_config_does_not_leak_to_cloud(monkeypatch):
"""Sedno bledu: skonfigurowany model lokalny przejmowal zadania do chmury."""
_clear(monkeypatch)
monkeypatch.setenv("LLM_PROVIDER", "local")
monkeypatch.setenv("LLM_BASE_URL", "http://ollama:11434/v1") # konfiguracja lokalnego
monkeypatch.setenv("LLM_MODEL", "llama3.1:8b")
monkeypatch.setenv("OPENAI_API_KEY", "sk-test")
local = factory.build_provider("local")
assert local.base_url == "http://ollama:11434/v1" and local.model == "llama3.1:8b"
openai = factory.build_provider("openai")
assert openai.base_url == "https://api.openai.com/v1", "zadanie do OpenAI poszloby do Ollamy"
assert openai.model == "gpt-4o-mini", "OpenAI dostalby nazwe modelu llama"
def test_provider_specific_settings_win(monkeypatch):
_clear(monkeypatch)
monkeypatch.setenv("LLM_PROVIDER", "local")
monkeypatch.setenv("ANTHROPIC_API_KEY", "sk-ant")
monkeypatch.setenv("ANTHROPIC_MODEL", "claude-opus-4-8")
p = factory.build_provider("anthropic")
assert p.model == "claude-opus-4-8" and p.api_key == "sk-ant"
def test_generic_vars_apply_only_to_default_provider(monkeypatch):
"""Zgodnosc wstecz: wspolne LLM_* konfiguruja dostawce domyslnego i tylko jego."""
_clear(monkeypatch)
monkeypatch.setenv("LLM_PROVIDER", "openai")
monkeypatch.setenv("LLM_API_KEY", "sk-generic")
monkeypatch.setenv("LLM_MODEL", "gpt-4o")
assert factory.build_provider("openai").model == "gpt-4o"
assert factory.build_provider("local").model == "llama3.1:8b" # nie dziedziczy
def test_cloud_without_key_is_rejected_clearly(monkeypatch):
_clear(monkeypatch)
monkeypatch.setenv("LLM_PROVIDER", "local")
for name in ("openai", "anthropic"):
with pytest.raises(LLMError, match=f"{name.upper()}_API_KEY"):
factory.build_provider(name)
def test_local_needs_no_key(monkeypatch):
_clear(monkeypatch)
assert factory.build_provider("local").api_key == ""
# ------------------------------------------- pusta odpowiedz modelu (cicha awaria)
# Regresja: model potrafi oddac pusta tresc (prompt zjadl caly kontekst ->
# finish_reason=length, completion_tokens=0). Wczesniej generate() zwracalo pusty
# tekst BEZ bledu, widok nic nie renderowal i uzytkownik dostawal pusta strone
# bez zadnego wyjasnienia. Pusta odpowiedz MUSI byc bledem.
def test_empty_completion_raises_instead_of_silent_blank(monkeypatch):
monkeypatch.setattr(httpx, "Client", _mock_client(lambda r: httpx.Response(200, json={
"model": "llama3.1:8b",
"choices": [{"message": {"content": ""}, "finish_reason": "length"}],
"usage": {"prompt_tokens": 8000, "completion_tokens": 0},
})))
with pytest.raises(LLMError, match="nie zwrócił żadnej treści"):
ChatCompletionsProvider("local", "http://x/v1", "m").generate("p", 2000)
def test_empty_completion_explains_context_window(monkeypatch):
"""Komunikat ma prowadzic do przyczyny, a nie tylko stwierdzac fakt."""
monkeypatch.setattr(httpx, "Client", _mock_client(lambda r: httpx.Response(200, json={
"choices": [{"message": {"content": " "}, "finish_reason": "length"}],
"usage": {"prompt_tokens": 8000, "completion_tokens": 0},
})))
with pytest.raises(LLMError) as ei:
ChatCompletionsProvider("local", "http://x/v1", "m").generate("p", 2000)
msg = str(ei.value)
assert "kontekstu" in msg and "budżet" in msg
assert "powód zakończenia: length" in msg # diagnostyka w tresci bledu
def test_whitespace_only_is_treated_as_empty(monkeypatch):
monkeypatch.setattr(httpx, "Client", _mock_client(lambda r: httpx.Response(200, json={
"choices": [{"message": {"content": "\n\n \t "}}],
})))
with pytest.raises(LLMError, match="nie zwrócił żadnej treści"):
ChatCompletionsProvider("local", "http://x/v1", "m").generate("p", 100)
def test_anthropic_empty_completion_raises(monkeypatch):
monkeypatch.setattr(httpx, "Client", _mock_client(lambda r: httpx.Response(200, json={
"model": "claude-x", "content": [], "stop_reason": "max_tokens",
"usage": {"input_tokens": 9000, "output_tokens": 0},
})))
with pytest.raises(LLMError, match="nie zwrócił żadnej treści"):
AnthropicProvider("https://api.anthropic.com", "claude-x", "klucz").generate("p", 100)
def test_normal_response_still_passes(monkeypatch):
"""Straznik nie moze psuc poprawnej odpowiedzi."""
monkeypatch.setattr(httpx, "Client", _mock_client(lambda r: httpx.Response(200, json={
"choices": [{"message": {"content": "Horoskop."}, "finish_reason": "stop"}],
"usage": {"prompt_tokens": 100, "completion_tokens": 20},
})))
assert ChatCompletionsProvider("local", "http://x/v1", "m").generate("p", 100).text == "Horoskop."
# ------------------------------------- kontynuacja: horoskop MA powstac zawsze
# Sedno wymagania: niezaleznie od objetosci promptu i limitu wyjscia, pelna tresc
# ma wrocic do uzytkownika. Model urwany na max_tokens jest proszony o dokonczenie
# w ramach tej samej rozmowy, a kawalki sa sklejane.
def _scripted(responses):
"""Transport oddajacy kolejne odpowiedzi z listy (po jednej na ture)."""
seq = list(responses)
seen = []
def handler(request):
seen.append(request)
return httpx.Response(200, json=seq.pop(0) if seq else seq_last)
seq_last = responses[-1]
return handler, seen
def _chat(text, finish):
return {"choices": [{"message": {"content": text}, "finish_reason": finish}],
"usage": {"completion_tokens": 10}}
def test_truncated_answer_is_continued_and_joined(monkeypatch):
handler, seen = _scripted([
_chat("Czesc pierwsza.", "length"),
_chat("Czesc druga. KONIEC", "stop"),
])
monkeypatch.setattr(httpx, "Client", _mock_client(handler))
out = ChatCompletionsProvider("local", "http://x/v1", "m").generate("prompt", 40000)
assert "Czesc pierwsza." in out.text and "Czesc druga." in out.text
assert "KONIEC" not in out.text # znacznik nie trafia do horoskopu
assert out.usage["turns"] == 2
def test_continuation_asks_in_same_conversation(monkeypatch):
"""Kontynuacja musi isc jako kolejna tura rozmowy, a ostatnia wiadomosc MUSI
byc od uzytkownika — Claude odrzuca prefill w turze asystenta (400)."""
handler, seen = _scripted([_chat("Poczatek", "length"), _chat("Reszta", "stop")])
monkeypatch.setattr(httpx, "Client", _mock_client(handler))
ChatCompletionsProvider("local", "http://x/v1", "m").generate("prompt", 40000)
msgs = json.loads(seen[1].content)["messages"]
assert msgs[-1]["role"] == "user", "ostatnia wiadomosc nie moze byc prefillem asystenta"
assert msgs[1]["role"] == "assistant" and "Poczatek" in msgs[1]["content"]
def test_anthropic_thinking_only_turn_is_continued(monkeypatch):
"""DOKLADNIE zgloszony objaw: cala tura poszla na myslenie, tekst pusty.
Wczesniej konczylo sie to pusta strona; teraz pytamy o tresc dalej."""
seq = [
{"content": [{"type": "thinking", "thinking": ""}], "stop_reason": "max_tokens",
"usage": {"output_tokens": 2000}},
{"content": [{"type": "text", "text": "Horoskop urodzeniowy..."}],
"stop_reason": "end_turn", "usage": {"output_tokens": 500}},
]
def handler(request):
return httpx.Response(200, json=seq.pop(0) if seq else seq[-1])
monkeypatch.setattr(httpx, "Client", _mock_client(handler))
out = AnthropicProvider("https://api.anthropic.com", "claude-opus-4-8", "k").generate("p", 40000)
assert out.text == "Horoskop urodzeniowy..."
assert out.usage["turns"] == 2
def test_anthropic_sends_thinking_config(monkeypatch):
"""Bez jawnego `thinking` Sonnet 5 wlacza myslenie sam — konfigurujemy to wprost."""
seen = []
def handler(request):
seen.append(json.loads(request.content))
return httpx.Response(200, json={"content": [{"type": "text", "text": "ok"}],
"stop_reason": "end_turn"})
monkeypatch.setattr(httpx, "Client", _mock_client(handler))
monkeypatch.delenv("ANTHROPIC_THINKING", raising=False)
AnthropicProvider("https://api.anthropic.com", "claude-opus-4-8", "k").generate("p", 5000)
assert seen[0]["thinking"] == {"type": "adaptive"}
seen.clear()
monkeypatch.setenv("ANTHROPIC_THINKING", "off")
AnthropicProvider("https://api.anthropic.com", "claude-opus-4-8", "k").generate("p", 5000)
assert seen[0]["thinking"] == {"type": "disabled"}
def test_complete_answer_does_not_loop(monkeypatch):
"""Straznik nie moze mnozyc zapytan, gdy model skonczyl normalnie."""
calls = {"n": 0}
def handler(request):
calls["n"] += 1
return httpx.Response(200, json=_chat("Gotowe.", "stop"))
monkeypatch.setattr(httpx, "Client", _mock_client(handler))
out = ChatCompletionsProvider("local", "http://x/v1", "m").generate("p", 40000)
assert out.text == "Gotowe." and calls["n"] == 1
# ------------------------------------------- wybor modelu przez uzytkownika (UI)
def test_model_from_request_wins_over_config(monkeypatch):
_clear(monkeypatch)
monkeypatch.setenv("ANTHROPIC_MODEL", "claude-opus-4-8")
monkeypatch.setenv("ANTHROPIC_API_KEY", "k")
p = factory.build_provider("anthropic", "claude-fable-5")
assert p.model == "claude-fable-5", "wybor z UI musi wygrac nad konfiguracja"
def test_blank_model_falls_back_to_configured_default(monkeypatch):
_clear(monkeypatch)
monkeypatch.setenv("ANTHROPIC_MODEL", "claude-sonnet-5")
monkeypatch.setenv("ANTHROPIC_API_KEY", "k")
assert factory.build_provider("anthropic", " ").model == "claude-sonnet-5"
def test_resolve_model_needs_no_api_key(monkeypatch):
"""Budzet promptu zalezy od okna kontekstu modelu — nie moze wymagac klucza.
Wczesniej liczenie budzetu szlo przez build_provider(), ktory bez klucza
rzuca bledem, wiec „maksymalny kontekst" cicho spadal do wartosci zapasowej.
"""
_clear(monkeypatch)
provider, model = factory.resolve_model("anthropic", "claude-haiku-4-5")
assert (provider, model) == ("anthropic", "claude-haiku-4-5")
with pytest.raises(LLMError): # samo zbudowanie nadal wymaga klucza
factory.build_provider("anthropic", "claude-haiku-4-5")
def test_max_budget_differs_between_models(monkeypatch):
"""Sedno funkcji: wieksze okno = wiekszy budzet promptu."""
from app.llm.limits import prompt_token_budget
_clear(monkeypatch)
opus = prompt_token_budget(*factory.resolve_model("anthropic", "claude-opus-4-8"))
haiku = prompt_token_budget(*factory.resolve_model("anthropic", "claude-haiku-4-5"))
local = prompt_token_budget(*factory.resolve_model("local", "llama3.1:8b"))
assert opus > haiku > local > 0
@@ -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,40 +65,37 @@ 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,
budget: str = "medium", from_date: str | None = None, to_date: str | None = None,
provider: str | None = None, model: str | None = None,
) -> dict[str, Any]:
"""Gotowy prompt do LLM z wyliczeń (LOG-29/30) — woła logic /chart/prompt."""
"""Gotowy prompt do LLM z wyliczeń (LOG-29/30) — woła logic /chart/prompt.
`provider` jest potrzebny dla budżetu maksymalny kontekst modelu": limit
znaków zależy wtedy od okna kontekstu konkretnego modelu."""
payload: dict[str, Any] = {
"profile": profile, "when_utc": when_utc_iso,
"lat": lat, "lon": lon, "budget": budget,
"lat": lat, "lon": lon, "budget": budget, "provider": provider,
"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, 60.0)) as client:
r = client.post(f"{self.base_url}/chart/prompt", json=payload)
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,
budget: str = "medium", provider: str | None = None,
budget: str = "medium", provider: str | None = None, model: str | None = None,
from_date: str | None = None, to_date: str | None = None,
) -> dict[str, Any]:
"""Napisany horoskop (LOG-31) — woła logic /chart/horoscope.
@@ -97,12 +108,17 @@ class LogicClient:
}
if provider:
payload["provider"] = provider
if model:
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)
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:
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,
@@ -113,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))
+466
View File
@@ -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 drogę: **AES-256-GCM** na ciele każdego żądania i odpowiedzi.
GCM daje jednocześnie poufność i uwierzytelnienie cudzy albo podmieniony bajt
nie odszyfruje się w ogóle, więc nie ma osobnego problemu zaszyfrowane, ale
podatne na modyfikację".
**Dwa niezależne klucze**, po jednym na parę rozmówców:
* ``LINK_KEY_PRESENTATION_LOGIC`` prezentacja logika,
* ``LINK_KEY_LOGIC_DATA`` logika dane.
Dzięki temu przejęcie klucza prezentacji nie daje dostępu do warstwy danych,
gdzie leżą całe bazy. Logika trzyma oba, bo rozmawia w obie strony.
Z każdego klucza łącza wyprowadzamy **osobne podklucze na kierunek** (HKDF).
Żądanie i odpowiedź nigdy nie szyfrują się tym samym kluczem, więc powtórzenie
losowej jednorazówki w jedną stronę nie osłabia drugiej.
Format ramki (bo strumień odpowiedzi może iść kawałkami patrz okno postępu):
[4 bajty długości][magia "AL1"][12 bajtów jednorazówki][szyfrogram + znacznik]
Do materiału uwierzytelnianego (AAD) wchodzą kierunek, ścieżka, znacznik czasu
i numer ramki. Skutek: ramki nie da się przekleić do innego endpointu, odtworzyć
po czasie (dopuszczalny poślizg ``MAX_SKEW``) ani przestawić w strumieniu.
Bez ustawionego klucza moduł **przepuszcza ruch otwartym tekstem** (dev, zgodność
wstecz) i krzyczy o tym przy starcie. Gdy klucz JEST ustawiony, warstwa serwerowa
działa fail-closed: nieszyfrowane żądanie dostaje odmowę, żeby przypadkowa
regresja po stronie klienta nie oznaczała cichego powrotu do jawnego ruchu.
"""
from __future__ import annotations
import base64
import binascii
import logging
import os
import struct
import time
from typing import Iterable, Iterator
from cryptography.exceptions import InvalidTag
from cryptography.hazmat.primitives import hashes
from cryptography.hazmat.primitives.ciphers.aead import AESGCM
from cryptography.hazmat.primitives.kdf.hkdf import HKDF
log = logging.getLogger("astrololo.link")
MAGIC = b"AL1"
VERSION = "v1"
NONCE_BYTES = 12
KEY_BYTES = 32 # AES-256
LENGTH_PREFIX = 4
MAX_FRAME = 64 * 1024 * 1024 # zapora przed alokacją z podanej długości
MAX_SKEW_SECONDS = 300.0
HEADER_ENC = "X-Astrololo-Enc"
HEADER_TS = "X-Astrololo-Enc-Ts"
CONTENT_TYPE = "application/vnd.astrololo.enc"
ENV_PRESENTATION_LOGIC = "LINK_KEY_PRESENTATION_LOGIC"
ENV_LOGIC_DATA = "LINK_KEY_LOGIC_DATA"
ENV_REQUIRED = "LINK_ENCRYPTION_REQUIRED"
REQUEST, RESPONSE = b"req", b"res"
# Sondy k8s pukają tu bez klucza i tak ma zostać — inaczej pierwsza literówka
# w sekrecie kładłaby pody zamiast pokazać błąd w aplikacji.
PUBLIC_PATHS = frozenset({"/health"})
class LinkError(Exception):
"""Cokolwiek poszło nie tak z kopertą — celowo bez szczegółów na zewnątrz."""
# --------------------------------------------------------------------- klucze
def parse_key(raw: str) -> bytes:
"""Klucz z konfiguracji: hex (64 znaki) albo base64. Zawsze 32 bajty."""
text = raw.strip()
if not text:
raise LinkError("pusty klucz łącza")
try:
key = bytes.fromhex(text)
except ValueError:
try:
key = base64.b64decode(text, validate=True)
except (binascii.Error, ValueError) as exc:
raise LinkError("klucz łącza nie jest ani hexem, ani base64") from exc
if len(key) != KEY_BYTES:
raise LinkError(
f"klucz łącza ma {len(key)} B zamiast {KEY_BYTES} — wygeneruj przez "
f"`openssl rand -hex 32`"
)
return key
def key_from_env(env_name: str) -> bytes | None:
"""Klucz albo None. Zły klucz to wyjątek OD RAZU — nie przy pierwszym żądaniu."""
raw = os.getenv(env_name, "")
return parse_key(raw) if raw.strip() else None
def encryption_required() -> bool:
"""Czy brak klucza ma być błędem, a nie cichym powrotem do jawnego ruchu.
Serwer sam z siebie broni się fail-closed, ale to za mało: klient BEZ klucza
wysyła pytanie otwartym tekstem i dopiero potem dostaje odmowę czyli treść
zapytania zdążyła już przelecieć przez sieć. Ta flaga zatrzymuje go, zanim
cokolwiek opuści proces. Ustawiana razem z kluczami we wdrożeniu.
"""
return os.getenv(ENV_REQUIRED, "").strip().lower() in {"1", "true", "yes", "on"}
def _subkey(link_key: bytes, direction: bytes) -> bytes:
return HKDF(
algorithm=hashes.SHA256(), length=KEY_BYTES, salt=None,
info=b"astrololo/link/" + direction,
).derive(link_key)
class Link:
"""Jedna para rozmówców: klucz plus wyprowadzone z niego podklucze."""
def __init__(self, link_key: bytes) -> None:
self._by_direction = {
REQUEST: AESGCM(_subkey(link_key, REQUEST)),
RESPONSE: AESGCM(_subkey(link_key, RESPONSE)),
}
# ---------------------------------------------------------- pojedyncza ramka
def _aad(self, direction: bytes, path: str, stamp: str, seq: int) -> bytes:
return b"|".join([MAGIC, direction, path.encode("utf-8"),
stamp.encode("ascii"), str(seq).encode("ascii")])
def seal(self, direction: bytes, path: str, stamp: str, seq: int,
plaintext: bytes) -> bytes:
nonce = os.urandom(NONCE_BYTES)
sealed = self._by_direction[direction].encrypt(
nonce, plaintext, self._aad(direction, path, stamp, seq))
return MAGIC + nonce + sealed
def open(self, direction: bytes, path: str, stamp: str, seq: int,
frame: bytes) -> bytes:
if not frame.startswith(MAGIC):
raise LinkError("ramka bez znacznika astrololo")
body = frame[len(MAGIC):]
if len(body) <= NONCE_BYTES:
raise LinkError("ramka za krótka")
nonce, sealed = body[:NONCE_BYTES], body[NONCE_BYTES:]
try:
return self._by_direction[direction].decrypt(
nonce, sealed, self._aad(direction, path, stamp, seq))
except InvalidTag as exc:
# Jeden komunikat na wszystkie przypadki: zły klucz, podmieniony bajt,
# przeklejenie z innej ścieżki, przestawiona ramka. Rozróżnianie ich
# na zewnątrz podpowiadałoby atakującemu, w co trafił.
raise LinkError("nie udało się odszyfrować — zły klucz albo naruszone dane") from exc
# ------------------------------------------------------------ strumień ramek
def seal_stream(self, direction: bytes, path: str, stamp: str,
chunks: Iterable[bytes]) -> Iterator[bytes]:
for seq, chunk in enumerate(chunks):
yield frame_out(self.seal(direction, path, stamp, seq, chunk))
def open_stream(self, direction: bytes, path: str, stamp: str,
raw: bytes) -> Iterator[bytes]:
for seq, frame in enumerate(frames_in(raw)):
yield self.open(direction, path, stamp, seq, frame)
def open_all(self, direction: bytes, path: str, stamp: str, raw: bytes) -> bytes:
return b"".join(self.open_stream(direction, path, stamp, raw))
# ---------------------------------------------------------------- ramkowanie
def frame_out(payload: bytes) -> bytes:
return struct.pack(">I", len(payload)) + payload
def frames_in(raw: bytes) -> Iterator[bytes]:
"""Rozbiera bufor na ramki. Ucięty strumień to błąd, nie cicha strata danych."""
offset = 0
while offset < len(raw):
if offset + LENGTH_PREFIX > len(raw):
raise LinkError("urwana ramka (brak nagłówka długości)")
(size,) = struct.unpack(">I", raw[offset:offset + LENGTH_PREFIX])
if size > MAX_FRAME:
raise LinkError("ramka ponad dopuszczalny rozmiar")
offset += LENGTH_PREFIX
if offset + size > len(raw):
raise LinkError("urwana ramka (za mało danych)")
yield raw[offset:offset + size]
offset += size
def unframe_incremental(buffer: bytearray) -> Iterator[bytes]:
"""Wyjmuje z bufora KOMPLETNE ramki i zjada je; resztę zostawia na później.
Dla odbioru na żywo: kawałki przychodzą podzielone dowolnie i ramka potrafi
rozjechać się między dwa odczyty.
"""
while True:
if len(buffer) < LENGTH_PREFIX:
return
(size,) = struct.unpack(">I", buffer[:LENGTH_PREFIX])
if size > MAX_FRAME:
raise LinkError("ramka ponad dopuszczalny rozmiar")
if len(buffer) < LENGTH_PREFIX + size:
return
frame = bytes(buffer[LENGTH_PREFIX:LENGTH_PREFIX + size])
del buffer[:LENGTH_PREFIX + size]
yield frame
# ------------------------------------------------------------- świeżość ruchu
def stamp_now() -> str:
return f"{time.time():.3f}"
def check_stamp(stamp: str) -> None:
"""Odrzuca ramki spoza okna czasowego — inaczej podsłuchane żądanie dałoby się
odtworzyć w dowolnym momencie w przyszłości."""
try:
sent = float(stamp)
except (TypeError, ValueError) as exc:
raise LinkError("brak albo błędny znacznik czasu") from exc
if abs(time.time() - sent) > MAX_SKEW_SECONDS:
raise LinkError("znacznik czasu poza dopuszczalnym oknem")
# =========================================================== strona serwerowa
class LinkCryptoMiddleware:
"""Rozszyfrowuje wchodzące żądania i zaszyfrowuje wychodzące odpowiedzi.
Napisane jako czyste ASGI, nie ``@app.middleware("http")``, bo trzeba
podmienić CIAŁO żądania jeszcze zanim zobaczy je FastAPI, oraz przepuścić
odpowiedź strumieniową kawałek po kawałku, bez zbierania jej w pamięci.
"""
def __init__(self, app, link: Link | None, layer: str) -> None:
self.app = app
self.link = link
self.layer = layer
async def __call__(self, scope, receive, send):
if scope["type"] != "http" or self.link is None or scope["path"] in PUBLIC_PATHS:
return await self.app(scope, receive, send)
path = scope["path"]
headers = {k.decode("latin-1").lower(): v.decode("latin-1") for k, v in scope["headers"]}
if headers.get(HEADER_ENC.lower()) != VERSION:
# Fail-closed. Klucz jest ustawiony, więc jawne żądanie oznacza albo
# pomyłkę w konfiguracji, albo kogoś obcego — w obu wypadkach nie
# chcemy po cichu wrócić do jawnego ruchu.
log.warning("warstwa %s: odrzucone żądanie bez szyfrowania łącza (%s)",
self.layer, path)
return await _refuse(send, "Łącze międzywarstwowe wymaga szyfrowania.")
stamp = headers.get(HEADER_TS.lower(), "")
try:
check_stamp(stamp)
plaintext = self.link.open_all(REQUEST, path, stamp, await _read_body(receive))
except LinkError as exc:
log.warning("warstwa %s: %s (%s)", self.layer, exc, path)
return await _refuse(send, "Nie udało się odczytać zaszyfrowanego żądania.")
scope = dict(scope)
scope["headers"] = _rewritten_headers(scope["headers"], len(plaintext))
await self.app(scope, _replay(plaintext, receive), self._sealing_send(send, path))
def _sealing_send(self, send, path: str):
state: dict = {"stamp": "", "seq": 0}
async def sealing(message):
if message["type"] == "http.response.start":
state["stamp"] = stamp_now()
keep = [(k, v) for k, v in message.get("headers", [])
if k.lower() not in (b"content-length", b"content-type")]
message = dict(message)
message["headers"] = keep + [
(b"content-type", CONTENT_TYPE.encode()),
(HEADER_ENC.lower().encode(), VERSION.encode()),
(HEADER_TS.lower().encode(), state["stamp"].encode()),
]
return await send(message)
if message["type"] == "http.response.body":
chunk = message.get("body", b"")
sealed = b""
if chunk:
sealed = frame_out(self.link.seal(
RESPONSE, path, state["stamp"], state["seq"], chunk))
state["seq"] += 1
return await send({"type": "http.response.body", "body": sealed,
"more_body": message.get("more_body", False)})
return await send(message)
return sealing
def _rewritten_headers(raw: Iterable[tuple[bytes, bytes]], length: int):
"""Po odszyfrowaniu ciało ma inną długość i zwykły typ — inaczej FastAPI
próbowałby sparsować JSON o cudzej deklarowanej wielkości."""
kept = [(k, v) for k, v in raw if k.lower() not in (b"content-length", b"content-type")]
kept.append((b"content-length", str(length).encode()))
if length:
kept.append((b"content-type", b"application/json"))
return kept
async def _read_body(receive) -> bytes:
body = bytearray()
while True:
message = await receive()
if message["type"] == "http.disconnect":
raise LinkError("rozłączenie w trakcie odbioru żądania")
body += message.get("body", b"")
if not message.get("more_body", False):
return bytes(body)
def _replay(body: bytes, original):
"""Podstawia odszyfrowane ciało jako jedyną porcję wejścia dla aplikacji.
Po oddaniu ciała oddajemy głos ORYGINALNEMU `receive`, zamiast od razu
zgłaszać rozłączenie. Odpowiedź strumieniowa nasłuchuje bowiem rozłączenia
równolegle do wysyłania i przerywa się, gdy je zobaczy na skróconej wersji
okno postępu dostawało pustą odpowiedź, choć zwykłe żądania działały.
"""
delivered = False
async def receive():
nonlocal delivered
if delivered:
return await original()
delivered = True
return {"type": "http.request", "body": body, "more_body": False}
return receive
async def _refuse(send, detail: str) -> None:
"""Odmowa leci JAWNIE — rozmówca właśnie pokazał, że nie umie odszyfrować,
więc zaszyfrowany komunikat o błędzie byłby dla niego nieczytelny."""
payload = f'{{"detail":"{detail}"}}'.encode("utf-8")
await send({"type": "http.response.start", "status": 400, "headers": [
(b"content-type", b"application/json"),
(b"content-length", str(len(payload)).encode()),
]})
await send({"type": "http.response.body", "body": payload})
def install(app, env_name: str, layer: str):
"""Podpina szyfrowanie łącza. Wołać PO `security.install`, żeby także odmowa
tokenowa (401) wracała zaszyfrowana inaczej klient by jej nie odczytał."""
link_key = key_from_env(env_name)
if link_key is None and encryption_required():
# Celowo wywracamy start. Ta sama zasada co przy sekrecie logowania:
# wolimy widoczną awarię niż usługę, która wstała i po cichu nie chroni
# niczego. Pod w CrashLoop widać od razu, jawny ruch — nie.
raise LinkError(
f"{ENV_REQUIRED} jest włączone, ale {env_name} nie ustawiony — "
f"warstwa {layer} nie wystartuje bez klucza łącza"
)
if link_key is None:
log.warning(
"UWAGA: %s nie ustawiony — warstwa %s rozmawia z sąsiadem JAWNYM tekstem, "
"więc treść baz interpretacyjnych jest widoczna dla każdego, kto podsłucha "
"ruch wewnątrz sieci.", env_name, layer,
)
return None
link = Link(link_key)
app.add_middleware(LinkCryptoMiddleware, link=link, layer=layer)
log.info("warstwa %s: łącze szyfrowane (AES-256-GCM, klucz z %s)", layer, env_name)
return link
# ============================================================ strona kliencka
def call(client, method: str, url: str, *, payload=None,
headers: dict[str, str] | None = None, link: Link | None) -> bytes:
"""Żądanie do sąsiedniej warstwy; zwraca odszyfrowane ciało odpowiedzi.
Ścieżkę do materiału uwierzytelnianego bierzemy Z URL-a, a nie z osobnego
argumentu gdyby klient i serwer liczyły inaczej, każde żądanie kończyłoby
się niejasnym błędem odszyfrowania.
"""
import json as _json
import httpx
request_headers = dict(headers or {})
if link is None:
if encryption_required():
# Zatrzymujemy się PRZED wysłaniem. Gdyby polecieć jawnie i dopiero
# zebrać odmowę, pytanie byłoby już na kablu — a to właśnie ono niesie
# sygnifikatory, o które pytamy bazę.
raise LinkError(
f"{ENV_REQUIRED} jest włączone, ale brak klucza łącza — żądanie "
f"NIE zostało wysłane, żeby jego treść nie poszła jawnym tekstem"
)
response = client.request(method, url, json=payload, headers=request_headers)
response.raise_for_status()
return response.content
path = httpx.URL(url).path
stamp = stamp_now()
plaintext = b"" if payload is None else _json.dumps(payload).encode("utf-8")
body = frame_out(link.seal(REQUEST, path, stamp, 0, plaintext))
request_headers.update({HEADER_ENC: VERSION, HEADER_TS: stamp,
"Content-Type": CONTENT_TYPE})
response = client.request(method, url, content=body, headers=request_headers)
if response.status_code >= 400 and response.headers.get(HEADER_ENC) != VERSION:
log.error("łącze %s odmówiło: %s", path, response.text[:200])
response.raise_for_status()
if response.headers.get(HEADER_ENC) != VERSION:
raise LinkError("odpowiedź przyszła nieszyfrowana, choć klucz łącza jest ustawiony")
reply_stamp = response.headers.get(HEADER_TS, "")
check_stamp(reply_stamp)
return link.open_all(RESPONSE, path, reply_stamp, response.content)
def call_json(client, method: str, url: str, *, payload=None,
headers: dict[str, str] | None = None, link: Link | None):
import json as _json
return _json.loads(call(client, method, url, payload=payload,
headers=headers, link=link))
def open_response_stream(response, link: Link | None) -> Iterator[bytes]:
"""Odbiór odpowiedzi płynącej kawałkami (okno postępu).
Ramka potrafi rozjechać się między dwa odczyty z gniazda, więc składamy
w buforze zamiast zakładać, że każdy kawałek to komplet.
"""
if link is None:
yield from response.iter_bytes()
return
if response.headers.get(HEADER_ENC) != VERSION:
raise LinkError("strumień przyszedł nieszyfrowany, choć klucz łącza jest ustawiony")
stamp = response.headers.get(HEADER_TS, "")
check_stamp(stamp)
path = response.request.url.path
buffer = bytearray()
seq = 0
for chunk in response.iter_bytes():
buffer += chunk
for frame in unframe_incremental(buffer):
yield link.open(RESPONSE, path, stamp, seq, frame)
seq += 1
if buffer:
raise LinkError("strumień urwał się w połowie ramki")
+27 -9
View File
@@ -39,6 +39,15 @@ def _build_utc(date: str, time: str, tz_offset: float) -> tuple[str, str]:
return utc.isoformat(), label
def _llm_catalog() -> dict:
"""Podpowiedzi modeli dla pola wyboru. Awaria logiki nie może wywrócić strony —
pole modelu jest tekstowe, więc bez katalogu nadal da się wpisać model ręcznie."""
try:
return logic.llm_models()
except httpx.HTTPError:
return {"providers": {}, "defaults": {}}
def _logic_error(e: Exception) -> str:
if isinstance(e, httpx.HTTPStatusError) and e.response.status_code == 404:
return (
@@ -53,7 +62,8 @@ def _logic_error(e: Exception) -> str:
def chart_form(request: Request):
return templates.TemplateResponse(
request, "chart.html",
{"result": None, "form": default_form(), "location_label": DEFAULT_LOCATION_LABEL},
{"result": None, "form": default_form(), "location_label": DEFAULT_LOCATION_LABEL,
"llm_catalog": _llm_catalog()},
)
@@ -115,7 +125,8 @@ def significators_search(
def interpret_form(request: Request):
return templates.TemplateResponse(
request, "interpret.html",
{"result": None, "form": default_form(), "location_label": DEFAULT_LOCATION_LABEL},
{"result": None, "form": default_form(), "location_label": DEFAULT_LOCATION_LABEL,
"llm_catalog": _llm_catalog()},
)
@@ -131,22 +142,25 @@ def interpret_run(
action: str = Form("report"),
prompt_budget: str = Form("medium"),
llm_provider: str = Form("local"),
llm_model: str = Form(""),
):
form = {"date": date, "time": time, "tz_offset": tz_offset,
"lat": lat, "lon": lon, "group": group, "prompt_budget": prompt_budget,
"llm_provider": llm_provider}
ctx: dict = {"form": form, "result": None, "error": None, "moment": None}
"llm_provider": llm_provider, "llm_model": llm_model}
ctx: dict = {"form": form, "result": None, "error": None, "moment": None,
"llm_catalog": _llm_catalog()}
try:
iso_utc, label = _build_utc(date, time, tz_offset)
ctx["moment"] = label
if action == "prompt":
ctx["prompt_result"] = logic.prompt(
profile="natal", when_utc_iso=iso_utc, lat=lat, lon=lon, budget=prompt_budget,
provider=llm_provider, model=llm_model,
)
elif action == "horoscope":
ctx["prompt_result"] = logic.horoscope(
profile="natal", when_utc_iso=iso_utc, lat=lat, lon=lon,
budget=prompt_budget, provider=llm_provider,
budget=prompt_budget, provider=llm_provider, model=llm_model,
)
else:
ctx["result"] = logic.report(when_utc_iso=iso_utc, lat=lat, lon=lon, group=group)
@@ -162,7 +176,8 @@ def interpret_run(
def timeline_form(request: Request):
return templates.TemplateResponse(
request, "timeline.html",
{"result": None, "form": default_form(), "location_label": DEFAULT_LOCATION_LABEL},
{"result": None, "form": default_form(), "location_label": DEFAULT_LOCATION_LABEL,
"llm_catalog": _llm_catalog()},
)
@@ -179,11 +194,13 @@ def timeline_run(
action: str = Form("timeline"),
prompt_budget: str = Form("medium"),
llm_provider: str = Form("local"),
llm_model: str = Form(""),
):
form = {"date": date, "time": time, "tz_offset": tz_offset, "lat": lat, "lon": lon,
"from_date": from_date, "to_date": to_date, "prompt_budget": prompt_budget,
"llm_provider": llm_provider}
ctx: dict = {"form": form, "result": None, "error": None, "moment": None}
"llm_provider": llm_provider, "llm_model": llm_model}
ctx: dict = {"form": form, "result": None, "error": None, "moment": None,
"llm_catalog": _llm_catalog()}
try:
iso_utc, label = _build_utc(date, time, tz_offset)
ctx["moment"] = label
@@ -191,11 +208,12 @@ def timeline_run(
ctx["prompt_result"] = logic.prompt(
profile="period", when_utc_iso=iso_utc, lat=lat, lon=lon,
budget=prompt_budget, from_date=from_date, to_date=to_date,
provider=llm_provider, model=llm_model,
)
elif action == "horoscope":
ctx["prompt_result"] = logic.horoscope(
profile="period", when_utc_iso=iso_utc, lat=lat, lon=lon,
budget=prompt_budget, provider=llm_provider,
budget=prompt_budget, provider=llm_provider, model=llm_model,
from_date=from_date, to_date=to_date,
)
else:
+34 -3
View File
@@ -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"},
@@ -0,0 +1,58 @@
// Podpowiedzi modeli zależne od wybranego dostawcy.
//
// Pole modelu jest CELOWO tekstowe (input + datalist), a nie zamkniętym <select>:
// katalog to tylko wygoda, a konto może mieć dostęp do modeli, o których kod nie
// wie. Puste pole = model domyślny dostawcy.
document.addEventListener('DOMContentLoaded', function () {
const providerEl = document.getElementById('llmProvider');
const modelEl = document.getElementById('llmModel');
const listEl = document.getElementById('llmModelList');
const hintEl = document.getElementById('llmModelHint');
const dataEl = document.getElementById('llmCatalog');
if (!providerEl || !modelEl || !listEl || !dataEl) return;
let catalog = {};
let defaults = {};
try {
const parsed = JSON.parse(dataEl.textContent || '{}');
catalog = parsed.providers || {};
defaults = parsed.defaults || {};
} catch (e) {
return; // brak katalogu — pole nadal działa jako wolny tekst
}
const fmt = n =>
n >= 1000000 ? (n / 1000000) + 'M' : (n >= 1000 ? Math.round(n / 1000) + 'k' : String(n));
function refresh(resetValue) {
const provider = providerEl.value;
const models = catalog[provider] || [];
listEl.innerHTML = '';
models.forEach(m => {
const opt = document.createElement('option');
opt.value = m.id;
opt.label = m.label + ' · kontekst ' + fmt(m.context_window);
listEl.appendChild(opt);
});
const fallback = defaults[provider] || '';
modelEl.placeholder = fallback ? 'domyślny: ' + fallback : 'domyślny dostawcy';
// po zmianie dostawcy stary model nie ma sensu (np. gpt-4o u Anthropica)
if (resetValue) modelEl.value = '';
if (hintEl) {
const chosen = models.find(m => m.id === modelEl.value.trim());
hintEl.textContent = chosen
? chosen.label + ' — okno kontekstu ' + fmt(chosen.context_window) +
' tokenów, maksymalna odpowiedź ' + fmt(chosen.max_output) + '.'
: 'Zostaw puste, by użyć modelu domyślnego. Możesz też wpisać dowolny ' +
'identyfikator modelu, do którego Twoje konto ma dostęp.';
}
}
providerEl.addEventListener('change', () => refresh(true));
modelEl.addEventListener('input', () => refresh(false));
refresh(false);
});
@@ -1,3 +1,6 @@
{# Katalog modeli wstrzykiwany przez handler — pole jest tekstowe, więc
można wpisać dowolny model, do którego konto ma dostęp. #}
<script type="application/json" id="llmCatalog">{{ llm_catalog | tojson }}</script>
{# Wspólny blok: sterowanie generowaniem + prompt (LOG-29/30) + horoskop (LOG-31).
Używany na ekranach Interpretacje i Kalendarz. #}
<div class="opts">
@@ -7,17 +10,26 @@
<option value="concise" {{ 'selected' if pb == 'concise' else '' }}>zwięzły (~4 tys. znaków)</option>
<option value="medium" {{ 'selected' if pb == 'medium' else '' }}>średni (~12 tys.)</option>
<option value="extensive" {{ 'selected' if pb == 'extensive' else '' }}>obszerny (~30 tys.)</option>
<option value="huge" {{ 'selected' if pb == 'huge' else '' }}>bardzo obszerny (~120 tys.)</option>
<option value="max" {{ 'selected' if pb == 'max' else '' }}>maksymalny kontekst modelu</option>
</select>
</label>
<label>Model
<select name="llm_provider">
<label>Dostawca
<select name="llm_provider" id="llmProvider">
{% set lp = form.llm_provider or 'local' %}
<option value="local" {{ 'selected' if lp == 'local' else '' }}>lokalny — nic nie opuszcza sieci</option>
<option value="anthropic" {{ 'selected' if lp == 'anthropic' else '' }}>Anthropic — dane wychodzą</option>
<option value="openai" {{ 'selected' if lp == 'openai' else '' }}>OpenAI — dane wychodzą</option>
</select>
</label>
<label>Model
<input type="text" name="llm_model" id="llmModel" list="llmModelList"
value="{{ form.llm_model or '' }}" placeholder="domyślny dostawcy"
autocomplete="off">
<datalist id="llmModelList"></datalist>
</label>
</div>
<p class="muted small" id="llmModelHint"></p>
<p class="muted small">
Prompt zawiera <strong>oryginalne opisy z baz</strong> oraz dane urodzeniowe. Model lokalny
przetwarza je u nas; wybór dostawcy w chmurze oznacza, że ta treść <strong>opuszcza naszą
@@ -51,6 +63,23 @@
</p>
{% endif %}
{% if prompt_result.warnings %}
{% for w in prompt_result.warnings %}
<p class="muted small"><strong>Uwaga:</strong> {{ w }}</p>
{% endfor %}
{% endif %}
{% if prompt_result.token_plan %}
{% set tp = prompt_result.token_plan %}
<p class="muted small">
Tokeny: prompt {{ tp.prompt_tokens }} · okno modelu {{ tp.context_window }} ·
zarezerwowane na odpowiedź {{ tp.max_output }}
{% if prompt_result.usage and prompt_result.usage.turns and prompt_result.usage.turns > 1 %}
· odpowiedź złożona z {{ prompt_result.usage.turns }} tur (model dokańczał urwany tekst)
{% endif %}
</p>
{% endif %}
{% if prompt_result.llm_error %}
<div class="error">
Nie udało się napisać horoskopu: {{ prompt_result.llm_error }}<br>
@@ -90,4 +90,5 @@
<script src="/static/now.js"></script>
<script src="/static/copy.js"></script>
<script src="/static/models.js"></script>
{% endblock %}
@@ -80,4 +80,5 @@
<script src="/static/now.js"></script>
<script src="/static/copy.js"></script>
<script src="/static/models.js"></script>
{% endblock %}
+2
View File
@@ -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
@@ -0,0 +1,155 @@
"""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.
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, 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
target = node.func.value
has_headers = any(kw.arg == "headers" for kw in node.keywords)
has_link = any(kw.arg == "link" for kw in node.keywords)
if (node.func.attr in HTTP_VERBS
and isinstance(target, ast.Name) and target.id == "client"):
# Surowe wywołanie httpx — omija szyfrowanie, więc dopuszczalne tylko
# dla ścieżek wyjętych spod niego (patrz link_crypto.PUBLIC_PATHS).
out.append((f"client.{node.func.attr}()", node.lineno, has_headers, True))
elif (node.func.attr in LINK_CALLS
and isinstance(target, ast.Name) and target.id == "link_crypto"):
out.append((f"link_crypto.{node.func.attr}()", node.lineno, has_headers, has_link))
return out
def test_client_module_exists():
assert CLIENT.is_file()
def test_every_outbound_call_sends_auth_header():
calls = _http_calls(CLIENT)
assert calls, "nie znaleziono żadnego wywołania HTTP — test przestał cokolwiek pilnować"
missing = [f"{CLIENT.name}:{line} {what}" for what, line, ok, _ in calls if not ok]
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."""
import os
from app.clients.logic_client import _auth_headers
old = os.environ.get("INTERNAL_TOKEN")
try:
os.environ["INTERNAL_TOKEN"] = "abc"
assert _auth_headers() == {"X-Astrololo-Token": "abc"}
os.environ.pop("INTERNAL_TOKEN")
assert _auth_headers() == {} # ochrona wyłączona = brak nagłówka
finally:
if old is not None:
os.environ["INTERNAL_TOKEN"] = old
else:
os.environ.pop("INTERNAL_TOKEN", None)
# --------------------------------------------- wybor modelu musi dojsc do logiki
# Ta sama klasa bledu co przy tokenie: dokladajac nowa sciezke latwo zapomniec
# przekazac parametr, a objaw (cichy powrot do modelu domyslnego) jest niewidoczny.
def _calls_to(path_fragment: str) -> list[int]:
"""Linie wywolan logic.<metoda>(...) w handlerach prezentacji."""
main = CLIENT.parent.parent / "main.py"
tree = ast.parse(main.read_text(encoding="utf-8"))
out = []
for node in ast.walk(tree):
if (isinstance(node, ast.Call) and isinstance(node.func, ast.Attribute)
and node.func.attr == path_fragment
and isinstance(node.func.value, ast.Name) and node.func.value.id == "logic"):
out.append(node.lineno)
return out
def _has_kwarg(main_src: str, lineno: int, name: str) -> bool:
tree = ast.parse(main_src)
for node in ast.walk(tree):
if isinstance(node, ast.Call) and node.lineno == lineno:
return any(kw.arg == name for kw in node.keywords)
return False
def test_every_llm_call_passes_selected_model():
main = CLIENT.parent.parent / "main.py"
src = main.read_text(encoding="utf-8")
missing = []
for method in ("prompt", "horoscope"):
for line in _calls_to(method):
if not _has_kwarg(src, line, "model"):
missing.append(f"main.py:{line} logic.{method}()")
assert not missing, (
"Wywołania bez wybranego modelu — po cichu użyją domyślnego: " + ", ".join(missing)
)
def test_every_llm_call_passes_provider():
main = CLIENT.parent.parent / "main.py"
src = main.read_text(encoding="utf-8")
missing = []
for method in ("prompt", "horoscope"):
for line in _calls_to(method):
if not _has_kwarg(src, line, "provider"):
missing.append(f"main.py:{line} logic.{method}()")
assert not missing, "Wywołania bez dostawcy: " + ", ".join(missing)
@@ -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"