Compare commits

2 Commits

Author SHA1 Message Date
team3
b5398f73d2 update 2026-07-01 22:01:32 +02:00
team3
fa718b7d6c update 2026-06-30 18:06:06 +02:00
20 changed files with 1720 additions and 650 deletions

View File

@@ -5,3 +5,9 @@ CLAUDE_CODE_OAUTH_TOKEN=
# MiniMax-Provider: API-Key aus der MiniMax-Console (Coding-Plan). # MiniMax-Provider: API-Key aus der MiniMax-Console (Coding-Plan).
MINIMAX_API_KEY= MINIMAX_API_KEY=
# Agent-Parallelität (optional). Zwei verschachtelte Limits, Default je 10 = bisheriges Verhalten.
# Global gilt über ALLE Themen, der Thema-Wert je Thema. Lokal das globale Limit hochsetzen,
# um mehrere Themen parallel mit je 10 Agenten zu fahren (z.B. global 50, Thema 10).
# MAX_CONCURRENT_AGENTS=10
# MAX_CONCURRENT_AGENTS_PER_TOPIC=10

View File

@@ -5,6 +5,7 @@ respective provider fails — the other keeps running unchanged.
""" """
import asyncio import asyncio
import heapq
import logging import logging
import os import os
import re import re
@@ -13,13 +14,26 @@ import signal
import tempfile import tempfile
import time import time
import urllib.request import urllib.request
from contextlib import asynccontextmanager
from pathlib import Path from pathlib import Path
from config import PROVIDERS, DEFAULT_PROVIDER, MAX_CONCURRENT_AGENTS, MAX_CONCURRENT_INTERACTIVE from config import (PROVIDERS, DEFAULT_PROVIDER, MAX_CONCURRENT_AGENTS,
MAX_CONCURRENT_AGENTS_PER_TOPIC, MAX_CONCURRENT_INTERACTIVE)
log = logging.getLogger("creator.agents") log = logging.getLogger("creator.agents")
_active_processes: dict[str, asyncio.subprocess.Process] = {} _active_processes: dict[str, asyncio.subprocess.Process] = {}
_active_started: dict[str, float] = {} # agent_key → wall-clock start (for the live runtime display)
def active_agents(scope_prefix: str | None = None) -> list[dict]:
"""Currently running agents and how long they've been running. Filter by key prefix
(e.g. f"blocks-{topic}-") for one topic. → [{key, runtime}] sorted longest-first."""
now = time.time()
out = [{"key": k, "runtime": round(now - t, 1)}
for k, t in list(_active_started.items())
if k in _active_processes and (not scope_prefix or k.startswith(scope_prefix))]
return sorted(out, key=lambda a: -a["runtime"])
# Cancelled scopes (key prefixes, symmetric to kill_process). An agent whose # Cancelled scopes (key prefixes, symmetric to kill_process). An agent whose
# key starts with one of these prefixes aborts BEFORE the spawn — so agents WAITING # key starts with one of these prefixes aborts BEFORE the spawn — so agents WAITING
@@ -41,8 +55,71 @@ def _scope_cancelled(agent_key: str) -> bool:
# Caps the real CLI processes — independent of the pipeline semaphore in # Caps the real CLI processes — independent of the pipeline semaphore in
# generator.py. The acquire happens BEFORE the spawn so that queue wait time # generator.py. The acquire happens BEFORE the spawn so that queue wait time
# does not count against the agent timeout. # does not count against the agent timeout.
_batch_sem = asyncio.Semaphore(MAX_CONCURRENT_AGENTS) class _PrioritySemaphore:
"""asyncio.Semaphore variant: when slots are scarce, the LOWEST priority number is served first
(FIFO within the same priority). Lets earlier pipeline columns grab agents before later ones."""
def __init__(self, value: int):
self._value = value
self._waiters: list = [] # heap of [priority, seq, future]
self._seq = 0
async def acquire(self, priority: int = 100):
if self._value > 0:
self._value -= 1
return
fut = asyncio.get_event_loop().create_future()
entry = [priority, self._seq, fut]
self._seq += 1
heapq.heappush(self._waiters, entry)
try:
await fut # release() hands us the slot directly (no value change)
except BaseException:
entry[2] = None # tombstone so release() skips this dead waiter
if fut.done() and not fut.cancelled():
self.release() # granted just before we were cancelled → pass it on
raise
def release(self):
while self._waiters:
entry = heapq.heappop(self._waiters)
if entry[2] is not None and not entry[2].done():
entry[2].set_result(None) # hand the slot straight to the highest-priority waiter
return
self._value += 1
_batch_sem = _PrioritySemaphore(MAX_CONCURRENT_AGENTS)
_interactive_sem = asyncio.Semaphore(MAX_CONCURRENT_INTERACTIVE) _interactive_sem = asyncio.Semaphore(MAX_CONCURRENT_INTERACTIVE)
# Per-topic caps (lazily created): each topic gets its own priority semaphore of size
# MAX_CONCURRENT_AGENTS_PER_TOPIC, nested INSIDE the global _batch_sem. Priority-based too, so the
# per-topic queue can't undo the global priority when one topic is the only load.
_topic_sems: dict[str, _PrioritySemaphore] = {}
# Earlier kanban columns get the scarce global slot first (smaller = higher priority).
_STAGE_PRIORITY = ("research", "verify", "naming", "small", "dep")
def _agent_priority(key: str) -> int:
for i, tag in enumerate(_STAGE_PRIORITY):
if f"-{tag}-" in key or key.endswith(f"-{tag}"):
return i
return len(_STAGE_PRIORITY) # downstream agents (subblocks/facts/…) after the inventory columns
@asynccontextmanager
async def _batch_gate(scope: str | None, priority: int):
"""Per-topic slot FIRST (fair), then the GLOBAL slot by priority (earlier columns win when
agents are scarce). Order matters — a waiter holds only its per-topic slot while queueing globally."""
topic_sem = _topic_sems.setdefault(scope, _PrioritySemaphore(MAX_CONCURRENT_AGENTS_PER_TOPIC)) if scope else None
if topic_sem is not None:
await topic_sem.acquire(priority)
await _batch_sem.acquire(priority)
try:
yield
finally:
_batch_sem.release()
if topic_sem is not None:
topic_sem.release()
# Serialize OpenCode starts: processes starting simultaneously collide on the # Serialize OpenCode starts: processes starting simultaneously collide on the
# internal session DB ("database is locked", exit after <1s). The short # internal session DB ("database is locked", exit after <1s). The short
@@ -102,6 +179,7 @@ def kill_process(agent_key_prefix: str) -> None:
for key, process in list(_active_processes.items()): for key, process in list(_active_processes.items()):
if process.returncode is not None: # clean up dead entries while iterating if process.returncode is not None: # clean up dead entries while iterating
_active_processes.pop(key, None) _active_processes.pop(key, None)
_active_started.pop(key, None)
continue continue
if key.startswith(agent_key_prefix): if key.startswith(agent_key_prefix):
log.debug("kill agent %s", key) log.debug("kill agent %s", key)
@@ -116,6 +194,8 @@ async def run_agent(
role: str = "fast", role: str = "fast",
capabilities: str = "none", capabilities: str = "none",
lane: str = "batch", lane: str = "batch",
scope: str | None = None,
on_line=None,
) -> tuple[int, str, str]: ) -> tuple[int, str, str]:
if _scope_cancelled(agent_key): # before queueing: don't even enter the queue if _scope_cancelled(agent_key): # before queueing: don't even enter the queue
return 1, "", "cancelled" return 1, "", "cancelled"
@@ -123,16 +203,16 @@ async def run_agent(
return 1, "", f"Unknown provider: {provider}" return 1, "", f"Unknown provider: {provider}"
if shutil.which(PROVIDERS[provider]["cli"]) is None: if shutil.which(PROVIDERS[provider]["cli"]) is None:
return 1, "", f"CLI '{PROVIDERS[provider]['cli']}' not installed (provider: {provider})" return 1, "", f"CLI '{PROVIDERS[provider]['cli']}' not installed (provider: {provider})"
sem = _interactive_sem if lane == "interactive" else _batch_sem gate = _interactive_sem if lane == "interactive" else _batch_gate(scope, _agent_priority(agent_key))
async with sem: async with gate:
if _scope_cancelled(agent_key): # after the acquire: cancelled in the queue → no spawn if _scope_cancelled(agent_key): # after the acquire: cancelled in the queue → no spawn
return 1, "", "cancelled" return 1, "", "cancelled"
if PROVIDERS[provider]["cli"] == "opencode": if PROVIDERS[provider]["cli"] == "opencode":
return await _run_opencode(agent_key, prompt, timeout, provider, role, capabilities) return await _run_opencode(agent_key, prompt, timeout, provider, role, capabilities, on_line=on_line)
return await _run_claude_cli(agent_key, prompt, timeout, role, capabilities) return await _run_claude_cli(agent_key, prompt, timeout, role, capabilities)
async def _communicate(agent_key: str, cmd: list[str], stdin_data: bytes | None, timeout: int, stagger: bool = False) -> tuple[int, str, str]: async def _communicate(agent_key: str, cmd: list[str], stdin_data: bytes | None, timeout: int, stagger: bool = False, on_line=None) -> tuple[int, str, str]:
start = time.monotonic() start = time.monotonic()
async def spawn(): async def spawn():
@@ -151,12 +231,29 @@ async def _communicate(agent_key: str, cmd: list[str], stdin_data: bytes | None,
else: else:
process = await spawn() process = await spawn()
_active_processes[agent_key] = process _active_processes[agent_key] = process
_active_started[agent_key] = time.time()
try: try:
try: try:
stdout, stderr = await asyncio.wait_for( if on_line is not None:
process.communicate(input=stdin_data), # Streaming path: read stdout line by line, hand each raw line to on_line LIVE.
timeout=timeout, out_chunks: list[str] = []
) async def _pump():
async for raw in process.stdout:
s = raw.decode("utf-8", errors="replace")
out_chunks.append(s)
try:
on_line(s)
except Exception:
log.debug("on_line callback failed", exc_info=True)
await asyncio.wait_for(_pump(), timeout=timeout)
await process.wait()
stderr_b = await process.stderr.read()
stdout, stderr = "".join(out_chunks).encode("utf-8"), stderr_b
else:
stdout, stderr = await asyncio.wait_for(
process.communicate(input=stdin_data),
timeout=timeout,
)
except asyncio.TimeoutError: except asyncio.TimeoutError:
_kill(process) _kill(process)
try: try:
@@ -175,6 +272,7 @@ async def _communicate(agent_key: str, cmd: list[str], stdin_data: bytes | None,
# the NEW process from tracking. # the NEW process from tracking.
if _active_processes.get(agent_key) is process: if _active_processes.get(agent_key) is process:
del _active_processes[agent_key] del _active_processes[agent_key]
_active_started.pop(agent_key, None)
async def _run_claude_cli(agent_key: str, prompt: str, timeout: int, role: str, capabilities: str) -> tuple[int, str, str]: async def _run_claude_cli(agent_key: str, prompt: str, timeout: int, role: str, capabilities: str) -> tuple[int, str, str]:
@@ -187,7 +285,7 @@ async def _run_claude_cli(agent_key: str, prompt: str, timeout: int, role: str,
return await _communicate(agent_key, cmd, prompt.encode("utf-8"), timeout) return await _communicate(agent_key, cmd, prompt.encode("utf-8"), timeout)
async def _run_opencode(agent_key: str, prompt: str, timeout: int, provider: str, role: str, capabilities: str) -> tuple[int, str, str]: async def _run_opencode(agent_key: str, prompt: str, timeout: int, provider: str, role: str, capabilities: str, on_line=None) -> tuple[int, str, str]:
cfg = PROVIDERS[provider] cfg = PROVIDERS[provider]
# Prompt via temp file instead of argv (ARG_MAX protection for large project prompts) # Prompt via temp file instead of argv (ARG_MAX protection for large project prompts)
with tempfile.NamedTemporaryFile("w", suffix=".md", delete=False, encoding="utf-8", dir=tempfile.gettempdir()) as f: with tempfile.NamedTemporaryFile("w", suffix=".md", delete=False, encoding="utf-8", dir=tempfile.gettempdir()) as f:
@@ -203,9 +301,11 @@ async def _run_opencode(agent_key: str, prompt: str, timeout: int, provider: str
"--dangerously-skip-permissions", "--dangerously-skip-permissions",
"-f", str(prompt_path), "-f", str(prompt_path),
] ]
if on_line is not None:
cmd += ["--format", "json"] # raw JSON events → parsed live by on_line
try: try:
rc, stdout, stderr = await _communicate(agent_key, cmd, None, timeout, stagger=True) rc, stdout, stderr = await _communicate(agent_key, cmd, None, timeout, stagger=True, on_line=on_line)
return rc, _clean_opencode_output(stdout), stderr return rc, (stdout if on_line is not None else _clean_opencode_output(stdout)), stderr
finally: finally:
prompt_path.unlink(missing_ok=True) prompt_path.unlink(missing_ok=True)

File diff suppressed because it is too large Load Diff

View File

@@ -1,3 +1,4 @@
import os
from pathlib import Path from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parent.parent PROJECT_ROOT = Path(__file__).resolve().parent.parent
@@ -41,13 +42,22 @@ EMBEDDING_BLOCK_CAP = 25 # max. titles per block (keep the LLM list short/s
# block context — from this cosine on two are the same statement (checked on aak: ≥0.88 are # block context — from this cosine on two are the same statement (checked on aak: ≥0.88 are
# without exception true duplicates). Conservative 0.90 so different aspects (∈NP ≠ NP-hard) stay separate. # without exception true duplicates). Conservative 0.90 so different aspects (∈NP ≠ NP-hard) stay separate.
EMBEDDING_SUB_DUP = 0.90 EMBEDDING_SUB_DUP = 0.90
# Cosine: two judge-subblocks state the SAME point → one cluster in the clarify majority vote.
# Lower than _DUP because it must merge paraphrases (not just typo-variants). Empirically 0.80 keeps
# distinct aspects (∈NP vs NP-hard) apart while clustering re-wordings of the same fact.
EMBEDDING_SUB_SAME = 0.80
# Cap for concurrent CLI agent processes (across all generations). # Caps for concurrent CLI agent processes (env-overridable). Two nested limits, both always active:
# Own lane for interactive calls (chat, elements) so they don't hang behind # a per-topic cap and a global cap across all topics. Defaults 10/10 = previous behavior (global
# running writers in the queue. # dominates). Locally raise the global cap to actually parallelize across topics (per-topic stays 10).
MAX_CONCURRENT_AGENTS = 10 # Own lane for interactive calls (chat, elements) so they don't hang behind running writers.
MAX_CONCURRENT_AGENTS = int(os.getenv("MAX_CONCURRENT_AGENTS", "10")) # global, all topics
MAX_CONCURRENT_AGENTS_PER_TOPIC = int(os.getenv("MAX_CONCURRENT_AGENTS_PER_TOPIC", "10")) # per topic
MAX_CONCURRENT_INTERACTIVE = 8 MAX_CONCURRENT_INTERACTIVE = 8
# Inventory engine: streaming kanban dataflow (kanban.py) is the default; set "0" for the legacy ER pipeline.
KANBAN_INVENTORY = os.getenv("KANBAN_INVENTORY", "1") != "0"
# Grace window of the consensus races (blocks, guide, OnePager): after the first # Grace window of the consensus races (blocks, guide, OnePager): after the first
# valid result the remaining agents may still become done for this many seconds # valid result the remaining agents may still become done for this many seconds
# (kill only once the minimum is already in). # (kill only once the minimum is already in).

View File

@@ -199,6 +199,59 @@ CREATE TABLE IF NOT EXISTS sub_artefakte (
) )
""" """
# Kanban streaming dataflow for the inventory phase. Cards (titles → chains → blocks) flow through
# columns; `stage` is the current/next column (the queue of a worker = WHERE stage = <predecessor>).
# `stage` is used instead of the reserved word `column`. chain_id/block_id are stable → upsert, not dup.
CREATE_KANBAN_TITLES = """
CREATE TABLE IF NOT EXISTS kanban_titles (
topic TEXT NOT NULL,
title_norm TEXT NOT NULL,
title TEXT NOT NULL,
source TEXT NOT NULL DEFAULT '',
content TEXT NOT NULL DEFAULT '',
stage TEXT NOT NULL DEFAULT 'merge',
updated_at TEXT NOT NULL,
PRIMARY KEY (topic, title_norm)
)
"""
CREATE_KANBAN_CHAINS = """
CREATE TABLE IF NOT EXISTS kanban_chains (
topic TEXT NOT NULL,
chain_id TEXT NOT NULL,
stage TEXT NOT NULL DEFAULT 'chain_verify',
main_title_norm TEXT,
dirty INTEGER NOT NULL DEFAULT 0,
updated_at TEXT NOT NULL,
PRIMARY KEY (topic, chain_id)
)
"""
CREATE_KANBAN_CHAIN_MEMBERS = """
CREATE TABLE IF NOT EXISTS kanban_chain_members (
topic TEXT NOT NULL,
chain_id TEXT NOT NULL,
title_norm TEXT NOT NULL,
PRIMARY KEY (topic, title_norm)
)
"""
CREATE_KANBAN_BLOCKS = """
CREATE TABLE IF NOT EXISTS kanban_blocks (
topic TEXT NOT NULL,
block_id TEXT NOT NULL,
chain_id TEXT,
title TEXT NOT NULL,
source TEXT NOT NULL DEFAULT '',
content TEXT NOT NULL DEFAULT '',
stage TEXT NOT NULL DEFAULT 'small_blocks',
is_small INTEGER NOT NULL DEFAULT 0,
parent_block_id TEXT,
updated_at TEXT NOT NULL,
PRIMARY KEY (topic, block_id)
)
"""
_db: aiosqlite.Connection | None = None _db: aiosqlite.Connection | None = None
@@ -230,6 +283,10 @@ async def init_db():
await db.execute(CREATE_SOURCE) await db.execute(CREATE_SOURCE)
await db.execute(CREATE_GUIDE_OUTLINE) await db.execute(CREATE_GUIDE_OUTLINE)
await db.execute(CREATE_SUB_ARTEFAKTE) await db.execute(CREATE_SUB_ARTEFAKTE)
await db.execute(CREATE_KANBAN_TITLES)
await db.execute(CREATE_KANBAN_CHAINS)
await db.execute(CREATE_KANBAN_CHAIN_MEMBERS)
await db.execute(CREATE_KANBAN_BLOCKS)
try: # migration for existing DBs without the step column try: # migration for existing DBs without the step column
await db.execute("ALTER TABLE guides ADD COLUMN step INTEGER") await db.execute("ALTER TABLE guides ADD COLUMN step INTEGER")
except aiosqlite.OperationalError: except aiosqlite.OperationalError:
@@ -706,6 +763,144 @@ async def delete_blocks(topic: str) -> None:
await db.commit() await db.commit()
# ── Kanban streaming dataflow (inventory) ───────────────────────────────────────
# Generic stage helpers. `stage` is the queue key: a worker pulls WHERE stage = <its input stage>.
_KANBAN_ID = {"kanban_titles": "title_norm", "kanban_chains": "chain_id", "kanban_blocks": "block_id"}
async def kanban_pull(topic: str, table: str, stage: str, limit: int) -> list[dict]:
"""Oldest `limit` cards sitting in `stage` (FIFO via updated_at)."""
idc = _KANBAN_ID[table] # validates table name
db = await get_db()
cursor = await db.execute(
f"SELECT * FROM {table} WHERE topic = ? AND stage = ? ORDER BY updated_at LIMIT ?", (topic, stage, limit))
rows = await cursor.fetchall()
return [_row_to_dict(row, cursor) for row in rows]
async def kanban_count(topic: str, table: str, stages) -> int:
"""How many cards sit in any of `stages` (str or list) — for queue length / quiescence."""
_ = _KANBAN_ID[table]
if isinstance(stages, str):
stages = [stages]
if not stages:
return 0
db = await get_db()
ph = ",".join("?" * len(stages))
cursor = await db.execute(f"SELECT count(*) FROM {table} WHERE topic = ? AND stage IN ({ph})", (topic, *stages))
return (await cursor.fetchone())[0]
async def kanban_advance(topic: str, table: str, id_val: str, stage: str) -> None:
"""Move a card to `stage` (advance to next column, or back for rework/retraction)."""
idc = _KANBAN_ID[table]
db = await get_db()
await db.execute(f"UPDATE {table} SET stage = ?, updated_at = ? WHERE topic = ? AND {idc} = ?",
(stage, _now(), topic, id_val))
await db.commit()
async def kanban_add_title(topic: str, title_norm: str, title: str, source: str = "", content: str = "") -> bool:
"""Research → titles queue (stage 'merge'). Exact dupes are dropped (PK conflict). → True if new."""
db = await get_db()
cursor = await db.execute(
"""INSERT INTO kanban_titles (topic, title_norm, title, source, content, stage, updated_at)
VALUES (?, ?, ?, ?, ?, 'merge', ?) ON CONFLICT(topic, title_norm) DO NOTHING""",
(topic, title_norm, title, source, content, _now()))
await db.commit()
return cursor.rowcount > 0
async def kanban_upsert_chain(topic: str, chain_id: str, stage: str, main_title_norm: str | None = None,
dirty: int = 0) -> None:
db = await get_db()
await db.execute(
"""INSERT INTO kanban_chains (topic, chain_id, stage, main_title_norm, dirty, updated_at)
VALUES (?, ?, ?, ?, ?, ?)
ON CONFLICT(topic, chain_id) DO UPDATE SET
stage = excluded.stage, main_title_norm = COALESCE(excluded.main_title_norm, kanban_chains.main_title_norm),
dirty = excluded.dirty, updated_at = excluded.updated_at""",
(topic, chain_id, stage, main_title_norm, dirty, _now()))
await db.commit()
async def kanban_set_chain_members(topic: str, chain_id: str, members: list[str]) -> None:
"""Replace the member set of a chain (one title belongs to exactly one chain)."""
db = await get_db()
await db.execute("DELETE FROM kanban_chain_members WHERE topic = ? AND chain_id = ?", (topic, chain_id))
for nm in members:
await db.execute(
"""INSERT INTO kanban_chain_members (topic, chain_id, title_norm) VALUES (?, ?, ?)
ON CONFLICT(topic, title_norm) DO UPDATE SET chain_id = excluded.chain_id""",
(topic, chain_id, nm))
await db.commit()
async def kanban_chain_members(topic: str, chain_id: str) -> list[str]:
db = await get_db()
cursor = await db.execute(
"SELECT title_norm FROM kanban_chain_members WHERE topic = ? AND chain_id = ?", (topic, chain_id))
return [r[0] for r in await cursor.fetchall()]
async def kanban_member_chain(topic: str, title_norm: str) -> str | None:
"""Which chain a title currently belongs to (or None)."""
db = await get_db()
cursor = await db.execute(
"SELECT chain_id FROM kanban_chain_members WHERE topic = ? AND title_norm = ?", (topic, title_norm))
row = await cursor.fetchone()
return row[0] if row else None
async def kanban_upsert_block(topic: str, block_id: str, chain_id: str | None, title: str, source: str = "",
content: str = "", stage: str = "small_blocks", is_small: int = 0,
parent_block_id: str | None = None) -> None:
"""Chain-id-stable block (Filter/Block). Upsert → growing chains overwrite, never duplicate."""
db = await get_db()
await db.execute(
"""INSERT INTO kanban_blocks (topic, block_id, chain_id, title, source, content, stage, is_small, parent_block_id, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(topic, block_id) DO UPDATE SET
chain_id = excluded.chain_id, title = excluded.title, source = excluded.source,
content = excluded.content, stage = excluded.stage, is_small = excluded.is_small,
parent_block_id = excluded.parent_block_id, updated_at = excluded.updated_at""",
(topic, block_id, chain_id, title, source, content, stage, is_small, parent_block_id, _now()))
await db.commit()
async def kanban_titles_by_norm(topic: str) -> dict[str, dict]:
"""All titles of a topic keyed by title_norm (the candidate universe for chaining)."""
db = await get_db()
cursor = await db.execute("SELECT * FROM kanban_titles WHERE topic = ?", (topic,))
rows = await cursor.fetchall()
return {(d := _row_to_dict(row, cursor))["title_norm"]: d for row in rows}
async def kanban_all_blocks(topic: str) -> list[dict]:
db = await get_db()
cursor = await db.execute("SELECT * FROM kanban_blocks WHERE topic = ?", (topic,))
rows = await cursor.fetchall()
return [_row_to_dict(row, cursor) for row in rows]
async def kanban_stage_counts(topic: str) -> dict[str, int]:
"""{stage: count} across all kanban tables — for the live board / quiescence."""
db = await get_db()
out: dict[str, int] = {}
for table in ("kanban_titles", "kanban_chains", "kanban_blocks"):
cursor = await db.execute(f"SELECT stage, count(*) FROM {table} WHERE topic = ? GROUP BY stage", (topic,))
for stage, n in await cursor.fetchall():
out[stage] = out.get(stage, 0) + n
return out
async def kanban_reset(topic: str) -> None:
db = await get_db()
for table in ("kanban_titles", "kanban_chains", "kanban_chain_members", "kanban_blocks"):
await db.execute(f"DELETE FROM {table} WHERE topic = ?", (topic,))
await db.commit()
async def upsert_subblock(topic: str, block_norm: str, sub_norm: str, block: str, sub_title: str) -> None: async def upsert_subblock(topic: str, block_norm: str, sub_norm: str, block: str, sub_title: str) -> None:
db = await get_db() db = await get_db()
await db.execute( await db.execute(

View File

@@ -469,7 +469,7 @@ async def _generate_sections(
topic=topic, assignment=assignments[i], facts=facts, topic=topic, assignment=assignments[i], facts=facts,
out_path=content_paths[i], extra=_extra(instructions), out_path=content_paths[i], extra=_extra(instructions),
), ),
_timeout("content", chunk_sizes[i]), provider=provider, role="guide", capabilities="full", _timeout("content", chunk_sizes[i]), provider=provider, role="guide", capabilities="full", scope=topic,
) )
for i in pending for i in pending
], writer_count, report, start=writer_count - len(pending)) ], writer_count, report, start=writer_count - len(pending))
@@ -509,7 +509,7 @@ async def _generate_sections(
topic=topic, assignment=_assignment_subs(followup_chunks[k], entries, subs_by_title), topic=topic, assignment=_assignment_subs(followup_chunks[k], entries, subs_by_title),
facts=facts, out_path=followup_paths[k], extra=_extra(instructions), facts=facts, out_path=followup_paths[k], extra=_extra(instructions),
), ),
_timeout("content", len(followup_chunks[k][0]["nums"])), provider=provider, role="guide", capabilities="full", _timeout("content", len(followup_chunks[k][0]["nums"])), provider=provider, role="guide", capabilities="full", scope=topic,
) )
for k in followup_pending for k in followup_pending
], len(followup_chunks), report_n, start=len(followup_chunks) - len(followup_pending)) ], len(followup_chunks), report_n, start=len(followup_chunks) - len(followup_pending))
@@ -583,7 +583,7 @@ async def _generate_sections(
), ),
out_path=fix_paths[i], extra=_extra(instructions), out_path=fix_paths[i], extra=_extra(instructions),
), ),
_timeout("content", len(fix_chunks[i])), provider=provider, role="guide", capabilities="full", _timeout("content", len(fix_chunks[i])), provider=provider, role="guide", capabilities="full", scope=topic,
) )
for i in fix_pending for i in fix_pending
], return_exceptions=True) ], return_exceptions=True)
@@ -620,7 +620,7 @@ async def _generate_sections(
contents=content_text(w_chunks[i]), contents=content_text(w_chunks[i]),
spec=spec, out_path=paths[i], extra=_extra(instructions), spec=spec, out_path=paths[i], extra=_extra(instructions),
), ),
_timeout("writer", 1), provider=provider, role="guide", capabilities="files", _timeout("writer", 1), provider=provider, role="guide", capabilities="files", scope=topic,
) )
for i in pending for i in pending
], len(w_chunks), report, start=len(w_chunks) - len(pending)) ], len(w_chunks), report, start=len(w_chunks) - len(pending))
@@ -668,7 +668,7 @@ async def _generate_sections(
topic=topic, format_name=format_name, assignment=_assignment_subs(nw_chunks[k], entries, subs_by_title), topic=topic, format_name=format_name, assignment=_assignment_subs(nw_chunks[k], entries, subs_by_title),
contents=content_text(nw_chunks[k]), spec=spec, out_path=nw_paths[k], extra=_extra(instructions), contents=content_text(nw_chunks[k]), spec=spec, out_path=nw_paths[k], extra=_extra(instructions),
), ),
_timeout("writer", 1), provider=provider, role="guide", capabilities="files", _timeout("writer", 1), provider=provider, role="guide", capabilities="files", scope=topic,
) )
for k in nw_pending for k in nw_pending
], len(nw_chunks), report_nw, start=len(nw_chunks) - len(nw_pending)) ], len(nw_chunks), report_nw, start=len(nw_chunks) - len(nw_pending))
@@ -770,7 +770,7 @@ async def _generate_sections(
tasks=tasks_text(fix_chunks[i], problems_by_num), tasks=tasks_text(fix_chunks[i], problems_by_num),
out_path=fix_paths[i], extra=_extra(instructions), out_path=fix_paths[i], extra=_extra(instructions),
), ),
_timeout("writer", len(fix_chunks[i])), provider=provider, role="guide", capabilities="full", _timeout("writer", len(fix_chunks[i])), provider=provider, role="guide", capabilities="full", scope=topic,
) )
for i in fix_pending for i in fix_pending
], return_exceptions=True) ], return_exceptions=True)

