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

log

USAGE_BUDGET_KEY

Classes

TaskStateStoreUsageBudget

Persists cumulative RunUsage across task attempts in the AIP-103 task state store.

Functions

dump_run_usage(usage)

Serialize a RunUsage to a JSON-safe dict.

load_run_usage(raw, *, key)

Deserialize a dict produced by dump_run_usage() back into a RunUsage.

copy_run_usage(usage)

Return an independent copy of usage.

subtract_run_usage(total, base)

Return the field-by-field usage in total that is not already in base.

Module Contents

airflow.providers.common.ai.utils.usage_budget.log[source]
airflow.providers.common.ai.utils.usage_budget.USAGE_BUDGET_KEY = '__commonai_usage__'[source]
airflow.providers.common.ai.utils.usage_budget.dump_run_usage(usage)[source]

Serialize a RunUsage to 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 a RunUsage.

Unknown keys in raw are ignored, so a record written by a newer version of this module still loads. Only the fields RunUsage currently declares are read.

Raises:

ValueError – raw is not a dict, or a field has the wrong shape (cost not a valid number, a count field not an int, a float field such as audio_seconds not a number, details not 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.copy shares the details dict, which RunUsage.incr mutates 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 through dump_run_usage() / load_run_usage() copies details too.

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 returns None for cost whenever it is unchanged – indistinguishable from “unknown” – cost here is None only when both sides are None; otherwise it is a numeric delta (0 when 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 RunUsage across task attempts in the AIP-103 task state store.

The stored record also carries the task instance’s max_tries at write time. Airflow bumps ti.max_tries when a task is cleared (clear_task_instances, airflow-core/src/airflow/models/taskinstance.py) but never changes it on an ordinary retry (handle_failure does not touch it). So a max_tries that differs from the current task instance’s means this row was written before the most recent clear – a new budget cycle – and load() starts over from zero instead of carrying a stale spend forward. A same-cycle retry, which never changes max_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.

load()[source]

Return the cumulative usage so far, or a fresh RunUsage() if none or stale.

save(usage)[source]

Best-effort write; a failure here must never fail the task, only the caller’s raise matters.

clear()[source]

Best-effort delete, called once the whole execute succeeds.

Was this entry helpful?