Lebenszyklus
Diese Seite ist die genaue Reihenfolge der Ereignisse für jeden Einstiegspunkt. Wenn du nur ein Diagramm liest, lies das für den Router; es ist die Obermenge.
- Agent.predict_batch
- Agent.system_one / predict
- Router.predict
- Router.predict_batch
- Modell-Lebenszyklus
- Caching mit skip
- Leere Eingaben
- Reihenfolgeregeln
- Nebenläufigkeit
Agent.predict_batch
predict_batch ist die einzige Implementierung; system_one und predict rufen sie mit einem
Zustand auf.
predict_batch(states, questions, batch_size=..., hooks=..., ...)
│
├─ active = installed hooks + per-call hooks (installed first)
├─ ctx = PredictContext(states, questions, model=self.model_id, agent=self)
│
├─ try:
│ │
│ ├─ on_predict_start ─────────────────────────────┐
│ │ a hook may: │
│ │ • rewrite ctx.states / ctx.questions │
│ │ • set ctx.max_len / ctx.head_max_len │
│ │ • ctx.skip(results) ─────────────┐ │
│ │ • raise (aborts; see errors) │ │
│ │ │ │
│ ├─ if ctx.results is not None: ◄─────────┘ │ cache hit
│ │ skip tokenization and forward │
│ ├─ else: │
│ │ validate states is a list │
│ │ for each batch chunk: │
│ │ _encode_state ─► collate ─► _forward │
│ │ _decode_answers │
│ │ ctx.results = [...] │
│ │ │
│ └─ (any failure here) ──► except BaseException: │
│ ctx.error = exc │
│ on_error │
│ re-raise │
│ │
│ finally: │
│ ctx.elapsed_ms = now - ctx.started_at │
│ if ctx.results: ctx.usage = aggregate_usage(...)│
│ on_predict_end ─────────────────────────────────┘
│
└─ return ctx.results
Dasselbe ctx-Objekt fließt durch Start, Fehler und Ende, sodass run_id sie korreliert und ein
End-Hook ctx.error lesen kann.
Agent.system_one / predict
system_one(state, questions, hooks=..., ...)
└─ predict_batch([state], questions, hooks=..., ...)[0]
Somit erbt system_one jeden Hook und denselben Lebenszyklus, mit ctx.states == [state].
Router.predict
Router.predict(state, questions, model=..., hooks=..., on_predict_start=..., on_predict_end=...)
│
├─ active = installed hooks + per-call hooks
│
├─ route(state, questions, ..., hooks=per-call, hooks_raise=...)
│ │
│ ├─ _route(...) detect script / language / workflow
│ ├─ on_route ──► ctx.decision a hook may replace the decision
│ └─ return ctx.decision
│
├─ load(decision["model"])
│ │
│ ├─ already resident? ──► return it
│ ├─ else build Agent(...) ──► on_load (after the Router lock is released)
│ └─ evict LRU checkpoints ──► on_evict (after the Router lock is released)
│
├─ ctx = PredictContext(states=[state], questions, decision, model=decision.model,
│ agent=agent, router=self)
├─ try:
│ ├─ on_predict_start
│ ├─ if ctx.results is None:
│ │ result = agent.system_one(ctx.states[0], ctx.questions)
│ │ └─ the Agent's own hooks run here (start / forward / end)
│ │ result["routing"] = decision
│ │ ctx.results = [result]
│ └─ else:
│ for each cached result: result.setdefault("routing", decision)
│ └─ (any failure) ──► except: on_error, re-raise
│ └─ finally: elapsed_ms, usage, on_predict_end
│
└─ return ctx.results[0]
Wichtige Punkte:
on_routeläuft, bevor das Modell geladen wird, sodass ein Hook einen Checkpoint festlegen und das Laden eines anderen vermeiden kann.- Predict-Hooks auf Router-Ebene umschließen den gesamten Aufruf. Sie werden nicht in den Agent weitergeleitet; ein angehängter Agent mit eigenen Hooks führt diese ebenfalls aus, was zu erwarten ist.
- Ein
ctx.skip()auf Router-Ebene fügt trotzdemroutinghinzu, sodass die Rückgabeform stabil ist.
Router.predict_batch
Jedes Ergebnis ist das, was predict für diese Anfrage zurückgibt, sodass Predict-Hooks auf
Router-Ebene auch hier pro Anfrage laufen: jede Anfrage bekommt ihr eigenes PredictContext,
run_id und elapsed_ms.
Router.predict_batch(requests, batch_size=..., hooks=...)
│
├─ active = installed hooks + per-call hooks (installed first; None and [] add nothing)
├─ route_batch(requests, hooks=per-call) ──► on_route, once per request (no checkpoint loaded yet)
│
└─ for each checkpoint, in order of first appearance:
│
├─ load(checkpoint) ──► on_load / on_evict
├─ for each request of this checkpoint, in input order:
│ ctx = PredictContext(states=[state], questions, decision, model, agent, router,
│ max_len=request.get("max_len"),
│ head_max_len=request.get("head_max_len"))
│ on_predict_start a hook may redact, rewrite, set a token budget or skip
├─ group the requests left to infer by (questions, ctx.max_len, ctx.head_max_len)
│ agent.predict_batch(states, questions, ...) ──► one shared forward pass per group
│ result["routing"] = decision; ctx.results = [result]
├─ ctx.usage, once per request
├─ (any failure) ──► for every started request, in reverse input order:
│ ctx.error = exc, on_error, on_predict_end; then re-raise
└─ on_predict_end, once per request of this checkpoint, in reverse input order
Wichtige Punkte:
- Jeder Start-Hook der Anfragen eines Checkpoints läuft vor allen ihren End-Hooks, weil sie sich
Vorwärtsdurchläufe teilen. Ein Cache, der in
on_predict_endgefüllt wird, kann daher innerhalb derselben Checkpoint-Gruppe keinen doppelten Zustand bedienen; über Aufrufe hinweg kann er das. - Aus demselben Grund enden die Anfragen in umgekehrter Reihenfolge zu ihrem Start, sodass ein Hook,
der im Start etwas setzt und im Ende zurücksetzt (ein
contextvars-Wert, ein OpenTelemetrycontext.attach/detach) zu dem Wert zurückkehrt, den er vorgefunden hat. - Ein Start-Hook, der
ctx.states,ctx.questionsoder das Token-Budget ersetzt, ändert nur seine eigene Anfrage: Die Anfragen werden für den Vorwärtsdurchlauf gruppiert, nachdem ihre Start-Hooks gelaufen sind. Ein geteiltes Fragen-dict an Ort und Stelle zu mutieren, ist etwas anderes und nicht, waspredicttut: Nichts aus der Gruppe wird inferiert, bis alle ihre Start-Hooks gelaufen sind, also erreicht die Änderung jede Anfrage, die das dict teilt, einschließlich derer, deren Start-Hooks früher liefen, und auch den Aufrufer. Weise stattdessenctx.questionsein neues dict zu. - Der Checkpoint-Name wird vor dem Gruppieren aufgelöst, sodass ein
on_route-Hook, der eine Anfrage über ein Alias ("ml") festlegt, den Vorwärtsdurchlauf dieses Checkpoints teilt, undctx.modelist der aufgelöste Name. - Eine Anfrage kann ihr eigenes
max_len/head_max_lentragen, die Form des Token-Budgets pro Anfrage, daspredictals Aufrufargumente nimmt. Ein Start-Hook, derctx.max_lensetzt, überschreibt es, weil der Hook läuft, nachdem der Kontext gebaut wurde. Anfragen, die unterschiedliche Budgets verlangen, können keinen Vorwärtsdurchlauf teilen, sodass ein Batch, der Budgets mischt, einenagent.predict_batch-Aufruf pro Budget macht. - Jede gestartete Anfrage bekommt genau ein
on_predict_end, auch wenn der End-Hook einer anderen Anfrage eine Ausnahme auslöst; der erste solche Fehler wird ausgelöst, nachdem alle gelaufen sind. - Wenn eine Checkpoint-Gruppe fehlschlägt, schlägt jede ihrer gestarteten Anfragen mit der Ausnahme
fehl: jede bekommt
on_errormitctx.errorauf diese gesetzt, dannon_predict_end. Das schließt einen Cache-Treffer und eine Anfrage ein, deren Fragengruppe bereits gelaufen war, weil der Aufrufer die Ausnahme und kein Ergebnis für irgendeine von ihnen bekommt, und es bedeutet, dassctx.errorder Fehler einer anderen Anfrage sein kann (ein Start-Hook, der für eine Anfrage eine Ausnahme auslöst, lässt seine Gruppe fehlschlagen). Anfragen von Checkpoint-Gruppen, die bereits fertig waren, sind mit ihren Ergebnissen beendet, wie es die früheren Aufrufe von[router.predict(...) for ...]getan hätten.
Modell-Lebenszyklus
on_load wird ausgelöst, wenn ein Checkpoint gebaut wird; on_evict, wenn einer freigegeben wird.
Beide laufen nachdem das interne Lock des Routers freigegeben wurde, sodass ein Hook sicher in
den Router zurückrufen kann.
load("multilingual")
│
├─ [lock]
│ build Agent(...) (seconds: download + weights)
│ register in _agents / _order
│ evict LRU if over max_loaded ──► evicted = ["english"]
├─ [unlock]
├─ on_evict("english")
└─ on_load("multilingual")
unload("english")
├─ [lock] remove from _agents / _order
├─ [unlock]
└─ on_evict("english")
attach(name, agent) registriert einen vorhandenen Agent und löst nicht on_load aus, weil kein
Checkpoint gebaut wurde.
Caching mit skip
on_predict_start
├─ cache hit? ctx.skip([cached_result])
│ └─ forward pass skipped
│ └─ on_predict_end still runs
│ └─ Router adds `routing` if missing
└─ cache miss? nothing
└─ forward pass runs
└─ on_predict_end can store the result
Siehe examples/hooks/cache.py
für einen funktionierenden Cache.
Leere Eingaben
Hooks werden trotzdem ausgelöst, damit Audit jeden Aufruf sieht:
| Eingabe | ctx.results am Ende |
|---|---|
predict_batch([]) |
[] |
predict_batch(states, {}) (keine Fragen) |
eine leere Antwort-Payload pro Zustand |
system_one(state, {}) |
eine einzelne leere Antwort-Payload |
In diesen Fällen findet keine Tokenisierung oder kein Vorwärtsdurchlauf statt, aber
on_predict_start und on_predict_end laufen.
Reihenfolgeregeln
- Installierte Hooks laufen immer vor Hooks pro Aufruf.
- Innerhalb einer Liste laufen Hooks in Listenreihenfolge.
- Für ein Ereignis läuft jeder Hook, der es implementiert, in dieser Reihenfolge, vor dem nächsten Ereignis.
on_errorläuft voron_predict_endauf dem Fehlerpfad.on_evictläuft voron_load, wenn ein einzelnesloadsowohl verdrängt als auch baut.
installed: [A, B] per-call: [C]
on_predict_start: A, B, C
on_predict_end: A, B, C
Nebenläufigkeit
Agent und Router können sicher aus vielen Threads aufgerufen werden. Jeder Aufruf erzeugt sein
eigenes PredictContext, sodass Kontexte nie über Anfragen hinweg durchsickern. Der einzige geteilte
Zustand sind die Hook-Objekte selbst, sodass ein Hook, der nicht thread-sicher ist, entweder seinen
eigenen Zustand schützen oder mit hooks_concurrent=False installiert werden muss.
hooks_concurrent=True (default) hooks_concurrent=False
thread 1 ─┐ thread 1 ─┐
thread 2 ─┼─ hooks run in parallel thread 2 ─┼─ one hook at a time
thread 3 ─┘ thread 3 ─┘ (RLock)
hooks_concurrent=False serialisiert jeden Hook-Aufruf, nicht ganze Aufrufe: zwei Aufrufe können
sich weiterhin zwischen Ereignissen verzahnen. Es verwendet ein reentrant Lock, sodass ein Hook in
denselben Agent/Router zurückrufen kann, ohne in einen Deadlock zu geraten.
Siehe auch
- Fehler: was passiert, wenn ein Hook eine Ausnahme auslöst, pro Ereignis.
- Muster und Anti-Muster: wie du den Lebenszyklus gut nutzt.