Documentation

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

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_route tourne 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ême routing, 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_end ne 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, un context.attach / detach OpenTelemetry) revient à la valeur qu’il a trouvée.
  • Un hook de start qui remplace ctx.states, ctx.questions ou 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 fait predict : 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_route qui épingle une requête par un alias ("ml") partage la passe avant de ce checkpoint, et ctx.model est 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 que predict prend comme arguments d’appel. Un hook de start qui définit ctx.max_len les 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 appel agent.predict_batch par 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_error avec ctx.error défini dessus, puis on_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 que ctx.error peut ê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

  1. Les hooks installés tournent avant les hooks par appel, toujours.
  2. Au sein d’une liste, les hooks tournent dans l’ordre de la liste.
  3. Pour un événement, chaque hook qui l’implémente tourne, dans cet ordre, avant l’événement suivant.
  4. on_error tourne avant on_predict_end sur le chemin d’échec.
  5. on_evict tourne avant on_load quand un seul load é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.