"""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"