Документация

Жизненный цикл

Эта страница — точный порядок событий для каждой точки входа. Если вы прочитаете только одну схему, читайте схему Router; она — надмножество.

Agent.predict_batch

predict_batch — единственная реализация; system_one и predict вызывают её с одним состоянием.

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

Один и тот же объект ctx проходит через начало, ошибку и конец, поэтому run_id их коррелирует, и hook конца может прочитать ctx.error.

Agent.system_one / predict

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

Так system_one наследует все hooks и тот же жизненный цикл, при 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]

Ключевые моменты:

  • on_route выполняется до загрузки модели, поэтому hook может зафиксировать чекпойнт и избежать загрузки другого.
  • Hooks предсказания уровня Router оборачивают весь вызов. Они не перенаправляются в Agent; присоединённый Agent со своими hooks выполняет и их, что ожидаемо.
  • ctx.skip() уровня Router всё равно добавляет routing, поэтому форма возврата стабильна.

Router.predict_batch

Каждый результат — это то, что predict возвращает для этого запроса, поэтому hooks предсказания уровня Router выполняются здесь тоже по запросу: каждый запрос получает свой PredictContext, run_id и 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

Ключевые моменты:

  • Каждый hook начала для запросов одного чекпойнта выполняется до любого из их hooks конца, потому что они разделяют прямые проходы. Поэтому кэш, который заполняется в on_predict_end, не может обслужить дублирующееся состояние внутри одной группы чекпойнта; между вызовами — может.
  • По той же причине запросы завершаются в порядке, обратном порядку начала, поэтому hook, который что-то устанавливает в начале и сбрасывает в конце (значение contextvars, context.attach / detach из OpenTelemetry), возвращается к значению, которое застал.
  • Hook начала, который заменяет ctx.states, ctx.questions или бюджет токенов, меняет только свой собственный запрос: запросы группируются для прямого прохода после того, как выполнились их hooks начала. Мутировать общий словарь questions на месте — это другое, и не то, что делает predict: ничего из группы не выводится, пока не выполнены все её hooks начала, поэтому изменение доходит до каждого запроса, разделяющего этот словарь, включая те, чьи hooks начала выполнились раньше, а также до вызывающего кода. Вместо этого присвойте новый словарь в ctx.questions.
  • Имя чекпойнта разрешается до группировки, поэтому hook on_route, который фиксирует запрос по псевдониму ("ml"), разделяет прямой проход этого чекпойнта, а ctx.model — это разрешённое имя.
  • Запрос может нести собственные max_len / head_max_len — форму бюджета токенов на запрос, который predict принимает как аргументы вызова. Hook начала, устанавливающий ctx.max_len, переопределяет его, потому что hook выполняется после создания контекста. Запросы, требующие разных бюджетов, не могут разделять прямой проход, поэтому пакет, смешивающий бюджеты, делает один вызов agent.predict_batch на каждый бюджет.
  • Каждый начатый запрос получает ровно один on_predict_end, даже когда hook конца другого запроса вызывает исключение; первая такая ошибка вызывается после того, как выполнились все они.
  • Если группа чекпойнта падает, каждый из её начатых запросов падает с этим исключением: каждый получает on_error с установленным в него ctx.error, а затем on_predict_end. Это включает попадание в кэш и запрос, чья группа вопросов уже выполнилась, потому что вызывающий получает исключение и ни одного результата для любого из них, и это значит, что ctx.error может быть сбоем другого запроса (hook начала, вызывающий исключение для одного запроса, роняет его группу). Запросы групп чекпойнтов, которые уже завершились, завершились со своими результатами, как это было бы в более ранних вызовах [router.predict(...) for ...].

Жизненный цикл модели

on_load срабатывает, когда строится чекпойнт; on_evict — когда один освобождается. Оба выполняются после освобождения внутренней блокировки Router, поэтому hook может безопасно вызывать Router обратно.

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) регистрирует существующий agent и не запускает on_load, потому что ни один чекпойнт не был построен.

Кэширование через 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

См. examples/hooks/cache.py, чтобы увидеть работающий кэш.

Пустые входные данные

Hooks всё равно срабатывают, чтобы аудит видел каждый вызов:

входные данные ctx.results в конце
predict_batch([]) []
predict_batch(states, {}) (без вопросов) одна полезная нагрузка с пустым ответом на состояние
system_one(state, {}) одна полезная нагрузка с пустым ответом

В этих случаях не происходит ни токенизации, ни прямого прохода, но on_predict_start и on_predict_end выполняются.

Правила порядка

  1. Установленные hooks выполняются перед hooks, переданными на вызов, всегда.
  2. Внутри списка hooks выполняются в порядке списка.
  3. Для одного события выполняются все hooks, которые его реализуют, в этом порядке, до следующего события.
  4. on_error выполняется до on_predict_end на пути сбоя.
  5. on_evict выполняется до on_load, когда один load и вытесняет, и строит.
installed: [A, B]   per-call: [C]
on_predict_start: A, B, C
on_predict_end:   A, B, C

Конкурентность

Agent и Router безопасно вызывать из многих потоков. Каждый вызов создаёт свой PredictContext, поэтому контексты никогда не протекают между запросами. Единственное общее состояние — сами объекты hooks, поэтому hook, не безопасный для потоков, должен либо защищать своё состояние, либо устанавливаться с 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 сериализует каждый вызов hook, а не целые вызовы: два вызова всё ещё могут чередоваться между событиями. Он использует реентрантную блокировку, поэтому hook может вызывать тот же Agent/Router обратно без взаимной блокировки.

См. также