Documentación

Ciclo de vida

Esta página es el orden exacto de los eventos de cada punto de entrada. Si solo vas a leer un diagrama, lee el de Router; es el superconjunto.

Agent.predict_batch

predict_batch es la única implementación; system_one y predict la llaman con un solo estado.

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

El mismo objeto ctx recorre el inicio, el error y el final, así que run_id los correlaciona y un hook de final puede leer ctx.error.

Agent.system_one / predict

system_one(state, questions, hooks=..., ...)
  └─ predict_batch([state], questions, hooks=..., ...)[0]

Así que system_one hereda todos los hooks y el mismo ciclo de vida, con 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]

Puntos clave:

  • on_route se ejecuta antes de que se cargue el modelo, así que un hook puede fijar un checkpoint y evitar que se cargue otro.
  • Los hooks de predicción a nivel de Router envuelven toda la llamada. No se reenvían al Agent; un Agent asociado con sus propios hooks también ejecuta esos, lo cual es lo esperado.
  • Un ctx.skip() a nivel de Router igualmente añade routing, así que la forma del valor devuelto es estable.

Router.predict_batch

Cada resultado es lo que predict devuelve para esa solicitud, así que los hooks de predicción a nivel de Router también se ejecutan por solicitud aquí: cada solicitud recibe su propio PredictContext, run_id y 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

Puntos clave:

  • Todos los hooks de inicio de las solicitudes de un checkpoint se ejecutan antes que cualquiera de sus hooks de final, porque comparten pasadas hacia adelante. Por eso, una caché que se llena en on_predict_end no puede servir un estado duplicado dentro del mismo grupo de checkpoint; sí puede hacerlo entre llamadas distintas.
  • Por la misma razón, las solicitudes terminan en orden inverso al de inicio, así que un hook que fija algo en el inicio y lo restablece en el final (un valor de contextvars, un context.attach / detach de OpenTelemetry) vuelve al valor que encontró.
  • Un hook de inicio que reemplaza ctx.states, ctx.questions o el presupuesto de tokens solo cambia su propia solicitud: las solicitudes se agrupan para la pasada hacia adelante después de que se hayan ejecutado sus hooks de inicio. Mutar en el lugar un dict de preguntas compartido es distinto, y no es lo que hace predict: nada del grupo se infiere hasta que se hayan ejecutado todos sus hooks de inicio, así que el cambio alcanza a todas las solicitudes que comparten el dict, incluidas aquellas cuyos hooks de inicio se ejecutaron antes, y también al llamador. En su lugar, asigna un dict nuevo a ctx.questions.
  • El nombre del checkpoint se resuelve antes de agrupar, así que un hook on_route que fija una solicitud mediante un alias ("ml") comparte la pasada hacia adelante de ese checkpoint, y ctx.model es el nombre resuelto.
  • Una solicitud puede llevar su propio max_len / head_max_len, la forma por solicitud del presupuesto de tokens que predict recibe como argumentos de llamada. Un hook de inicio que fija ctx.max_len lo anula, porque el hook se ejecuta después de que se construye el contexto. Las solicitudes que piden presupuestos distintos no pueden compartir una pasada hacia adelante, así que un lote que mezcla presupuestos hace una llamada a agent.predict_batch por cada presupuesto.
  • Cada solicitud iniciada recibe exactamente un on_predict_end, incluso cuando el hook de final de otra solicitud lanza una excepción; el primer error de ese tipo se lanza después de que todas se hayan ejecutado.
  • Si falla un grupo de checkpoint, todas sus solicitudes iniciadas fallan con la excepción: cada una recibe on_error con ctx.error puesto a esa excepción, y luego on_predict_end. Eso incluye un acierto de caché y una solicitud cuyo grupo de preguntas ya se había ejecutado, porque el llamador recibe la excepción y ningún resultado para ninguna de ellas, y significa que ctx.error puede ser el fallo de otra solicitud (un hook de inicio que lanza una excepción para una solicitud hace fallar a su grupo). Las solicitudes de grupos de checkpoint que ya habían terminado finalizaron con sus resultados, como habría ocurrido en las llamadas anteriores de [router.predict(...) for ...].

Ciclo de vida del modelo

on_load se dispara cuando se construye un checkpoint; on_evict, cuando se libera uno. Ambos se ejecutan después de que se libere el bloqueo interno del Router, así que un hook puede llamar de vuelta al Router con seguridad.

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) registra un agent existente y no dispara on_load, porque no se construyó ningún checkpoint.

Caché con 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

Consulta examples/hooks/cache.py para ver una caché que funciona.

Entradas vacías

Los hooks igualmente se disparan para que la auditoría vea cada llamada:

entrada ctx.results al final
predict_batch([]) []
predict_batch(states, {}) (sin preguntas) una carga de respuesta vacía por estado
system_one(state, {}) una sola carga de respuesta vacía

En estos casos no ocurre ninguna tokenización ni pasada hacia adelante, pero on_predict_start y on_predict_end se ejecutan.

Reglas de orden

  1. Los hooks instalados se ejecutan antes que los hooks por llamada, siempre.
  2. Dentro de una lista, los hooks se ejecutan en el orden de la lista.
  3. Para un evento, todos los hooks que lo implementan se ejecutan, en ese orden, antes del evento siguiente.
  4. on_error se ejecuta antes de on_predict_end en la ruta de fallo.
  5. on_evict se ejecuta antes de on_load cuando un solo load tanto desaloja como construye.
installed: [A, B]   per-call: [C]
on_predict_start: A, B, C
on_predict_end:   A, B, C

Concurrencia

Agent y Router se pueden llamar con seguridad desde muchos hilos. Cada llamada crea su propio PredictContext, así que los contextos nunca se filtran entre solicitudes. El único estado compartido son los propios objetos hook, así que un hook que no sea seguro para hilos debe proteger su propio estado o instalarse con hooks_concurrent=False.

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 serializa cada invocación de hook, no las llamadas completas: dos llamadas aún pueden intercalarse entre eventos. Usa un bloqueo reentrante, así que un hook puede llamar de vuelta al mismo Agent/Router sin producir un interbloqueo.

Véase también