Files
planer/phase0/beispiel-einheiten-pipeline-m3.md
2026-07-22 16:12:23 +02:00

231 lines
13 KiB
Markdown

# Einheiten: backend/pipeline.py
stand: UNBEKANNT
## backend/pipeline.py::cancel_guide
anker: backend/pipeline.py::is_guide_cancelled
anker: backend/pipeline.py::clear_guide_cancelled
beschreibung: Verwaltet den Cancel-Zustand eines Guides (Markieren, Abfragen, Zurücksetzen) und markiert den Guide-Datensatz als fehlgeschlagen unter Erhalt des Fortschritts.
fakten:
- Beim Markieren wird der Guide zur internen Cancel-Menge hinzugefügt, der Agent-Scope gelöscht und laufende Subprozesse beendet.
beleg: "_cancelled.add(guide_id)"
- Nach dem Abbruch wird der Guide-Datensatz auf status="error" gesetzt, der Fortschritt bleibt erhalten und es wird eine UTC-Zeitmarke geschrieben.
beleg: "await update_guide(guide_id, status=\"error\", progress=None, error_msg=\"Cancelled — progress is preserved\""
- `is_guide_cancelled` meldet ausschließlich über die Modul-interne Menge, nicht aus dem Datensatz.
beleg: "return guide_id in _cancelled"
- `clear_guide_cancelled` entfernt den Guide aus der Cancel-Menge und räumt den zugehörigen Agent-Scope auf, sodass ein Neustart nicht blockiert wird.
beleg: "clear_scope(f\"{guide_id}-\") # clear scope → restart not blocked"
kanten:
- ruft-auf: backend/pipeline.py::is_guide_cancelled
- ruft-auf: backend/pipeline.py::clear_guide_cancelled
- nutzt: backend/database.py::update_guide
- nutzt: backend/agents.py::cancel_scope
- nutzt: backend/agents.py::kill_process
- nutzt: backend/agents.py::clear_scope
## backend/pipeline.py::_set_progress
anker: backend/pipeline.py::_set_step
anker: backend/pipeline.py::_fail
beschreibung: Schreibt Fortschritts- und Schrittinformationen sowie Fehlermeldungen eines Guides persistiert in den Datensatz.
fakten:
- `_set_progress` aktualisiert ausschließlich das Fortschrittsfeld zusammen mit einer UTC-Zeitmarke.
beleg: "await update_guide(guide_id, progress=progress, updated_at=now)"
- `_set_step` schreibt zusätzlich den numerischen Schritt zusammen mit dem Fortschritt und einem Zeitstempel.
beleg: "await update_guide(guide_id, step=step, progress=progress, updated_at=now)"
- `_fail` setzt den Guide auf den Fehlerstatus, leert den Fortschritt und hinterlegt die übergebene Fehlermeldung.
beleg: "await update_guide(guide_id, status=\"error\", progress=None, error_msg=msg, updated_at=now)"
kanten:
- nutzt: backend/database.py::update_guide
## backend/pipeline.py::_prompt
anker: backend/pipeline.py::_extra
anker: backend/pipeline.py::_log
anker: backend/pipeline.py::_claude_error
anker: backend/pipeline.py::_is_infra
anker: backend/pipeline.py::_gather_error
anker: backend/pipeline.py::_timeout
anker: backend/pipeline.py::_problems_schema
anker: backend/pipeline.py::_str_list
anker: backend/pipeline.py::_runde_schema
anker: backend/pipeline.py::_enum_map_schema
anker: backend/pipeline.py::_yesno_schema
beschreibung: Stellt Hilfsfunktionen für Prompt-Bau, Logging, Fehlerklassifikation, Timeout-Berechnung und JSON-Schemavalidierung bereit.
fakten:
- `_prompt` lädt ein Prompt-Template aus `TEMPLATES_DIR/Prompt/<name>.md` und formatiert es mit den übergebenen Schlüsselwortargumenten.
beleg: "template = (TEMPLATES_DIR / \"Prompt\" / f\"{name}.md\").read_text(encoding=\"utf-8\")"
- `_extra` hängt zusätzliche Nutzeranweisungen in einem klar markierten Block an, sofern welche vorhanden sind.
beleg: "return f\"\\n\\nADDITIONAL INSTRUCTIONS FROM THE USER:\\n{instructions}\\n\" if instructions else \"\""
- `_is_infra` erkennt Transport- bzw. Infrastrukturfehler anhand einer festen Markerliste, um sie von inhaltlichen Fehlern zu unterscheiden.
beleg: "return any(m in (err or \"\") for m in _INFRA_MARKERS)"
- `_claude_error` formatiert eine Fehlermeldung aus Returncode, stdout und stderr, mit Fallback auf das Ende der stdout-Ausgabe.
beleg: "return f\"{label} (exit {returncode}, no output)\""
- `_timeout` berechnet die Timeout-Dauer eines Schritts als Basis plus skalierten Zuschlag pro Eintrag.
beleg: "return base + per * n"
- `_problems_schema` liefert eine leere Liste bei `{"ok": True}`, eine bereinigte Problemliste bei vorhandenen Problemen, sonst `None`.
beleg: "if data.get(\"ok\") is True:"
- `_runde_schema` liefert im Finalmodus nur dann ein Ergebnis, wenn keine Restfragen offen sind.
beleg: "if include is None or rest is None or (final and rest):"
- `_enum_map_schema` liefert eine Parser-Factory, die bei ungültigen IDs oder Werten strikt `None` zurückgibt.
beleg: "if value not in allowed:"
- `_yesno_schema` ist die konkrete Ausprägung der Enum-Map-Factory für das Triage-Feld "relevant" mit den Werten "ja"/"nein".
beleg: "_yesno_schema = _enum_map_schema(\"relevant\", _YESNO) # triage gate ∈ ja/nein"
kanten:
- nutzt: backend/config.py::TEMPLATES_DIR
- nutzt: backend/config.py::TIMEOUTS
## backend/pipeline.py::_race
anker: backend/pipeline.py::AgentInfraError
beschreibung: Startet parallele Agent-Slots, sammelt eine konfigurierbare Quorum-Anzahl gültiger Ergebnisse und behandelt dabei Timeouts, Infrastrukturfehler, Hedging, Grace-Periode und Stornierung.
fakten:
- Slots, die länger als `hedge_s` ohne Ergebnis laufen, bekommen genau einen parallelen Zwilling mit dem Suffix `-h`.
beleg: "hedged.add(i)"
- Sobald das Quorum steht und die Grace-Frist abgelaufen ist, werden laufende Agents beendet und die bis dahin gesammelten Ergebnisse zurückgegeben.
beleg: "if deadline is not None and len(results) >= quorum and loop.time() >= deadline:"
- Infrastrukturfehler werden mit eigenem Zähler und wachsendem Backoff bis zu `_INFRA_MAX_RETRIES` wiederholt; bei Erschöpfung wird `AgentInfraError` ausgelöst.
beleg: "raise AgentInfraError("
- Eine aktive Stornierung führt sofort zum Abbruch ohne Neustart und Rückgabe von `None`.
beleg: "if cancelled and cancelled():"
- Im `finally`-Block werden alle noch laufenden Tasks storniert und ihre Subprozesse beendet, auch beim vorzeitigen Abbruch.
beleg: "for task, i in tasks.items():"
kanten:
- nutzt: backend/agents.py::run_agent
- nutzt: backend/agents.py::kill_process
- nutzt: backend/pipeline.py::_is_infra
- nutzt: backend/pipeline.py::_claude_error
- nutzt: backend/pipeline.py::_log
- nutzt: backend/pipeline.py::_timeout
- nutzt: backend/config.py::MAX_CONCURRENT_GENERATIONS
## backend/pipeline.py::GenContext
anker: backend/pipeline.py::run_single_slot
beschreibung: Bündelt Pipeline-Parameter (Topic, Provider, Stornierungsprüfung, Guide-ID) zu einem Wert, der durch die Pipeline-Glieder weitergereicht wird.
fakten:
- `GenContext` kapselt die langen Argumentlisten der Pipeline-Aufrufe in einem Wert.
beleg: "\"\"\"Pipeline parameters passed through — saves long argument signatures.\"\"\""
- `is_cancelled` ist eine parameterlose Funktion, die den aktuellen Stornierungszustand abfragt.
beleg: "is_cancelled: Callable[[], bool]"
- `run_single_slot` führt genau einen Agent-Aufruf als Rennen mit Quorum 1 aus und übersetzt das Ergebnis in einen `(status, wert)`-Tripel.
beleg: "res = await _race(ctx.topic, label, slots, 1, timeout, ctx.provider, cancelled=ctx.is_cancelled)"
- Bei gesetztem Quorum und Stornierung gibt `run_single_slot` `(CANCELLED, None)` zurück; bei verfehltem Quorum `(FAILED, None)`.
beleg: "if res is None:"
kanten:
- ruft-auf: backend/pipeline.py::_race
## backend/pipeline.py::_gather_progress
beschreibung: Führt Coroutinen nebenläufig aus, meldet nach jeder Vervollständigung den Live-Fortschritt und liefert die Ergebnisse in Eingabereihenfolge samt Ausnahmen.
fakten:
- Vor dem Start der Coroutinen wird der Initialwert `done` an den Reporter übergeben.
beleg: "await report(done, total)"
- Jeder abgeschlossene Job ruft den Reporter auch dann auf, wenn die Coroutine eine Ausnahme wirft, dank `try/finally`.
beleg: "finally:"
- `asyncio.gather` wird mit `return_exceptions=True` aufgerufen, damit ein Fehler in einer Coroutine die übrigen nicht abbricht.
beleg: "return await asyncio.gather(*[wrap(c) for c in coros], return_exceptions=True)"
## backend/pipeline.py::is_guide_cancelled
beschreibung: Beantwortet die Frage, ob ein Guide abgesagt wurde, ausschließlich anhand der Modul-internen Cancel-Menge.
fakten:
- Die Prüfung liest ausschließlich die Modul-interne Cancel-Menge und nicht den Datensatz.
beleg: "return guide_id in _cancelled"
kanten:
## backend/pipeline.py::clear_guide_cancelled
beschreibung: Entfernt einen Guide aus der internen Cancel-Menge und räumt den zugehörigen Agent-Scope auf, damit ein Neustart nicht blockiert wird.
fakten:
- Entfernt die Guide-ID aus der Modul-internen Cancel-Menge.
beleg: "_cancelled.discard(guide_id)"
- Löscht den Agent-Scope unter dem Prefix "{guide_id}-" mit Hinweis auf Restart-Freigabe.
beleg: "clear_scope(f\"{guide_id}-\") # clear scope → restart not blocked"
kanten:
- nutzt: backend/agents.py::clear_scope
## backend/pipeline.py::_set_step
beschreibung: Schreibt zusätzlich zum Fortschritt den numerischen Schritt eines Guides in den Datensatz und stempelt die Aktualisierungszeit.
fakten:
- Persistiert step, progress und updated_at in einem update_guide-Aufruf.
beleg: "await update_guide(guide_id, step=step, progress=progress, updated_at=now)"
- Setzt den Zeitstempel als UTC-ISO-String.
beleg: "now = datetime.now(timezone.utc).isoformat()"
kanten:
- nutzt: backend/database.py::update_guide
## backend/pipeline.py::_fail
beschreibung: Setzt den Guide-Datensatz konsistent auf den Fehlerstatus mit übergebener Meldung und löscht den Fortschritt.
fakten:
- Schreibt status="error", leert progress und hinterlegt die übergebene Fehlermeldung.
beleg: "await update_guide(guide_id, status=\"error\", progress=None, error_msg=msg, updated_at=now)"
- Stempelt den Datensatz mit einer UTC-Zeitmarke.
beleg: "now = datetime.now(timezone.utc).isoformat()"
kanten:
- nutzt: backend/database.py::update_guide
## backend/pipeline.py::_extra
beschreibung: Erzeugt einen klar markierten Block mit zusätzlichen Nutzeranweisungen für einen Prompt oder liefert einen leeren String, falls keine Anweisungen vorliegen.
fakten:
- Mit Anweisungen wird ein Block mit fester Überschrift eingefügt.
beleg: "return f\"\\n\\nADDITIONAL INSTRUCTIONS FROM THE USER:\\n{instructions}\\n\" if instructions else \"\""
- Ohne Anweisungen ist das Ergebnis der leere String.
beleg: "if instructions else \"\""
kanten:
## backend/pipeline.py::_log
beschreibung: Schreibt eine einheitlich formatierte Info-Meldung mit Topic-Präfix über den Pipeline-Logger.
fakten:
- Nutzt den Modul-Logger mit dem Namen "creator.pipeline".
beleg: "log = logging.getLogger(\"creator.pipeline\")"
- Setzt Topic in eckige Klammern vor die eigentliche Meldung.
beleg: "log.info(\"[%s] %s\", topic, msg)"
kanten:
- wird-genutzt-von: backend/pipeline.py::_race
## backend/pipeline.py::_claude_error
beschreibung: Baut eine kompakte Fehlermeldung aus Returncode, stdout und stderr eines Agent-Aufrufs mit abgestufter Fallback-Strategie.
fakten:
- Zieht zuerst den gestrippten stderr heran.
beleg: "stderr = (stderr or \"\").strip()"
- Bevorzugt stderr bis 1000 Zeichen als Fehlermeldung.
beleg: "return f\"{label}: {stderr[:1000]}\""
- Bei leerem stderr wird das Ende der stdout als Fallback genutzt.
beleg: "tail = (stdout or \"\").strip()[-500:]"
- Bei komplett fehlender Ausgabe erscheint nur Label und Exitcode mit "no output".
beleg: "return f\"{label} (exit {returncode}, no output)\""
kanten:
- wird-genutzt-von: backend/pipeline.py::_race
- wird-genutzt-von: backend/pipeline.py::_gather_error
## backend/pipeline.py::_is_infra
beschreibung: Klassifiziert eine Fehlermeldung als Transport- bzw. Infrastrukturfehler anhand einer festen Markerliste.
fakten:
- Die Markerliste enthält HTTP 429, HTTP 5xx, rate_limit, Timeout after und diverse Netzfehler.
beleg: "_INFRA_MARKERS = (\"HTTP 429\", \"HTTP 5\", \"rate_limit\", \"Timeout after\","
- Ein Marker-Treffer genügt, um die Meldung als Infra-Fehler zu werten.
beleg: "return any(m in (err or \"\") for m in _INFRA_MARKERS)"
- Der Docstring grenzt Infra-Fehler von inhaltlichem Fehlschlag ab.
beleg: "\"\"\"Transport-/Infra-Fehler (retry + pause) statt inhaltlichem Fehlschlag.\"\"\""
kanten:
- wird-genutzt-von: backend/pipeline.py::_race
## backend/pipeline.py::_gather_error
beschreibung: Wählt aus einer Liste von Slot-Ergebnissen den ersten Fehler aus und formatiert ihn als kompakten Fehlertext.
fakten:
- Iteriert die Ergebnisliste in Reihenfolge und meldet den ersten Treffer.
beleg: "for r in results:"
- Ausnahmen werden als Klassenname plus Meldung ausgegeben.
beleg: "return f\"{label}: {type(r).__name__}: {r}\""
- Nicht-null-Returncodes werden an `_claude_error` zur Formatierung delegiert.
beleg: "return _claude_error(label, returncode, stdout, stderr)"
- Ohne verwertbares Ergebnis wird ein pauschaler Fehlertext geliefert.
beleg: "return f\"{label}: no usable result\""
kanten:
- nutzt: backend/pipeline.py::_claude_error
## backend/pipeline.py::_timeout
beschreibung: Berechnet die Timeout-Dauer eines Pipeline-Schritts als Basis plus skalierten Zuschlag pro Eintrag aus der zentralen TIMEOUTS-Tabelle.
fakten:
- Liest Basis und Zuschlag pro Eintrag aus TIMEOUTS für den angegebenen Schritt.
beleg: "base, per = TIMEOUTS[step]"
- Liefert die Summe aus Basis und Zuschlag multipliziert mit n.
beleg: "return base + per * n"
kanten:
- nutzt: backend/config.py::TIMEOUTS
## backend/pipeline.py::_problems_schema
beschreibung: Validiert Agent-Antworten für das Probleme-Feld und liefert je nach Form eine leere Liste, die bereinigte