airflow.providers.common.ai.utils.usage_budget¶
Cross-attempt pydantic_ai.usage.RunUsage accounting, backed by the task state store.
RunUsage is a plain (non-pydantic) dataclass, so it needs manual JSON
(de)serialization – dump_run_usage / load_run_usage do that, iterating
dataclasses.fields rather than hard-coding the field list so a future
pydantic-ai field is carried through automatically. This module has no
top-level import of any Airflow >= 3.3-only symbol: it must stay importable on
older Airflow versions, even though TaskStateStoreUsageBudget is only
constructed on 3.3+ (see AgentOperator._build_usage_budget).
Attributes¶
Classes¶
Persists cumulative |
Functions¶
|
Serialize a |
|
Deserialize a dict produced by |
|
Return an independent copy of usage. |
|
Return the field-by-field usage in total that is not already in base. |
Module Contents¶
- airflow.providers.common.ai.utils.usage_budget.dump_run_usage(usage)[source]¶
Serialize a
RunUsageto a JSON-safe dict.
- airflow.providers.common.ai.utils.usage_budget.load_run_usage(raw, *, key)[source]¶
Deserialize a dict produced by
dump_run_usage()back into aRunUsage.Unknown keys in raw are ignored, so a record written by a newer version of this module still loads. Only the fields
RunUsagecurrently declares are read.- Raises:
ValueError – raw is not a dict, or a field has the wrong shape (
costnot a valid number, a count field not an int, a float field such asaudio_secondsnot a number,detailsnot a dict). The message names key so the error points at which task state store key to delete to reset the budget.
- airflow.providers.common.ai.utils.usage_budget.copy_run_usage(usage)[source]¶
Return an independent copy of usage.
copy.copyshares thedetailsdict, whichRunUsage.incrmutates in place (usage.py_incr_usage_tokens), so a shallow copy would let a later increment of the original leak into the copy. Round-tripping throughdump_run_usage()/load_run_usage()copiesdetailstoo.
- airflow.providers.common.ai.utils.usage_budget.subtract_run_usage(total, base)[source]¶
Return the field-by-field usage in total that is not already in base.
Unlike
RunUsage.__sub__(usage.py), which returnsNoneforcostwhenever it is unchanged – indistinguishable from “unknown” –costhere isNoneonly when both sides areNone; otherwise it is a numeric delta (0when unchanged), because the result is meant to be reported (XCom, logs), where a numeric zero and an unknown cost are not the same thing.
- class airflow.providers.common.ai.utils.usage_budget.TaskStateStoreUsageBudget(accessor, *, max_tries)[source]¶
Persists cumulative
RunUsageacross task attempts in the AIP-103 task state store.The stored record also carries the task instance’s
max_triesat write time. Airflow bumpsti.max_trieswhen a task is cleared (clear_task_instances,airflow-core/src/airflow/models/taskinstance.py) but never changes it on an ordinary retry (handle_failuredoes not touch it). So amax_triesthat differs from the current task instance’s means this row was written before the most recent clear – a new budget cycle – andload()starts over from zero instead of carrying a stale spend forward. A same-cycle retry, which never changesmax_tries, keeps accumulating.- Parameters:
accessor (airflow.sdk.execution_time.context.TaskStateStoreAccessor) – The task state store accessor for the current task instance (
context["task_state_store"]).max_tries (int) – The current task instance’s
max_tries.