This commit is contained in:
team3
2026-06-22 13:44:33 +02:00
parent 0adc223c0a
commit 613a4a2b1f
4 changed files with 51 additions and 30 deletions

View File

@@ -25,7 +25,7 @@ from lernen import FRAGETYPEN
from paths import arbeit_dir, bausteine_path, frage_muster_path, project_dir, subbausteine_path, quelle_path, quelle_crawl_dir, safe_ordner from paths import arbeit_dir, bausteine_path, frage_muster_path, project_dir, subbausteine_path, quelle_path, quelle_crawl_dir, safe_ordner
from crawl import crawl from crawl import crawl
from pipeline import ( from pipeline import (
CANCELLED, FAILED, GenContext, _extra, _log, _prompt, _race, _relevanz_schema, _rest_schema, CANCELLED, FAILED, GenContext, _extra, _gather_fortschritt, _log, _prompt, _race, _relevanz_schema, _rest_schema,
_runde_schema, _semaphore, _str_liste, _stufen_schema, _timeout, run_single_slot, _runde_schema, _semaphore, _str_liste, _stufen_schema, _timeout, run_single_slot,
) )
from textkit import ( from textkit import (
@@ -158,6 +158,14 @@ def _step_idx(topic: str, name: str) -> int:
return _bausteine_steps(topic).index(name) return _bausteine_steps(topic).index(name)
def _melde_p(set_p, topic: str, schritt: str):
"""Async-Melde-Callback für _gather_fortschritt: setzt „<Schritt> d/t…" + Schritt-Index."""
idx = _step_idx(topic, schritt)
async def melde(d, t):
set_p(f"{schritt} {d}/{t}", step=idx)
return melde
# Grobe Anzeige-Phasen: bündeln die Feinschritte (intern bleibt alles feingranular). # Grobe Anzeige-Phasen: bündeln die Feinschritte (intern bleibt alles feingranular).
# Sonderschritte (Quelle laden, Ergänzung) gehören zur Phase „Inventar". # Sonderschritte (Quelle laden, Ergänzung) gehören zur Phase „Inventar".
PHASEN = ( PHASEN = (
@@ -572,8 +580,7 @@ async def _subbausteine_block(ctx: GenContext, set_p, files: dict, entries: dict
neu = await _race(topic, f"Subbausteine Paket {c}", slots, 2 - vorhanden, _timeout("subbaustein", len(chunk)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) neu = await _race(topic, f"Subbausteine Paket {c}", slots, 2 - vorhanden, _timeout("subbaustein", len(chunk)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE)
return not is_cancelled() and neu is not None return not is_cancelled() and neu is not None
set_p(f"Subbausteine finden ({n} Pakete)…", step=_step_idx(topic, "Subbausteine finden")) oks = await _gather_fortschritt([_finde(c, chunk) for c, chunk in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Subbausteine finden"))
oks = await asyncio.gather(*[_finde(c, chunk) for c, chunk in enumerate(chunks, 1)], return_exceptions=True)
if is_cancelled(): if is_cancelled():
return None return None
if not all(ok is True for ok in oks): if not all(ok is True for ok in oks):
@@ -607,8 +614,7 @@ async def _subbausteine_block(ctx: GenContext, set_p, files: dict, entries: dict
if status == FAILED: if status == FAILED:
_log(topic, f"Subbaustein-Klärung Paket {c} fehlgeschlagen — nur Konsens übernommen") _log(topic, f"Subbaustein-Klärung Paket {c} fehlgeschlagen — nur Konsens übernommen")
set_p(f"Subbausteine klären ({n} Pakete)…", step=_step_idx(topic, "Subbausteine klären")) await _gather_fortschritt([_klaere(c, chunk) for c, chunk in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Subbausteine klären"))
await asyncio.gather(*[_klaere(c, chunk) for c, chunk in enumerate(chunks, 1)], return_exceptions=True)
if is_cancelled(): if is_cancelled():
return None return None
@@ -664,8 +670,7 @@ async def _stufen_block(ctx: GenContext, set_p, files: dict, roh: dict, instruct
neu = await _race(topic, f"Stufen Paket {c}", slots, 2 - vorhanden, _timeout("stufe", len(item_idxs)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) neu = await _race(topic, f"Stufen Paket {c}", slots, 2 - vorhanden, _timeout("stufe", len(item_idxs)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE)
return not is_cancelled() and neu is not None return not is_cancelled() and neu is not None
set_p(f"Stufen finden ({n} Pakete)…", step=_step_idx(topic, "Stufen finden")) oks = await _gather_fortschritt([_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Stufen finden"))
oks = await asyncio.gather(*[_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True)
if is_cancelled(): if is_cancelled():
return None return None
if not all(ok is True for ok in oks): if not all(ok is True for ok in oks):
@@ -719,8 +724,7 @@ async def _stufen_block(ctx: GenContext, set_p, files: dict, roh: dict, instruct
ergebnis = {**{k: "mittel" for k in strittig}, **ergebnis, **entsch} ergebnis = {**{k: "mittel" for k in strittig}, **ergebnis, **entsch}
return {item_idxs[k - 1] + 1: stufe for k, stufe in ergebnis.items()} return {item_idxs[k - 1] + 1: stufe for k, stufe in ergebnis.items()}
set_p(f"Stufen klären ({n} Pakete)…", step=_step_idx(topic, "Stufen klären")) parts = await _gather_fortschritt([_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Stufen klären"))
parts = await asyncio.gather(*[_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True)
if is_cancelled(): if is_cancelled():
return None return None
stufe_by_id: dict[int, str] = {} stufe_by_id: dict[int, str] = {}
@@ -783,8 +787,7 @@ async def _relevanz_block(ctx: GenContext, set_p, files: dict, sidecar: dict, in
neu = await _race(topic, f"Relevanz Paket {c}", slots, 2 - vorhanden, _timeout("relevanz", len(item_idxs)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE) neu = await _race(topic, f"Relevanz Paket {c}", slots, 2 - vorhanden, _timeout("relevanz", len(item_idxs)), provider, cancelled=is_cancelled, grace=KONSENS_GRACE)
return not is_cancelled() and neu is not None return not is_cancelled() and neu is not None
set_p(f"Relevanz finden ({n} Pakete)…", step=_step_idx(topic, "Relevanz finden")) oks = await _gather_fortschritt([_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Relevanz finden"))
oks = await asyncio.gather(*[_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True)
if is_cancelled(): if is_cancelled():
return None return None
if not all(ok is True for ok in oks): if not all(ok is True for ok in oks):
@@ -838,8 +841,7 @@ async def _relevanz_block(ctx: GenContext, set_p, files: dict, sidecar: dict, in
ergebnis = {**{k: "relevant" for k in strittig}, **ergebnis, **entsch} ergebnis = {**{k: "relevant" for k in strittig}, **ergebnis, **entsch}
return {item_idxs[k - 1] + 1: rel for k, rel in ergebnis.items()} return {item_idxs[k - 1] + 1: rel for k, rel in ergebnis.items()}
set_p(f"Relevanz klären ({n} Pakete)…", step=_step_idx(topic, "Relevanz klären")) parts = await _gather_fortschritt([_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Relevanz klären"))
parts = await asyncio.gather(*[_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True)
if is_cancelled(): if is_cancelled():
return None return None
relevanz_by_id: dict[int, str] = {} relevanz_by_id: dict[int, str] = {}
@@ -915,8 +917,7 @@ async def _frage_muster_block(ctx: GenContext, set_p, files: dict, sidecar: dict
if status == FAILED: if status == FAILED:
_log(topic, f"Frage-Muster Baustein {c} fehlgeschlagen — kein Muster (Fallback Live)") _log(topic, f"Frage-Muster Baustein {c} fehlgeschlagen — kein Muster (Fallback Live)")
set_p(f"Fragen finden ({len(bausteine)} Bausteine)…", step=_step_idx(topic, "Fragen finden")) await _gather_fortschritt([_finde(c, t, r) for c, (t, r) in enumerate(bausteine, 1)], len(bausteine), _melde_p(set_p, topic, "Fragen finden"))
await asyncio.gather(*[_finde(c, t, r) for c, (t, r) in enumerate(bausteine, 1)], return_exceptions=True)
if is_cancelled(): if is_cancelled():
return None return None
@@ -954,8 +955,7 @@ async def _frage_muster_block(ctx: GenContext, set_p, files: dict, sidecar: dict
if status == FAILED: if status == FAILED:
_log(topic, f"Frage-Muster-Klärung Baustein {c} fehlgeschlagen — Roh-Muster übernommen") _log(topic, f"Frage-Muster-Klärung Baustein {c} fehlgeschlagen — Roh-Muster übernommen")
set_p(f"Fragen klären ({len(bausteine)} Bausteine)…", step=_step_idx(topic, "Fragen klären")) await _gather_fortschritt([_klaere(c, t) for c, (t, _) in enumerate(bausteine, 1)], len(bausteine), _melde_p(set_p, topic, "Fragen klären"))
await asyncio.gather(*[_klaere(c, t) for c, (t, _) in enumerate(bausteine, 1)], return_exceptions=True)
if is_cancelled(): if is_cancelled():
return None return None
@@ -967,19 +967,19 @@ async def _frage_muster_block(ctx: GenContext, set_p, files: dict, sidecar: dict
ergebnis = {titel: _finalisiere(c, rel) for c, (titel, rel) in enumerate(bausteine, 1)} ergebnis = {titel: _finalisiere(c, rel) for c, (titel, rel) in enumerate(bausteine, 1)}
# Phase „Fragen prüfen": Bausteine mit relevanten Subs aber 0 Mustern eine Runde nachholen. # Phase „Fragen prüfen": Bausteine mit relevanten Subs aber 0 Mustern eine Runde nachholen.
set_p(f"Fragen prüfen ({len(bausteine)} Bausteine)", step=_step_idx(topic, "Fragen prüfen")) set_p("Fragen prüfen…", step=_step_idx(topic, "Fragen prüfen"))
leer = [(c, t, r) for c, (t, r) in enumerate(bausteine, 1) if not ergebnis.get(t)] leer = [(c, t, r) for c, (t, r) in enumerate(bausteine, 1) if not ergebnis.get(t)]
if leer: if leer:
_log(topic, f"Frage-Muster: {len(leer)} Baustein(e) ohne Muster — Nachrunde") _log(topic, f"Frage-Muster: {len(leer)} Baustein(e) ohne Muster — Nachrunde")
for c, t, r in leer: for c, t, r in leer:
roh_path(c).unlink(missing_ok=True) roh_path(c).unlink(missing_ok=True)
final_path(c).unlink(missing_ok=True) final_path(c).unlink(missing_ok=True)
await asyncio.gather(*[_finde(c, t, r) for c, t, r in leer], return_exceptions=True) await _gather_fortschritt([_finde(c, t, r) for c, t, r in leer], len(leer), _melde_p(set_p, topic, "Fragen prüfen"))
if is_cancelled(): if is_cancelled():
return None return None
for c, t, r in leer: for c, t, r in leer:
roh_by_c[c] = _waehle(c, r) roh_by_c[c] = _waehle(c, r)
await asyncio.gather(*[_klaere(c, t) for c, t, r in leer], return_exceptions=True) await _gather_fortschritt([_klaere(c, t) for c, t, r in leer], len(leer), _melde_p(set_p, topic, "Fragen prüfen"))
if is_cancelled(): if is_cancelled():
return None return None
for c, t, r in leer: for c, t, r in leer:

View File

@@ -24,7 +24,7 @@ LESBARKEIT_HART_ANTEIL = 0.30 # … ODER wenn dieser Anteil der Sätze hart ist
# Deckel für gleichzeitige CLI-Agenten-Prozesse (über alle Generierungen hinweg). # Deckel für gleichzeitige CLI-Agenten-Prozesse (über alle Generierungen hinweg).
# Eigene Spur für interaktive Aufrufe (Chat, Elemente), damit sie nicht hinter # Eigene Spur für interaktive Aufrufe (Chat, Elemente), damit sie nicht hinter
# laufenden Writern in der Warteschlange hängen. # laufenden Writern in der Warteschlange hängen.
MAX_CONCURRENT_AGENTS = 30 MAX_CONCURRENT_AGENTS = 100
MAX_CONCURRENT_INTERACTIVE = 4 MAX_CONCURRENT_INTERACTIVE = 4
# Grace-Fenster der Konsens-Races (Bausteine, Guide, OnePager): Nach dem ersten # Grace-Fenster der Konsens-Races (Bausteine, Guide, OnePager): Nach dem ersten

View File

@@ -27,7 +27,7 @@ from jsonio import read_json_file as _json_datei
from paths import bausteine_path, guide_content_path, project_dir, subbausteine_path from paths import bausteine_path, guide_content_path, project_dir, subbausteine_path
from pipeline import ( from pipeline import (
CANCELLED, FAILED, GenContext, _claude_error, _extra, CANCELLED, FAILED, GenContext, _claude_error, _extra,
_fail, _gather_error, _log, _prompt, _race, _rest_schema, _runde_schema, _fail, _gather_error, _gather_fortschritt, _log, _prompt, _race, _rest_schema, _runde_schema,
_semaphore, _set_progress, _set_step, _timeout, clear_guide_cancelled, _semaphore, _set_progress, _set_step, _timeout, clear_guide_cancelled,
is_guide_cancelled, run_single_slot, is_guide_cancelled, run_single_slot,
) )
@@ -575,8 +575,8 @@ async def _generate_sections(
inhalt_paths = [content_path.parent / f"{content_path.stem}.inhalt-chunk-{i}.md" for i in range(1, writer_count + 1)] inhalt_paths = [content_path.parent / f"{content_path.stem}.inhalt-chunk-{i}.md" for i in range(1, writer_count + 1)]
offen = [i for i, p in enumerate(inhalt_paths) if not p.exists()] offen = [i for i, p in enumerate(inhalt_paths) if not p.exists()]
if offen: if offen:
await _set_step(guide_id, 2, f"Sammle Inhalte ({writer_count} Agenten)…" if writer_count > 1 else "Sammle Inhalte") async def melde(d, t): await _set_step(guide_id, 2, f"Sammle Inhalte {d}/{t}")
results = await asyncio.gather(*[ results = await _gather_fortschritt([
run_agent( run_agent(
f"{guide_id}-inhalt-{i + 1}", f"{guide_id}-inhalt-{i + 1}",
_prompt( _prompt(
@@ -587,7 +587,7 @@ async def _generate_sections(
_timeout("inhalt", chunk_sizes[i]), provider=provider, role="guide", capabilities="full", _timeout("inhalt", chunk_sizes[i]), provider=provider, role="guide", capabilities="full",
) )
for i in offen for i in offen
], return_exceptions=True) ], writer_count, melde, start=writer_count - len(offen))
if is_cancelled(): if is_cancelled():
return None return None
if not any(p.exists() for p in inhalt_paths): if not any(p.exists() for p in inhalt_paths):
@@ -624,7 +624,9 @@ async def _generate_sections(
"role": "judge", "capabilities": "files", "role": "judge", "capabilities": "files",
"payload": (lambda result, p=check_paths[i]: _lese_probleme_schema(_json_datei(p))), "payload": (lambda result, p=check_paths[i]: _lese_probleme_schema(_json_datei(p))),
} for i in offen_checks] } for i in offen_checks]
await _race(topic, "Inhalts-Prüfung", slots, len(slots), _timeout("inhalt_check", max(chunk_sizes)), provider, cancelled=is_cancelled) n_checks = len(slots)
upd = lambda n: asyncio.create_task(_set_step(guide_id, 3, f"Prüfe Inhalte {n}/{n_checks}"))
await _race(topic, "Inhalts-Prüfung", slots, len(slots), _timeout("inhalt_check", max(chunk_sizes)), provider, on_update=upd, cancelled=is_cancelled)
if is_cancelled(): if is_cancelled():
return None return None
@@ -676,8 +678,8 @@ async def _generate_sections(
paths = [content_path.parent / f"{content_path.stem}.chunk-{i}.md" for i in range(1, writer_count + 1)] paths = [content_path.parent / f"{content_path.stem}.chunk-{i}.md" for i in range(1, writer_count + 1)]
offen = [i for i, p in enumerate(paths) if not p.exists()] offen = [i for i, p in enumerate(paths) if not p.exists()]
if offen: if offen:
await _set_step(guide_id, 4, f"Schreibe Sections ({writer_count} Writer)…" if writer_count > 1 else "Schreibe Sections") async def melde(d, t): await _set_step(guide_id, 4, f"Schreibe Sections {d}/{t}")
results = await asyncio.gather(*[ results = await _gather_fortschritt([
run_agent( run_agent(
f"{guide_id}-w{i + 1}", f"{guide_id}-w{i + 1}",
_prompt( _prompt(
@@ -689,7 +691,7 @@ async def _generate_sections(
_timeout("writer", chunk_sizes[i]), provider=provider, role="guide", capabilities="full", _timeout("writer", chunk_sizes[i]), provider=provider, role="guide", capabilities="full",
) )
for i in offen for i in offen
], return_exceptions=True) ], writer_count, melde, start=writer_count - len(offen))
if is_cancelled(): if is_cancelled():
return None return None
for i, r in zip(offen, results): for i, r in zip(offen, results):
@@ -748,7 +750,9 @@ async def _generate_sections(
"role": "judge", "capabilities": "files", "role": "judge", "capabilities": "files",
"payload": (lambda result, p=check_paths[i]: _lese_probleme_schema(_json_datei(p))), "payload": (lambda result, p=check_paths[i]: _lese_probleme_schema(_json_datei(p))),
} for i in offen_checks] } for i in offen_checks]
res = await _race(topic, f"Lese-Prüfung r{runde}", slots, len(slots), _timeout("lese_check", max(chunk_sizes)), provider, cancelled=is_cancelled) n_checks = len(slots)
upd = lambda n: asyncio.create_task(_set_step(guide_id, 5, f"Prüfe Lesbarkeit {n}/{n_checks}"))
res = await _race(topic, f"Lese-Prüfung r{runde}", slots, len(slots), _timeout("lese_check", max(chunk_sizes)), provider, on_update=upd, cancelled=is_cancelled)
if is_cancelled(): if is_cancelled():
return None return None
if res is None: if res is None:

View File

@@ -309,3 +309,20 @@ async def run_single_slot(
return OK, res[0] return OK, res[0]
async def _gather_fortschritt(coros, total, melde, start=0):
"""Läuft `coros` nebenläufig und meldet Live-Fortschritt: `await melde(fertig, total)`
nach jedem Abschluss (und einmal initial). Ergebnisse in Reihenfolge, return_exceptions=True."""
done = start
async def wrap(c):
nonlocal done
try:
return await c
finally:
done += 1
await melde(done, total)
await melde(done, total)
return await asyncio.gather(*[wrap(c) for c in coros], return_exceptions=True)