748
backend/kanban.py Normal file
View File

@@ -0,0 +1,748 @@
"""Streaming kanban dataflow for the inventory phase.
Each column is a worker that pulls cards from its input `stage` (the queue), processes up to
KANBAN_BATCH at a time, and advances them to the next stage. Cards: titles → chains → blocks.
Streaming columns run continuously; barrier columns start only at QUIESCENCE of everything before
them (no active worker + empty queues). Verify columns push failures back (rework). The Chain column
re-clusters live: a chain that gains a member is marked dirty and flows back to chain_verify.
Reused from blocks.py (imported lazily-safe — kanban is only imported after blocks is loaded):
embedding clustering, `_pairs_schema`/`_cliques`, `_canonical`, research prompt + file payload.
"""
import asyncio
import json
import uuid
import database as db
import embedding
import blocks
from config import RESEARCH_GRACE, MAX_CONCURRENT_AGENTS_PER_TOPIC
from pipeline import GenContext, run_single_slot, _prompt, _timeout, _log, OK
from textkit import _norm_title, _title, _parse_selection
from jsonio import read_json_file as _json_file
KANBAN_BATCH = 5 # cards a worker pulls per package (micro-batching)
# 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
# Stages whose processor mutates shared cross-card state and MUST run one package at a time.
# Online chain-clustering reads the whole universe + membership; parallel packages would race.
_SERIAL_STAGES = {"chain"}
_ID_COL = {"kanban_titles": "title_norm", "kanban_chains": "chain_id", "kanban_blocks": "block_id"}
CHAIN_CAP = 12 # max members per chain — caps the O(n²) pair-verification blow-up
_POLL = 0.3 # seconds between empty-queue polls
# Stage order. A card's `stage` = the column it waits in (its worker's input).
TITLE_STAGES = ["merge", "chain", "chained"] # 'chained' = consumed into a chain
CHAIN_STAGES = ["chain_verify", "naming", "naming_verify", "chain_filter", "filter_verify", "block_assemble"]
BLOCK_STAGES = ["small_blocks", "small_verify", "dependency", "dependency_verify", "main"]
DONE_CHAIN = "done_chain"
DONE_BLOCK = "done_block"
REJECTED = "rejected" # block dropped by filter_verify (off-topic / noise) — terminal, never mirrored
# Predecessor stages for each barrier (must ALL be quiescent before the barrier worker runs).
_BEFORE_CHAIN_FILTER = ["merge", "chain", "chain_verify", "naming", "naming_verify"]
_BEFORE_BLOCK = _BEFORE_CHAIN_FILTER + ["chain_filter", "filter_verify"]
_BEFORE_MAIN = ["small_blocks", "small_verify", "dependency", "dependency_verify"]
class _Flow:
"""Shared runtime state: active-task counters per stage + a wakeup event. `producers` counts the
running research agents (initial + any added live via the generate button); 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):
self.topic = topic
self.work_dir = work_dir
self.active: dict[str, int] = {}
self.producers = 1 # the initial research agent
self.research_tag = 0
self.stop = False
self.wake = asyncio.Event()
self.spawn_research = None # set by run_kanban: () → coroutine that adds one more research agent
@property
def research_done(self) -> bool:
return self.producers <= 0
def add_producer(self):
self.producers += 1
self.wake.set()
def done_producer(self):
self.producers -= 1
self.wake.set()
def next_tag(self) -> int:
self.research_tag += 1
return self.research_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)
async def queued_in(self, table: str, stages) -> bool:
return await db.kanban_count(self.topic, table, list(stages)) > 0
# ── Research producer ────────────────────────────────────────────────────────────
async def _ingest_titles(topic: str, text: str) -> int:
"""Parse a reader file into kanban_titles (stage 'merge'). Exact dupes drop on the PK. → new count."""
n, seen = 0, set()
for record in _parse_selection(text).values():
title = _title(record)
norm = _norm_title(title)
if not norm or norm in seen:
continue
seen.add(norm)
parts = [t.strip() for t in record.split("")]
source = parts[2] if len(parts) >= 3 else ""
desc = parts[1] if len(parts) >= 2 else ""
if await db.kanban_add_title(topic, norm, title, source, desc):
n += 1
return n
RESEARCH_RUNTIME = 900 # one research agent, one round, ~15 min hard cap — the tail ingests live while it writes
_POLL_RESEARCH = 3 # seconds between live reads of a running research file
# Live registry of running flows, so the "+ research" button can attach another agent to a live run.
_active_flows: dict[str, "_Flow"] = {}
def _extract_text(raw_line: str) -> str:
"""Best-effort: pull assistant/tool text out of ONE opencode `--format json` event line.
Recursively collects every `text`/`content` string — robust to the exact event schema."""
try:
obj = json.loads(raw_line)
except Exception:
return ""
parts: list[str] = []
def _walk(o):
if isinstance(o, dict):
for k, v in o.items():
if k in ("text", "content") and isinstance(v, str):
parts.append(v)
else:
_walk(v)
elif isinstance(o, list):
for v in o:
_walk(v)
_walk(obj)
return "".join(parts)
async def _research_once(ctx: GenContext, files: dict, q: dict, folder, instructions: str, tag: str, flow: "_Flow"):
"""ONE agent searches the topic; its titles go into the merge queue LIVE. Two sources feed the
ingest: the JSON event stream (on_line → text buffer) AND the file the agent writes — whichever
the agent uses, cards stream in immediately (not only after it finishes)."""
work_dir = files["arbeit"]
caps = "files" if folder else "full"
p = work_dir / f"research-{tag}.md"
p.unlink(missing_ok=True)
stop = asyncio.Event()
buf: list[str] = [] # assistant text streamed live from the JSON events
def _on_line(raw: str): # sync, called per stdout line by the agent runner
if (t := _extract_text(raw)):
buf.append(t)
async def _drain() -> bool: # ingest from BOTH event buffer and file (idempotent, dupes drop on PK)
text = "".join(buf)
if (ft := blocks._file_payload(p)):
text += "\n" + ft
return bool(text) and await _ingest_titles(ctx.topic, text)
async def _tail(): # live-ingest loop while the agent runs
while not stop.is_set():
try:
await asyncio.wait_for(stop.wait(), timeout=_POLL_RESEARCH)
except asyncio.TimeoutError:
pass
if await _drain():
flow.wake.set() # new cards → wake the workers
tail = asyncio.create_task(_tail())
try:
await run_single_slot(
ctx, f"research-{tag}", key=f"blocks-{ctx.topic}-research-{tag}",
prompt=blocks._build_research_prompt(ctx.topic, p, instructions, q["type"], folder),
role="quick", capabilities=caps,
payload=(lambda result, p=p: blocks._file_payload(p)),
timeout=RESEARCH_RUNTIME, on_line=_on_line,
)
finally:
stop.set()
await tail
if await _drain(): # final catch-up
flow.wake.set()
_log(ctx.topic, f"Research {tag}: titles → merge queue")
async def _research(ctx: GenContext, files: dict, q: dict, folder, instructions: str, flow: _Flow):
"""The initial research producer (already counted in flow.producers=1)."""
try:
await _research_once(ctx, files, q, folder, instructions, "1", flow)
finally:
flow.done_producer()
async def _extra_research(ctx: GenContext, files: dict, q: dict, folder, instructions: str, flow: _Flow):
"""One more research agent, added live via the generate button. Keeps the flow awake until done."""
flow.add_producer()
try:
await _research_once(ctx, files, q, folder, instructions, f"x{flow.next_tag()}", flow)
finally:
flow.done_producer()
def add_research_agent(topic: str) -> bool:
"""Attach one more research agent to a running flow. → True if a run was live to attach to."""
flow = _active_flows.get(topic)
if flow is None or flow.stop or flow.spawn_research is None:
return False
asyncio.create_task(flow.spawn_research())
return True
# ── Generic worker loop ────────────────────────────────────────────────────────────
async def _quiescent(flow: _Flow, stages) -> bool:
"""True iff no worker is active in `stages` AND no card is queued in any of them (all tables).
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
for tb in ("kanban_titles", "kanban_chains", "kanban_blocks"):
if await db.kanban_count(flow.topic, tb, list(stages)):
return False
return True
async def _worker(flow: _Flow, table: str, in_stage: str, process, upstream, *, barrier=False, inflight=WORKER_INFLIGHT):
"""Pull cards from `in_stage`, run `process` — keeping up to `inflight` packages running CONCURRENTLY
so a busy column fills the agent slots instead of doing one package at a time. `upstream` = all stages
before this one. A barrier worker only pulls when `upstream` is fully quiescent. ANY worker exits only
when research is done, its own queue is empty, AND `upstream` is quiescent (nothing can still arrive).
Double-pull safety: each stage has exactly ONE worker, so an in-memory `claimed` set of card-ids (held
while a package runs) is enough to keep concurrent pulls from grabbing the same cards."""
topic = flow.topic
idc = _ID_COL[table]
claimed: set[str] = set()
tasks: set[asyncio.Task] = set()
async def _run(cards):
ids = [c[idc] for c in cards]
flow.enter(in_stage)
try:
await process(cards)
except Exception as e: # one bad package must not kill the worker
_log(topic, f"worker {in_stage}: {type(e).__name__}: {e}")
finally:
flow.leave(in_stage)
for i in ids:
claimed.discard(i)
flow.wake.set()
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 not barrier or await _quiescent(flow, upstream):
while len(tasks) < inflight:
rows = await db.kanban_pull(topic, table, in_stage, KANBAN_BATCH + len(claimed))
fresh = [r for r in rows if r[idc] not in claimed][:KANBAN_BATCH]
if not fresh:
break
for r in fresh:
claimed.add(r[idc])
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
up_quiet = await _quiescent(flow, upstream)
if (flow.research_done and up_quiet and not flow.active_in([in_stage])
and await db.kanban_count(topic, table, in_stage) == 0):
return # nothing left and nothing upstream can produce
await _sleep_wake(flow)
finally:
for t in tasks:
t.cancel()
if tasks:
await asyncio.gather(*tasks, return_exceptions=True)
async def _sleep_wake(flow: _Flow):
try:
await asyncio.wait_for(flow.wake.wait(), timeout=_POLL)
except asyncio.TimeoutError:
pass
flow.wake.clear()
# ── Column processors ──────────────────────────────────────────────────────────────
async def _proc_merge(flow: _Flow, cards):
"""Exact dedup happened at ingest (PK). Merge just advances titles to the chain column."""
for c in cards:
await db.kanban_advance(flow.topic, "kanban_titles", c["title_norm"], "chain")
flow.wake.set()
async def _proc_chain(flow: _Flow, cards):
"""Embedding blocking: for each new title, find the most similar existing title (cosine ≥ floor).
Join its chain (or open a new one), mark the chain dirty → chain_verify. Live-growing clusters."""
topic = flow.topic
by_norm = await db.kanban_titles_by_norm(topic)
# Universe = titles already chained + the new batch (for nearest-neighbour search).
universe = [nm for nm, r in by_norm.items() if r["stage"] in ("chain", "chained")]
if len(universe) < 1:
return
texts = [f"{by_norm[nm]['title']}{by_norm[nm]['content']}" if by_norm[nm]["content"] else by_norm[nm]["title"]
for nm in universe]
sims = await asyncio.to_thread(embedding.embed_sims, texts) if (
embedding and await asyncio.to_thread(embedding.available)) else None
idx = {nm: i for i, nm in enumerate(universe)}
# existing membership
member_chain = {}
for nm in universe:
cid = await _chain_of(topic, nm)
if cid:
member_chain[nm] = cid
touched = set()
for c in cards:
nm = c["title_norm"]
target = None
if sims is not None and nm in idx:
best, bestcos = None, blocks.DEDUP_PAIR_FLOOR
for other in universe:
if other == nm or other not in member_chain and other not in idx:
continue
cos = float(sims[idx[nm]][idx[other]]) if other in idx else -1
if cos >= bestcos and other != nm:
best, bestcos = other, cos
if best is not None:
target = member_chain.get(best)
cid = target or f"c-{uuid.uuid4().hex[:12]}"
members = set(await db.kanban_chain_members(topic, cid))
if target and len(members) >= CHAIN_CAP: # neighbour's chain is full → start a fresh chain
cid = f"c-{uuid.uuid4().hex[:12]}"
members = set()
members.add(nm)
await db.kanban_set_chain_members(topic, cid, sorted(members))
await db.kanban_upsert_chain(topic, cid, "chain_verify", dirty=1)
member_chain[nm] = cid
await db.kanban_advance(topic, "kanban_titles", nm, "chained")
touched.add(cid)
flow.wake.set()
async def _chain_of(topic: str, title_norm: str) -> str | None:
return await db.kanban_member_chain(topic, title_norm)
async def _members_dicts(topic: str, members: list[str]) -> list[dict]:
by = await db.kanban_titles_by_norm(topic)
return [by[m] for m in members if m in by]
async def _proc_chain_verify(ctx: GenContext, flow: _Flow, cards):
"""Pairwise-verify the batch's chains IN PARALLEL (one agent per chain). Failures split off."""
await asyncio.gather(*[_verify_one(ctx, flow, c) for c in cards], return_exceptions=True)
flow.wake.set()
async def _verify_one(ctx: GenContext, flow: _Flow, c):
topic = flow.topic
cid = c["chain_id"]
members = await db.kanban_chain_members(topic, cid)
dicts = await _members_dicts(topic, members)
if len(dicts) <= 1:
await db.kanban_upsert_chain(topic, cid, "naming", dirty=0)
return
# Only the embedding-NEAR candidate pairs (cosine ≥ floor) — NOT all O(n²) pairs. A 12-member
# chain shrinks from 66 pairs to a handful. Transitivity (connected components) does the rest.
nm = [d["title_norm"] for d in dicts]
texts = [f"{d['title']}{d['content']}" if d["content"] else d["title"] for d in dicts]
sims = await asyncio.to_thread(embedding.embed_sims, texts) if (
embedding and await asyncio.to_thread(embedding.available)) else None
if sims is not None:
pairs = [(nm[i], nm[j]) for i in range(len(nm)) for j in range(i + 1, len(nm))
if float(sims[i][j]) >= blocks.DEDUP_PAIR_FLOOR]
else:
pairs = [(a, b) for x, a in enumerate(members) for b in members[x + 1:]]
keep_edges = await _verify_pairs(ctx, flow.work_dir, topic, cid, dicts, pairs)
groups = _components(members, keep_edges) # transitive groups over confirmed near-pairs
groups.sort(key=len, reverse=True)
main = groups[0] if groups else members
await db.kanban_set_chain_members(topic, cid, sorted(main))
await db.kanban_upsert_chain(topic, cid, "naming", dirty=0)
for g in groups[1:]: # the rest split into fresh chains
ncid = f"c-{uuid.uuid4().hex[:12]}"
await db.kanban_set_chain_members(topic, ncid, sorted(g))
await db.kanban_upsert_chain(topic, ncid, "naming", dirty=0)
def _components(members: list[str], edges) -> list[list[str]]:
"""Connected components (union-find) over confirmed pairs. Members without an edge stay alone."""
parent = {m: m for m in members}
def find(x):
while parent[x] != x:
parent[x] = parent[parent[x]]
x = parent[x]
return x
for a, b in edges:
if a in parent and b in parent:
parent[find(a)] = find(b)
comp: dict[str, list[str]] = {}
for m in members:
comp.setdefault(find(m), []).append(m)
return list(comp.values())
async def _verify_pairs(ctx, work_dir, topic, cid, dicts, pairs):
"""Judge the candidate pairs in DEDUP_PAIRS_CHUNK packages, all packages IN PARALLEL → confirmed edges."""
by = {d["title_norm"]: d for d in dicts}
chunks = [pairs[k:k + blocks.DEDUP_PAIRS_CHUNK] for k in range(0, len(pairs), blocks.DEDUP_PAIRS_CHUNK)]
async def _chunk(ci, chunk):
path = work_dir / f"verify-{cid}-{ci}.json"
lines = "\n\n".join(
f"{j + 1}.\nA: {by[a]['title']}{by[a]['content']}\nB: {by[b]['title']}{by[b]['content']}"
for j, (a, b) in enumerate(chunk))
await run_single_slot(
ctx, f"Chain verify {cid}", key=f"blocks-{topic}-verify-{cid}-{ci}",
prompt=_prompt("Blocks-Paar-Filter", topic=topic, pairs=lines, out_path=path),
role="judge", capabilities="files",
payload=lambda result, p=path: blocks._pairs_schema(_json_file(p)),
timeout=_timeout("selection_mapping", len(chunk)))
verdict = blocks._pairs_schema(_json_file(path)) or {}
return [(a, b) for j, (a, b) in enumerate(chunk) if verdict.get(j + 1)]
results = await asyncio.gather(*[_chunk(ci, ch) for ci, ch in enumerate(chunks)], return_exceptions=True)
return [e for r in results if isinstance(r, list) for e in r]
async def _proc_naming(ctx: GenContext, flow: _Flow, cards):
"""Pick the best member title per chain — batch runs IN PARALLEL (one agent per chain)."""
await asyncio.gather(*[_name_one(ctx, flow, c) for c in cards], return_exceptions=True)
flow.wake.set()
async def _name_one(ctx: GenContext, flow: _Flow, c):
topic = flow.topic
cid = c["chain_id"]
members = await db.kanban_chain_members(topic, cid)
dicts = await _members_dicts(topic, members)
if len(dicts) <= 1:
await db.kanban_upsert_chain(topic, cid, "naming_verify",
main_title_norm=(members[0] if members else None), dirty=0)
return
winner = await _choose_title(ctx, flow.work_dir, topic, cid, members, dicts, "Blocks-Naming")
await db.kanban_upsert_chain(topic, cid, "naming_verify", main_title_norm=winner, dirty=0)
async def _proc_naming_verify(ctx: GenContext, flow: _Flow, cards):
"""Second judge checks each title — batch runs IN PARALLEL."""
await asyncio.gather(*[_namecheck_one(ctx, flow, c) for c in cards], return_exceptions=True)
flow.wake.set()
async def _namecheck_one(ctx: GenContext, flow: _Flow, c):
topic = flow.topic
cid = c["chain_id"]
members = await db.kanban_chain_members(topic, cid)
dicts = await _members_dicts(topic, members)
winner = c.get("main_title_norm") or (members[0] if members else None)
if len(dicts) > 1:
winner = await _choose_title(ctx, flow.work_dir, topic, cid, members, dicts, "Blocks-Naming-Check",
current=(members.index(winner) + 1 if winner in members else 1))
await db.kanban_upsert_chain(topic, cid, "chain_filter", main_title_norm=winner, dirty=0)
async def _choose_title(ctx, work_dir, topic, cid, members, dicts, template, current=None):
by = {d["title_norm"]: d for d in dicts}
path = work_dir / f"naming-{cid}.json"
lines = "\n".join(f"{k + 1}. {by[m]['title']}{by[m]['content']}" for k, m in enumerate(members) if m in by)
kw = dict(topic=topic, members=lines, out_path=path)
if current is not None:
kw["current"] = current
await run_single_slot(
ctx, f"Naming {cid}", key=f"blocks-{topic}-naming-{cid}",
prompt=_prompt(template, **kw), role="judge", capabilities="files",
payload=lambda result, p=path: blocks._naming_schema(_json_file(p), len(members)),
timeout=_timeout("selection_mapping", len(members)))
best = blocks._naming_schema(_json_file(path), len(members))
if best is None:
rep = blocks._canonical(dicts, list(range(len(dicts))), set())
w = _norm_title(rep["title"])
return w if w in members else members[0]
return members[best - 1]
async def _proc_chain_filter(flow: _Flow, cards):
"""BARRIER. Reduce each chain to its winner → upsert a block (chain_id stable). → filter_verify."""
topic = flow.topic
by = await db.kanban_titles_by_norm(topic)
for c in cards:
cid = c["chain_id"]
members = await db.kanban_chain_members(topic, cid)
winner = c.get("main_title_norm") if c.get("main_title_norm") in members else (members[0] if members else None)
if not winner or winner not in by:
await db.kanban_upsert_chain(topic, cid, DONE_CHAIN, dirty=0)
continue
w = by[winner]
await db.kanban_upsert_block(topic, f"b-{cid}", cid, w["title"], w["source"], w["content"], stage="filter_verify_b")
await db.kanban_upsert_chain(topic, cid, "filter_verify", dirty=0)
flow.wake.set()
async def _proc_filter_verify(ctx: GenContext, flow: _Flow, cards):
"""An agent confirms each reduced block is a valid, on-topic, self-contained concept. Off-topic /
noise / empty blocks are dropped (→ REJECTED, chain done). IN PARALLEL (one agent per block)."""
await asyncio.gather(*[_filtercheck_one(ctx, flow, c) for c in cards], return_exceptions=True)
flow.wake.set()
async def _filtercheck_one(ctx: GenContext, flow: _Flow, c):
topic = flow.topic
cid = c["chain_id"]
bid = f"b-{cid}"
by = await db.kanban_titles_by_norm(topic)
members = await db.kanban_chain_members(topic, cid)
winner = c.get("main_title_norm") if c.get("main_title_norm") in by else (members[0] if members else None)
async def _drop():
await db.kanban_advance(topic, "kanban_blocks", bid, REJECTED)
await db.kanban_upsert_chain(topic, cid, DONE_CHAIN, dirty=0)
async def _pass():
await db.kanban_advance(topic, "kanban_blocks", bid, "block_assemble_b")
await db.kanban_upsert_chain(topic, cid, "block_assemble", dirty=0)
if not winner or winner not in by: # nothing to verify → drop the empty chain
await _drop()
return
w = by[winner]
path = flow.work_dir / f"filtercheck-{cid}.json"
await run_single_slot(
ctx, f"Filter verify {cid}", key=f"blocks-{topic}-verify-filter-{cid}",
prompt=_prompt("Blocks-Filter-Check", topic=topic, title=w["title"], content=w["content"], out_path=path),
role="judge", capabilities="files",
payload=lambda result, p=path: _keep_schema(_json_file(p)),
timeout=_timeout("selection_mapping", 1))
keep = _keep_schema(_json_file(path))
await (_drop() if keep is False else _pass()) # None (parse fail) → keep, conservative
async def _proc_block(flow: _Flow, cards):
"""BARRIER. Assemble the final block row → small_blocks. (chain card consumed → done.)"""
topic = flow.topic
for c in cards:
cid = c["chain_id"]
await db.kanban_advance(topic, "kanban_blocks", f"b-{cid}", "small_blocks")
await db.kanban_upsert_chain(topic, cid, DONE_CHAIN, dirty=0)
flow.wake.set()
def _small_schema(data, count):
"""{"small": {"1": true, ...}} → {block_index: bool} · else None."""
if not isinstance(data, dict) or not isinstance(data.get("small"), dict):
return None
out = {}
for k, v in data["small"].items():
try:
n = int(k)
except (ValueError, TypeError):
continue
if 1 <= n <= count:
out[n] = str(v).strip().casefold() in ("true", "ja", "yes", "1")
return out or None
def _keep_schema(data):
"""{"keep": true/false} → bool · None when absent/unparseable (caller keeps on None, conservative)."""
if not isinstance(data, dict) or "keep" not in data:
return None
return str(data["keep"]).strip().casefold() in ("true", "ja", "yes", "1")
def _dep_schema(data, count):
"""{"parent": N} → 0..count (0 = standalone) · else None."""
if not isinstance(data, dict):
return None
try:
n = int(data.get("parent"))
except (ValueError, TypeError):
return None
return n if 0 <= n <= count else None
async def _proc_small(ctx: GenContext, flow: _Flow, cards):
"""Judge marks fragment-like blocks (batch). → small_verify."""
topic = flow.topic
path = flow.work_dir / f"small-{cards[0]['block_id']}.json"
lines = "\n".join(f"{i + 1}. {c['title']}{c['content']}" for i, c in enumerate(cards))
await run_single_slot(
ctx, "Small blocks", key=f"blocks-{topic}-small-{cards[0]['block_id']}",
prompt=_prompt("Blocks-Small", topic=topic, blocks=lines, out_path=path),
role="judge", capabilities="files",
payload=lambda result, p=path: _small_schema(_json_file(p), len(cards)),
timeout=_timeout("selection_mapping", len(cards)))
verdict = _small_schema(_json_file(path), len(cards)) or {}
for i, c in enumerate(cards):
is_small = 1 if verdict.get(i + 1) else 0
await db.kanban_upsert_block(topic, c["block_id"], c["chain_id"], c["title"], c["source"], c["content"],
stage="small_verify", is_small=is_small, parent_block_id=c.get("parent_block_id"))
flow.wake.set()
async def _proc_small_verify(ctx: GenContext, flow: _Flow, cards):
"""Second judge re-checks the small flag (consensus): a block stays `small` only if it was marked
small AND this judge also calls it a fragment. Disagreement → keep as a main block (conservative).
→ dependency."""
topic = flow.topic
path = flow.work_dir / f"smallcheck-{cards[0]['block_id']}.json"
lines = "\n".join(f"{i + 1}. {c['title']}{c['content']}" for i, c in enumerate(cards))
await run_single_slot(
ctx, "Small verify", key=f"blocks-{topic}-verify-small-{cards[0]['block_id']}",
prompt=_prompt("Blocks-Small", topic=topic, blocks=lines, out_path=path),
role="judge", capabilities="files",
payload=lambda result, p=path: _small_schema(_json_file(p), len(cards)),
timeout=_timeout("selection_mapping", len(cards)))
verdict = _small_schema(_json_file(path), len(cards)) or {}
for i, c in enumerate(cards):
is_small = 1 if (c["is_small"] and verdict.get(i + 1)) else 0
await db.kanban_upsert_block(topic, c["block_id"], c["chain_id"], c["title"], c["source"], c["content"],
stage="dependency", is_small=is_small, parent_block_id=c.get("parent_block_id"))
flow.wake.set()
async def _proc_dependency(ctx: GenContext, flow: _Flow, cards):
"""For each SMALL block, a judge picks its parent from the full list — batch runs IN PARALLEL."""
topic = flow.topic
parents = [b for b in await db.kanban_all_blocks(topic) if not b["is_small"] and b["stage"] != REJECTED]
plist = "\n".join(f"{i + 1}. {b['title']}" for i, b in enumerate(parents))
await asyncio.gather(*[_dep_one(ctx, flow, c, parents, plist) for c in cards], return_exceptions=True)
flow.wake.set()
async def _dep_one(ctx: GenContext, flow: _Flow, c, parents, plist):
topic = flow.topic
if not c["is_small"] or not parents:
await db.kanban_advance(topic, "kanban_blocks", c["block_id"], "dependency_verify")
return
path = flow.work_dir / f"dep-{c['block_id']}.json"
await run_single_slot(
ctx, "Dependency", key=f"blocks-{topic}-dep-{c['block_id']}",
prompt=_prompt("Blocks-Dependency", topic=topic,
small=f"{c['title']}{c['content']}", parents=plist, out_path=path),
role="judge", capabilities="files",
payload=lambda result, p=path: _dep_schema(_json_file(p), len(parents)),
timeout=_timeout("selection_mapping", len(parents)))
pick = _dep_schema(_json_file(path), len(parents))
parent_id = parents[pick - 1]["block_id"] if pick else None
await db.kanban_upsert_block(topic, c["block_id"], c["chain_id"], c["title"], c["source"], c["content"],
stage="dependency_verify", is_small=c["is_small"], parent_block_id=parent_id)
async def _proc_dependency_verify(flow: _Flow, cards):
"""A small block without a parent is demarked → becomes a main block. → main."""
topic = flow.topic
for c in cards:
is_small = c["is_small"]
if is_small and not c.get("parent_block_id"):
is_small = 0
await db.kanban_upsert_block(topic, c["block_id"], c["chain_id"], c["title"], c["source"], c["content"],
stage="main", is_small=is_small, parent_block_id=c.get("parent_block_id"))
flow.wake.set()
async def _proc_main(flow: _Flow, cards):
"""BARRIER. Finalize non-small blocks → mirror into the legacy `blocks` table as consensus."""
topic = flow.topic
for c in cards:
if not c["is_small"]:
norm = _norm_title(c["title"])
await db.upsert_block(topic, norm, c["title"], c["content"], [c["source"]] if c["source"] else [])
await db.set_block_status(topic, norm, "consensus")
await db.kanban_advance(topic, "kanban_blocks", c["block_id"], DONE_BLOCK)
flow.wake.set()
# ── Orchestration ──────────────────────────────────────────────────────────────────
async def run_kanban(ctx: GenContext, set_p, files: dict, q: dict, folder, instructions: str,
research: bool = True) -> bool:
"""Run the streaming inventory. Returns True when the whole flow reaches quiescence at 'main'.
research=False ("Continue"): process the EXISTING queue without searching new titles. No initial
research producer, producers=0 → research_done is true at once; workers drain the queue and exit.
The +Research button can still attach an agent later via flow.spawn_research."""
topic = ctx.topic
flow = _Flow(topic, files["arbeit"])
flow.spawn_research = lambda: _extra_research(ctx, files, q, folder, instructions, flow)
if not research:
flow.producers = 0 # continue the existing queue, search no new titles
_active_flows[topic] = flow
set_p("Kanban inventory…")
ORDER = ["merge", "chain", "chain_verify", "naming", "naming_verify", "chain_filter",
"filter_verify", "block_assemble", "small_blocks", "small_verify",
"dependency", "dependency_verify", "main"]
up = {s: ORDER[:i] for i, s in enumerate(ORDER)} # upstream = all stages before this one
barriers = {"chain_filter", "block_assemble", "main"}
specs = [
("kanban_titles", "merge", lambda cs: _proc_merge(flow, cs)),
("kanban_titles", "chain", lambda cs: _proc_chain(flow, cs)),
("kanban_chains", "chain_verify", lambda cs: _proc_chain_verify(ctx, flow, cs)),
("kanban_chains", "naming", lambda cs: _proc_naming(ctx, flow, cs)),
("kanban_chains", "naming_verify", lambda cs: _proc_naming_verify(ctx, flow, cs)),
("kanban_chains", "chain_filter", lambda cs: _proc_chain_filter(flow, cs)),
("kanban_chains", "filter_verify", lambda cs: _proc_filter_verify(ctx, flow, cs)),
("kanban_chains", "block_assemble", lambda cs: _proc_block(flow, cs)),
("kanban_blocks", "small_blocks", lambda cs: _proc_small(ctx, flow, cs)),
("kanban_blocks", "small_verify", lambda cs: _proc_small_verify(ctx, flow, cs)),
("kanban_blocks", "dependency", lambda cs: _proc_dependency(ctx, flow, cs)),
("kanban_blocks", "dependency_verify", lambda cs: _proc_dependency_verify(flow, cs)),
("kanban_blocks", "main", lambda cs: _proc_main(flow, cs)),
]
workers = [_research(ctx, files, q, folder, instructions, flow)] if research else []
for table, stage, proc in specs:
workers.append(_worker(flow, table, stage, proc, up[stage], barrier=(stage in barriers),
inflight=(1 if stage in _SERIAL_STAGES else WORKER_INFLIGHT)))
progress = asyncio.create_task(_progress(flow, set_p))
try:
await asyncio.gather(*workers, return_exceptions=True)
finally:
flow.stop = True
progress.cancel()
_active_flows.pop(topic, None)
if ctx.is_cancelled():
return False
n = await db.kanban_count(topic, "kanban_blocks", DONE_BLOCK)
_log(topic, f"Kanban: done — {n} blocks finalized")
return True
async def _progress(flow: _Flow, set_p):
while not flow.stop:
try:
counts = await db.kanban_stage_counts(flow.topic)
total = sum(counts.values())
set_p(f"Kanban: {total} cards in flow")
except Exception:
pass
await asyncio.sleep(1.0)

