Documentação

Ciclo de vida

Esta página é a ordem precisa dos eventos para cada ponto de entrada. Se só leres um diagrama, lê o do Router; é o superconjunto.

Agent.predict_batch

predict_batch é a única implementação; system_one e predict chamam-na com um só 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

O mesmo objeto ctx flui através do início, do erro e do fim, por isso run_id correlaciona-os e um hook de fim pode ler ctx.error.

Agent.system_one / predict

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

Portanto, system_one herda cada hook e o mesmo ciclo de vida, com 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]

Pontos-chave:

  • on_route corre antes de o modelo ser carregado, por isso um hook pode fixar um checkpoint e evitar carregar outro.
  • Os hooks de predição ao nível do Router envolvem toda a chamada. Não são encaminhados para dentro do Agent; um Agent ligado com hooks próprios corre também esses, o que é esperado.
  • Um ctx.skip() ao nível do Router continua a acrescentar routing, por isso a forma do retorno é estável.

Router.predict_batch

Cada resultado é o que predict devolve para esse pedido, por isso os hooks de predição ao nível do Router correm aqui também por pedido: cada pedido recebe o seu próprio PredictContext, run_id e 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

Pontos-chave:

  • Cada hook de início dos pedidos de um checkpoint corre antes de qualquer dos seus hooks de fim, porque partilham passagens diretas. Uma cache que enche no on_predict_end não consegue, portanto, servir um estado duplicado dentro do mesmo grupo de checkpoint; consegue-o entre chamadas.
  • Pela mesma razão, os pedidos terminam na ordem inversa àquela em que começaram, por isso um hook que define algo no início e o repõe no fim (um valor contextvars, um context.attach / detach do OpenTelemetry) desenrola para o valor que encontrou.
  • Um hook de início que substitui ctx.states, ctx.questions ou o orçamento de tokens altera apenas o seu próprio pedido: os pedidos são agrupados para a passagem direta depois de os seus hooks de início terem corrido. Mutar um dict de perguntas partilhado no local é diferente, e não é o que predict faz: nada do grupo é inferido até todos os seus hooks de início terem corrido, por isso a alteração alcança cada pedido que partilhe o dict, incluindo aqueles cujos hooks de início correram antes, e o autor da chamada também. Atribui antes um novo dict a ctx.questions.
  • O nome do checkpoint é resolvido antes do agrupamento, por isso um hook on_route que fixa um pedido por um alias ("ml") partilha a passagem direta desse checkpoint, e ctx.model é o nome resolvido.
  • Um pedido pode transportar os seus próprios max_len / head_max_len, a forma por pedido do orçamento de tokens que predict aceita como argumentos de chamada. Um hook de início que define ctx.max_len substitui-o, porque o hook corre depois de o contexto ser construído. Os pedidos que pedem orçamentos diferentes não podem partilhar uma passagem direta, por isso um lote que mistura orçamentos faz uma chamada agent.predict_batch por orçamento.
  • Cada pedido iniciado recebe exatamente um on_predict_end, mesmo quando o hook de fim de outro pedido levanta exceção; o primeiro desses erros é levantado depois de todos eles terem corrido.
  • Se um grupo de checkpoint falhar, cada um dos seus pedidos iniciados falha com a exceção: cada um recebe on_error com ctx.error definido para ela, e depois on_predict_end. Isso inclui uma cache hit e um pedido cujo grupo de perguntas já tinha corrido, porque o autor da chamada recebe a exceção e nenhum resultado para qualquer deles, e significa que ctx.error pode ser a falha de outro pedido (um hook de início que levanta exceção por um pedido falha o seu grupo). Os pedidos de grupos de checkpoint que já terminaram terminaram com os seus resultados, como teriam as chamadas anteriores de [router.predict(...) for ...].

Ciclo de vida do modelo

on_load dispara quando um checkpoint é construído; on_evict quando um é libertado. Ambos correm depois de o lock interno do Router ser libertado, por isso um hook pode chamar de volta para o Router em segurança.

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 um agente existente e não dispara on_load, porque nenhum checkpoint foi construído.

Colocação em cache com 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

Vê examples/hooks/cache.py para uma cache funcional.

Inputs vazios

Os hooks disparam na mesma, para que a auditoria veja cada chamada:

input ctx.results no fim
predict_batch([]) []
predict_batch(states, {}) (sem perguntas) um payload de resposta vazia por estado
system_one(state, {}) um único payload de resposta vazia

Não acontece tokenização nem passagem direta nestes casos, mas on_predict_start e on_predict_end correm.

Regras de ordenação

  1. Os hooks instalados correm antes dos hooks por chamada, sempre.
  2. Dentro de uma lista, os hooks correm pela ordem da lista.
  3. Para um evento, cada hook que o implementa corre, nessa ordem, antes do evento seguinte.
  4. on_error corre antes de on_predict_end no caminho de falha.
  5. on_evict corre antes de on_load quando um único load tanto liberta como constrói.
installed: [A, B]   per-call: [C]
on_predict_start: A, B, C
on_predict_end:   A, B, C

Concorrência

Agent e Router são seguros de chamar a partir de muitas threads. Cada chamada cria o seu próprio PredictContext, por isso os contextos nunca se infiltram entre pedidos. O único estado partilhado são os próprios objetos hook, por isso um hook que não seja thread-safe tem de guardar o seu próprio estado ou ser instalado com 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 invocação de hook, não chamadas inteiras: duas chamadas podem ainda intercalar-se entre eventos. Usa um lock reentrante, por isso um hook pode chamar de volta para o mesmo Agent/Router sem impasse.

Ver também