From 613a4a2b1f210c3e624070205b973036f658109c Mon Sep 17 00:00:00 2001 From: team3 Date: Mon, 22 Jun 2026 13:44:33 +0200 Subject: [PATCH] update --- backend/bausteine.py | 40 ++++++++++++++++++++-------------------- backend/config.py | 2 +- backend/guide.py | 22 +++++++++++++--------- backend/pipeline.py | 17 +++++++++++++++++ 4 files changed, 51 insertions(+), 30 deletions(-) diff --git a/backend/bausteine.py b/backend/bausteine.py index 44b61ba..7e3fcd5 100644 --- a/backend/bausteine.py +++ b/backend/bausteine.py @@ -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 crawl import crawl 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, ) from textkit import ( @@ -158,6 +158,14 @@ def _step_idx(topic: str, name: str) -> int: return _bausteine_steps(topic).index(name) +def _melde_p(set_p, topic: str, schritt: str): + """Async-Melde-Callback für _gather_fortschritt: setzt „ 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). # Sonderschritte (Quelle laden, Ergänzung) gehören zur Phase „Inventar". 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) return not is_cancelled() and neu is not None - set_p(f"Subbausteine finden ({n} Pakete)…", step=_step_idx(topic, "Subbausteine finden")) - oks = await asyncio.gather(*[_finde(c, chunk) for c, chunk in enumerate(chunks, 1)], return_exceptions=True) + oks = await _gather_fortschritt([_finde(c, chunk) for c, chunk in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Subbausteine finden")) if is_cancelled(): return None 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: _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 asyncio.gather(*[_klaere(c, chunk) for c, chunk in enumerate(chunks, 1)], return_exceptions=True) + await _gather_fortschritt([_klaere(c, chunk) for c, chunk in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Subbausteine klären")) if is_cancelled(): 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) return not is_cancelled() and neu is not None - set_p(f"Stufen finden ({n} Pakete)…", step=_step_idx(topic, "Stufen finden")) - oks = await asyncio.gather(*[_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True) + oks = await _gather_fortschritt([_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Stufen finden")) if is_cancelled(): return None 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} 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 asyncio.gather(*[_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True) + parts = await _gather_fortschritt([_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Stufen klären")) if is_cancelled(): return None 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) return not is_cancelled() and neu is not None - set_p(f"Relevanz finden ({n} Pakete)…", step=_step_idx(topic, "Relevanz finden")) - oks = await asyncio.gather(*[_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True) + oks = await _gather_fortschritt([_rate(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Relevanz finden")) if is_cancelled(): return None 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} 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 asyncio.gather(*[_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], return_exceptions=True) + parts = await _gather_fortschritt([_klaere(c, idxs) for c, idxs in enumerate(chunks, 1)], len(chunks), _melde_p(set_p, topic, "Relevanz klären")) if is_cancelled(): return None 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: _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 asyncio.gather(*[_finde(c, t, r) for c, (t, r) in enumerate(bausteine, 1)], return_exceptions=True) + await _gather_fortschritt([_finde(c, t, r) for c, (t, r) in enumerate(bausteine, 1)], len(bausteine), _melde_p(set_p, topic, "Fragen finden")) if is_cancelled(): return None @@ -954,8 +955,7 @@ async def _frage_muster_block(ctx: GenContext, set_p, files: dict, sidecar: dict if status == FAILED: _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 asyncio.gather(*[_klaere(c, t) for c, (t, _) in enumerate(bausteine, 1)], return_exceptions=True) + await _gather_fortschritt([_klaere(c, t) for c, (t, _) in enumerate(bausteine, 1)], len(bausteine), _melde_p(set_p, topic, "Fragen klären")) if is_cancelled(): 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)} # 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)] if leer: _log(topic, f"Frage-Muster: {len(leer)} Baustein(e) ohne Muster — Nachrunde") for c, t, r in leer: roh_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(): return None for c, t, r in leer: 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(): return None for c, t, r in leer: diff --git a/backend/config.py b/backend/config.py index 662de5b..e7fe878 100644 --- a/backend/config.py +++ b/backend/config.py @@ -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). # Eigene Spur für interaktive Aufrufe (Chat, Elemente), damit sie nicht hinter # laufenden Writern in der Warteschlange hängen. -MAX_CONCURRENT_AGENTS = 30 +MAX_CONCURRENT_AGENTS = 100 MAX_CONCURRENT_INTERACTIVE = 4 # Grace-Fenster der Konsens-Races (Bausteine, Guide, OnePager): Nach dem ersten diff --git a/backend/guide.py b/backend/guide.py index 1a2d699..3ab0e0f 100644 --- a/backend/guide.py +++ b/backend/guide.py @@ -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 pipeline import ( 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, 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)] offen = [i for i, p in enumerate(inhalt_paths) if not p.exists()] if offen: - await _set_step(guide_id, 2, f"Sammle Inhalte ({writer_count} Agenten)…" if writer_count > 1 else "Sammle Inhalte…") - results = await asyncio.gather(*[ + async def melde(d, t): await _set_step(guide_id, 2, f"Sammle Inhalte {d}/{t}…") + results = await _gather_fortschritt([ run_agent( f"{guide_id}-inhalt-{i + 1}", _prompt( @@ -587,7 +587,7 @@ async def _generate_sections( _timeout("inhalt", chunk_sizes[i]), provider=provider, role="guide", capabilities="full", ) for i in offen - ], return_exceptions=True) + ], writer_count, melde, start=writer_count - len(offen)) if is_cancelled(): return None if not any(p.exists() for p in inhalt_paths): @@ -624,7 +624,9 @@ async def _generate_sections( "role": "judge", "capabilities": "files", "payload": (lambda result, p=check_paths[i]: _lese_probleme_schema(_json_datei(p))), } 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(): 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)] offen = [i for i, p in enumerate(paths) if not p.exists()] if offen: - await _set_step(guide_id, 4, f"Schreibe Sections ({writer_count} Writer)…" if writer_count > 1 else "Schreibe Sections…") - results = await asyncio.gather(*[ + async def melde(d, t): await _set_step(guide_id, 4, f"Schreibe Sections {d}/{t}…") + results = await _gather_fortschritt([ run_agent( f"{guide_id}-w{i + 1}", _prompt( @@ -689,7 +691,7 @@ async def _generate_sections( _timeout("writer", chunk_sizes[i]), provider=provider, role="guide", capabilities="full", ) for i in offen - ], return_exceptions=True) + ], writer_count, melde, start=writer_count - len(offen)) if is_cancelled(): return None for i, r in zip(offen, results): @@ -748,7 +750,9 @@ async def _generate_sections( "role": "judge", "capabilities": "files", "payload": (lambda result, p=check_paths[i]: _lese_probleme_schema(_json_datei(p))), } 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(): return None if res is None: diff --git a/backend/pipeline.py b/backend/pipeline.py index 777854d..a8b2efa 100644 --- a/backend/pipeline.py +++ b/backend/pipeline.py @@ -309,3 +309,20 @@ async def run_single_slot( 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) + +