Merge origin/master do PRE-16: naprawa regresji strumienia postepu
Testy / Testy warstwy logicznej (silnik) (push) Successful in 11m33s
Testy / Testy warstwy prezentacji (dostęp do baz) (push) Successful in 9m43s
Testy / Build obrazu silnika B (swisseph) (push) Successful in 33s
Testy / Kontrola składni wszystkich warstw (push) Successful in 15s
Testy / Testy warstwy logicznej (silnik) (pull_request) Successful in 10m39s
Testy / Testy warstwy prezentacji (dostęp do baz) (pull_request) Successful in 9m45s
Testy / Build obrazu silnika B (swisseph) (pull_request) Successful in 26s
Testy / Kontrola składni wszystkich warstw (pull_request) Successful in 15s

Konflikt w logic_client.py::positions() rozwiazany biorac OBIE zmiany:
routing przez szyfrowane _post() (PRE-16) + dluzszy timeout takze dla tables
(LOG-23) — warunek (stations or tables).

Wazniejsza rzecz, ktora scalenie ujawnilo: okno postepu (#19/#22) i szyfrowanie
lacz (#21) powstaly na rownoleglych galeziach, ktore sie nie widzialy. Po
zejsciu razem strumien horoskopu szedl SUROWYM httpx, z pominieciem szyfrowania.
Przy wlaczonym LINK_ENCRYPTION_REQUIRED serwer odrzucalby to zadanie (400), a
nawet bez wymagania odpowiedz wracalaby jako nieczytelne ramki — okno postepu
przestaloby dzialac na produkcji.

Naprawa:
- nowy link_crypto.stream_lines(): strumieniowe POST przez szyfrowane lacze;
  pieczetuje zadanie i odszyfrowuje odpowiedz ramka po ramce, sklejajac bufor bo
  granice ramek nie pokrywaja sie z granicami linii NDJSON. Dostarczanie na zywo
  zachowane. Bez klucza — jak dotad (dev).
- horoscope_stream() w kliencie idzie teraz przez stream_lines zamiast surowego
  client.stream.
- fail-closed takze dla strumienia: bez klucza przy wymaganym szyfrowaniu klient
  nie wysyla NIC (wczesniej cialo — dane urodzenia — szloby w eter, dopiero potem
  serwer odmawial). Ujednolica kontrakt z call().

Weryfikacja e2e na prawdziwym uvicornie z podsluchem gniazda: z kluczem strumien
dziala (5 etapow + result, na zywo), na kablu ZERO tresci bazy (grep=0; jedyne
'horoscope' to sciezka URL w naglowku, ktory z zalozenia jest jawny); bez klucza
klient zatrzymuje sie przed wyslaniem.

Testy: +4 na stream_lines (round-trip, sciezka jawna, fail-closed serwera i
klienta), niezmiennik strukturalny rozszerzony o stream_lines jako droge w dol
(sprawdzone celowym zepsuciem — czerwienieje). Calosc: logika 234 passed /
1 skipped, prezentacja 25 passed.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-07-23 16:10:28 +02:00
20 changed files with 1412 additions and 85 deletions
+368
View File
@@ -0,0 +1,368 @@
"""Tabele pomocnicze horoskopu (LOG-23).
Zbiór wyliczeń, które astrolog czyta „obok" pozycji: bilans żywiołów i jakości,
faza Księżyca, stopnie krytyczne, dzień i godziny planetarne, syzygia prenatalna
oraz podziały (dwunastniki i nawamsa).
Dwie rzeczy wymagają prawdziwego liczenia, nie tabelki:
* **godziny planetarne** — są NIERÓWNE: dzień od wschodu do zachodu Słońca dzieli
się na 12 części, noc osobno. Bez faktycznego wschodu/zachodu wynik byłby
zmyślony, więc szukamy ich numerycznie (przejście wysokości Słońca przez 0°50);
* **syzygia prenatalna** — ostatni nów albo pełnia PRZED urodzeniem; szukamy
wstecz momentu, w którym elongacja Księżyca przechodzi przez 0° lub 180°.
Moduł jest silnik-agnostyczny: potrzebuje tylko `positions()` i `sidereal()`.
"""
from __future__ import annotations
import math
from datetime import datetime, timedelta
from app.engine.formats import SIGNS, in_sign, norm360, sign_index
from app.engine.models import ChartMoment
from app.engine.zodiac import to_equatorial
# --- żywioły i jakości ------------------------------------------------------
ELEMENTS = ["Fire", "Earth", "Air", "Water"]
QUALITIES = ["Cardinal", "Fixed", "Mutable"]
ELEMENT_PL = {"Fire": "Ogień", "Earth": "Ziemia", "Air": "Powietrze", "Water": "Woda"}
QUALITY_PL = {"Cardinal": "Kardynalny", "Fixed": "Stały", "Mutable": "Zmienny"}
CLASSICAL = ["Sun", "Moon", "Mercury", "Venus", "Mars", "Jupiter", "Saturn"]
MODERN = CLASSICAL + ["Uranus", "Neptune", "Pluto"]
# --- dzień i godziny planetarne --------------------------------------------
# Kolejność chaldejska: od najwolniejszej do najszybszej planety
CHALDEAN = ["Saturn", "Jupiter", "Mars", "Sun", "Venus", "Mercury", "Moon"]
# Władca dnia wg dnia tygodnia (0 = poniedziałek, jak w datetime.weekday())
WEEKDAY_RULER = ["Moon", "Mars", "Mercury", "Jupiter", "Venus", "Saturn", "Sun"]
# wysokość środka tarczy Słońca przy wschodzie/zachodzie (refrakcja + promień tarczy)
SUNRISE_ALTITUDE = -0.833
def element_of(sign: str) -> str:
return ELEMENTS[SIGNS.index(sign) % 4]
def quality_of(sign: str) -> str:
return QUALITIES[SIGNS.index(sign) % 3]
def tally(positions: list[dict], asc_sign: str | None = None,
modern: bool = True) -> dict:
"""Bilans żywiołów i jakości (LOG-23).
Liczymy w dwóch wariantach naraz, bo szkoły się różnią: 7 planet klasycznych
i 10 z nowożytnymi. Ascendent doliczany osobno — bywa traktowany jak punkt
równorzędny planetom.
"""
wanted = MODERN if modern else CLASSICAL
by_name = {p.get("name"): p for p in positions}
def count(names: list[str], with_asc: bool) -> dict:
elements = dict.fromkeys(ELEMENTS, 0)
qualities = dict.fromkeys(QUALITIES, 0)
used = []
for name in names:
p = by_name.get(name)
if not p or not p.get("sign"):
continue
elements[element_of(p["sign"])] += 1
qualities[quality_of(p["sign"])] += 1
used.append(name)
if with_asc and asc_sign:
elements[element_of(asc_sign)] += 1
qualities[quality_of(asc_sign)] += 1
used.append("Asc")
return {"elements": elements, "qualities": qualities,
"counted": used, "total": len(used)}
classical = count(CLASSICAL, False)
result = {
"classical_7": classical,
"with_modern_10": count(wanted, False),
"classical_7_plus_asc": count(CLASSICAL, True),
"with_modern_10_plus_asc": count(wanted, True),
}
# brakujące żywioły — klasyczne „no air" itd., podstawa pod scoring (LOG-21)
base = result["with_modern_10_plus_asc"]
result["missing_elements"] = [e for e, n in base["elements"].items() if n == 0]
result["missing_qualities"] = [q for q, n in base["qualities"].items() if n == 0]
result["labels"] = {"elements": ELEMENT_PL, "qualities": QUALITY_PL}
return result
# --- faza Księżyca ----------------------------------------------------------
_PHASES = [
(0.0, "New Moon", "Nów"),
(45.0, "Waxing Crescent", "Sierp przybywający"),
(90.0, "First Quarter", "Pierwsza kwadra"),
(135.0, "Waxing Gibbous", "Garb przybywający"),
(180.0, "Full Moon", "Pełnia"),
(225.0, "Waning Gibbous", "Garb ubywający"),
(270.0, "Last Quarter", "Ostatnia kwadra"),
(315.0, "Waning Crescent", "Sierp ubywający"),
]
def moon_phase(sun_lon: float, moon_lon: float) -> dict:
"""Faza Księżyca z elongacji (Księżyc Słońce)."""
angle = norm360(moon_lon - sun_lon)
idx = int(((angle + 22.5) % 360.0) // 45.0)
_, name, name_pl = _PHASES[idx]
illumination = (1.0 - math.cos(math.radians(angle))) / 2.0
return {
"angle": round(angle, 4),
"phase": name,
"phase_pl": name_pl,
"illumination": round(illumination, 4),
"waxing": angle < 180.0,
}
# --- stopnie krytyczne ------------------------------------------------------
# klasyczne stopnie krytyczne zależą od jakości znaku
_CRITICAL = {"Cardinal": (0, 13, 26), "Fixed": (8, 21), "Mutable": (4, 17)}
CRITICAL_ORB = 1.0
def critical_degrees(positions: list[dict]) -> list[dict]:
"""Obiekty stojące na stopniach krytycznych, 0° albo 29° (anaretycznym)."""
out = []
for p in positions:
lon = p.get("decimal")
sign = p.get("sign")
if lon is None or not sign:
continue
deg = norm360(lon) - sign_index(lon) * 30.0
flags = []
for critical in _CRITICAL[quality_of(sign)]:
if abs(deg - critical) <= CRITICAL_ORB:
flags.append(f"stopień krytyczny {critical}° ({QUALITY_PL[quality_of(sign)].lower()})")
if deg >= 29.0:
flags.append("29° — stopień anaretyczny (koniec znaku)")
elif deg < 1.0:
flags.append("0° — wejście w znak")
if flags:
out.append({"name": p.get("name"), "sign": sign,
"in_sign": p.get("in_sign"), "flags": flags})
return out
# --- podziały: dwunastnik i nawamsa ----------------------------------------
def dwadasamsa(lon: float) -> float:
"""12. część (dwadasamsa): znak dzielony na 12 po 2°30, licząc od siebie."""
lon = norm360(lon)
start = sign_index(lon) * 30.0
return norm360(start + (lon - start) * 12.0)
def navamsa(lon: float) -> float:
"""9. część (nawamsa): 108 podziałów po 3°20 liczonych od 0° Barana."""
lon = norm360(lon)
part = int(lon // (30.0 / 9.0))
return norm360((part % 12) * 30.0 + (lon % (30.0 / 9.0)) * 9.0)
def divisional(positions: list[dict]) -> list[dict]:
"""Pozycje w podziałach 12. i 9. — obie tabele naraz."""
out = []
for p in positions:
lon = p.get("decimal")
if lon is None:
continue
d12, d9 = dwadasamsa(lon), navamsa(lon)
out.append({
"name": p.get("name"),
"d12_sign": SIGNS[sign_index(d12)], "d12_in_sign": in_sign(d12),
"d9_sign": SIGNS[sign_index(d9)], "d9_in_sign": in_sign(d9),
})
return out
# --- wschód/zachód Słońca i godziny planetarne ------------------------------
def sun_altitude(engine, moment: ChartMoment) -> float:
"""Wysokość Słońca nad horyzontem [°] dla momentu i miejsca."""
ramc, eps = engine.sidereal(moment)
sun = engine.positions(moment, ["Sun"])[0]
ra, dec = to_equatorial(sun.longitude, sun.latitude, eps)
hour_angle = math.radians(norm360(ramc - ra))
phi, d = math.radians(moment.lat), math.radians(dec)
sin_alt = math.sin(d) * math.sin(phi) + math.cos(d) * math.cos(phi) * math.cos(hour_angle)
return math.degrees(math.asin(max(-1.0, min(1.0, sin_alt))))
def _at(moment: ChartMoment, when: datetime) -> ChartMoment:
return ChartMoment(when_utc=when, lat=moment.lat, lon=moment.lon)
def _crossings(engine, moment: ChartMoment, start: datetime, end: datetime,
step_minutes: int = 20) -> list[tuple[datetime, str]]:
"""Momenty przejścia Słońca przez horyzont w oknie [start, end].
Skan zgrubny + bisekcja — ten sam wzorzec co przy stacjach planet (LOG-03).
"""
out: list[tuple[datetime, str]] = []
step = timedelta(minutes=step_minutes)
t0 = start
f0 = sun_altitude(engine, _at(moment, t0)) - SUNRISE_ALTITUDE
while t0 < end:
t1 = min(t0 + step, end)
f1 = sun_altitude(engine, _at(moment, t1)) - SUNRISE_ALTITUDE
if f0 == 0.0 or (f0 < 0.0) != (f1 < 0.0):
lo, hi, flo = t0, t1, f0
for _ in range(40): # ~sekundowa dokładność
mid = lo + (hi - lo) / 2
fmid = sun_altitude(engine, _at(moment, mid)) - SUNRISE_ALTITUDE
if (flo < 0.0) != (fmid < 0.0):
hi = mid
else:
lo, flo = mid, fmid
out.append((lo + (hi - lo) / 2, "sunrise" if f1 > f0 else "sunset"))
t0, f0 = t1, f1
return out
def planetary_hours(engine, moment: ChartMoment) -> dict | None:
"""Dzień i godziny planetarne w porządku chaldejskim (LOG-23).
Godziny są NIERÓWNE: dzień (wschód→zachód) i noc (zachód→wschód) dzielą się
na 12 części każde. Doba planetarna zaczyna się o WSCHODZIE, nie o północy —
dlatego władcę dnia bierzemy z dnia tygodnia tego wschodu, który otworzył
bieżący okres.
Zwraca None dla dnia polarnego/nocy polarnej, gdzie wschód nie występuje.
"""
now = moment.when_utc
events = _crossings(engine, moment, now - timedelta(hours=30), now + timedelta(hours=30))
if not events:
return None # brak wschodu/zachodu w oknie
before = [e for e in events if e[0] <= now]
after = [e for e in events if e[0] > now]
if not before or not after:
return None
last_time, last_kind = before[-1]
next_time, _ = after[0]
daytime = last_kind == "sunrise"
period_start, period_end = last_time, next_time
# doba planetarna startuje o wschodzie: w nocy to wschód POPRZEDZAJĄCY zachód
day_start = last_time if daytime else next((t for t, k in reversed(before)
if k == "sunrise"), last_time)
length = (period_end - period_start) / 12
index = int((now - period_start) / length)
index = max(0, min(11, index))
day_ruler = WEEKDAY_RULER[day_start.weekday()]
hour_number = index if daytime else index + 12 # 0..23 od wschodu
ruler = CHALDEAN[(CHALDEAN.index(day_ruler) + hour_number) % 7]
hours = []
for i in range(12):
start = period_start + length * i
hours.append({
"index": i + 1,
"ruler": CHALDEAN[(CHALDEAN.index(day_ruler) + (i if daytime else i + 12)) % 7],
"start": start.isoformat(timespec="seconds"),
"end": (start + length).isoformat(timespec="seconds"),
"current": i == index,
})
return {
"day_ruler": day_ruler,
"hour_ruler": ruler,
"hour_number": hour_number + 1,
"daytime": daytime,
"period": "dzień" if daytime else "noc",
"hour_length_minutes": round(length.total_seconds() / 60.0, 2),
"period_start": period_start.isoformat(timespec="seconds"),
"period_end": period_end.isoformat(timespec="seconds"),
"hours": hours,
}
# --- syzygia prenatalna -----------------------------------------------------
def prenatal_syzygy(engine, moment: ChartMoment, max_days: float = 32.0) -> dict | None:
"""Ostatni nów albo pełnia PRZED podanym momentem (LOG-23).
Szukamy wstecz przejścia elongacji przez 0° (nów) lub 180° (pełnia); bierzemy
to, które wypadło później. Cykl trwa ~29,5 dnia, więc okno 32 dni wystarcza.
"""
def elongation(when: datetime) -> float:
pts = {p.name: p.longitude for p in
engine.positions(_at(moment, when), ["Sun", "Moon"])}
return norm360(pts["Moon"] - pts["Sun"])
def signed(when: datetime, target: float) -> float:
"""Odległość od celu w [180, 180] — zeruje się dokładnie w syzygii."""
return ((elongation(when) - target + 180.0) % 360.0) - 180.0
best: tuple[datetime, str] | None = None
for target, kind in ((0.0, "new_moon"), (180.0, "full_moon")):
step = timedelta(hours=6)
t1 = moment.when_utc
f1 = signed(t1, target)
scanned = timedelta()
while scanned < timedelta(days=max_days):
t0 = t1 - step
f0 = signed(t0, target)
if (f0 < 0.0) != (f1 < 0.0) and abs(f0 - f1) < 180.0:
lo, hi, flo = t0, t1, f0
for _ in range(40):
mid = lo + (hi - lo) / 2
fmid = signed(mid, target)
if (flo < 0.0) != (fmid < 0.0):
hi = mid
else:
lo, flo = mid, fmid
found = lo + (hi - lo) / 2
if best is None or found > best[0]:
best = (found, kind)
break
t1, f1 = t0, f0
scanned += step
if best is None:
return None
when, kind = best
pts = {p.name: p.longitude for p in engine.positions(_at(moment, when), ["Sun", "Moon"])}
lon = pts["Sun"] if kind == "new_moon" else pts["Moon"]
return {
"type": kind,
"type_pl": "nów" if kind == "new_moon" else "pełnia",
"when_utc": when.isoformat(timespec="seconds"),
"days_before_birth": round((moment.when_utc - when).total_seconds() / 86400.0, 3),
"sign": SIGNS[sign_index(lon)],
"in_sign": in_sign(lon),
"decimal": round(norm360(lon), 6),
}
# --- złożenie wszystkiego ---------------------------------------------------
def build_tables(engine, moment: ChartMoment, chart: dict,
heavy: bool = True) -> dict:
"""Komplet tabel dla policzonego horoskopu.
`heavy=False` pomija wyliczenia wymagające szukania numerycznego (godziny
planetarne, syzygia) — przydatne, gdy liczy się czas odpowiedzi.
"""
positions = chart.get("positions") or []
by_name = {p.get("name"): p for p in positions}
asc_sign = (chart.get("angles") or {}).get("Asc", {}).get("sign")
out: dict = {
"tally": tally(positions, asc_sign),
"critical_degrees": critical_degrees(positions),
"divisional": divisional(positions),
}
if "Sun" in by_name and "Moon" in by_name:
out["moon_phase"] = moon_phase(by_name["Sun"]["decimal"], by_name["Moon"]["decimal"])
if heavy:
out["planetary_hours"] = planetary_hours(engine, moment)
out["prenatal_syzygy"] = prenatal_syzygy(engine, moment)
return out
+55
View File
@@ -464,3 +464,58 @@ def open_response_stream(response, link: Link | None) -> Iterator[bytes]:
seq += 1
if buffer:
raise LinkError("strumień urwał się w połowie ramki")
def stream_lines(client, url: str, *, payload, headers: dict[str, str] | None = None,
link: Link | None) -> Iterator[str]:
"""Strumieniowe POST zwracające kolejne NIEPUSTE linie NDJSON — na żywo.
Dla okna postępu: linie muszą docierać w trakcie pracy, nie na końcu, więc
czytamy strumień, a nie całe ciało. Gdy łącze ma klucz, żądanie jest
pieczętowane, a odpowiedź odszyfrowywana ramka po ramce; granice ramek NIE
pokrywają się z granicami linii, więc sklejamy bajty w buforze i tniemy je
dopiero na znakach nowej linii.
Bez klucza zachowuje się jak dotąd (surowy strumień), żeby dev bez sekretów
działał bez zmian.
"""
import json as _json
import httpx as _httpx
request_headers = dict(headers or {})
if link is None:
if encryption_required():
# Ten sam kontrakt co w `call`: nie wypuszczamy jawnego żądania, gdy
# szyfrowanie jest wymagane. Bez tego serwer owszem odrzuca (400), ale
# ciało żądania — tu dane urodzenia — zdążyłoby już pójść w eter.
raise LinkError(
f"{ENV_REQUIRED} jest włączone, ale brak klucza łącza — strumień "
f"NIE został wysłany, żeby jego treść nie poszła jawnym tekstem"
)
with client.stream("POST", url, json=payload, headers=request_headers) as response:
response.raise_for_status()
for text_line in response.iter_lines():
if text_line:
yield text_line
return
path = _httpx.URL(url).path
stamp = stamp_now()
body = frame_out(link.seal(REQUEST, path, stamp, 0, _json.dumps(payload).encode("utf-8")))
request_headers.update({HEADER_ENC: VERSION, HEADER_TS: stamp, "Content-Type": CONTENT_TYPE})
with client.stream("POST", url, content=body, headers=request_headers) as response:
response.raise_for_status()
buffer = bytearray()
for plain in open_response_stream(response, link):
buffer += plain
while True:
nl = buffer.find(b"\n")
if nl < 0:
break
text_line = bytes(buffer[:nl])
del buffer[:nl + 1]
if text_line:
yield text_line.decode("utf-8")
if buffer:
yield bytes(buffer).decode("utf-8")
+31 -2
View File
@@ -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,
+71
View File
@@ -44,6 +44,7 @@ class PositionsRequest(BaseModel):
objects: list[str] | None = None
house_system: str = "whole_sign" # whole_sign | equal | porphyry
stations: bool = False # licz stacje (LOG-03; wolniejsze — root-findy)
tables: bool = False # tabele dodatkowe (LOG-23; szuka wschodu/zachodu)
zodiac: str = "tropical" # LOG-04: tropical | sidereal_{lahiri,fagan_bradley,krishnamurti} | draconic
@@ -75,6 +76,10 @@ def chart_positions(req: PositionsRequest) -> dict:
st = find_stations(engine, moment, p["name"])
if st:
p["stations"] = st
if req.tables:
from app.engine.tables import build_tables
chart["tables"] = build_tables(engine, moment, chart)
return chart
@@ -268,6 +273,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.
+70
View File
@@ -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"