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

13 KiB

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