View File

@@ -33,6 +33,7 @@ class BlocksCreateRequest(BaseModel):
ab_phase: int | None = Field(default=None, ge=1, le=9) # re-run from a coarse phase (position in _phasen(topic), 1-based; up to 9: …outline/questions/artifacts); None = resume/continue without deleting ab_phase: int | None = Field(default=None, ge=1, le=9) # re-run from a coarse phase (position in _phasen(topic), 1-based; up to 9: …outline/questions/artifacts); None = resume/continue without deleting
ab_step: int | None = Field(default=None, ge=0) # re-run from a fine sub-step (0-based index into _blocks_steps); takes precedence over ab_phase ab_step: int | None = Field(default=None, ge=0) # re-run from a fine sub-step (0-based index into _blocks_steps); takes precedence over ab_phase
to_step: int | None = Field(default=None, ge=0) # stop AFTER this fine sub-step (0-based index into _blocks_steps); None = run to the end to_step: int | None = Field(default=None, ge=0) # stop AFTER this fine sub-step (0-based index into _blocks_steps); None = run to the end
research: bool = True # False = "Continue": process the existing queue, search no new titles
class BlocksResetStepRequest(BaseModel): class BlocksResetStepRequest(BaseModel):

View File

@@ -167,7 +167,7 @@ _yesno_schema = _enum_map_schema("relevant", _YESNO) # triage gate
_MAX_RESTARTS = 2 _MAX_RESTARTS = 2
async def _race(topic: str, label: str, slots: list[dict], quorum: int, timeout: int, provider: str, on_update=None, cancelled=None, *, grace: int | None = None) -> list | None: async def _race(topic: str, label: str, slots: list[dict], quorum: int, timeout: int, provider: str, on_update=None, cancelled=None, *, grace: int | None = None, min_runtime: int | None = None, max_runtime: int | None = None) -> list | None:
"""Starts all slots in parallel and collects `quorum` valid results. """Starts all slots in parallel and collects `quorum` valid results.
Slot spec: {key, prompt, role, capabilities, payload}. `payload(result)` Slot spec: {key, prompt, role, capabilities, payload}. `payload(result)`
@@ -180,10 +180,18 @@ async def _race(topic: str, label: str, slots: list[dict], quorum: int, timeout:
a timer of `grace` seconds. After it expires, running agents are only a timer of `grace` seconds. After it expires, running agents are only
killed if the minimum stands — otherwise the race, including restarts, killed if the minimum stands — otherwise the race, including restarts,
keeps running until it stands. Returns: `quorum` to `len(slots)` results. keeps running until it stands. Returns: `quorum` to `len(slots)` results.
`min_runtime` (wall-clock from start): the race does not return before it
elapses while agents are still running — gives them time to search thoroughly.
`max_runtime` (wall-clock from start): hard cap — returns whatever is collected
(or None if nothing), killing the rest. Both default off; only Research sets them.
""" """
attempts = {i: 0 for i in range(len(slots))} attempts = {i: 0 for i in range(len(slots))}
tasks: dict[asyncio.Task, int] = {} tasks: dict[asyncio.Task, int] = {}
loop = asyncio.get_running_loop() loop = asyncio.get_running_loop()
start = loop.time()
min_deadline = start + min_runtime if min_runtime else None
max_deadline = start + max_runtime if max_runtime else None
deadline: float | None = None deadline: float | None = None
def spawn(i: int) -> None: def spawn(i: int) -> None:
@@ -191,6 +199,7 @@ async def _race(topic: str, label: str, slots: list[dict], quorum: int, timeout:
task = asyncio.create_task(run_agent( task = asyncio.create_task(run_agent(
slot["key"], slot["prompt"], timeout, slot["key"], slot["prompt"], timeout,
provider=provider, role=slot["role"], capabilities=slot["capabilities"], provider=provider, role=slot["role"], capabilities=slot["capabilities"],
scope=topic, on_line=slot.get("on_line"),
)) ))
tasks[task] = i tasks[task] = i
@@ -202,12 +211,22 @@ async def _race(topic: str, label: str, slots: list[dict], quorum: int, timeout:
while tasks: while tasks:
if cancelled and cancelled(): if cancelled and cancelled():
return None return None
if deadline is not None and len(results) >= quorum and loop.time() >= deadline: # Hard wall-clock cap: return whatever we have (None if empty), kill the rest.
if max_deadline is not None and loop.time() >= max_deadline:
_log(topic, f"{label}: max runtime {max_runtime}s reached ({len(results)} valid)")
return results or None
min_ok = min_deadline is None or loop.time() >= min_deadline
if deadline is not None and len(results) >= quorum and loop.time() >= deadline and min_ok:
return results return results
# Grace set and minimum reached → only wait for the remaining deadline # Wake up for the earliest relevant deadline (grace, min, or max).
wait_timeout = None waits = []
if deadline is not None and len(results) >= quorum: if deadline is not None and len(results) >= quorum:
wait_timeout = max(0.0, deadline - loop.time()) waits.append(deadline - loop.time())
if min_deadline is not None:
waits.append(min_deadline - loop.time())
if max_deadline is not None:
waits.append(max_deadline - loop.time())
wait_timeout = max(0.0, min(waits)) if waits else None
done, _ = await asyncio.wait(tasks.keys(), return_when=asyncio.FIRST_COMPLETED, timeout=wait_timeout) done, _ = await asyncio.wait(tasks.keys(), return_when=asyncio.FIRST_COMPLETED, timeout=wait_timeout)
if not done: if not done:
continue continue
@@ -234,7 +253,8 @@ async def _race(topic: str, label: str, slots: list[dict], quorum: int, timeout:
_log(topic, f"{label}: first result — grace {grace}s running") _log(topic, f"{label}: first result — grace {grace}s running")
if on_update: if on_update:
on_update(len(results)) on_update(len(results))
if len(results) >= quorum and (grace is None or loop.time() >= deadline): if (len(results) >= quorum and (grace is None or loop.time() >= deadline)
and (min_deadline is None or loop.time() >= min_deadline)):
return results return results
continue continue
@@ -272,13 +292,13 @@ OK, CANCELLED, FAILED = "ok", "cancelled", "failed"
async def run_single_slot( async def run_single_slot(
ctx: GenContext, label: str, *, ctx: GenContext, label: str, *,
key: str, prompt: str, role: str, capabilities: str, payload, timeout: int, key: str, prompt: str, role: str, capabilities: str, payload, timeout: int, on_line=None,
) -> tuple[str, object]: ) -> tuple[str, object]:
"""One agent, one valid result (race with quorum 1). """One agent, one valid result (race with quorum 1).
→ (OK, value) | (CANCELLED, None) | (FAILED, None) → (OK, value) | (CANCELLED, None) | (FAILED, None)
""" """
slots = [{"key": key, "prompt": prompt, "role": role, "capabilities": capabilities, "payload": payload}] slots = [{"key": key, "prompt": prompt, "role": role, "capabilities": capabilities, "payload": payload, "on_line": on_line}]
res = await _race(ctx.topic, label, slots, 1, timeout, ctx.provider, cancelled=ctx.is_cancelled) res = await _race(ctx.topic, label, slots, 1, timeout, ctx.provider, cancelled=ctx.is_cancelled)
if ctx.is_cancelled(): if ctx.is_cancelled():
return CANCELLED, None return CANCELLED, None

View File

@@ -164,7 +164,7 @@ async def create_blocks(req: BlocksCreateRequest):
raise HTTPException(400, "Link must start with http:// or https://.") raise HTTPException(400, "Link must start with http:// or https://.")
qp.parent.mkdir(parents=True, exist_ok=True) qp.parent.mkdir(parents=True, exist_ok=True)
atomic_write_json(qp, {"type": type, "location": location, "spec": req.instructions.strip()}) atomic_write_json(qp, {"type": type, "location": location, "spec": req.instructions.strip()})
asyncio.create_task(generate_blocks(topic, req.instructions.strip(), req.provider, ab_phase=req.ab_phase, ab_step=req.ab_step, to_step=req.to_step)) asyncio.create_task(generate_blocks(topic, req.instructions.strip(), req.provider, ab_phase=req.ab_phase, ab_step=req.ab_step, to_step=req.to_step, research=req.research))
return {"ok": True} return {"ok": True}
@@ -231,6 +231,28 @@ async def get_blocks_uebersicht(topic: str):
return await load_overview(topic) return await load_overview(topic)
@router.get("/blocks/kanban")
async def get_kanban_board(topic: str):
"""Live card counts per kanban column (empty dict when the streaming inventory is not in use)."""
import database as db
return await db.kanban_stage_counts(topic)
@router.get("/blocks/agents")
async def get_active_agents(topic: str):
"""Currently running agents for this topic + their runtime (seconds). Label = key minus prefix."""
from agents import active_agents
prefix = f"blocks-{topic}-"
return [{"label": a["key"][len(prefix):], "runtime": a["runtime"]} for a in active_agents(prefix)]
@router.post("/blocks/research")
async def add_research(topic: str):
"""Attach one more research agent to the running kanban flow (live breadth boost)."""
from kanban import add_research_agent
return {"started": add_research_agent(topic)}
@router.get("/blocks/question-pattern") @router.get("/blocks/question-pattern")
async def get_question_pattern(topic: str, block: str): async def get_question_pattern(topic: str, block: str):
"""Unlocked question patterns of a block (up to the current level; empty = live).""" """Unlocked question patterns of a block (up to the current level; empty = live)."""

View File

@@ -224,12 +224,12 @@ async function handleResetFromStep(step) {
await loadBlocks() await loadBlocks()
} }
async function handleBlocksClick({ instructions, abPhase = null, abStep = null, toStep = null }) { async function handleBlocksClick({ instructions, abPhase = null, abStep = null, toStep = null, research = true }) {
if (!selectedTopic.value) return if (!selectedTopic.value) return
uiError.value = null uiError.value = null
try { try {
// Source is already fixed here; abPhase/abStep set the start, toStep an optional end limit. // Source is already fixed here; abPhase/abStep set the start, toStep an optional end limit.
await apiCreateBausteine(selectedTopic.value, instructions, provider.value, undefined, undefined, abPhase, abStep, toStep) await apiCreateBausteine(selectedTopic.value, instructions, provider.value, undefined, undefined, abPhase, abStep, toStep, research)
} catch (e) { } catch (e) {
uiError.value = e.message uiError.value = e.message
return return
@@ -424,7 +424,7 @@ onMounted(async () => {
@close="mainView = 'detail'" @close="mainView = 'detail'"
@restartFrom="(r) => handleBlocksClick({ instructions: '', abStep: r.from, toStep: r.to })" @restartFrom="(r) => handleBlocksClick({ instructions: '', abStep: r.from, toStep: r.to })"
@resetFrom="handleResetFromStep" @resetFrom="handleResetFromStep"
@restartAll="() => handleBlocksClick({ abPhase: blocks.ready ? 1 : null })" @restartAll="(o) => handleBlocksClick({ research: o?.research ?? false })"
@removeAll="handleResetBlocks" @removeAll="handleResetBlocks"
@cancel="handleCancelBlocks" @cancel="handleCancelBlocks"
/> />

View File

@@ -47,11 +47,11 @@ export async function fetchBlocksStatus(topic) {
return res.json() return res.json()
} }
export async function createBlocks(topic, instructions = '', provider = 'claude', sourceType = 'thema', sourceOrt = '', abPhase = null, abStep = null, toStep = null) { export async function createBlocks(topic, instructions = '', provider = 'claude', sourceType = 'thema', sourceOrt = '', abPhase = null, abStep = null, toStep = null, research = true) {
const res = await fetch(`${BASE}/blocks`, { const res = await fetch(`${BASE}/blocks`, {
method: 'POST', method: 'POST',
headers: { 'Content-Type': 'application/json' }, headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ topic, instructions, provider, source_type: sourceType, source_location: sourceOrt, ab_phase: abPhase, ab_step: abStep, to_step: toStep }), body: JSON.stringify({ topic, instructions, provider, source_type: sourceType, source_location: sourceOrt, ab_phase: abPhase, ab_step: abStep, to_step: toStep, research }),
}) })
return jsonOrThrow(res) return jsonOrThrow(res)
} }
@@ -147,6 +147,24 @@ export async function fetchBlocksOverview(topic) {
return jsonOrThrow(res) return jsonOrThrow(res)
} }
// Live card counts per kanban column ({} when the streaming inventory is not in use).
export async function fetchKanban(topic) {
const res = await fetch(`${BASE}/blocks/kanban?topic=${encodeURIComponent(topic)}`)
return jsonOrThrow(res)
}
// Currently running agents for a topic + their runtime in seconds.
export async function fetchAgents(topic) {
const res = await fetch(`${BASE}/blocks/agents?topic=${encodeURIComponent(topic)}`)
return jsonOrThrow(res)
}
// Attach one more research agent to the running kanban flow.
export async function addResearch(topic) {
const res = await fetch(`${BASE}/blocks/research?topic=${encodeURIComponent(topic)}`, { method: 'POST' })
return jsonOrThrow(res)
}
export async function cancelGuide(id) { export async function cancelGuide(id) {
await fetch(`${BASE}/guides/${id}/cancel`, { method: 'POST' }) await fetch(`${BASE}/guides/${id}/cancel`, { method: 'POST' })
} }

View File

@@ -1,6 +1,6 @@
<script setup> <script setup>
import { ref, computed, watch } from 'vue' import { ref, computed, watch, onUnmounted } from 'vue'
import { fetchBlocksOverview } from '../api.js' import { fetchBlocksOverview, fetchKanban, fetchAgents, addResearch } from '../api.js'
const props = defineProps({ const props = defineProps({
topic: { type: String, required: true }, topic: { type: String, required: true },
@@ -12,6 +12,39 @@ const props = defineProps({
}) })
const emit = defineEmits(['close', 'restartFrom', 'resetFrom', 'restartAll', 'removeAll', 'cancel']) const emit = defineEmits(['close', 'restartFrom', 'resetFrom', 'restartAll', 'removeAll', 'cancel'])
// Live kanban board (streaming inventory). Ordered columns + their card counts.
const KANBAN_COLS = [
['merge', 'Merge'], ['chain', 'Chain'], ['chain_verify', 'Verify'], ['naming', 'Naming'],
['naming_verify', 'Name✓'], ['chain_filter', 'Filter'], ['filter_verify', 'Filter✓'],
['block_assemble', 'Block'], ['small_blocks', 'Small'], ['small_verify', 'Small✓'],
['dependency', 'Dep'], ['dependency_verify', 'Dep✓'], ['main', 'Main'], ['done_block', 'Done'],
]
const kanban = ref({})
const agents = ref([])
const kanbanCols = computed(() => KANBAN_COLS.map(([k, label]) => ({ key: k, label, n: kanban.value[k] || 0 })))
const kanbanActive = computed(() => Object.values(kanban.value).some((n) => n > 0))
function fmtRuntime(s) {
const m = Math.floor(s / 60), sec = Math.floor(s % 60)
return `${m}:${String(sec).padStart(2, '0')}`
}
const researchBusy = ref(false)
async function moreResearch() {
researchBusy.value = true
try { await addResearch(props.topic) } catch { /* ignore */ }
setTimeout(() => { researchBusy.value = false }, 800) // brief debounce against double-clicks
}
let kanbanTimer = null
async function pollKanban() {
try { kanban.value = await fetchKanban(props.topic) } catch { /* ignore */ }
try { agents.value = await fetchAgents(props.topic) } catch { /* ignore */ }
}
watch(() => [props.topic, props.generating], () => {
clearInterval(kanbanTimer)
pollKanban()
if (props.generating) kanbanTimer = setInterval(pollKanban, 1000)
}, { immediate: true })
onUnmounted(() => clearInterval(kanbanTimer))
// Group sub-steps by phase, carrying the global index for the re-run. // Group sub-steps by phase, carrying the global index for the re-run.
const phaseGroups = computed(() => { const phaseGroups = computed(() => {
const out = [] const out = []
@@ -47,7 +80,11 @@ function arm(action, fn) {
if (confirm.value === action) { confirm.value = null; fn() } if (confirm.value === action) { confirm.value = null; fn() }
else confirm.value = action else confirm.value = action
} }
function regenerateFromHere() { const from = startSel.value, to = endSel.value; clearSel(); emit('restartFrom', { from, to }) } // regenerate: re-run ONLY the picked step(s). No end → bound to the start itself (single step);
// with an end → the whole [start, end] range. Later steps stay intact (bounded reset).
function regenerateFromHere() { const from = startSel.value, to = endSel.value ?? startSel.value; clearSel(); emit('restartFrom', { from, to }) }
// continue: run from the start point straight through to the end of the pipeline (full cascade).
function continueFromHere() { const from = startSel.value; clearSel(); emit('restartFrom', { from, to: null }) }
function deleteFromHere() { const from = startSel.value; clearSel(); emit('resetFrom', from) } function deleteFromHere() { const from = startSel.value; clearSel(); emit('resetFrom', from) }
const items = ref([]) const items = ref([])
@@ -109,7 +146,8 @@ const subTotal = computed(() => items.value.reduce((n, b) => n + (b.subblocks?.l
<div class="bk-steps-top"> <div class="bk-steps-top">
<div v-if="progress" class="bk-progress"><span class="bk-progress-dot"></span>{{ progress }}</div> <div v-if="progress" class="bk-progress"><span class="bk-progress-dot"></span>{{ progress }}</div>
<div v-if="!generating" class="bk-global-actions"> <div v-if="!generating" class="bk-global-actions">
<button class="bk-act play" @click="emit('restartAll')">{{ partial ? 'Continue' : ready ? 'Regenerate' : 'Generate' }}</button> <button class="bk-act play" @click="emit('restartAll', { research: false })" title="Process the existing queue — search no new topics">Continue</button>
<button class="bk-act play" @click="emit('restartAll', { research: true })" title="Start one research agent and process the queue">+ Research</button>
<button <button
v-if="ready || partial" v-if="ready || partial"
class="bk-act danger" class="bk-act danger"
@@ -118,9 +156,22 @@ const subTotal = computed(() => items.value.reduce((n, b) => n + (b.subblocks?.l
>{{ confirm === 'remove' ? 'Sure?' : 'Remove' }}</button> >{{ confirm === 'remove' ? 'Sure?' : 'Remove' }}</button>
</div> </div>
<div v-else class="bk-global-actions"> <div v-else class="bk-global-actions">
<button class="bk-act play" :disabled="researchBusy" @click="moreResearch" title="Start one more research agent">+ Research</button>
<button class="bk-act danger" @click="emit('cancel')">Cancel</button> <button class="bk-act danger" @click="emit('cancel')">Cancel</button>
</div> </div>
</div> </div>
<div v-if="generating || kanbanActive" class="bk-kanban">
<div v-for="c in kanbanCols" :key="c.key" class="bk-kcol" :class="{ 'bk-kactive': c.n > 0 }">
<span class="bk-kcount">{{ c.n }}</span>
<span class="bk-klabel">{{ c.label }}</span>
</div>
</div>
<div v-if="agents.length" class="bk-agents">
<span class="bk-agents-label">{{ agents.length }} Agenten aktiv:</span>
<span v-for="a in agents" :key="a.label" class="bk-agent">
{{ a.label }} <span class="bk-agent-time">{{ fmtRuntime(a.runtime) }}</span>
</span>
</div>
<div class="bk-phasen"> <div class="bk-phasen">
<div v-for="g in phaseGroups" :key="g.phase" class="bk-phase"> <div v-for="g in phaseGroups" :key="g.phase" class="bk-phase">
<span class="bk-phase-label">{{ g.phase }}</span> <span class="bk-phase-label">{{ g.phase }}</span>
@@ -139,7 +190,8 @@ const subTotal = computed(() => items.value.reduce((n, b) => n + (b.subblocks?.l
</div> </div>
<div v-if="startSel !== null && !generating" class="bk-step-actions"> <div v-if="startSel !== null && !generating" class="bk-step-actions">
<span class="bk-step-actions-label">From «{{ startLabel }}»<span v-if="endSel !== null"> to «{{ endLabel }}»</span>:</span> <span class="bk-step-actions-label">From «{{ startLabel }}»<span v-if="endSel !== null"> to «{{ endLabel }}»</span>:</span>
<button class="bk-act play" @click="regenerateFromHere"> regenerate</button> <button class="bk-act play" @click="regenerateFromHere" :title="endSel === null ? 'Re-run only this step; later steps stay' : 'Re-run every step in the range; later steps stay'"> regenerate {{ endSel === null ? 'step' : 'range' }}</button>
<button v-if="endSel === null" class="bk-act" @click="continueFromHere" title="Run from here through to the end"> continue</button>
<button class="bk-act danger" :class="{ armed: confirm === 'reset' }" @click="arm('reset', deleteFromHere)">{{ confirm === 'reset' ? 'Sure?' : ' delete all' }}</button> <button class="bk-act danger" :class="{ armed: confirm === 'reset' }" @click="arm('reset', deleteFromHere)">{{ confirm === 'reset' ? 'Sure?' : ' delete all' }}</button>
<button class="bk-act ghost" @click="clearSel">Cancel</button> <button class="bk-act ghost" @click="clearSel">Cancel</button>
</div> </div>
@@ -269,6 +321,28 @@ const subTotal = computed(() => items.value.reduce((n, b) => n + (b.subblocks?.l
/* Header: progress left, global buttons right */ /* Header: progress left, global buttons right */
.bk-steps-top { display: flex; align-items: center; gap: 1rem; min-height: 1.9rem; margin-bottom: 0.7rem; } .bk-steps-top { display: flex; align-items: center; gap: 1rem; min-height: 1.9rem; margin-bottom: 0.7rem; }
.bk-steps-top .bk-progress { margin-bottom: 0; } .bk-steps-top .bk-progress { margin-bottom: 0; }
/* Live kanban board (streaming inventory) */
.bk-kanban { display: flex; flex-wrap: wrap; gap: 0.3rem; margin-bottom: 0.7rem; }
.bk-kcol {
display: flex; flex-direction: column; align-items: center; gap: 1px;
min-width: 3.1rem; padding: 0.3rem 0.4rem;
border: 1px solid var(--border-strong); border-radius: 6px; background: var(--panel);
}
.bk-kcol.bk-kactive { border-color: var(--accent); background: var(--accent-soft); }
.bk-kcount { font-size: 0.95rem; font-weight: 700; color: var(--text); }
.bk-kactive .bk-kcount { color: var(--accent); }
.bk-klabel { font-size: 0.6rem; text-transform: uppercase; letter-spacing: 0.03em; color: var(--text-faint); }
/* Running agents + live runtime */
.bk-agents { display: flex; flex-wrap: wrap; align-items: center; gap: 0.3rem 0.5rem; margin-bottom: 0.7rem; font-size: 0.78rem; }
.bk-agents-label { color: var(--text-muted); font-weight: 600; }
.bk-agent {
display: inline-flex; align-items: center; gap: 0.35rem;
padding: 0.12rem 0.5rem; border: 1px solid var(--accent); border-radius: 10px;
background: var(--accent-soft); color: var(--text);
}
.bk-agent-time { font-variant-numeric: tabular-nums; font-weight: 700; color: var(--accent); }
.bk-global-actions { margin-left: auto; display: flex; gap: 0.4rem; } .bk-global-actions { margin-left: auto; display: flex; gap: 0.4rem; }
/* Action bar for the selected start point */ /* Action bar for the selected start point */

View File

@@ -0,0 +1,17 @@
A small block for the topic "{topic}" may be a sub-topic of a bigger block. Pick the ONE block it belongs under — or 0 if it stands on its own after all.
SMALL BLOCK:
{small}
CANDIDATE PARENT BLOCKS:
{parents}
Rules:
- Pick the number of the block the small block is a detail/sub-step/property of.
- 0 = it is NOT a sub-topic of any of them (it is actually standalone).
- Only pick a parent if the small block clearly belongs INSIDE it.
Write ONLY the JSON file to: {out_path}
Format (the chosen parent number, or 0):
{{"parent": 0}}

View File

@@ -0,0 +1,23 @@
You are filtering the block inventory of a learning guide for the topic "{topic}". Judge the ONE block below.
BLOCK:
{title} — {content}
## Question
Is this a **valid, self-contained learning block** that genuinely belongs to "{topic}"?
- **keep = true** → a real, teachable concept of this topic: a method, definition, syntax element, problem, rule, or feature. The default.
- **keep = false** → drop it, ONLY if it clearly is one of:
- **Off-topic**: not actually about "{topic}" (a stray crawl artifact, a different subject).
- **Noise / meta**: navigation, "Table of Contents", "Dos and Don'ts", a tool/website name, a page section — not a concept you would learn.
- **Empty / degenerate**: title says nothing teachable, or the content is a non-statement.
## Rules
- When in doubt → **keep** (true). Only drop a CLEAR off-topic / noise / empty case.
- Judge by the CONTENT (after the "—"), not only the title.
- Do NOT drop something just because it is narrow or overlaps another block — that is handled elsewhere. Drop only off-topic / noise / empty.
Write ONLY the JSON file to: {out_path}
Format:
{{"keep": true}}

View File

@@ -0,0 +1,15 @@
The numbered entries below all describe the SAME block for the topic "{topic}". Entry number {current} was chosen as the canonical title. Check whether that is the best choice — if another entry is a clearly better canonical name, pick it instead.
MEMBERS:
{members}
Rules:
- Pick an EXISTING entry number — do NOT invent a title.
- Best = most concrete, precise, self-explanatory, established term for the shared concept.
- If the current choice ({current}) is already the best, return it unchanged.
- When in doubt, keep the current choice.
Write ONLY the JSON file to: {out_path}
Format (the best member number, nothing else):
{{"best": {current}}}

View File

@@ -0,0 +1,15 @@
The numbered entries below all describe the SAME block (concept) for the topic "{topic}", just worded differently. Pick the ONE entry whose title is the best canonical name for this block.
MEMBERS:
{members}
Rules:
- Pick an EXISTING entry — do NOT invent a new title or umbrella term.
- Prefer the most CONCRETE, precise, self-explanatory title for the shared concept.
- Prefer the established/standard term (correct spelling, full form over cryptic abbreviation) — but stay concrete, never over-general.
- Avoid reference/placeholder titles ("Satz 7.18", "Punkt 3", "(**)") if a meaningful one exists.
Write ONLY the JSON file to: {out_path}
Format (the chosen member number, nothing else):
{{"best": 1}}

View File

@@ -1,18 +1,23 @@
Two research passes have noted blocks for the topic "{topic}". For EACH pair, decide whether A and B denote the SAME block (the same concept, just worded differently) → **ja**, or whether they are TWO DIFFERENT blocks → **nein**. Research has noted block candidates for the topic "{topic}". For EACH pair, decide whether A and B denote the SAME block (the same concept, just worded differently) → **ja**, or whether they are TWO DIFFERENT blocks → **nein**.
PAIRS: PAIRS:
{pairs} {pairs}
Rules: Rules:
- **Watch the CORE ENTITY first** (the problem/object in question): Clique, Vertex Cover, Independent Set, Dominating Set, Set Cover, FVS, Knapsack … If the entities are DIFFERENT → **nein**, no matter how identical the phrasing. - **Watch the CORE ENTITY first** (the object/concept/operation in question). If the entities are DIFFERENT → **nein**, no matter how identical the phrasing. Examples of **nein** despite near-identical wording:
- Identical phrasing is deceptive. These pairs are **nein** (different entity despite nearly identical wording): - "Plugin-**Lebenszyklus**" ↔ "App-**Lebenszyklus**" (different object: Plugin vs App)
- "Lower Bound **Clique** bzgl. Knoten" ↔ "Lower Bound **Vertex Cover** bzgl. Knoten" - "**install()**-Methode" ↔ "**uninstall()**-Methode" (different operation)
- "Lower Bound Clique bzgl. **Knoten**" ↔ "Lower Bound Clique bzgl. **Kanten**" - "Lower Bound **Clique**" ↔ "Lower Bound **Vertex Cover**" (different problem)
- "Verifizierer für **FVS**" ↔ "Verifizierer für **Knapsack**" - "Present **Perfect**" ↔ "Past **Perfect**" (different tense)
- "**Cliquenproblem**" ↔ "**Vertex-Cover-Problem**" - **ja** ONLY on genuine semantic equivalence: the same concept/operation, just different wording, naming, language or abbreviation. Examples:
- **ja** only on genuine semantic equivalence: same solution to the same problem, the same entity, just different wording/naming (e.g. "SET COVER" ↔ "Mengenüberdeckungsproblem", "Cliquenproblem" ↔ "k-CLIQUE", "List Scheduling" ↔ "LPT-Algorithmus"). - "EntityRepository" ↔ "Repository für CRUD-Operationen am DAL"
- **nein** also for different aspects of the same problem: "Set Cover (Problem)" ↔ "Set Cover ETH-Schranke"; a problem ↔ its reduction to another; a problem ↔ its verifier. - "Set Cover" ↔ "Mengenüberdeckungsproblem"
- When in doubt **nein** — better two separate blocks than wrongly merging two concepts. - "List Scheduling" ↔ "LPT-Algorithmus"
- **nein** for different ASPECTS, properties, parts or methods of the same thing — they are their own blocks (a later step folds true fragments back in):
- a concept ↔ one of its properties/details ("DAL" ↔ "DAL-Versionierung")
- a concept ↔ a single method/step of it ("Plugin-Lebenszyklus" ↔ "install()-Methode")
- a thing ↔ its sub-component, its verifier, its reduction to another thing
- When in doubt → **nein**. Better two separate blocks than wrongly merging two concepts.
Write ONLY the JSON file to: {out_path} Write ONLY the JSON file to: {out_path}

View File

@@ -17,7 +17,8 @@ Rules:
- Write title and description in GERMAN (technical terms/code identifiers stay original). - Write title and description in GERMAN (technical terms/code identifiers stay original).
- Description at most ~12 words. - Description at most ~12 words.
Write ONLY the Markdown file to: {blocks_path} Write the Markdown file to: {blocks_path}
**Stream INCREMENTALLY — your output is read LIVE while you work.** The MOMENT you find a block: (1) print its line in your reply, AND (2) re-write the file with all blocks so far. One line per block, immediately, do NOT wait until the end. Cards appear as soon as a line lands.
Format: EXACTLY one line per block: `N. Title — Kurzbeschreibung — Source` Format: EXACTLY one line per block: `N. Title — Kurzbeschreibung — Source`
The source (3rd segment) MUST be the exact file name or URL of the crawl page the block comes from — it drives the coverage check. The source (3rd segment) MUST be the exact file name or URL of the crawl page the block comes from — it drives the coverage check.

View File

@@ -0,0 +1,15 @@
Below are block candidates for the topic "{topic}". For EACH, decide whether it is a STANDALONE learning block or a SMALL fragment.
BLOCKS:
{blocks}
Rules:
- small = true → a property, detail, sub-step or notation that only makes sense inside another block
(e.g. "install() method" belongs to "Plugin lifecycle"; "Knapsack ∈ NP" belongs to "Knapsack").
- small = false → a self-contained learning unit (its own problem/method/concept).
- When in doubt → false (keep as a main block).
Write ONLY the JSON file to: {out_path}
Format (each block number → true/false):
{{"small": {{"1": true, "2": false}}}}