Documentação

Ciclo de vida

Esta página é a ordem precisa dos eventos para cada ponto de entrada. Se você só ler um diagrama, leia o Router; ele é o superconjunto.

Agent.predict_batch

predict_batch é a implementação única; system_one e predict a chamam com um único 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 pelo início, pelo erro e pelo fim, então run_id os correlaciona e um hook de fim consegue ler ctx.error.

Agent.system_one / predict

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

Então system_one herda todo 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 roda antes de o modelo ser carregado, então um hook pode fixar um checkpoint e evitar carregar outro.
  • Os hooks de predição no nível do Router envolvem a chamada inteira. Eles não são repassados para dentro do Agent; um Agent anexado com seus próprios hooks também roda esses, o que é esperado.
  • Um ctx.skip() no nível do Router ainda adiciona routing, então o formato do retorno é estável.

Router.predict_batch

Cada resultado é o que predict retorna para aquela solicitação, então os hooks de predição no nível do Router rodam por solicitação aqui também: toda solicitação ganha 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:

  • Todo hook de início das solicitações de um checkpoint roda antes de qualquer hook de fim delas, porque elas compartilham passadas diretas. Um cache que se preenche em on_predict_end, portanto, não consegue servir um estado duplicado dentro do mesmo grupo de checkpoint; consegue entre chamadas.
  • Pelo mesmo motivo as solicitações terminam na ordem inversa à que começaram, então um hook que define algo no início e o restaura no fim (um valor de contextvars, um context.attach / detach do OpenTelemetry) volta ao valor que encontrou.
  • Um hook de início que substitui ctx.states, ctx.questions ou o orçamento de tokens muda apenas a própria solicitação: as solicitações são agrupadas para a passada direta depois de seus hooks de início terem rodado. Mutar um dict de perguntas compartilhado no lugar é diferente, e não é o que predict faz: nada do grupo é inferido até todos os seus hooks de início terem rodado, então a mudança alcança toda solicitação que compartilha o dict, incluindo as cujos hooks de início rodaram antes, e o chamador também. Atribua um novo dict a ctx.questions em vez disso.
  • O nome do checkpoint é resolvido antes do agrupamento, então um hook on_route que fixa uma solicitação por um alias ("ml") compartilha a passada direta daquele checkpoint, e ctx.model é o nome resolvido.
  • Uma solicitação pode carregar seu próprio max_len / head_max_len, a forma por solicitação do orçamento de tokens que predict aceita como argumentos de chamada. Um hook de início que define ctx.max_len o sobrescreve, porque o hook roda depois de o contexto ser construído. Solicitações que pedem orçamentos diferentes não conseguem compartilhar uma passada direta, então um lote que mistura orçamentos faz uma chamada de agent.predict_batch por orçamento.
  • Toda solicitação iniciada ganha exatamente um on_predict_end, mesmo quando o hook de fim de outra solicitação lança; o primeiro erro desse tipo é lançado depois que todos eles rodaram.
  • Se um grupo de checkpoint falha, cada uma de suas solicitações iniciadas falha com a exceção: cada uma recebe on_error com ctx.error definido para ela, depois on_predict_end. Isso inclui um acerto de cache e uma solicitação cujo grupo de perguntas já rodou, porque o chamador recebe a exceção e nenhum resultado de nenhuma delas, e significa que ctx.error pode ser a falha de outra solicitação (um hook de início que lança para uma solicitação falha todo o seu grupo). Solicitações de grupos de checkpoint que já terminaram encerraram com seus resultados, como as chamadas anteriores de [router.predict(...) for ...] teriam feito.

Ciclo de vida do modelo

on_load dispara quando um checkpoint é construído; on_evict, quando um é liberado. Ambos rodam depois que o lock interno do Router é liberado, então um hook pode chamar de volta o Router com 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.

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

Veja examples/hooks/cache.py para um cache funcional.

Entradas vazias

Os hooks ainda disparam para que a auditoria veja cada chamada:

entrada ctx.results no fim
predict_batch([]) []
predict_batch(states, {}) (sem perguntas) uma carga útil de resposta vazia por estado
system_one(state, {}) uma única carga útil de resposta vazia

Nenhuma tokenização ou passada direta acontece nesses casos, mas on_predict_start e on_predict_end rodam.

Regras de ordenação

  1. Hooks instalados rodam antes de hooks por chamada, sempre.
  2. Dentro de uma lista, os hooks rodam na ordem da lista.
  3. Para um mesmo evento, todo hook que o implementa roda, nessa ordem, antes do próximo evento.
  4. on_error roda antes de on_predict_end no caminho de falha.
  5. on_evict roda antes de on_load quando um único load faz evict e build ao mesmo tempo.
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 para chamar a partir de muitas threads. Cada chamada cria seu próprio PredictContext, então os contextos nunca vazam entre solicitações. O único estado compartilhado são os próprios objetos de hook, então um hook que não é thread-safe precisa ou proteger 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 ainda podem se intercalar entre eventos. Ele usa um lock reentrante, então um hook pode chamar de volta o mesmo Agent/Router sem deadlock.

Veja também