Dokumentation

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

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_route lä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 trotzdem routing hinzu, 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_end gefü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 OpenTelemetry context.attach / detach) zu dem Wert zurückkehrt, den er vorgefunden hat.
  • Ein Start-Hook, der ctx.states, ctx.questions oder 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, was predict tut: 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 stattdessen ctx.questions ein 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, und ctx.model ist der aufgelöste Name.
  • Eine Anfrage kann ihr eigenes max_len / head_max_len tragen, die Form des Token-Budgets pro Anfrage, das predict als Aufrufargumente nimmt. Ein Start-Hook, der ctx.max_len setzt, ü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, einen agent.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_error mit ctx.error auf diese gesetzt, dann on_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, dass ctx.error der 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

  1. Installierte Hooks laufen immer vor Hooks pro Aufruf.
  2. Innerhalb einer Liste laufen Hooks in Listenreihenfolge.
  3. Für ein Ereignis läuft jeder Hook, der es implementiert, in dieser Reihenfolge, vor dem nächsten Ereignis.
  4. on_error läuft vor on_predict_end auf dem Fehlerpfad.
  5. on_evict läuft vor on_load, wenn ein einzelnes load sowohl 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