airflow.providers.common.ai.operators.agent

Operator for running pydantic-ai agents with tools and multi-turn reasoning.

Classes

HITLReviewLink

Link that opens the live chat window for a running feedback session.

AgentOperator

Run a pydantic-ai Agent with tools and multi-turn reasoning.

Module Contents

Bases: airflow.providers.common.compat.sdk.BaseOperatorLink

Link 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.

name = 'HITL Review'[source]

Name of the link. This will be the button name on the task UI.

Link to external system.

Parameters:
Returns:

link to external system

Return type:

str

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.HITLReviewMixin

Run a pydantic-ai Agent with tools and multi-turn reasoning.

Provide llm_conn_id and optional toolsets to 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_id and token usage are pushed to XCom under the run_id and usage keys, so a downstream task can reference the run and its cost. usage is this attempt’s own usage, not the cross-attempt cumulative total described under usage_limits below; it is pushed on a failed attempt too, so a downstream all_done task 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. The run_id also ties the task to its GenAI trace (see the provider’s observability docs). With enable_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_ids set in the connection’s extra field. None (default) reads the connection’s own extra field; an explicit [] disables a chain configured there. See PydanticAIHook for 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 Pydantic BaseModel subclass 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 of SQLToolset, MCPToolset and HookToolset (its hook’s conn_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 is SandboxToolset.attach_to, which is how SandboxToolset(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 than params or dag_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. A Toolset capability holding one of the toolsets above has its connection IDs templated the same way as toolsets=. 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. Passing capabilities inside agent_params still 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 a LoggingToolset that logs tool calls with timing at INFO level and arguments at DEBUG level. Set to False to disable.

  • agent_params (dict[str, Any] | None) – Additional keyword arguments passed to the pydantic-ai Agent constructor (e.g. retries, model_settings).

  • usage_limits (pydantic_ai.usage.UsageLimits | dict[str, Any] | None) –

    Optional pydantic-ai UsageLimits enforced 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 other template_fields entry, then coerced to that field’s type (Decimal, int, or bool). 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 a ValueError naming the field and the rendered value, instead of silently disabling the limit. A UsageLimits instance 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 default request_limit of 50 requests.

    A dict that omits request_limit gets the same default of 50 requests – pass "request_limit": None explicitly 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 set usage_limits=None, to keep the old per-attempt headroom. On Airflow < 3.3, or when usage_limits is None, 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 how durable replay 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. Default False. A replayed step adds nothing to the usage counted against usage_limits or reported in the usage XCom – 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_path to be set. Tools are durably cached when provided via toolsets= or via a concrete pydantic-ai Toolset capability. Tools reaching the agent through any other capability – MCP, PrefixTools, CombinedCapability, a Toolset backed by a callable factory, or capabilities loaded from a spec_file – are not cached and re-run on retry; put tools you need replayed in toolsets=. Provider-native capabilities such as WebSearch and Thinking execute inside the model call and are covered by model-response caching. Cannot be combined with a SandboxToolset (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-harness CodeMode capability (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 in agent_params["model_settings"] or a spec file take precedence: setting any anthropic_cache* key leaves Anthropic caching entirely to you, and a CachePoint in the prompt or message history leaves all of it to you. Set False where 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 list of pydantic-ai ModelMessage objects, or their JSON form as str / bytes – e.g. "{{ ti.xcom_pull(task_ids='ask', key='message_history', default='[]') }}" (pass default='[]' 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 key message_history so 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 with enable_hitl_review (raises) – the post-review transcript is not yet recoverable.

HITL Review parameters (requires the hitl_review plugin):

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-review or through the HITL Review extra link on the task instance. Default False. Cannot be combined with a SandboxToolset that 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. A SandboxToolset attached 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’s sandbox_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 HITLMaxIterationsError without calling the LLM. E.g. 5 allows changes at iterations 1–4. Default 5.

  • hitl_timeout (datetime.timedelta | None) – Maximum wall-clock time to wait for all review rounds combined. None means no timeout (the operator blocks until a terminal action).

  • 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(...), or requires_approval=True on 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 the awaiting_input state 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_limits applies to both sides of the pause. Not available together with durable, enable_hitl_review, a CodeMode capability, or a SandboxToolset that provisions its own sandbox; there, a tool that requires approval fails the task as before, except one called from inside CodeMode’s run_code, which does not run and is reported back to the model. A SandboxToolset attached 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 a tool_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 True and output_type is a Pydantic BaseModel subclass, the model instance is dumped to a dict via model_dump() before being pushed to XCom. Default False – the Pydantic instance flows through XCom unchanged. Set to True when a downstream consumer needs the dict shape.

deserialization_allowed_class_fields: ClassVar[tuple[str, ...]] = ('output_type',)[source]
template_fields: collections.abc.Sequence[str] = ('prompt', 'llm_conn_id', 'model_id', 'fallback_conn_ids', 'system_prompt', 'agent_params',...[source]
prompt[source]
llm_conn_id[source]
model_id = None[source]
fallback_conn_ids = None[source]
system_prompt = ''[source]
output_type[source]
serialize_output = False[source]
toolsets = None[source]
enable_tool_logging = True[source]
agent_params[source]
capabilities = None[source]
usage_limits = None[source]
message_history = None[source]
durable = False[source]
cache_prompt = True[source]
enable_hitl_review = False[source]
max_hitl_iterations = 5[source]
hitl_timeout = None[source]
hitl_poll_interval = 10.0[source]
tool_approval_timeout = None[source]
on_tool_approval_timeout = 'fail'[source]
tool_approval_assigned_users[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 RunUsage with the run that produced the output being reviewed – so a usage_limits cap bounds the initial run plus every regeneration combined – only when usage_limits is set. With usage_limits=None, each regeneration starts from a fresh RunUsage(), 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-tracked RunUsage (via the task state store) is gated on >= 3.3, in execute().

Was this entry helpful?