airflow.providers.common.ai.durable.replay_usage

Keeps durable replays out of the usage a run is counted and limited against.

Classes

ReplayUsageLedger

Nets durable replays out of the RunUsage a run is counted and limited against.

Functions

subtract_request_usage(run_usage, usage)

Subtract one response's token usage and cost from a run's usage, in place.

fill_replayed_cost(response)

Price a replayed response in place, the way pydantic-ai's graph is about to.

Module Contents

airflow.providers.common.ai.durable.replay_usage.subtract_request_usage(run_usage, usage)[source]

Subtract one response’s token usage and cost from a run’s usage, in place.

airflow.providers.common.ai.durable.replay_usage.fill_replayed_cost(response)[source]

Price a replayed response in place, the way pydantic-ai’s graph is about to.

CachingModel stores a live response before the graph prices it, so a cached response comes back with usage.cost=None even for a priced model. The graph fills the cost right after CachingModel.request returns and only when it is still unset, so pricing it here first means the cost this module subtracts is exactly the cost the graph then adds. Uses the public ModelResponse.cost(), which runs the same price lookup as the graph’s own best-effort pricing; a response it cannot price stays unpriced, as it would in the graph.

class airflow.providers.common.ai.durable.replay_usage.ReplayUsageLedger(*, run_usage, usage_limits)[source]

Nets durable replays out of the RunUsage a run is counted and limited against.

pydantic-ai counts a replayed step exactly like a live one: after CachingModel.request returns, the graph does requests += 1 and adds the response’s usage, and after CachingToolset.call_tool returns, the tool manager does tool_calls += 1. The ledger cancels each replay when it happens, so the run’s counts only move for live work. Nothing here is persisted: the result is the same whether the cache entry was written by the previous attempt, an attempt that failed halfway through its own replay, or an attempt before a clear (which resets the budget but keeps the durable cache).

Two of pydantic-ai’s limit checks run before the durable layer sees the step it guards: check_before_request before each model request, and the up-front check_before_tool_call projection over a whole step’s function-tool calls. For those, the ledger takes a credit ahead of time – when a replayed model response is followed in the cache by the next model step or by cached tool results – and resolves each credit when the step runs: a replay keeps it, a live call gives it back and re-runs the check pydantic-ai made against the credited count. Credits that are never resolved (the run took a different path, or stopped) are given back by settle(), so they never reach the persisted total; until then, a credit for a stale cache entry the run never reaches has only loosened that one up-front check.

Parameters:
  • run_usage (pydantic_ai.usage.RunUsage) – The RunUsage passed to the run as usage=.

  • usage_limits (pydantic_ai.usage.UsageLimits | None) – The limits passed to the run; None means pydantic-ai’s defaults, as in Agent.run.

run_usage[source]
credit_request()[source]

Credit the next model request, which the cache says will be a replay.

credit_tool_steps(steps, *, batch_calls)[source]

Credit the tool calls at these step indices, which have cached results.

Parameters:

batch_calls (int) – The number of function-tool calls this replayed response issued, i.e. the count pydantic-ai’s up-front check_before_tool_call projects forward before any of them run.

settle()[source]

Give back every unresolved credit.

Returns:

Whether a request credit was outstanding.

Return type:

bool

record_model_replay(response, *, continuation)[source]

Cancel what the graph is about to add for a replayed model response.

A continuation segment (Anthropic pause_turn, OpenAI background mode) is not a request of its own: the graph counts one request per chain and commits the merged response’s usage once. Merging sums the segments’ usage, except when a segment re-polls the same provider response id, where the merged usage is that segment’s cumulative snapshot – so the previous segment’s subtraction is given back before this one is taken.

Between this subtraction and the graph adding the response’s usage back, the graph awaits after_model_request; a non-ModelRetry exception raised there, or an AirflowTaskTimeout landing in that window, leaves this response’s tokens and cost under-counted.

record_live_model_request(*, had_request_credit)[source]

Re-run the pre-request check for a credited request that turned out to be live.

pydantic-ai checked this request against the credited count, which assumed it would replay for free.

record_tool_replay(step)[source]

Cancel the tool_calls += 1 the tool manager does after a replayed call returns.

record_live_tool_call(step)[source]

Give back a credit whose cached result did not replay, and re-check the limit.

pydantic-ai’s up-front check projected the whole batch this credited call belongs to as if every credited call in it would replay for free. With this one turning out to be live, that assumption is wrong for one more call: re-run the check against pydantic-ai’s original projection plus one for each credited call in the batch that has turned out to be live so far (including this one). A credit that is still unresolved keeps counting as a replay – if it does replay, the original projection was already right for it; if it turns out live later, its own call to this method re-checks again. So the last credited call in a batch to turn live is checked against the sum of every live call in the batch, regardless of the order concurrent calls finish in, and the batch as a whole never exceeds the limit.

batch_calls (passed to credit_tool_steps()) is derived from the response’s tool-call parts matched against model_request_parameters’s function_tools, which also lists tools that need approval or are external – so in a non-resume batch, where those calls are deferred rather than executed, batch_calls can count more than pydantic-ai’s own projection does. That over-count flows into _tool_batch_projection and this method’s recheck, so it can only make the recheck stricter – a possible false block, never a missed one. pydantic-ai’s projection, in turn, also covers 'unknown'-kind calls that batch_calls does not; that under-count never reaches this method or record_tool_replay(), because nothing calls record_tool_call for calls of 'unknown' kind, so it can’t push the final count over the limit.

Was this entry helpful?