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
- Agent.system_one / predict
- Router.predict
- Router.predict_batch
- Ciclo de vida do modelo
- Cache com skip
- Entradas vazias
- Regras de ordenação
- Concorrência
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_routeroda 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 adicionarouting, 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, umcontext.attach/detachdo OpenTelemetry) volta ao valor que encontrou. - Um hook de início que substitui
ctx.states,ctx.questionsou 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 quepredictfaz: 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 actx.questionsem vez disso. - O nome do checkpoint é resolvido antes do agrupamento, então um hook
on_routeque fixa uma solicitação por um alias ("ml") compartilha a passada direta daquele checkpoint, ectx.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 quepredictaceita como argumentos de chamada. Um hook de início que definectx.max_leno 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 deagent.predict_batchpor 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_errorcomctx.errordefinido para ela, depoison_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 quectx.errorpode 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
- Hooks instalados rodam antes de hooks por chamada, sempre.
- Dentro de uma lista, os hooks rodam na ordem da lista.
- Para um mesmo evento, todo hook que o implementa roda, nessa ordem, antes do próximo evento.
on_errorroda antes deon_predict_endno caminho de falha.on_evictroda antes deon_loadquando um únicoloadfaz 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
- Erros: o que acontece quando um hook lança, por evento.
- Padrões e antipadrões: como usar bem o ciclo de vida.