Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 114b7eebdf |
@@ -1,296 +0,0 @@
|
||||
# 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.
|
||||
@@ -1,466 +0,0 @@
|
||||
"""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")
|
||||
@@ -10,7 +10,7 @@ from contextlib import asynccontextmanager
|
||||
|
||||
from fastapi import FastAPI
|
||||
|
||||
from app import link_crypto, security
|
||||
from app import security
|
||||
from app.config import settings
|
||||
from app.models import HealthInfo, SearchQuery, SearchResult
|
||||
from app.providers.factory import build_provider
|
||||
@@ -26,9 +26,6 @@ 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)
|
||||
|
||||
@@ -7,5 +7,3 @@ 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
|
||||
|
||||
@@ -10,7 +10,6 @@ from typing import Any
|
||||
|
||||
import httpx
|
||||
|
||||
from app import link_crypto
|
||||
from app.config import settings
|
||||
|
||||
|
||||
@@ -20,13 +19,6 @@ 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("/")
|
||||
@@ -41,13 +33,11 @@ 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:
|
||||
return link_crypto.call_json(client, "POST", f"{self.base_url}/search",
|
||||
payload=payload, headers=_auth_headers(),
|
||||
link=_link())
|
||||
r = client.post(f"{self.base_url}/search", json=payload, headers=_auth_headers())
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
|
||||
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()
|
||||
|
||||
@@ -1,466 +0,0 @@
|
||||
"""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")
|
||||
@@ -124,7 +124,14 @@ class _Driver:
|
||||
"""(tekst, czy_ucięta, zużycie, nazwa_modelu, powód_zakończenia)."""
|
||||
raise NotImplementedError
|
||||
|
||||
def generate(self, prompt: str, max_tokens: int) -> Completion:
|
||||
def generate(self, prompt: str, max_tokens: int, on_event=None) -> Completion:
|
||||
"""`on_event(dict)` dostaje zdarzenia postępu — UI pokazuje z nich log.
|
||||
Raportujemy KAŻDĄ turę, bo to ona trwa; bez tego pasek postępu byłby
|
||||
ozdobnikiem, a nie informacją."""
|
||||
def emit(kind: str, message: str, **extra):
|
||||
if on_event:
|
||||
on_event({"type": kind, "message": message, **extra})
|
||||
|
||||
messages: list[dict] = [{"role": "user", "content": prompt}]
|
||||
parts: list[str] = []
|
||||
usage: dict = {}
|
||||
@@ -137,9 +144,29 @@ class _Driver:
|
||||
while turns <= MAX_CONTINUATIONS:
|
||||
turns += 1
|
||||
budget = max(256, min(remaining, TURN_TOKENS_CAP))
|
||||
emit("turn_start",
|
||||
f"Tura {turns}: wysyłam do modelu {self.model} (limit {budget} tokenów)…",
|
||||
turn=turns)
|
||||
started = time.monotonic()
|
||||
text, truncated, turn_usage, model_name, stop = self._turn(messages, budget)
|
||||
took = time.monotonic() - started
|
||||
_merge_usage(usage, turn_usage)
|
||||
remaining -= budget
|
||||
|
||||
# Odejmujemy tokeny FAKTYCZNIE wyprodukowane, nie zamówiony limit tury.
|
||||
# Inaczej pierwsza tura zjadałaby cały budżet i urwana odpowiedź nigdy
|
||||
# nie doczekałaby się kontynuacji — wracałby do użytkownika fragment
|
||||
# udający całość.
|
||||
produced = (turn_usage or {}).get("completion_tokens")
|
||||
if produced is None:
|
||||
produced = (turn_usage or {}).get("output_tokens")
|
||||
if produced is None:
|
||||
produced = max(1, int(len(text) / 3.6))
|
||||
remaining -= max(1, int(produced))
|
||||
|
||||
emit("turn_end",
|
||||
f"Tura {turns}: odebrano {len(text.strip())} znaków w {took:.1f}s"
|
||||
+ (" — odpowiedź urwana, poproszę o dokończenie" if truncated else ""),
|
||||
turn=turns, chars=len(text.strip()), truncated=truncated)
|
||||
|
||||
chunk = text.strip()
|
||||
if chunk:
|
||||
@@ -172,6 +199,8 @@ class _Driver:
|
||||
final = _join(parts)
|
||||
if not final:
|
||||
raise LLMError(_explain_empty(turns, usage, stop))
|
||||
emit("generated", f"Gotowe: {len(final)} znaków w {turns} turach.",
|
||||
chars=len(final), turns=turns)
|
||||
|
||||
usage["turns"] = turns
|
||||
return Completion(text=final, model=model_name, provider=self.name,
|
||||
|
||||
@@ -12,7 +12,7 @@ import httpx
|
||||
from fastapi import FastAPI, HTTPException
|
||||
from pydantic import BaseModel
|
||||
|
||||
from app import link_crypto, security
|
||||
from app import security
|
||||
from app.clients.data_client import DataClient
|
||||
from app.models import QueryRequest, QueryResponse
|
||||
from app.service import QueryService
|
||||
@@ -20,9 +20,6 @@ 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
|
||||
@@ -268,6 +265,72 @@ def chart_horoscope(req: HoroscopeRequest) -> dict:
|
||||
return out
|
||||
|
||||
|
||||
@app.post("/chart/horoscope/stream")
|
||||
def chart_horoscope_stream(req: HoroscopeRequest):
|
||||
"""To samo co /chart/horoscope, ale strumieniuje POSTĘP w trakcie pracy.
|
||||
|
||||
Pisanie horoskopu trwa minutami — bez sygnału aplikacja wygląda na zawieszoną.
|
||||
Strumień (NDJSON, jedna linia = jedno zdarzenie) niesie RZECZYWISTE etapy:
|
||||
budowę promptu, limity modelu i każdą turę generowania. Ostatnie zdarzenie
|
||||
(`result`) ma identyczny kształt co odpowiedź zwykłego endpointu.
|
||||
"""
|
||||
from fastapi.responses import StreamingResponse
|
||||
|
||||
from app.progress import stream
|
||||
|
||||
def work(emit) -> dict:
|
||||
from app.llm.base import LLMError
|
||||
from app.llm.factory import build_provider
|
||||
from app.llm.limits import plan
|
||||
|
||||
emit({"type": "stage", "message": "Liczę horoskop i szukam wskazań w bazach…"})
|
||||
out = chart_prompt(req)
|
||||
st = out.get("stats", {})
|
||||
emit({"type": "stage", "message":
|
||||
f"Prompt gotowy: {st.get('chars', 0)} znaków, "
|
||||
f"wskazań {st.get('included', 0)}"
|
||||
+ (f", pominięto {st['omitted']}" if st.get("omitted") else "")})
|
||||
if out.get("data_error"):
|
||||
emit({"type": "warn", "message": out["data_error"]})
|
||||
|
||||
try:
|
||||
provider = build_provider(req.provider, req.model)
|
||||
emit({"type": "stage", "message":
|
||||
f"Dostawca: {provider.name}, model: {provider.model}"
|
||||
+ ("" if not provider.leaves_lan else " — dane opuszczają sieć")})
|
||||
|
||||
emit({"type": "stage", "message": "Liczę tokeny promptu…"})
|
||||
prompt_tokens = provider.count_tokens(out["prompt"])
|
||||
budget = plan(provider.name, provider.model, prompt_tokens, req.max_tokens)
|
||||
out["token_plan"] = budget
|
||||
emit({"type": "stage", "message":
|
||||
f"Prompt {prompt_tokens} tok. · okno modelu {budget['context_window']} · "
|
||||
f"na odpowiedź {budget['max_output']}"})
|
||||
for warning in budget["warnings"]:
|
||||
emit({"type": "warn", "message": warning})
|
||||
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"], on_event=emit)
|
||||
except LLMError as e:
|
||||
emit({"type": "warn", "message": f"Model zawiódł: {e}"})
|
||||
out["llm_error"] = str(e)
|
||||
return out
|
||||
|
||||
out.update(horoscope=result.text, provider=result.provider, model=result.model,
|
||||
leaves_lan=result.leaves_lan, usage=result.usage)
|
||||
return out
|
||||
|
||||
return StreamingResponse(
|
||||
stream(work),
|
||||
media_type="application/x-ndjson",
|
||||
headers={"Cache-Control": "no-store", "X-Accel-Buffering": "no"},
|
||||
)
|
||||
|
||||
|
||||
@app.get("/llm/models")
|
||||
def llm_models() -> dict:
|
||||
"""Podpowiedzi modeli per dostawca — UI buduje z tego listę wyboru.
|
||||
|
||||
@@ -0,0 +1,70 @@
|
||||
"""Strumień postępu długiej operacji (NDJSON).
|
||||
|
||||
Po co: pisanie horoskopu trwa — czasem minuty. Bez sygnału aplikacja wygląda na
|
||||
zawieszoną. Zamiast udawanego paska postępu strumieniujemy **rzeczywiste**
|
||||
zdarzenia z kolejnych etapów, żeby log pokazywał to, co faktycznie się dzieje.
|
||||
|
||||
Dlaczego NDJSON, a nie SSE: `EventSource` w przeglądarce obsługuje wyłącznie GET,
|
||||
a to jest POST z ciałem. Strumień „jedna linia = jeden obiekt JSON" czyta się
|
||||
zwykłym `fetch()` i jest trywialny do sparsowania.
|
||||
|
||||
Dlaczego wątek: właściwa praca (silnik, baza, model) jest synchroniczna. Puszczamy
|
||||
ją w wątku roboczym, a generator odpompowuje kolejkę zdarzeń — dzięki temu
|
||||
zdarzenia docierają w trakcie pracy, a nie dopiero na końcu.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import queue
|
||||
import threading
|
||||
import traceback
|
||||
from collections.abc import Iterator
|
||||
from typing import Any, Callable
|
||||
|
||||
_HEARTBEAT_SECONDS = 10.0
|
||||
_DONE = object()
|
||||
|
||||
|
||||
def line(kind: str, message: str, **extra: Any) -> str:
|
||||
return json.dumps({"type": kind, "message": message, **extra}, ensure_ascii=False) + "\n"
|
||||
|
||||
|
||||
def stream(work: Callable[[Callable[[dict], None]], dict]) -> Iterator[str]:
|
||||
"""Uruchamia `work(emit)` w wątku i strumieniuje zdarzenia w czasie rzeczywistym.
|
||||
|
||||
`work` dostaje funkcję `emit(zdarzenie)` i zwraca końcowy wynik, który leci
|
||||
jako ostatnie zdarzenie typu `result`. Wyjątek zamienia się w zdarzenie `error`
|
||||
— połączenie nigdy nie urywa się bez wyjaśnienia.
|
||||
"""
|
||||
events: queue.Queue = queue.Queue()
|
||||
|
||||
def emit(event: dict) -> None:
|
||||
events.put(event)
|
||||
|
||||
def run() -> None:
|
||||
try:
|
||||
result = work(emit)
|
||||
events.put({"type": "result", "message": "Gotowe.", "result": result})
|
||||
except Exception as e: # noqa: BLE001 — zgłaszamy KAŻDY błąd
|
||||
events.put({
|
||||
"type": "error",
|
||||
"message": f"{type(e).__name__}: {e}",
|
||||
"detail": traceback.format_exc(limit=3),
|
||||
})
|
||||
finally:
|
||||
events.put(_DONE)
|
||||
|
||||
worker = threading.Thread(target=run, daemon=True)
|
||||
worker.start()
|
||||
|
||||
while True:
|
||||
try:
|
||||
event = events.get(timeout=_HEARTBEAT_SECONDS)
|
||||
except queue.Empty:
|
||||
# cisza dłuższa niż heartbeat: dajemy znak życia, żeby pośredniki
|
||||
# (proxy, load balancer) nie uznały połączenia za martwe
|
||||
yield line("ping", "…")
|
||||
continue
|
||||
if event is _DONE:
|
||||
break
|
||||
yield json.dumps(event, ensure_ascii=False) + "\n"
|
||||
@@ -4,5 +4,3 @@ 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
|
||||
|
||||
@@ -1,351 +0,0 @@
|
||||
"""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")
|
||||
@@ -405,3 +405,33 @@ def test_max_budget_differs_between_models(monkeypatch):
|
||||
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
|
||||
|
||||
|
||||
def test_turn_budget_counts_produced_not_requested(monkeypatch):
|
||||
"""Regresja: odejmowanie ZAMOWIONEGO limitu tury zamiast wyprodukowanych
|
||||
tokenow konczylo petle po jednej turze — urwany fragment wracal jako calosc."""
|
||||
seq = [_chat("Fragment 1. ", "length"), _chat("Fragment 2. ", "length"),
|
||||
_chat("Zakonczenie. KONIEC", "stop")]
|
||||
|
||||
def handler(request):
|
||||
return httpx.Response(200, json=seq.pop(0) if seq else seq[-1])
|
||||
|
||||
monkeypatch.setattr(httpx, "Client", _mock_client(handler))
|
||||
# budzet 8000 < TURN_TOKENS_CAP: przy starej logice byla dokladnie jedna tura
|
||||
out = ChatCompletionsProvider("local", "http://x/v1", "m").generate("p", 8000)
|
||||
assert out.usage["turns"] == 3, "urwana odpowiedz musi byc kontynuowana"
|
||||
assert "Fragment 1." in out.text and "Zakonczenie." in out.text
|
||||
|
||||
|
||||
def test_generate_reports_progress_events(monkeypatch):
|
||||
"""Log w UI ma pokazywac RZECZYWISTE tury, nie udawany pasek postepu."""
|
||||
seq = [_chat("Czesc. ", "length"), _chat("Reszta. KONIEC", "stop")]
|
||||
monkeypatch.setattr(httpx, "Client", _mock_client(
|
||||
lambda r: httpx.Response(200, json=seq.pop(0) if seq else seq[-1])))
|
||||
events = []
|
||||
ChatCompletionsProvider("local", "http://x/v1", "m").generate(
|
||||
"p", 8000, on_event=events.append)
|
||||
kinds = [e["type"] for e in events]
|
||||
assert kinds.count("turn_start") == 2 and kinds.count("turn_end") == 2
|
||||
assert kinds[-1] == "generated"
|
||||
assert any("urwana" in e["message"] for e in events if e["type"] == "turn_end")
|
||||
|
||||
@@ -10,7 +10,6 @@ from typing import Any
|
||||
|
||||
import httpx
|
||||
|
||||
from app import link_crypto
|
||||
from app.config import settings
|
||||
|
||||
|
||||
@@ -20,29 +19,16 @@ 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}
|
||||
return self._post("/api/query", payload, settings.http_timeout)
|
||||
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()
|
||||
|
||||
def positions(
|
||||
self,
|
||||
@@ -65,15 +51,20 @@ class LogicClient:
|
||||
"zodiac": zodiac,
|
||||
}
|
||||
# stacje wymagają root-findów — dłuższy timeout
|
||||
timeout = max(settings.http_timeout, 60.0) if stations else settings.http_timeout
|
||||
return self._post("/chart/positions", payload, 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()
|
||||
|
||||
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}
|
||||
return self._post("/chart/report", payload, max(settings.http_timeout, 30.0))
|
||||
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()
|
||||
|
||||
def prompt(
|
||||
self, profile: str, when_utc_iso: str, lat: float, lon: float,
|
||||
@@ -91,7 +82,10 @@ class LogicClient:
|
||||
}
|
||||
if from_date and to_date:
|
||||
payload["from_date"], payload["to_date"] = from_date, to_date
|
||||
return self._post("/chart/prompt", payload, max(settings.http_timeout, 60.0))
|
||||
with httpx.Client(timeout=max(settings.http_timeout, 60.0)) as client:
|
||||
r = client.post(f"{self.base_url}/chart/prompt", json=payload, headers=_auth_headers())
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
|
||||
def horoscope(
|
||||
self, profile: str, when_utc_iso: str, lat: float, lon: float,
|
||||
@@ -112,13 +106,31 @@ class LogicClient:
|
||||
payload["model"] = model
|
||||
if from_date and to_date:
|
||||
payload["from_date"], payload["to_date"] = from_date, to_date
|
||||
return self._post("/chart/horoscope", payload, max(settings.http_timeout, 300.0))
|
||||
with httpx.Client(timeout=max(settings.http_timeout, 300.0)) as client:
|
||||
r = client.post(f"{self.base_url}/chart/horoscope", json=payload, headers=_auth_headers())
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
|
||||
def horoscope_stream(self, payload: dict[str, Any]):
|
||||
"""Strumień postępu pisania horoskopu (NDJSON) — przekazywany do przeglądarki.
|
||||
|
||||
Timeout jest długi, bo generowanie trwa; strumień i tak niesie heartbeat,
|
||||
więc cisza na łączu nie zostanie wzięta za zerwanie.
|
||||
"""
|
||||
with httpx.Client(timeout=httpx.Timeout(None, connect=15.0)) as client:
|
||||
with client.stream("POST", f"{self.base_url}/chart/horoscope/stream",
|
||||
json=payload, headers=_auth_headers()) as r:
|
||||
r.raise_for_status()
|
||||
for chunk in r.iter_lines():
|
||||
if chunk:
|
||||
yield chunk
|
||||
|
||||
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())
|
||||
r = client.get(f"{self.base_url}/llm/models", headers=_auth_headers())
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
|
||||
def timeline(
|
||||
self, when_utc_iso: str, lat: float, lon: float,
|
||||
@@ -129,4 +141,7 @@ class LogicClient:
|
||||
"when_utc": when_utc_iso, "lat": lat, "lon": lon,
|
||||
"from_date": from_date, "to_date": to_date, "interpret": interpret,
|
||||
}
|
||||
return self._post("/chart/timeline", payload, max(settings.http_timeout, 60.0))
|
||||
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()
|
||||
|
||||
@@ -1,466 +0,0 @@
|
||||
"""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")
|
||||
@@ -9,11 +9,12 @@ Strona główna „/" = wprowadzenie danych horoskopu i podgląd policzonych poz
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
import httpx
|
||||
from fastapi import FastAPI, Form, HTTPException, Query, Request
|
||||
from fastapi.responses import HTMLResponse
|
||||
from fastapi.responses import HTMLResponse, JSONResponse
|
||||
from fastapi.staticfiles import StaticFiles
|
||||
from fastapi.templating import Jinja2Templates
|
||||
|
||||
@@ -228,6 +229,61 @@ def timeline_run(
|
||||
return templates.TemplateResponse(request, "timeline.html", ctx)
|
||||
|
||||
|
||||
# ---------------- Postęp pisania horoskopu (strumień do okna z logiem) ----------------
|
||||
@app.post("/horoscope/stream")
|
||||
def horoscope_stream(
|
||||
profile: str = Form("natal"),
|
||||
date: str = Form(...),
|
||||
time: str = Form(...),
|
||||
tz_offset: float = Form(0.0),
|
||||
lat: float = Form(0.0),
|
||||
lon: float = Form(0.0),
|
||||
prompt_budget: str = Form("medium"),
|
||||
llm_provider: str = Form("local"),
|
||||
llm_model: str = Form(""),
|
||||
from_date: str = Form(""),
|
||||
to_date: str = Form(""),
|
||||
):
|
||||
"""Przekazuje strumień postępu z logiki i DOKLEJA gotowy HTML wyniku.
|
||||
|
||||
Dzięki temu okno postępu wstawia dokładnie ten sam widok, który wyrenderowałoby
|
||||
przeładowanie strony — jedno źródło prawdy dla wyglądu wyniku.
|
||||
"""
|
||||
from fastapi.responses import StreamingResponse
|
||||
|
||||
try:
|
||||
iso_utc, _ = _build_utc(date, time, tz_offset)
|
||||
except ValueError as e:
|
||||
return JSONResponse({"detail": f"Niepoprawne dane wejściowe: {e}"}, status_code=422)
|
||||
|
||||
payload: dict = {
|
||||
"profile": profile, "when_utc": iso_utc, "lat": lat, "lon": lon,
|
||||
"budget": prompt_budget, "provider": llm_provider, "model": llm_model,
|
||||
}
|
||||
if profile == "period" and from_date and to_date:
|
||||
payload["from_date"], payload["to_date"] = from_date, to_date
|
||||
|
||||
def relay():
|
||||
try:
|
||||
for raw in logic.horoscope_stream(payload):
|
||||
try:
|
||||
event = json.loads(raw)
|
||||
except ValueError:
|
||||
continue
|
||||
if event.get("type") == "result":
|
||||
html = templates.get_template("_prompt_result.html").render(
|
||||
prompt_result=event.get("result") or {}
|
||||
)
|
||||
event["html"] = html
|
||||
yield json.dumps(event, ensure_ascii=False) + "\n"
|
||||
except httpx.HTTPError as e:
|
||||
yield json.dumps({"type": "error", "message": _logic_error(e)},
|
||||
ensure_ascii=False) + "\n"
|
||||
|
||||
return StreamingResponse(relay(), media_type="application/x-ndjson",
|
||||
headers={"Cache-Control": "no-store", "X-Accel-Buffering": "no"})
|
||||
|
||||
|
||||
# ---------------- Geokoder (proxy OSM/Nominatim dla wyszukiwarki lokalizacji) ----------------
|
||||
@app.get("/geocode")
|
||||
def geocode_search(q: str = Query("", description="Nazwa / adres / POI do wyszukania")):
|
||||
|
||||
@@ -7,9 +7,7 @@ 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ń);
|
||||
rozliczany per adres klienta, a za odwrotnym proxy — po TRUST_PROXY=true —
|
||||
per adres z nagłówka, nie per adres proxy (patrz `client_ip`).
|
||||
* **limit żądań** — hamuje masowe odpytywanie (eksfiltrację przez pętlę zapytań).
|
||||
|
||||
Świadomie NIE logujemy treści żądań ani promptów — logi to kolejny nośnik wycieku.
|
||||
|
||||
@@ -48,11 +46,6 @@ 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/",)
|
||||
|
||||
@@ -81,31 +74,6 @@ 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:
|
||||
@@ -137,7 +105,8 @@ def install(app) -> None:
|
||||
if _is_public(request.url.path):
|
||||
return await call_next(request)
|
||||
|
||||
if _rate_limited(client_ip(request)):
|
||||
client = request.client.host if request.client else "?"
|
||||
if _rate_limited(client):
|
||||
return JSONResponse(
|
||||
{"detail": "Zbyt wiele żądań — spróbuj za chwilę."},
|
||||
status_code=429, headers={"Retry-After": "60"},
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
// Okno postępu przy pisaniu horoskopu.
|
||||
//
|
||||
// Problem: generowanie trwa minutami, a zwykły POST formularza nie daje żadnego
|
||||
// sygnału — aplikacja wygląda na zawieszoną. Zamiast udawanego paska postępu
|
||||
// czytamy strumień RZECZYWISTYCH zdarzeń z serwera (NDJSON) i wypisujemy je
|
||||
// jako log: budowa promptu, limity modelu, każda tura generowania.
|
||||
//
|
||||
// Degradacja: jeśli przeglądarka nie umie strumieniować `fetch`, nie przechwytujemy
|
||||
// wysyłki — formularz idzie klasycznie i wszystko działa jak wcześniej, tylko bez okna.
|
||||
document.addEventListener('DOMContentLoaded', function () {
|
||||
const canStream = typeof fetch === 'function' && typeof ReadableStream === 'function' &&
|
||||
typeof TextDecoder === 'function';
|
||||
const form = document.querySelector('form[action="/interpret"], form[action="/timeline"]');
|
||||
if (!canStream || !form) return;
|
||||
|
||||
const profile = form.getAttribute('action') === '/timeline' ? 'period' : 'natal';
|
||||
const btn = form.querySelector('button[value="horoscope"]');
|
||||
if (!btn) return;
|
||||
|
||||
// --- okno ---------------------------------------------------------------
|
||||
const overlay = document.createElement('div');
|
||||
overlay.className = 'progress-overlay';
|
||||
overlay.hidden = true;
|
||||
overlay.innerHTML =
|
||||
'<div class="progress-box" role="dialog" aria-modal="true" aria-label="Postęp generowania">' +
|
||||
'<div class="progress-head">' +
|
||||
'<span class="progress-spinner" aria-hidden="true"></span>' +
|
||||
'<strong id="progressTitle">Piszę horoskop…</strong>' +
|
||||
'<span class="progress-clock" id="progressClock">0:00</span>' +
|
||||
'</div>' +
|
||||
'<ol class="progress-log" id="progressLog"></ol>' +
|
||||
'<p class="muted small">Nie zamykaj tej karty — generowanie trwa na serwerze.</p>' +
|
||||
'<div class="actions"><button type="button" class="ghost" id="progressClose" hidden>Zamknij</button></div>' +
|
||||
'</div>';
|
||||
document.body.appendChild(overlay);
|
||||
|
||||
const logEl = overlay.querySelector('#progressLog');
|
||||
const clockEl = overlay.querySelector('#progressClock');
|
||||
const titleEl = overlay.querySelector('#progressTitle');
|
||||
const closeEl = overlay.querySelector('#progressClose');
|
||||
let timer = null;
|
||||
|
||||
function addLine(text, kind) {
|
||||
const li = document.createElement('li');
|
||||
if (kind) li.className = 'log-' + kind;
|
||||
const now = new Date();
|
||||
li.textContent = String(now.getHours()).padStart(2, '0') + ':' +
|
||||
String(now.getMinutes()).padStart(2, '0') + ':' +
|
||||
String(now.getSeconds()).padStart(2, '0') + ' ' + text;
|
||||
logEl.appendChild(li);
|
||||
logEl.scrollTop = logEl.scrollHeight;
|
||||
}
|
||||
|
||||
function startClock() {
|
||||
const t0 = Date.now();
|
||||
timer = setInterval(function () {
|
||||
const s = Math.floor((Date.now() - t0) / 1000);
|
||||
clockEl.textContent = Math.floor(s / 60) + ':' + String(s % 60).padStart(2, '0');
|
||||
}, 1000);
|
||||
}
|
||||
|
||||
function finish(title, allowClose) {
|
||||
if (timer) { clearInterval(timer); timer = null; }
|
||||
titleEl.textContent = title;
|
||||
overlay.querySelector('.progress-spinner').style.visibility = 'hidden';
|
||||
if (allowClose) closeEl.hidden = false;
|
||||
}
|
||||
|
||||
closeEl.addEventListener('click', function () { overlay.hidden = true; });
|
||||
|
||||
// --- przechwycenie wysyłki ----------------------------------------------
|
||||
btn.addEventListener('click', function (event) {
|
||||
event.preventDefault();
|
||||
|
||||
if (!form.reportValidity()) return; // te same reguły co przy zwykłej wysyłce
|
||||
|
||||
const data = new FormData(form);
|
||||
data.set('profile', profile);
|
||||
|
||||
logEl.innerHTML = '';
|
||||
closeEl.hidden = true;
|
||||
overlay.querySelector('.progress-spinner').style.visibility = '';
|
||||
titleEl.textContent = 'Piszę horoskop…';
|
||||
overlay.hidden = false;
|
||||
startClock();
|
||||
addLine('Wysyłam żądanie…');
|
||||
|
||||
fetch('/horoscope/stream', { method: 'POST', body: data })
|
||||
.then(function (response) {
|
||||
if (!response.ok || !response.body) throw new Error('HTTP ' + response.status);
|
||||
const reader = response.body.getReader();
|
||||
const decoder = new TextDecoder();
|
||||
let buffer = '';
|
||||
|
||||
function pump() {
|
||||
return reader.read().then(function (chunk) {
|
||||
if (chunk.done) return;
|
||||
buffer += decoder.decode(chunk.value, { stream: true });
|
||||
const lines = buffer.split('\n');
|
||||
buffer = lines.pop(); // ostatni może być niepełny
|
||||
lines.forEach(function (raw) {
|
||||
if (!raw.trim()) return;
|
||||
let ev;
|
||||
try { ev = JSON.parse(raw); } catch (e) { return; }
|
||||
if (ev.type === 'ping') return; // sam heartbeat, nie logujemy
|
||||
if (ev.type === 'result') {
|
||||
addLine(ev.message || 'Gotowe.', 'ok');
|
||||
if (ev.html) {
|
||||
const host = document.getElementById('promptResult');
|
||||
if (host) host.innerHTML = ev.html;
|
||||
}
|
||||
finish('Gotowe', true);
|
||||
overlay.hidden = true;
|
||||
const host = document.getElementById('promptResult');
|
||||
if (host) host.scrollIntoView({ behavior: 'smooth', block: 'start' });
|
||||
return;
|
||||
}
|
||||
addLine(ev.message || ev.type, ev.type === 'error' ? 'err'
|
||||
: ev.type === 'warn' ? 'warn' : null);
|
||||
if (ev.type === 'error') finish('Nie udało się', true);
|
||||
});
|
||||
return pump();
|
||||
});
|
||||
}
|
||||
return pump();
|
||||
})
|
||||
.catch(function (e) {
|
||||
addLine('Połączenie przerwane: ' + e.message, 'err');
|
||||
addLine('Możesz spróbować ponownie — nic nie zostało utracone.');
|
||||
finish('Nie udało się', true);
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -88,3 +88,27 @@ textarea.prompt { width: 100%; margin-top: .5rem; padding: .7rem .8rem; box-sizi
|
||||
border-radius: 10px; font-family: ui-monospace, SFMono-Regular, Menlo, monospace;
|
||||
font-size: .82rem; line-height: 1.45; resize: vertical; white-space: pre; }
|
||||
textarea.prompt:focus { outline: 2px solid var(--accent); outline-offset: 1px; }
|
||||
|
||||
/* Okno postępu przy generowaniu horoskopu */
|
||||
.progress-overlay { position: fixed; inset: 0; background: rgba(8,9,20,.72);
|
||||
display: flex; align-items: center; justify-content: center;
|
||||
z-index: 50; padding: 1rem; }
|
||||
.progress-overlay[hidden] { display: none; }
|
||||
.progress-box { background: var(--panel); border: 1px solid var(--line); border-radius: 14px;
|
||||
padding: 1rem 1.1rem; width: min(38rem, 100%); box-shadow: 0 18px 48px rgba(0,0,0,.45); }
|
||||
.progress-head { display: flex; align-items: center; gap: .6rem; margin-bottom: .6rem; }
|
||||
.progress-clock { margin-left: auto; font-family: ui-monospace, Menlo, monospace;
|
||||
color: var(--muted); font-size: .9rem; }
|
||||
.progress-spinner { width: 14px; height: 14px; border-radius: 50%; flex: none;
|
||||
border: 2px solid var(--line); border-top-color: var(--accent);
|
||||
animation: spin .8s linear infinite; }
|
||||
@keyframes spin { to { transform: rotate(360deg); } }
|
||||
@media (prefers-reduced-motion: reduce) { .progress-spinner { animation: none; } }
|
||||
.progress-log { list-style: none; margin: 0 0 .6rem; padding: .6rem .7rem;
|
||||
background: #12132a; border: 1px solid var(--line); border-radius: 10px;
|
||||
max-height: 15rem; overflow-y: auto;
|
||||
font-family: ui-monospace, Menlo, monospace; font-size: .78rem; line-height: 1.6; }
|
||||
.progress-log li { color: var(--ink); white-space: pre-wrap; }
|
||||
.progress-log li.log-ok { color: #7fd18b; }
|
||||
.progress-log li.log-warn { color: #e0c060; }
|
||||
.progress-log li.log-err { color: #ef6b6b; }
|
||||
|
||||
@@ -36,80 +36,4 @@
|
||||
sieć</strong>. Możesz najpierw obejrzeć prompt, a dopiero potem wysłać.
|
||||
</p>
|
||||
|
||||
{% if prompt_result %}
|
||||
{% set st = prompt_result.stats %}
|
||||
|
||||
{# --- wynik: napisany horoskop (LOG-31) + transparentność (PRE-15) --- #}
|
||||
{% if prompt_result.horoscope %}
|
||||
<div class="meta">
|
||||
Horoskop napisany przez: <strong>{{ prompt_result.provider }}</strong> ·
|
||||
model: {{ prompt_result.model }}
|
||||
{% if prompt_result.usage and prompt_result.usage.completion_tokens %}
|
||||
· tokeny odpowiedzi: {{ prompt_result.usage.completion_tokens }}
|
||||
{% endif %}
|
||||
{% if prompt_result.leaves_lan %}
|
||||
· <strong class="retro">dane opuściły sieć</strong>
|
||||
{% else %}
|
||||
· dane nie opuściły sieci
|
||||
{% endif %}
|
||||
</div>
|
||||
<div class="actions">
|
||||
<button type="button" class="ghost" data-copy="#horoscopeText">Kopiuj horoskop</button>
|
||||
</div>
|
||||
<textarea id="horoscopeText" class="prompt" rows="20" readonly>{{ prompt_result.horoscope }}</textarea>
|
||||
<p class="muted small">
|
||||
Treść wygenerował model językowy na podstawie {{ st.included }} wskazań z baz.
|
||||
<strong>Nie stanowi porady medycznej, prawnej ani finansowej.</strong>
|
||||
</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>
|
||||
Prompt poniżej jest gotowy — możesz go skopiować i użyć ręcznie.
|
||||
</div>
|
||||
{% endif %}
|
||||
|
||||
{# --- prompt: zawsze dostępny do podglądu i skopiowania --- #}
|
||||
<div class="meta">
|
||||
Prompt · {{ st.chars }} znaków (~{{ st.est_tokens }} tokenów) ·
|
||||
budżet: {{ st.budget }} ·
|
||||
wskazań: <strong>{{ st.included }}</strong>
|
||||
{% if st.omitted %}· pominięto: <strong>{{ st.omitted }}</strong>{% endif %}
|
||||
{% if st.deduplicated %}· scalono powtórek: {{ st.deduplicated }}{% endif %}
|
||||
</div>
|
||||
|
||||
{% if st.omitted %}
|
||||
<p class="muted small">
|
||||
Pominięto {{ st.omitted }} najsłabszych wskazań (próg wagi {{ st.min_score_included }}).
|
||||
Chcesz komplet — wybierz obszerniejszy budżet i wygeneruj ponownie.
|
||||
</p>
|
||||
{% endif %}
|
||||
|
||||
{% if prompt_result.data_error %}
|
||||
<div class="error">{{ prompt_result.data_error }} — prompt złożony z samych wyliczeń.</div>
|
||||
{% endif %}
|
||||
|
||||
<div class="actions">
|
||||
<button type="button" class="ghost" data-copy="#promptText">Kopiuj prompt</button>
|
||||
<span class="muted small">Możesz też wkleić go samodzielnie do ChatGPT lub Claude.</span>
|
||||
</div>
|
||||
<textarea id="promptText" class="prompt" rows="14" readonly>{{ prompt_result.prompt }}</textarea>
|
||||
{% endif %}
|
||||
<div id="promptResult">{% include "_prompt_result.html" %}</div>
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
{# Blok WYNIKU generowania (prompt/horoskop). Wydzielony, bo wstawia go także
|
||||
okno postępu po zakończeniu strumienia — dzięki temu jest jedno źródło
|
||||
prawdy dla wyglądu wyniku, niezależnie od drogi, którą przyszedł. #}
|
||||
{% if prompt_result %}
|
||||
{% set st = prompt_result.stats %}
|
||||
|
||||
{# --- wynik: napisany horoskop (LOG-31) + transparentność (PRE-15) --- #}
|
||||
{% if prompt_result.horoscope %}
|
||||
<div class="meta">
|
||||
Horoskop napisany przez: <strong>{{ prompt_result.provider }}</strong> ·
|
||||
model: {{ prompt_result.model }}
|
||||
{% if prompt_result.usage and prompt_result.usage.completion_tokens %}
|
||||
· tokeny odpowiedzi: {{ prompt_result.usage.completion_tokens }}
|
||||
{% endif %}
|
||||
{% if prompt_result.leaves_lan %}
|
||||
· <strong class="retro">dane opuściły sieć</strong>
|
||||
{% else %}
|
||||
· dane nie opuściły sieci
|
||||
{% endif %}
|
||||
</div>
|
||||
<div class="actions">
|
||||
<button type="button" class="ghost" data-copy="#horoscopeText">Kopiuj horoskop</button>
|
||||
</div>
|
||||
<textarea id="horoscopeText" class="prompt" rows="20" readonly>{{ prompt_result.horoscope }}</textarea>
|
||||
<p class="muted small">
|
||||
Treść wygenerował model językowy na podstawie {{ st.included }} wskazań z baz.
|
||||
<strong>Nie stanowi porady medycznej, prawnej ani finansowej.</strong>
|
||||
</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>
|
||||
Prompt poniżej jest gotowy — możesz go skopiować i użyć ręcznie.
|
||||
</div>
|
||||
{% endif %}
|
||||
|
||||
{# --- prompt: zawsze dostępny do podglądu i skopiowania --- #}
|
||||
<div class="meta">
|
||||
Prompt · {{ st.chars }} znaków (~{{ st.est_tokens }} tokenów) ·
|
||||
budżet: {{ st.budget }} ·
|
||||
wskazań: <strong>{{ st.included }}</strong>
|
||||
{% if st.omitted %}· pominięto: <strong>{{ st.omitted }}</strong>{% endif %}
|
||||
{% if st.deduplicated %}· scalono powtórek: {{ st.deduplicated }}{% endif %}
|
||||
</div>
|
||||
|
||||
{% if st.omitted %}
|
||||
<p class="muted small">
|
||||
Pominięto {{ st.omitted }} najsłabszych wskazań (próg wagi {{ st.min_score_included }}).
|
||||
Chcesz komplet — wybierz obszerniejszy budżet i wygeneruj ponownie.
|
||||
</p>
|
||||
{% endif %}
|
||||
|
||||
{% if prompt_result.data_error %}
|
||||
<div class="error">{{ prompt_result.data_error }} — prompt złożony z samych wyliczeń.</div>
|
||||
{% endif %}
|
||||
|
||||
<div class="actions">
|
||||
<button type="button" class="ghost" data-copy="#promptText">Kopiuj prompt</button>
|
||||
<span class="muted small">Możesz też wkleić go samodzielnie do ChatGPT lub Claude.</span>
|
||||
</div>
|
||||
<textarea id="promptText" class="prompt" rows="14" readonly>{{ prompt_result.prompt }}</textarea>
|
||||
{% endif %}
|
||||
@@ -91,4 +91,5 @@
|
||||
<script src="/static/now.js"></script>
|
||||
<script src="/static/copy.js"></script>
|
||||
<script src="/static/models.js"></script>
|
||||
<script src="/static/progress.js"></script>
|
||||
{% endblock %}
|
||||
|
||||
@@ -81,4 +81,5 @@
|
||||
<script src="/static/now.js"></script>
|
||||
<script src="/static/copy.js"></script>
|
||||
<script src="/static/models.js"></script>
|
||||
<script src="/static/progress.js"></script>
|
||||
{% endblock %}
|
||||
|
||||
@@ -3,5 +3,3 @@ 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
|
||||
|
||||
@@ -1,44 +1,30 @@
|
||||
"""Niezmiennik: KAŻDE wyjście HTTP w dół niesie token międzywarstwowy (LOG-32)
|
||||
oraz klucz szyfrujący łącze (PRE-16).
|
||||
"""Niezmiennik: KAŻDE wyjście HTTP w dół niesie token międzywarstwowy (LOG-32).
|
||||
|
||||
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.
|
||||
strukturalnie: nie ma wywołania bez `headers=`.
|
||||
"""
|
||||
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ół."""
|
||||
def _http_calls(path: pathlib.Path) -> list[tuple[str, int, bool]]:
|
||||
"""(nazwa_metody_http, linia, czy_ma_headers) dla każdego client.post/get."""
|
||||
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
|
||||
if node.func.attr not in ("post", "get", "put", "patch", "delete"):
|
||||
continue
|
||||
if not (isinstance(node.func.value, ast.Name) and node.func.value.id == "client"):
|
||||
continue
|
||||
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))
|
||||
out.append((node.func.attr, node.lineno, has_headers))
|
||||
return out
|
||||
|
||||
|
||||
@@ -49,43 +35,13 @@ def test_client_module_exists():
|
||||
def test_every_outbound_call_sends_auth_header():
|
||||
calls = _http_calls(CLIENT)
|
||||
assert calls, "nie znaleziono żadnego wywołania HTTP — test przestał cokolwiek pilnować"
|
||||
missing = [f"{CLIENT.name}:{line} {what}" for what, line, ok, _ in calls if not ok]
|
||||
missing = [f"{CLIENT.name}:{line} client.{verb}()" for verb, 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."""
|
||||
|
||||
@@ -131,71 +131,3 @@ 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"
|
||||
|
||||
Reference in New Issue
Block a user