Cycle de vie
Cette page est l’ordre précis des événements pour chaque point d’entrée. Si tu ne lis qu’un seul diagramme, lis celui du Router ; c’est le sur-ensemble.
- Agent.predict_batch
- Agent.system_one / predict
- Router.predict
- Router.predict_batch
- Cycle de vie du modèle
- Mise en cache avec skip
- Entrées vides
- Règles d’ordre
- Concurrence
Agent.predict_batch
predict_batch est l’unique implémentation ; system_one et predict l’appellent avec un seul état.
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
Le même objet ctx circule du start à l’erreur et à la fin, donc run_id les corrèle et un hook de
fin peut lire ctx.error.
Agent.system_one / predict
system_one(state, questions, hooks=..., ...)
└─ predict_batch([state], questions, hooks=..., ...)[0]
Donc system_one hérite de chaque hook et du même cycle de vie, avec 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]
Points clés :
on_routetourne avant le chargement du modèle, donc un hook peut épingler un checkpoint et éviter d’en charger un autre.- Les hooks predict au niveau Router enveloppent tout l’appel. Ils ne sont pas transmis à l’Agent ; un Agent attaché avec ses propres hooks exécute aussi ceux-là, ce qui est attendu.
- Un
ctx.skip()au niveau Router ajoute quand mêmerouting, donc la forme de retour est stable.
Router.predict_batch
Chaque résultat est ce que predict renvoie pour cette requête, donc les hooks predict au niveau
Router tournent eux aussi par requête ici : chaque requête a son propre PredictContext, run_id et
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
Points clés :
- Chaque hook de start des requêtes d’un checkpoint tourne avant l’un de leurs hooks de fin, parce
qu’elles partagent des passes avant. Un cache qui se remplit dans
on_predict_endne peut donc pas servir un état dupliqué au sein du même groupe de checkpoint ; il peut le faire entre appels. - Pour la même raison, les requêtes se terminent dans l’ordre inverse de leur démarrage, donc un hook
qui définit quelque chose au start et le réinitialise à la fin (une valeur
contextvars, uncontext.attach/detachOpenTelemetry) revient à la valeur qu’il a trouvée. - Un hook de start qui remplace
ctx.states,ctx.questionsou le budget de jetons ne change que sa propre requête : les requêtes sont groupées pour la passe avant après l’exécution de leurs hooks de start. Muter sur place un dict de questions partagé est différent, et ce n’est pas ce que faitpredict: rien du groupe n’est inféré avant que tous ses hooks de start aient tourné, donc le changement atteint chaque requête qui partage le dict, y compris celles dont les hooks de start ont tourné plus tôt, et l’appelant aussi. Assigne plutôt un nouveau dict àctx.questions. - Le nom du checkpoint est résolu avant le groupement, donc un hook
on_routequi épingle une requête par un alias ("ml") partage la passe avant de ce checkpoint, etctx.modelest le nom résolu. - Une requête peut porter ses propres
max_len/head_max_len, la forme par requête du budget de jetons quepredictprend comme arguments d’appel. Un hook de start qui définitctx.max_lenles surcharge, parce que le hook tourne après la construction du contexte. Les requêtes qui demandent des budgets différents ne peuvent pas partager une passe avant, donc un lot qui mélange les budgets fait un appelagent.predict_batchpar budget. - Chaque requête démarrée reçoit exactement un
on_predict_end, même quand le hook de fin d’une autre requête lève une erreur ; la première erreur de ce type est levée après que toutes ont tourné. - Si un groupe de checkpoint échoue, chacune de ses requêtes démarrées échoue avec l’exception : chacune
reçoit
on_erroravecctx.errordéfini dessus, puison_predict_end. Cela inclut un hit de cache et une requête dont le groupe de questions avait déjà tourné, parce que l’appelant reçoit l’exception et aucun résultat pour aucune d’elles, et cela signifie quectx.errorpeut être l’échec d’une autre requête (un hook de start qui lève pour une requête fait échouer son groupe). Les requêtes des groupes de checkpoints déjà terminés se sont terminées avec leurs résultats, comme l’auraient fait les appels antérieurs de[router.predict(...) for ...].
Cycle de vie du modèle
on_load se déclenche quand un checkpoint est construit ; on_evict quand l’un est libéré. Les deux
tournent après la libération du verrou interne du Router, donc un hook peut rappeler le Router en
toute sécurité.
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) enregistre un agent existant et ne déclenche pas on_load, parce qu’aucun
checkpoint n’a été construit.
Mise en cache avec 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
Voir examples/hooks/cache.py pour un cache fonctionnel.
Entrées vides
Les hooks se déclenchent quand même, pour que l’audit voie chaque appel :
| entrée | ctx.results à la fin |
|---|---|
predict_batch([]) |
[] |
predict_batch(states, {}) (sans questions) |
une charge utile de réponse vide par état |
system_one(state, {}) |
une seule charge utile de réponse vide |
Aucune tokenisation ni passe avant n’a lieu dans ces cas, mais on_predict_start et on_predict_end
tournent.
Règles d’ordre
- Les hooks installés tournent avant les hooks par appel, toujours.
- Au sein d’une liste, les hooks tournent dans l’ordre de la liste.
- Pour un événement, chaque hook qui l’implémente tourne, dans cet ordre, avant l’événement suivant.
on_errortourne avanton_predict_endsur le chemin d’échec.on_evicttourne avanton_loadquand un seulloadévince et construit à la fois.
installed: [A, B] per-call: [C]
on_predict_start: A, B, C
on_predict_end: A, B, C
Concurrence
Agent et Router peuvent être appelés sans risque depuis de nombreux threads. Chaque appel crée son
propre PredictContext, donc les contextes ne fuient jamais entre requêtes. Le seul état partagé, ce
sont les objets hook eux-mêmes, donc un hook qui n’est pas sûr entre threads doit soit protéger son
propre état, soit être installé avec 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 sérialise chaque invocation de hook, pas des appels entiers : deux appels
peuvent quand même s’entrelacer entre les événements. Il utilise un verrou réentrant, donc un hook peut
rappeler le même Agent/Router sans interbloquer.
Voir aussi
- Erreurs : ce qui se passe quand un hook lève une erreur, par événement.
- Motifs et anti-motifs : comment bien utiliser le cycle de vie.