Files
planer/phase0/beispiel-einheiten-kanban-m27hs.md
2026-07-22 16:12:23 +02:00

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