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
- Agent.system_one / predict
- Router.predict
- Router.predict_batch
- Ciclo de vida do modelo
- Colocação em cache com skip
- Inputs vazios
- Regras de ordenação
- Concorrência
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_routecorre 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 acrescentarrouting, 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_endnã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, umcontext.attach/detachdo OpenTelemetry) desenrola para o valor que encontrou. - Um hook de início que substitui
ctx.states,ctx.questionsou 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 quepredictfaz: 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 actx.questions. - O nome do checkpoint é resolvido antes do agrupamento, por isso um hook
on_routeque fixa um pedido por um alias ("ml") partilha a passagem direta desse checkpoint, ectx.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 quepredictaceita como argumentos de chamada. Um hook de início que definectx.max_lensubstitui-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 chamadaagent.predict_batchpor 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_errorcomctx.errordefinido para ela, e depoison_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 quectx.errorpode 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
- Os hooks instalados correm antes dos hooks por chamada, sempre.
- Dentro de uma lista, os hooks correm pela ordem da lista.
- Para um evento, cada hook que o implementa corre, nessa ordem, antes do evento seguinte.
on_errorcorre antes deon_predict_endno caminho de falha.on_evictcorre antes deon_loadquando um únicoloadtanto 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
- Erros: o que acontece quando um hook levanta exceção, por evento.
- Padrões e antipadrões: como usar bem o ciclo de vida.