airflow.providers.common.ai.durable.replay_usage¶
Keeps durable replays out of the usage a run is counted and limited against.
Classes¶
Nets durable replays out of the |
Functions¶
|
Subtract one response's token usage and cost from a run's usage, in place. |
|
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.
CachingModelstores a live response before the graph prices it, so a cached response comes back withusage.cost=Noneeven for a priced model. The graph fills the cost right afterCachingModel.requestreturns 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 publicModelResponse.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
RunUsagea run is counted and limited against.pydantic-ai counts a replayed step exactly like a live one: after
CachingModel.requestreturns, the graph doesrequests += 1and adds the response’s usage, and afterCachingToolset.call_toolreturns, the tool manager doestool_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_requestbefore each model request, and the up-frontcheck_before_tool_callprojection 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 bysettle(), 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
RunUsagepassed to the run asusage=.usage_limits (pydantic_ai.usage.UsageLimits | None) – The limits passed to the run;
Nonemeans pydantic-ai’s defaults, as inAgent.run.
- 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_callprojects forward before any of them run.
- settle()[source]¶
Give back every unresolved credit.
- Returns:
Whether a request credit was outstanding.
- Return type:
- 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-ModelRetryexception raised there, or anAirflowTaskTimeoutlanding 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 += 1the 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 tocredit_tool_steps()) is derived from the response’s tool-call parts matched againstmodel_request_parameters’sfunction_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_callscan count more than pydantic-ai’s own projection does. That over-count flows into_tool_batch_projectionand 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 thatbatch_callsdoes not; that under-count never reaches this method orrecord_tool_replay(), because nothing callsrecord_tool_callfor calls of'unknown'kind, so it can’t push the final count over the limit.