"""Generic streaming kanban engine (no concrete stages — boards define those). Each column is a worker that pulls cards from its input `stage` (the queue = kanban_cards rows WHERE stage = ), processes up to KANBAN_BATCH at a time, and advances them. Streaming columns run continuously; BARRIER columns start only at QUIESCENCE of every stage before them (no active worker + no queued card). SERIAL columns process one package at a time (their processor mutates shared cross-card state). Failure handling: a processor exception (including parse-fails it raises) sends the package's unadvanced cards into exponential backoff (retries++, not_before); after MAX_CARD_RETRIES the card goes to stage 'dead' (dead-letter — visible on the board, requeue-able via API). No card is ever deleted by the engine. Board definitions live in board_inventory.py / board_artefacts.py; run via run_flow(). """ import asyncio import logging import database as db from config import KANBAN_BATCH, MAX_CARD_RETRIES, MAX_CONCURRENT_AGENTS_PER_TOPIC, RETRY_BACKOFF from agents import AgentInfraError log = logging.getLogger("creator.kanban") # How many packages ONE worker keeps in flight at once. A worker no longer blocks on a single # package — it keeps pulling and dispatching until this many run concurrently, so a busy column # fills the agent slots (the per-topic semaphore is the real cap; over-dispatch just queues cheaply). WORKER_INFLIGHT = MAX_CONCURRENT_AGENTS_PER_TOPIC _POLL = 0.3 # seconds between empty-queue polls # Live registry of running flows (topic → Flow), so routes can attach research agents, # report `generating`, and cancel. active_flows: dict[str, "Flow"] = {} class Flow: """Shared runtime state of one topic run: active-task counters per stage + a wakeup event. `producers` counts running research agents (initial + any added live); research counts as done only when ALL producers have finished, so the flow stays awake while extras still search.""" def __init__(self, topic: str, work_dir=None): self.topic = topic self.work_dir = work_dir self.active: dict[str, int] = {} self.producers = 0 self.producer_tag = 0 self.stop = False self.wake = asyncio.Event() self.spawn_research = None # set by the board: () → coroutine adding one more research agent self.state: dict = {} # board-private shared state (embedding caches, one-shot flags …) self.active_cards: set[str] = set() # "board:card_id" currently inside a processor (live display) @property def research_done(self) -> bool: return self.producers <= 0 def add_producer(self): """MUST be called synchronously BEFORE create_task of the producer — otherwise workers can pass their exit check in the gap and never see the new producer (quiescence race).""" self.producers += 1 self.wake.set() def done_producer(self): self.producers -= 1 self.wake.set() def next_tag(self) -> int: self.producer_tag += 1 return self.producer_tag def enter(self, stage: str): self.active[stage] = self.active.get(stage, 0) + 1 def leave(self, stage: str): self.active[stage] = max(0, self.active.get(stage, 0) - 1) self.wake.set() def active_in(self, stages) -> bool: return any(self.active.get(s, 0) > 0 for s in stages) class Stage: """One column: board + stage name + processor. `upstream` (all stages before it, across boards) is filled by chain_stages(). process(cards) gets the pulled package (list of card dicts with decoded payload). barrier: pull only when every upstream stage is quiescent (relational judgements need the full set). gate: extra callable that must be truthy before the stage pulls (works without barrier too — e.g. the consensus gate holds cards until research is done so late reader votes still count). drain: pull the WHOLE queue as one package (global passes like the fragment filter); implies serial.""" def __init__(self, board: str, stage: str, process, *, barrier: bool = False, serial: bool = False, gate=None, drain: bool = False): self.board = board self.stage = stage self.process = process self.barrier = barrier self.serial = serial or drain self.gate = gate self.drain = drain self.upstream: list[str] = [] def chain_stages(stages: list[Stage]) -> list[Stage]: """Fill each stage's upstream = every stage listed before it (list order = flow order). Producers are upstream of everything implicitly via flow.research_done.""" seen: list[str] = [] for s in stages: s.upstream = list(seen) seen.append(s.stage) return stages async def quiescent(flow: Flow, stages) -> bool: """True iff no worker is active in `stages` AND no card is queued in any of them. The barrier/exit condition — must include QUEUED cards, not just active workers, or a worker could exit in a momentary lull while an upstream worker still has work to push down.""" if not stages: return True if flow.active_in(stages): return False return await db.kanban_count(flow.topic, list(stages)) == 0 async def _sleep_wake(flow: Flow): try: await asyncio.wait_for(flow.wake.wait(), timeout=_POLL) except asyncio.TimeoutError: pass flow.wake.clear() async def _fail_package(flow: Flow, spec: Stage, cards: list[dict], error: str): """Backoff/dead-letter for the cards the processor did NOT advance (their stage is unchanged — advanced cards must not be punished for a failure after their move).""" for c in cards: cur = await db.kanban_get_card(flow.topic, spec.board, c["card_id"]) if cur is None or cur["stage"] != spec.stage: continue dead = await db.kanban_fail_card(flow.topic, spec.board, c["card_id"], error, MAX_CARD_RETRIES, RETRY_BACKOFF) if dead: log.warning("kanban %s/%s: card %s → dead (%s)", flow.topic, spec.stage, c["card_id"], error) async def _worker(flow: Flow, spec: Stage, inflight: int, all_stages: list[str]): """Pull cards from spec.stage, run spec.process — keeping up to `inflight` packages running CONCURRENTLY so a busy column fills the agent slots. A barrier worker only pulls when upstream is fully quiescent (and its gate, if any, is open). ANY worker exits only when research is done and the WHOLE flow is quiescent — global instead of per-stage, so a downstream stage that feeds cards back upstream (gap-check → ingest) never strands work. Double-checked over one grace sleep (a producer attached in the lull keeps the flow alive). Double-pull safety: each stage has exactly ONE worker, so an in-memory `claimed` set of card-ids (held while a package runs) keeps concurrent pulls from grabbing the same cards.""" topic = flow.topic claimed: set[str] = set() tasks: set[asyncio.Task] = set() batch = 100_000 if spec.drain else KANBAN_BATCH async def _run(cards): ids = [c["card_id"] for c in cards] flow.enter(spec.stage) flow.active_cards.update(f"{spec.board}:{i}" for i in ids) try: await spec.process(cards) except AgentInfraError as e: # 429/Timeout erschöpft → Lauf pausieren, NICHT dead-letter log.warning("kanban %s/%s: infra-pause: %s", topic, spec.stage, e) flow.state["infra_paused"] = True flow.state["infra_error"] = str(e) flow.stop = True flow.wake.set() except Exception as e: # one bad package must not kill the worker → backoff/dead-letter log.info("kanban %s/%s: %s: %s", topic, spec.stage, type(e).__name__, e) try: await _fail_package(flow, spec, cards, f"{type(e).__name__}: {e}") except Exception: log.exception("kanban %s/%s: fail-handling broke", topic, spec.stage) finally: flow.leave(spec.stage) for i in ids: claimed.discard(i) flow.active_cards.discard(f"{spec.board}:{i}") flow.wake.set() async def _idle_exit() -> bool: return (flow.research_done and not flow.active_in(all_stages) and await db.kanban_count(topic, all_stages) == 0) async def _may_pull() -> bool: if spec.gate is not None and not spec.gate(): return False if not spec.barrier: return True return await quiescent(flow, spec.upstream) try: while not flow.stop: tasks = {t for t in tasks if not t.done()} # Fill the pipeline: pull fresh cards and dispatch until `inflight` packages run. if await _may_pull(): while len(tasks) < inflight: rows = await db.kanban_pull(topic, spec.board, spec.stage, batch + len(claimed)) fresh = [r for r in rows if r["card_id"] not in claimed][:batch] if not fresh: break for r in fresh: claimed.add(r["card_id"]) tasks.add(asyncio.create_task(_run(list(fresh)))) if tasks: # busy → wait for a package to finish, then refill await asyncio.wait(tasks, timeout=_POLL, return_when=asyncio.FIRST_COMPLETED) continue # idle: nothing in flight and nothing pulled if await _idle_exit(): # Real grace sleep (NOT _sleep_wake — the wake event is usually already set # by the last package and would collapse the window to 0ms). A producer # attached during the lull flips research_done and keeps us alive. await asyncio.sleep(_POLL) if await _idle_exit(): return # nothing left and nothing upstream can produce continue await _sleep_wake(flow) finally: for t in tasks: t.cancel() if tasks: await asyncio.gather(*tasks, return_exceptions=True) async def run_flow(flow: Flow, stages: list[Stage], producers=(), set_p=None) -> None: """Run producers + one worker per stage until global quiescence. `producers` are coroutines already counted via flow.add_producer() BEFORE this call (quiescence race). Registers the flow in active_flows for live attach/cancel.""" active_flows[flow.topic] = flow names = [s.stage for s in stages] def _spawn_workers(): return [asyncio.ensure_future(_worker(flow, s, 1 if s.serial else WORKER_INFLIGHT, names)) for s in stages] workers = [asyncio.ensure_future(p) for p in producers] + _spawn_workers() progress = asyncio.create_task(_progress(flow, set_p)) if set_p else None try: while True: await asyncio.gather(*workers, return_exceptions=True) # Restart round: a producer attached exactly as the workers exited (missed even the # grace sleep) leaves live producers or queued cards behind → run the workers again. if flow.stop or (flow.research_done and await quiescent(flow, names)): break workers = _spawn_workers() finally: flow.stop = True if progress: progress.cancel() if active_flows.get(flow.topic) is flow: active_flows.pop(flow.topic, None) async def _progress(flow: Flow, set_p): while not flow.stop: try: counts = await db.kanban_stage_counts(flow.topic) total = sum(n for stages in counts.values() for n in stages.values()) set_p(f"Kanban: {total} Karten im Fluss") except Exception: pass await asyncio.sleep(1.0)