6.4 KiB
Einheiten: backend/kanban.py
stand: UNBEKANNT
backend/kanban.py::Flow
beschreibung: Verwaltet den Laufzeit-Kontext eines Topics: aktive Zähler, Producer-Status und Koordinations-Events. fakten:
- Zählt aktive Producer und meldet done, wenn alle fertig sind beleg: "return self.producers <= 0"
- Notifybar, wenn sich Producer-Zähler oder aktive Karten ändern beleg: "self.wake.set()"
- Führt Topic-spezifischen Zustand (z.B. infra_paused) beleg: "self.state: dict = {}"
- Registriert laufende Flows global für Cancel/Attach beleg: "active_flows: dict[str, "Flow"] = {}"
- Verfolgt Karten-IDs, die gerade in einem Processor laufen beleg: "self.active_cards: set[str] = set()" kanten:
- nutzt: backend/database.py::kanban_get_card
- wird-genutzt-von: backend/kanban.py::_worker
- wird-genutzt-von: backend/kanban.py::run_flow
backend/kanban.py::Stage
beschreibung: Kapselt eine Spalte des Kanban-Boards: Board/Stage-Name, Processor und Barrier/Serial/Gate/Drain-Parameter. fakten:
- Prozessor wird pro Stage nur einmal gespeichert, von allen Rows gemeinsam genutzt beleg: "self.process = process"
- Serial erzwingt sequentielle Abarbeitung (auch bei drain) beleg: "self.serial = serial or drain"
- Barrier-Wert wird unverändert übernommen beleg: "self.barrier = barrier"
- Upstream wird von chain_stages befüllt, nicht im Konstruktor beleg: "self.upstream: list[str] = []" kanten:
- nutzt: backend/kanban.py::chain_stages
- wird-genutzt-von: backend/kanban.py::_worker
- wird-genutzt-von: backend/kanban.py::run_flow
backend/kanban.py::chain_stages
beschreibung: Baut die upstream-Kette aller Stages — jede Stage sieht alle vorherigen Stages als upstream. fakten:
- Befüllt upstream jeder Stage mit allen davor liegenden Stage-Namen beleg: "s.upstream = list(seen)"
- Gibt die Eingabeliste zurück (Method-Chaining-fähig) beleg: "return stages"
- producers gelten implizit als upstream von allem beleg: "Producers are upstream of everything implicitly" kanten:
- wird-genutzt-von: backend/kanban.py::Stage
backend/kanban.py::quiescent
beschreibung: Prüft, ob alle gegebenen Stages ruhen (keine aktiven Worker UND keine Cards in der DB). fakten:
- Leere Stage-Liste gilt als quiescent beleg: "if not stages: return True"
- Zählt nur Cards in der DB, nicht in-flight im Speicher beleg: "return await db.kanban_count(flow.topic, list(stages)) == 0"
- Berücksichtigt auch Producer-Status via flow.active_in beleg: "if flow.active_in(stages): return False" kanten:
- nutzt: backend/database.py::kanban_count
- nutzt: backend/kanban.py::Flow
- wird-genutzt-von: backend/kanban.py::_worker
backend/kanban.py::_sleep_wake
beschreibung: Blockiert bis Wake-Event oder Poll-Timeout, dann cleart das Event. fakten:
- Wartet maximal _POLL Sekunden beleg: "await asyncio.wait_for(flow.wake.wait(), timeout=_POLL)"
- Timeout bedeutet kein Wake innerhalb von _POLL beleg: "except asyncio.TimeoutError: pass"
- Event wird nach Return gecleart beleg: "flow.wake.clear()" kanten:
- nutzt: backend/kanban.py::Flow
backend/kanban.py::_fail_package
beschreibung: Setzt nicht-advanced Cards auf Backoff oder Dead-Letter nach Processor-Fehler. fakten:
- Überspringt Cards, die bereits weitergingen (ihre Stage hat sich geändert) beleg: "if cur is None or cur["stage"] != spec.stage: continue"
- Nutzt db.kanban_fail_card für Backoff und potentielles Dead-Lettering beleg: "dead = await db.kanban_fail_card(flow.topic, spec.board, c["card_id""
- Loggt wann eine Card nach dead geht beleg: "log.warning("kanban %s/%s: card %s → dead (%s)"" kanten:
- nutzt: backend/database.py::kanban_get_card
- nutzt: backend/database.py::kanban_fail_card
- wird-genutzt-von: backend/kanban.py::_worker
backend/kanban.py::_worker
beschreibung: Dauerläufer pro Stage: zieht Cards, dispatcht bis INFLIGHT concurrent, handhabt Errors und Exit. fakten:
- Zieht max KANBAN_BATCH Cards pro Pull (bzw. 100_000 bei drain) beleg: "batch = 100_000 if spec.drain else KANBAN_BATCH"
- Hält Claims für alle in-flight Cards um Doppel-Pull zu verhindern beleg: "claimed: set[str] = set()"
- AgentInfraError (429/Timeout) pausiert den Flow komplett beleg: "flow.state["infra_paused"] = True"
- Any Exception wird nicht-propagiert — nur _fail_package aufgerufen beleg: "except Exception as e: log.info("kanban %s/%s: %s: %s""
- Prüft Gate (falls gesetzt) VOR dem Pull beleg: "if spec.gate is not None and not spec.gate(): return False"
- Barrier-Worker wartet auf upstream-Quieszenz vor dem Pull beleg: "return await quiescent(flow, spec.upstream)"
- Exit erst wenn research_done UND alle Stages quiescent beleg: "return (flow.research_done and not flow.active_in(all_stages)" kanten:
- ruft-auf: backend/kanban.py::quiescent
- ruft-auf: backend/kanban.py::_fail_package
- ruft-auf: backend/kanban.py::_sleep_wake
- nutzt: backend/database.py::kanban_pull
- nutzt: backend/kanban.py::Flow
- nutzt: backend/kanban.py::Stage
- wird-genutzt-von: backend/kanban.py::run_flow
backend/kanban.py::run_flow
beschreibung: Startet Producer + einen Worker pro Stage, läuft bis globaler Quieszenz oder stop-Flag. fakten:
- Registriert Flow global für Cancel/Attach beleg: "active_flows[flow.topic] = flow"
- Ein Worker pro Stage (serial=1, sonst WORKER_INFLIGHT) beleg: "1 if s.serial else WORKER_INFLIGHT"
- Restart-Schleife fängt Producer an, die genau beim Exit dazukamen beleg: "workers = _spawn_workers()"
- Stoppt Flows bei infra_paused oder bei globaler Quieszenz beleg: "if flow.stop or (flow.research_done and await quiescent(flow, names))"
- Räumt active_flows im finally-Block auf beleg: "active_flows.pop(flow.topic, None)" kanten:
- ruft-auf: backend/kanban.py::_worker
- ruft-auf: backend/kanban.py::_progress
- ruft-auf: backend/kanban.py::quiescent
- nutzt: backend/kanban.py::Flow
- nutzt: backend/kanban.py::Stage
- nutzt: backend/database.py::kanban_stage_counts
- wird-genutzt-von: backend/kanban.py::Flow (spawn_research callback)
backend/kanban.py::_progress
beschreibung: Reportet periodisch die Gesamtzahl aller Karten im Flow an den Callback. fakten:
- Zählt Cards pro Stage via db.kanban_stage_counts beleg: "counts = await db.kanban_stage_counts(flow.topic)"
- Summiert über alle Stages und Boards beleg: "total = sum(n for stages in counts.values() for n in stages.values())"
- Ruft set_p alle 1.0 Sekunden auf beleg: "await asyncio.sleep(1.0)" kanten:
- nutzt: backend/database.py::kanban_stage_counts
- nutzt: backend/kanban.py::Flow
- wird-ausgelöst-von: backend/kanban.py::run_flow