airflow.providers.common.ai.operators.agent¶
Operator for running pydantic-ai agents with tools and multi-turn reasoning.
Classes¶
Link that opens the live chat window for a running feedback session. |
|
Run a pydantic-ai Agent with tools and multi-turn reasoning. |
Module Contents¶
- class airflow.providers.common.ai.operators.agent.HITLReviewLink[source]¶
Bases:
airflow.providers.common.compat.sdk.BaseOperatorLinkLink that opens the live chat window for a running feedback session.
The URL is constructed directly from the task instance key so that the link is available immediately — even while the task is still running — without waiting for an XCom value to be committed.
- get_link(operator, *, ti_key)[source]¶
Link to external system.
- Parameters:
operator (airflow.providers.common.compat.sdk.BaseOperator) – The Airflow operator object this link is associated to.
ti_key (airflow.providers.common.compat.sdk.TaskInstanceKey) – TaskInstance ID to return link for.
- Returns:
link to external system
- Return type:
- class airflow.providers.common.ai.operators.agent.AgentOperator(*, prompt, llm_conn_id, model_id=None, fallback_conn_ids=None, system_prompt='', output_type=str, toolsets=None, capabilities=None, enable_tool_logging=True, agent_params=None, usage_limits=None, durable=False, cache_prompt=True, message_history=None, enable_hitl_review=False, max_hitl_iterations=5, hitl_timeout=None, hitl_poll_interval=10.0, serialize_output=False, tool_approval_timeout=None, on_tool_approval_timeout='fail', tool_approval_assigned_users=None, **kwargs)[source]¶
Bases:
airflow.providers.common.ai.mixins.cancellable_run.CancellableAgentRunMixin,airflow.providers.common.compat.sdk.BaseOperator,airflow.providers.common.ai.mixins.hitl_review.HITLReviewMixinRun a pydantic-ai Agent with tools and multi-turn reasoning.
Provide
llm_conn_idand optionaltoolsetsto let the operator build and run the agent. The agent reasons about the prompt, calls tools in a multi-turn loop, and returns a final answer.Alongside the returned agent output, the run’s
run_idand tokenusageare pushed to XCom under therun_idandusagekeys, so a downstream task can reference the run and its cost.usageis this attempt’s own usage, not the cross-attempt cumulative total described underusage_limitsbelow; it is pushed on a failed attempt too, so a downstreamall_donetask or failure callback can read what the last attempt spent – XCom is cleared at the start of every attempt, so only the most recent attempt’s value survives, not each historical attempt’s. Therun_idalso ties the task to its GenAI trace (see the provider’s observability docs). Withenable_hitl_review, these reflect the initial model run, not the human-feedback regenerations.- Parameters:
prompt (str) – The prompt to send to the agent.
llm_conn_id (str) – Connection ID for the LLM provider.
model_id (str | None) – Model identifier (e.g.
"openai:gpt-5"). Overrides the model stored in the connection’s extra field.fallback_conn_ids (list[str] | None) – Connection IDs to fail over to, in order, when the primary provider is unavailable. Overrides the
fallback_conn_idsset in the connection’s extra field.None(default) reads the connection’s own extra field; an explicit[]disables a chain configured there. SeePydanticAIHookfor how blank entries in the list are dropped.system_prompt (str) – System-level instructions for the agent.
output_type (type) – Expected output type. Default
str. Set to a PydanticBaseModelsubclass for structured output; the model instance is returned to XCom unchanged so downstream tasks can type-hint it directly. The class must be defined at module scope – nested classes cannot be deserialized from XCom.toolsets (list[pydantic_ai.toolsets.abstract.AbstractToolset] | None) – List of pydantic-ai toolsets the agent can use (e.g.
SQLToolset,HookToolset). The connection IDs ofSQLToolset,MCPToolsetandHookToolset(its hook’sconn_name_attr) are templated, e.g.SQLToolset(db_conn_id="warehouse_{{ var.value.environment }}")per environment, or"tenant_{{ task.op_kwargs.customer }}"per map index of a mapped@task.agent, and so isSandboxToolset.attach_to, which is howSandboxToolset(attach_to="{{ ti.xcom_pull('provision') }}")receives the sandbox an upstream task created. Each task instance renders its own copy and logs the rendered toolset id; the toolset object in the Dag file is not modified. Derive the connection ID from values the Dag controls rather thanparamsordag_run.conf, which whoever triggers the Dag controls.capabilities (list[pydantic_ai.capabilities.AgentCapability[Any]] | None) – pydantic-ai capabilities for the agent, e.g.
[Thinking(effort="high"), WebSearch()]. A capability bundles tools, instructions, model settings and lifecycle hooks; pydantic-ai wraps their hooks in list order, first outermost, unless a capability declares its own position. AToolsetcapability holding one of the toolsets above has its connection IDs templated the same way astoolsets=. Capabilities passed here are not stored in the serialized Dag (the worker builds them from the Dag file), except on a mapped task, where they are stored as their repr. Passingcapabilitiesinsideagent_paramsstill works, but stores each capability’s repr in the serialized Dag, and cannot be combined with this argument (the task fails when it runs).enable_tool_logging (bool) – When
True(default), wraps each toolset in aLoggingToolsetthat logs tool calls with timing at INFO level and arguments at DEBUG level. Set toFalseto disable.agent_params (dict[str, Any] | None) – Additional keyword arguments passed to the pydantic-ai
Agentconstructor (e.g.retries,model_settings).usage_limits (pydantic_ai.usage.UsageLimits | dict[str, Any] | None) –
Optional pydantic-ai
UsageLimitsenforced on every agent run (initial run, durable replay, and HITL regeneration), or a dict of the same fields (e.g.{"cost_limit": "{{ params.budget }}", "request_limit": 5}). The dict form is templated: each value is rendered by Jinja like any othertemplate_fieldsentry, then coerced to that field’s type (Decimal,int, orbool). A value that cannot be coerced – a Variable that exists but is empty renders to"", a typo renders to a non-numeric string – fails the task with aValueErrornaming the field and the rendered value, instead of silently disabling the limit. AUsageLimitsinstance passed directly is used as-is and is not templated or validated.None(default) sets no token, cost, or tool-call limits, but pydantic-ai still caps each run at its defaultrequest_limitof50requests.A dict that omits
request_limitgets the same default of50requests – pass"request_limit": Noneexplicitly for no request cap.On Airflow >= 3.3, this counts usage across every attempt combined – initial run, retries, and HITL regenerations all add to one running total instead of resetting each attempt. Scale each limit by
retries + 1, or setusage_limits=None, to keep the old per-attempt headroom. On Airflow < 3.3, or whenusage_limitsisNone, each attempt is checked and counted on its own, unchanged. See Single prompts: LLMOperator and @task.llm for the full caveats, and the cross-attempt usage budget for how it is persisted, reset, and howdurablereplay and HITL regeneration interact with it.durable (bool) – Experimental. When
True, enables step-level caching of model responses and tool results for durable execution. On retry, cached steps are replayed instead of re-executing. Each cached step is verified against the current request before replay: if the prompt, model, settings, tools, or message history changed since the failed attempt, the affected steps re-run live (with a warning) instead of replaying stale results. DefaultFalse. A replayed step adds nothing to the usage counted againstusage_limitsor reported in theusageXCom – not its request, tokens, cost, or tool calls – so every attempt counts only the model and tool calls it actually makes. This holds the same way after clearing a failed task instance: it starts a fresh budget but keeps the durable cache its attempts left behind, and whatever the rerun replays from that cache is free. On Airflow >= 3.3 the cache is kept in the AIP-103 task state store, so no extra configuration is needed. On older Airflow versions it is persisted to ObjectStorage and requires[common.ai] durable_cache_pathto be set. Tools are durably cached when provided viatoolsets=or via a concrete pydantic-aiToolsetcapability. Tools reaching the agent through any other capability –MCP,PrefixTools,CombinedCapability, aToolsetbacked by a callable factory, or capabilities loaded from aspec_file– are not cached and re-run on retry; put tools you need replayed intoolsets=. Provider-native capabilities such asWebSearchandThinkingexecute inside the model call and are covered by model-response caching. Cannot be combined with aSandboxToolset(raises), attached or not: a replayed tool result describes a workspace state the replay did not reproduce, and the first call that misses the cache runs against whatever the sandbox holds now. Cannot be combined with a pydantic-ai-harnessCodeModecapability (raises).cache_prompt (bool) – When
True(default), asks the provider to cache the tool definitions, system prompt and conversation so far, so the next request in the run – and a mapped task’s other instances within the cache lifetime – reads them back at a fraction of the input price instead of paying for them again. Turns on prompt caching for Anthropic models and for Bedrock and OpenRouter models that support it; a no-op for OpenAI and Gemini, which cache long prompts on their own. A provider’s own cache settings inagent_params["model_settings"]or a spec file take precedence: setting anyanthropic_cache*key leaves Anthropic caching entirely to you, and aCachePointin the prompt or message history leaves all of it to you. SetFalsewhere a cache write is rarely read back, such as a single long request that is not mapped. See Prompt caching for when caching costs more than it saves.message_history (list[pydantic_ai.messages.ModelMessage] | str | bytes | None) – Prior conversation to seed the run with, for multi-turn sessions that span task runs. Accepts a
listof pydantic-aiModelMessageobjects, or their JSON form asstr/bytes– e.g."{{ ti.xcom_pull(task_ids='ask', key='message_history', default='[]') }}"(passdefault='[]'so the first run, with no XCom yet, starts a fresh session instead of failing to parse the string"None").None(default) is a single-turn run – no behavior change. When set (an empty[]/""starts a fresh session), the full transcript after the run –result.all_messages()– is pushed to XCom under the keymessage_historyso the next run can resume. Persisting that transcript under a session key (e.g. in object storage) is the DAG’s responsibility. The transcript is cumulative and grows each turn; for long sessions use an object-storage XCom backend or trim old turns. Not supported together withenable_hitl_review(raises) – the post-review transcript is not yet recoverable.
HITL Review parameters (requires the
hitl_reviewplugin):- Parameters:
enable_hitl_review (bool) – When
True, the operator enters an iterative review loop after the first generation. A human reviewer can approve, reject, or request changes via the plugin’s REST API at/hitl-reviewor through the HITL Review extra link on the task instance. DefaultFalse. Cannot be combined with aSandboxToolsetthat provisions its own sandbox (raises): regeneration after feedback is a second run, which would start from an empty sandbox while its history describes the first run’s files. ASandboxToolsetattached to a sandbox another task owns (attach_to) is fine, since both runs find the same files, as long as the reviewer answers inside that sandbox’s lifetime: the wait spends the provisioning backend’ssandbox_timeout.max_hitl_iterations (int) – Maximum outputs shown to the reviewer (1 = initial output). When the reviewer requests changes at iteration >= this limit, the task fails with
HITLMaxIterationsErrorwithout calling the LLM. E.g. 5 allows changes at iterations 1–4. Default5.hitl_timeout (datetime.timedelta | None) – Maximum wall-clock time to wait for all review rounds combined.
Nonemeans no wall-clock timeout; the review still ends on a terminal action (approve or reject),max_hitl_iterations, or when polling the human action XCom fails too many times in a row.hitl_poll_interval (float) – Seconds between XCom polls while waiting for a human response. Default
10.
Per-tool approval (Airflow 3.3+, experimental):
Mark the tools a human must approve with pydantic-ai’s own API –
toolset.approval_required(...), orrequires_approval=Trueon a function tool – and the task pauses before running them. The pending calls, with their arguments, appear on the Required Actions page; the task waits in theawaiting_inputstate without holding a worker slot. On Approve the calls run and the agent carries on. On Reject the agent is told the call was denied (with the reviewer’s reason, when given) and carries on without it. A task instance asks at most once per Dag run, across retries and clears; a second request fails the task.usage_limitsapplies to both sides of the pause. Not available together withdurable,enable_hitl_review, aCodeModecapability, or aSandboxToolsetthat provisions its own sandbox; there, a tool that requires approval fails the task as before, except one called from insideCodeMode’srun_code, which does not run and is reported back to the model. ASandboxToolsetattached to a sandbox another task owns is fine: the sandbox outlives the pause.- Parameters:
tool_approval_timeout (datetime.timedelta | None) – Experimental. How long the pause waits for a decision.
None(default) waits indefinitely. Must be positive.on_tool_approval_timeout (Literal['fail', 'deny']) – Experimental. What a timed-out pause does:
"fail"(default) fails the task,"deny"rejects the pending calls so the agent carries on without them, and needs atool_approval_timeout. There is no approve-on-timeout.tool_approval_assigned_users (airflow.sdk.execution_time.hitl.HITLUser | collections.abc.Iterable[airflow.sdk.execution_time.hitl.HITLUser] | None) – Experimental. Users allowed to decide.
None(default) leaves it to anyone who can act on the task’s Required Actions.serialize_output (bool) – If
Trueandoutput_typeis a PydanticBaseModelsubclass, the model instance is dumped to adictviamodel_dump()before being pushed to XCom. DefaultFalse– the Pydantic instance flows through XCom unchanged. Set toTruewhen a downstream consumer needs the dict shape.
- template_fields: collections.abc.Sequence[str] = ('prompt', 'llm_conn_id', 'model_id', 'fallback_conn_ids', 'system_prompt', 'agent_params',...[source]¶
- property llm_hook: airflow.providers.common.ai.hooks.pydantic_ai.PydanticAIHook[source]¶
Return PydanticAIHook for the configured LLM connection.
- execute(context)[source]¶
Derive when creating an operator.
The main method to execute the task. Context is the same dictionary used as when rendering jinja templates.
Refer to get_template_context for more context.
- resume_after_tool_approval(context, tool_call_ids, usage, transcript_sha256, toolset_ids, event, attempt_usage=None)[source]¶
Continue a run paused by
_pause_for_tool_approval()with the reviewer’s decision.
- regenerate_with_feedback(*, feedback, message_history)[source]¶
Re-run the agent with feedback appended to the conversation history.
Shares the cross-run
RunUsagewith the run that produced the output being reviewed – so ausage_limitscap bounds the initial run plus every regeneration combined – only whenusage_limitsis set. Withusage_limits=None, each regeneration starts from a freshRunUsage(), matching the behaviour before the cross-attempt budget existed: nothing shares usage. This applies on every Airflow version; only the cross-attempt persistence of a shared, budget-trackedRunUsage(via the task state store) is gated on >= 3.3, inexecute().