320 lines
16 KiB
Python
320 lines
16 KiB
Python
"""Scan-Board: zerlegt ein Projekt-Repo in Einheiten und Sichten (.planer/-Dateien).
|
|
|
|
Stages (Engine: kanban.py):
|
|
Producer (Chunker, mechanisch) → schneiden (LLM je Chunk) → nachfass (LLM, nur bei
|
|
Abdeckungslücken) → sammeln (mechanisch, drain: Fragmente je Quelldatei mergen,
|
|
Sicht-Karten anlegen) → sichten (LLM je Sicht-Art) → schreiben (mechanisch, drain:
|
|
.planer/ ins Ziel-Repo).
|
|
|
|
Payload-Disziplin: Chunk-Text und Fragmente bleiben in der Karte (~15 KB), ab `sammeln`
|
|
tragen Karten nur noch Staging-Pfade (storage/laeufe/<run>/).
|
|
"""
|
|
import asyncio
|
|
import logging
|
|
import time
|
|
from pathlib import Path
|
|
|
|
import artefakte
|
|
import chunker
|
|
import config
|
|
import database as db
|
|
import gates
|
|
from agents import AgentInfraError, run_agent
|
|
from kanban import Flow, Stage, chain_stages, run_flow
|
|
|
|
log = logging.getLogger("planer.scan")
|
|
|
|
BOARD = "scan"
|
|
SICHT_ARTEN = ("features", "architektur", "flows")
|
|
MAX_NACHFASS_SYMBOLE = 40 # Notbremse gegen endlos wachsende Nachfass-Prompts
|
|
|
|
|
|
def _template(name: str) -> str:
|
|
return (config.TEMPLATES_DIR / f"{name}.md").read_text(encoding="utf-8")
|
|
|
|
|
|
def _prompt(name: str, **ersetzungen: str) -> str:
|
|
text = _template(name)
|
|
for k, v in ersetzungen.items():
|
|
text = text.replace(f"<{k}>", v)
|
|
return text
|
|
|
|
|
|
async def _llm(flow: Flow, stufe: str, card_id: str, rolle: str, prompt: str, label: str) -> str:
|
|
"""Ein Gate-gesicherter LLM-Call: rc!=0 → ValueError (Engine-Backoff);
|
|
AgentInfraError propagiert (Lauf pausiert)."""
|
|
rc, out, err = await run_agent(
|
|
f"scan-{flow.topic}-{stufe}-{card_id}", prompt, config.AGENT_TIMEOUT,
|
|
role=rolle, scope=flow.topic, label=label)
|
|
if rc != 0:
|
|
raise ValueError(f"{stufe} rc={rc}: {(err or out)[:200]}")
|
|
return out
|
|
|
|
|
|
class ScanKontext:
|
|
def __init__(self, wurzel: Path, stand: str, lauf_dir: Path):
|
|
self.wurzel = wurzel
|
|
self.stand = stand # HEAD-Hash des Ziel-Repos beim Scan-Start
|
|
self.lauf_dir = lauf_dir # Staging: storage/laeufe/<run>/
|
|
self.quelle_cache: dict = {} # Datei → Text (Gates lesen viel)
|
|
self.bestand: dict[str, str] = {} # unveränderte Einheiten-Dateien (Re-Scan)
|
|
|
|
|
|
# --- Prozessoren ---------------------------------------------------------------------
|
|
|
|
async def _je_karte(flow: Flow, cards: list[dict], verarbeite) -> None:
|
|
"""Karten eines Pakets parallel verarbeiten; Erfolge advancen, dann ersten Fehler
|
|
werfen (nur unadvancte Karten gehen in den Backoff). AgentInfraError hat Vorrang."""
|
|
ergebnisse = await asyncio.gather(*(verarbeite(c) for c in cards), return_exceptions=True)
|
|
moves = [(c["card_id"], ziel) for c, ziel in zip(cards, ergebnisse) if isinstance(ziel, str)]
|
|
if moves:
|
|
await db.kanban_advance_many(flow.topic, BOARD, moves)
|
|
flow.wake.set()
|
|
fehler = [e for e in ergebnisse if isinstance(e, BaseException)]
|
|
for e in fehler:
|
|
if isinstance(e, AgentInfraError):
|
|
raise e
|
|
if fehler:
|
|
raise fehler[0]
|
|
|
|
|
|
def _mit_befunden(prompt: str, payload: dict) -> str:
|
|
"""Repair-Muster (Creator-Lehre): der Retry bekommt die Gate-Befunde des letzten
|
|
Versuchs — blinde Wiederholung heilt systematische Fehler nicht."""
|
|
befunde = payload.get("befunde")
|
|
if not befunde:
|
|
return prompt
|
|
liste = "\n".join(f"- {b}" for b in befunde[:6])
|
|
return (prompt + "\n\n## WICHTIG: Korrektur\n\nDein vorheriger Versuch scheiterte an "
|
|
"diesen Prüf-Befunden. Behebe GENAU diese Punkte (Zitate müssen wörtlich aus "
|
|
"der Quelle stammen — kopiere sie zeichengenau):\n" + liste + "\n")
|
|
|
|
|
|
async def _befunde_merken(flow: Flow, card: dict, fehler: list[str]) -> None:
|
|
card["payload"]["befunde"] = fehler[:6]
|
|
await db.kanban_set_payload(flow.topic, BOARD, card["card_id"], card["payload"])
|
|
|
|
|
|
async def _proc_schneiden(ctx: ScanKontext, flow: Flow, cards: list[dict]) -> None:
|
|
async def verarbeite(c: dict) -> str:
|
|
p = c["payload"]
|
|
prompt = _mit_befunden(
|
|
_prompt("einheiten", DATEI=p["datei"], QUELLE=p["text"],
|
|
SYMBOLE="\n".join(f"- {s}" for s in p["symbole"]) or "- (keine)"), p)
|
|
out = await _llm(flow, "schneiden", c["card_id"], "schneiden", prompt, p["datei"])
|
|
out, gestrichen = gates.repariere_kanten(out, ctx.wurzel, p["datei"], ctx.quelle_cache)
|
|
out, korrigiert = gates.repariere_belege(out, ctx.wurzel, p["datei"], ctx.quelle_cache)
|
|
fehler, warnungen, einheiten = gates.pruefe_fragment(out, ctx.wurzel, p["datei"],
|
|
ctx.quelle_cache)
|
|
if fehler:
|
|
await _befunde_merken(flow, c, fehler)
|
|
raise ValueError(f"Gates: {'; '.join(fehler[:3])}")
|
|
warnungen += [f"Kante gestrichen: {z}" for z in gestrichen]
|
|
warnungen += [f"Beleg geschnappt: {z}" for z in korrigiert]
|
|
_, zusatz, _ = gates.parse_einheiten(out)
|
|
fehlend = gates.fehlende_symbole(p["symbole"], einheiten, zusatz)
|
|
p["fragmente"] = [out]
|
|
p["fehlend"] = fehlend
|
|
p["warnungen"] = warnungen
|
|
p.pop("befunde", None)
|
|
await db.kanban_set_payload(flow.topic, BOARD, c["card_id"], p)
|
|
return "nachfass" if fehlend else "sammeln"
|
|
await _je_karte(flow, cards, verarbeite)
|
|
|
|
|
|
async def _proc_nachfass(ctx: ScanKontext, flow: Flow, cards: list[dict]) -> None:
|
|
async def verarbeite(c: dict) -> str:
|
|
p = c["payload"]
|
|
fehlend = p["fehlend"][:MAX_NACHFASS_SYMBOLE]
|
|
prompt = _mit_befunden(
|
|
_prompt("nachfass", DATEI=p["datei"], QUELLE=p["text"],
|
|
FEHLEND="\n".join(f"- {s}" for s in fehlend),
|
|
VORHANDEN="\n\n".join(p["fragmente"])), p)
|
|
out = await _llm(flow, "nachfass", c["card_id"], "nachfass", prompt, p["datei"])
|
|
out, _ = gates.repariere_kanten(out, ctx.wurzel, p["datei"], ctx.quelle_cache)
|
|
out, _ = gates.repariere_belege(out, ctx.wurzel, p["datei"], ctx.quelle_cache)
|
|
fehler, _, _ = gates.pruefe_fragment(out, ctx.wurzel, p["datei"], ctx.quelle_cache) \
|
|
if out.strip().startswith("## ") else ([], [], [])
|
|
if fehler:
|
|
await _befunde_merken(flow, c, fehler)
|
|
raise ValueError(f"Gates (Nachfass): {'; '.join(fehler[:3])}")
|
|
# Abdeckung über ALLE Fragmente des Chunks prüfen
|
|
gesamt = "\n\n".join(p["fragmente"] + [out])
|
|
einheiten, zusatz, _ = gates.parse_einheiten(gesamt)
|
|
immer_noch = gates.fehlende_symbole(p["symbole"], einheiten, zusatz)
|
|
if immer_noch:
|
|
raise ValueError(f"Nachfass ließ Symbole offen: {', '.join(immer_noch[:5])}")
|
|
p["fragmente"].append(out)
|
|
p["fehlend"] = []
|
|
p.pop("befunde", None)
|
|
await db.kanban_set_payload(flow.topic, BOARD, c["card_id"], p)
|
|
return "sammeln"
|
|
await _je_karte(flow, cards, verarbeite)
|
|
|
|
|
|
async def _proc_sammeln(ctx: ScanKontext, flow: Flow, cards: list[dict]) -> None:
|
|
"""Drain: alle Chunk-Karten liegen vor. Fragmente je Quelldatei mergen, Staging
|
|
schreiben, Übersichten bauen, eine Karte je Sicht-Art anlegen."""
|
|
je_datei: dict[str, list[dict]] = {}
|
|
for c in cards:
|
|
je_datei.setdefault(c["payload"]["datei"], []).append(c)
|
|
(ctx.lauf_dir / "einheiten").mkdir(parents=True, exist_ok=True)
|
|
ids: set[str] = set()
|
|
kurz_zeilen: list[str] = [] # ID + beschreibung (features)
|
|
kanten_zeilen: list[str] = [] # ID + beschreibung + kanten (architektur/flows)
|
|
for datei, karten in sorted(je_datei.items()):
|
|
karten.sort(key=lambda c: c["payload"]["start"])
|
|
rumpf = artefakte.merge_fragmente([f for c in karten for f in c["payload"]["fragmente"]])
|
|
text = artefakte.einheiten_datei(datei, rumpf, ctx.stand)
|
|
(ctx.lauf_dir / artefakte.datei_pfad(datei)).write_text(text, encoding="utf-8")
|
|
ids |= artefakte.alle_ids(text)
|
|
einheiten, _, _ = gates.parse_einheiten(text)
|
|
for e in einheiten:
|
|
kurz_zeilen.append(f"- {e['id']}: {e['felder'].get('beschreibung', '')}")
|
|
kanten = "; ".join(f"{typ} {ziel}" for _, typ, ziel in e["kanten"])
|
|
kanten_zeilen.append(kurz_zeilen[-1] + (f"\n kanten: {kanten}" if kanten else ""))
|
|
for datei, text in ctx.bestand.items(): # Re-Scan: unveränderte Dateien zählen mit
|
|
ids |= artefakte.alle_ids(text)
|
|
einheiten, _, _ = gates.parse_einheiten(text)
|
|
for e in einheiten:
|
|
kurz_zeilen.append(f"- {e['id']}: {e['felder'].get('beschreibung', '')}")
|
|
kanten = "; ".join(f"{typ} {ziel}" for _, typ, ziel in e["kanten"])
|
|
kanten_zeilen.append(kurz_zeilen[-1] + (f"\n kanten: {kanten}" if kanten else ""))
|
|
(ctx.lauf_dir / "uebersicht-kurz.md").write_text("\n".join(sorted(kurz_zeilen)), encoding="utf-8")
|
|
(ctx.lauf_dir / "uebersicht-kanten.md").write_text("\n".join(sorted(kanten_zeilen)), encoding="utf-8")
|
|
(ctx.lauf_dir / "ids.txt").write_text("\n".join(sorted(ids)), encoding="utf-8")
|
|
for art in SICHT_ARTEN:
|
|
await db.kanban_upsert_card(flow.topic, BOARD, f"sicht-{art}", "sicht", "sichten",
|
|
{"art": art})
|
|
await db.kanban_advance_many(flow.topic, BOARD,
|
|
[(c["card_id"], "chunk-fertig") for c in cards])
|
|
flow.wake.set()
|
|
|
|
|
|
async def _proc_sichten(ctx: ScanKontext, flow: Flow, cards: list[dict]) -> None:
|
|
ids = set((ctx.lauf_dir / "ids.txt").read_text(encoding="utf-8").splitlines())
|
|
|
|
async def verarbeite(c: dict) -> str:
|
|
art = c["payload"]["art"]
|
|
uebersicht = "uebersicht-kurz.md" if art == "features" else "uebersicht-kanten.md"
|
|
prompt = _mit_befunden(
|
|
_prompt(f"sichten-{art}", PROJEKT=flow.topic,
|
|
EINHEITEN=(ctx.lauf_dir / uebersicht).read_text(encoding="utf-8")), c["payload"])
|
|
out = await _llm(flow, "sichten", c["card_id"], "sichten", prompt, art)
|
|
out, gestrichen = gates.repariere_referenzen(out, ids)
|
|
if gestrichen:
|
|
log.info("scan %s: %s — %d unbekannte Referenzen gestrichen", flow.topic, art,
|
|
len(gestrichen))
|
|
fehler = gates.pruefe_referenzen(out, ids)
|
|
if fehler:
|
|
await _befunde_merken(flow, c, fehler)
|
|
raise ValueError(f"Referenzen ({art}): {'; '.join(fehler[:3])}")
|
|
(ctx.lauf_dir / f"sicht-{art}.md").write_text(out, encoding="utf-8")
|
|
return "schreiben"
|
|
await _je_karte(flow, cards, verarbeite)
|
|
|
|
|
|
async def _proc_schreiben(ctx: ScanKontext, flow: Flow, cards: list[dict]) -> None:
|
|
"""Drain nach allen Sichten: .planer/ im Ziel-Repo schreiben."""
|
|
dateien: dict[str, str] = {}
|
|
for p in (ctx.lauf_dir / "einheiten").glob("*.md"):
|
|
dateien[f"einheiten/{p.name}"] = p.read_text(encoding="utf-8")
|
|
for datei, text in ctx.bestand.items():
|
|
dateien.setdefault(artefakte.datei_pfad(datei), text)
|
|
fehlend: list[str] = [] # gescheiterte Sichten: auslassen statt Kaskade (sichtbar degradieren)
|
|
if (ctx.lauf_dir / "sicht-features.md").is_file():
|
|
features = (ctx.lauf_dir / "sicht-features.md").read_text(encoding="utf-8")
|
|
kern, _, bereiche = features.partition("# Bereich:")
|
|
dateien["kern.md"] = kern.strip() + "\n"
|
|
for i, block in enumerate(("# Bereich:" + bereiche).split("# Bereich:")):
|
|
if not block.strip():
|
|
continue
|
|
name = block.splitlines()[0].strip().lower().replace(" ", "-").replace("/", "-")
|
|
dateien[f"features/{name or f'bereich-{i}'}.md"] = "# Bereich:" + block.rstrip() + "\n"
|
|
else:
|
|
fehlend.append("features")
|
|
for art, zieldatei in (("architektur", "architektur.md"), ("flows", "flows.md")):
|
|
if (ctx.lauf_dir / f"sicht-{art}.md").is_file():
|
|
dateien[zieldatei] = (ctx.lauf_dir / f"sicht-{art}.md").read_text(encoding="utf-8")
|
|
else:
|
|
fehlend.append(art)
|
|
if fehlend:
|
|
flow.state["sichten_fehlend"] = fehlend
|
|
log.warning("scan %s: Sichten fehlen (Karten dead): %s", flow.topic, ", ".join(fehlend))
|
|
if not (ctx.wurzel / ".planer" / "pruefung.md").is_file(): # nie Nutzer-Edits überschreiben
|
|
dateien["pruefung.md"] = artefakte.pruefung_geruest(ctx.wurzel, flow.topic)
|
|
artefakte.schreibe_planer(ctx.wurzel, dateien)
|
|
flow.state["geschrieben"] = sorted(dateien)
|
|
await db.kanban_advance_many(flow.topic, BOARD, [(c["card_id"], "fertig") for c in cards])
|
|
flow.wake.set()
|
|
|
|
|
|
# --- Flow-Start ------------------------------------------------------------------------
|
|
|
|
def scan_stages(ctx: ScanKontext, flow: Flow) -> list[Stage]:
|
|
# gate=research_done an den Barriere-Stufen ist PFLICHT (Abnahme-Befund): die Barriere
|
|
# prüft nur die Upstream-Stufen — ohne Gate drained `sammeln`, während der Producer
|
|
# noch Karten einreiht, und die Sichten sehen eine Teil-Übersicht.
|
|
fertig_produziert = lambda: flow.research_done # noqa: E731
|
|
return chain_stages([
|
|
Stage(BOARD, "schneiden", lambda cs: _proc_schneiden(ctx, flow, cs)),
|
|
Stage(BOARD, "nachfass", lambda cs: _proc_nachfass(ctx, flow, cs)),
|
|
Stage(BOARD, "sammeln", lambda cs: _proc_sammeln(ctx, flow, cs),
|
|
barrier=True, drain=True, gate=fertig_produziert),
|
|
Stage(BOARD, "sichten", lambda cs: _proc_sichten(ctx, flow, cs)),
|
|
Stage(BOARD, "schreiben", lambda cs: _proc_schreiben(ctx, flow, cs),
|
|
barrier=True, drain=True, gate=fertig_produziert),
|
|
])
|
|
|
|
|
|
async def run_scan(wurzel: Path, projekt: str, stand: str, set_p=None,
|
|
nur_dateien: set[str] | None = None,
|
|
bestand: dict[str, str] | None = None,
|
|
fortsetzen: bool = False) -> dict:
|
|
"""Scan eines Projekt-Repos. `nur_dateien` beschränkt auf geänderte Dateien (Re-Scan);
|
|
`bestand` = unveränderte Einheiten-Dateien {quelldatei: text}. `fortsetzen=True`
|
|
setzt einen abgebrochenen/fehlgeschlagenen Lauf fort: tote Chunks laufen neu durchs
|
|
LLM, fertige nur durchs mechanische Sammeln (Fragmente stecken im Payload) — setzt
|
|
unveränderten HEAD voraus. → {status, geschrieben, dead}."""
|
|
run_id = f"scan-{time.time_ns()}" # eindeutig auch bei zwei Läufen in derselben Sekunde
|
|
lauf_dir = config.LAEUFE_DIR / projekt / run_id
|
|
lauf_dir.mkdir(parents=True, exist_ok=True)
|
|
ctx = ScanKontext(wurzel, stand, lauf_dir)
|
|
ctx.bestand = bestand or {}
|
|
flow = Flow(projekt)
|
|
db.set_current_run(projekt, run_id)
|
|
|
|
async def producer():
|
|
try:
|
|
if fortsetzen:
|
|
n_neu = n_fertig = 0
|
|
for c in await db.kanban_cards(flow.topic, BOARD, kind="chunk"):
|
|
if c["stage"] == "dead":
|
|
await db.kanban_upsert_card(flow.topic, BOARD, c["card_id"],
|
|
"chunk", "schneiden", None)
|
|
n_neu += 1
|
|
elif c["stage"] == "chunk-fertig":
|
|
await db.kanban_upsert_card(flow.topic, BOARD, c["card_id"],
|
|
"chunk", "sammeln", None)
|
|
n_fertig += 1
|
|
log.info("scan %s: Fortsetzung — %d neu, %d nur sammeln", projekt, n_neu, n_fertig)
|
|
return
|
|
chunks = await asyncio.to_thread(chunker.chunks_fuer_repo, wurzel, nur_dateien)
|
|
for i, chunk in enumerate(chunks):
|
|
await db.kanban_upsert_card(flow.topic, BOARD, f"chunk-{i:04d}", "chunk",
|
|
"schneiden", chunk)
|
|
log.info("scan %s: %d Chunks angelegt", projekt, len(chunks))
|
|
finally:
|
|
flow.done_producer()
|
|
|
|
flow.add_producer() # SYNCHRON vor create_task (Quieszenz-Race)
|
|
try:
|
|
await run_flow(flow, scan_stages(ctx, flow), [producer()], set_p)
|
|
finally:
|
|
db.set_current_run(projekt, None)
|
|
dead = await db.kanban_dead(projekt)
|
|
status = "pausiert" if flow.state.get("infra_paused") else ("fehler" if dead else "ok")
|
|
return {"status": status, "infra_error": flow.state.get("infra_error", ""),
|
|
"geschrieben": flow.state.get("geschrieben", []),
|
|
"dead": [c["card_id"] for c in dead]}
|