Agents with tools: AgentOperator and @task.agent

Use AgentOperator or the @task.agent decorator to run an LLM agent with tools: the agent reasons about the prompt, calls tools (database queries, API calls, etc.) in a multi-turn loop, and returns a final answer.

This is different from LLMOperator, which sends a single prompt and returns the output. AgentOperator manages a stateful tool-call loop where the LLM decides which tools to call and when to stop.

SQL agent

The most common pattern: give an agent access to a database so it can answer questions by writing and executing SQL.

airflow/providers/common/ai/example_dags/example_agent.py[source]

if SQLToolset is not None:

    @dag(tags=["example"])
    def example_agent_operator_sql():
        AgentOperator(
            task_id="analyst",
            prompt="What are the top 5 customers by order count?",
            llm_conn_id="pydanticai_default",
            system_prompt=(
                "You are a SQL analyst. Use the available tools to explore "
                "the schema and answer the question with data."
            ),
            toolsets=[
                # ``allowed_tables`` scopes the agent's intent, but it is an
                # application-level guardrail, not a security boundary. Point
                # ``postgres_default`` at a least-privilege role whose SELECT grants
                # are limited to these tables -- that is the boundary that holds even
                # if the agent (which may be under prompt injection) reaches for data
                # through a function the parser cannot see. See the "Security" section
                # of the toolsets docs.
                SQLToolset(
                    db_conn_id="postgres_default",
                    allowed_tables=["customers", "orders"],
                    # Functions sqlglot cannot type are rejected while allowed_tables is
                    # set; list any the agent legitimately needs (e.g. to shape output).
                    allowed_functions=["json_build_object"],
                    max_rows=20,
                )
            ],
        )

The SQLToolset provides four tools to the agent:

Tool

Description

list_tables

Lists available table names (filtered by allowed_tables if set)

get_schema

Returns column names and types for a table

query

Executes a SQL query and returns rows as JSON

check_query

Validates SQL syntax without executing it

Hook-based tools

Wrap any Airflow Hook’s methods as agent tools using HookToolset. Only methods you explicitly list are exposed; there is no auto-discovery.

airflow/providers/common/ai/example_dags/example_agent.py[source]

@dag(tags=["example"])
def example_agent_operator_hook():
    from airflow.providers.http.hooks.http import HttpHook

    http_hook = HttpHook(http_conn_id="my_api")

    AgentOperator(
        task_id="api_explorer",
        prompt="What endpoints are available and what does /status return?",
        llm_conn_id="pydanticai_default",
        system_prompt="You are an API explorer. Use the tools to discover and call endpoints.",
        toolsets=[
            HookToolset(
                http_hook,
                allowed_methods=["run"],
                tool_name_prefix="http_",
            )
        ],
    )


TaskFlow decorator

The @task.agent decorator wraps AgentOperator. The function returns the prompt string; all other parameters are passed to the operator.

airflow/providers/common/ai/example_dags/example_agent.py[source]

if SQLToolset is not None:

    @dag(tags=["example"])
    def example_agent_decorator():
        @task.agent(
            llm_conn_id="pydanticai_default",
            system_prompt="You are a data analyst. Use tools to answer questions.",
            toolsets=[
                SQLToolset(
                    db_conn_id="postgres_default",
                    allowed_tables=["orders"],
                )
            ],
        )
        def analyze(question: str):
            return f"Answer this question about our orders data: {question}"

        analyze("What was our total revenue last month?")

Multimodal prompts

The decorated callable may also return a Sequence[UserContent] – for example, a list mixing strings with ImageUrl, BinaryContent, or other pydantic-ai user-content types – to send vision, audio, or document inputs to the model. This mirrors the input types accepted by pydantic-ai’s Agent.run_sync.

from pydantic_ai.messages import ImageUrl


@task.agent(llm_conn_id="pydanticai_default", system_prompt="You are an image analyst.")
def analyze_review(image_url: str):
    return ["Describe what you see:", ImageUrl(url=image_url)]

Note

Combining a non-string prompt with enable_hitl_review=True is not currently supported – the HITL session model stores the prompt as a string, so a Sequence prompt will raise at the review boundary.

Structured output

Set output_type to a Pydantic BaseModel subclass to get structured data back. The model instance is pushed to XCom unchanged so downstream tasks can type-hint the class directly (def downstream(result: MyModel)) and use attribute access (result.field).

Structured output and XCom explains the XCom deserialization rules, the cross-Dag gap and serialize_output.

airflow/providers/common/ai/example_dags/example_agent.py[source]

# Pydantic output classes must be defined at module scope so downstream
# tasks can re-import them when deserializing the XCom payload.
class Analysis(BaseModel):
    """Structured analysis output for the agent example."""

    summary: str
    top_items: list[str]
    row_count: int


airflow/providers/common/ai/example_dags/example_agent.py[source]

if SQLToolset is not None:

    @dag(tags=["example"])
    def example_agent_structured_output():
        @task.agent(
            llm_conn_id="pydanticai_default",
            system_prompt="You are a data analyst. Return structured results.",
            output_type=Analysis,
            toolsets=[SQLToolset(db_conn_id="postgres_default")],
        )
        def analyze(question: str):
            return f"Analyze: {question}"

        analyze("What are the trending products this week?")

Chaining with downstream tasks

The agent’s output is pushed to XCom like any other operator, so downstream tasks can consume it.

airflow/providers/common/ai/example_dags/example_agent.py[source]

if SQLToolset is not None:

    @dag(tags=["example"])
    def example_agent_chain():
        @task.agent(
            llm_conn_id="pydanticai_default",
            system_prompt="You are a SQL analyst.",
            toolsets=[SQLToolset(db_conn_id="postgres_default", allowed_tables=["orders"])],
        )
        def investigate(question: str):
            return f"Investigate: {question}"

        @task
        def send_report(analysis: str):
            """Send the agent's analysis to a downstream system."""
            print(f"Report: {analysis}")
            return analysis

        result = investigate("Summarize order trends for last quarter")
        send_report(result)

Dynamic system prompt

system_prompt is a templated field, so instead of a static string it can be a Jinja expression that reads a value an earlier task already computed – for example, tailoring the agent’s instructions to a classification produced upstream.

airflow/providers/common/ai/example_dags/example_agent.py[source]

@dag(tags=["example"])
def example_agent_dynamic_system_prompt():
    @task
    def classify(ticket: str) -> dict:
        category = "shipping" if "order" in ticket.lower() else "other"
        return {"priority": "high", "category": category}

    @task.agent(
        llm_conn_id="pydanticai_default",
        # system_prompt is a templated field -- Jinja renders it at task-run
        # time, pulling the classification an upstream task already computed.
        system_prompt=(
            "You are handling a {{ ti.xcom_pull(task_ids='classify')['priority'] }}-priority "
            "'{{ ti.xcom_pull(task_ids='classify')['category'] }}' ticket. "
            "Draft a concise, friendly reply."
        ),
    )
    def draft_reply(ticket: str, triage: dict) -> str:
        # `triage` creates the task dependency; its content also flows into
        # system_prompt via Jinja above. The returned string is the *prompt*
        # sent to the agent -- the drafted reply is this task's XCom output.
        return f"Draft a reply for: {ticket}"

    ticket = "Where is my order? It still hasn't shipped."
    draft_reply(ticket, classify(ticket))


Open the Rendered Template tab on the task instance to see the substituted system_prompt after Jinja fills in classify’s XCom values.

Reuse one agent across tasks

When several tasks, or several Dags, run the same agent, define it once and import it. AgentOperator and @task.agent take the whole agent definition as keyword arguments, so a dict in a module next to your Dags is enough. A task that needs something different overrides single keys.

# dags/shared_agents/__init__.py
from airflow.providers.common.ai.toolsets.sql import SQLToolset

ORDERS_ANALYST = {
    "llm_conn_id": "pydanticai_default",
    "system_prompt": "You are the orders analyst. Answer only from the orders database.",
    "toolsets": [SQLToolset(db_conn_id="orders_db", allowed_tables=["orders"])],
    "agent_params": {"name": "orders_analyst"},
}
# dags/orders.py
from shared_agents import ORDERS_ANALYST

from airflow.sdk import dag, task


@dag(schedule=None)
def orders():
    @task.agent(**ORDERS_ANALYST)
    def weekly_summary() -> str:
        return "Summarize this week's orders."

    @task.agent(**{**ORDERS_ANALYST, "system_prompt": "Answer in one sentence."})
    def one_liner() -> str:
        return "How many orders are there?"

    weekly_summary()
    one_liner()


orders()

When span export is on (see Observability (OpenTelemetry tracing)), the name in agent_params becomes the gen_ai.agent.name attribute on each agent run’s span, so traces from every task that uses the definition group under one agent.

To keep the definition out of Python, for example to share it with a program that does not run on Airflow, write a pydantic-ai agent spec file and pass its path through agent_params. A model set on the connection wins over a model in the file, and system_prompt is added to the file’s instructions:

# dags/shared_agents/orders_analyst.yaml
name: orders_analyst
instructions: >
  You are the orders analyst. Answer only from the orders database.
retries: 2
from pathlib import Path

AgentOperator(
    task_id="orders_question",
    llm_conn_id="pydanticai_default",
    prompt="How many orders are there?",
    agent_params={"spec_file": Path(__file__).parent / "shared_agents" / "orders_analyst.yaml"},
)

Build the path from __file__: a relative path resolves against the worker’s working directory, not the Dag file.

With durable=True, tools from capabilities declared in the spec file are not replayed on retry; they run again. Pass tools you need replayed in toolsets=.

Agent features

Four features have pages of their own:

  • Multi-turn sessions and message history: pass message_history to carry a conversation across runs.

  • Durable execution: set durable=True to replay completed model and tool steps on retry instead of paying for them again.

  • Guardrails and capabilities: pass pydantic-ai capabilities and pydantic-ai-shields guardrails through agent_params.

  • Code mode: set code_mode=True to collapse the agent’s tools into a single run_code tool the model drives by writing Python.

Durable execution

Moved to Durable execution.

Parameters

  • prompt: The prompt to send to the agent (operator) or the return value of the decorated function (decorator).

  • llm_conn_id: Airflow connection ID for the LLM provider.

  • model_id: Model identifier (e.g. "openai:gpt-5"). Overrides the connection’s extra field.

  • system_prompt: System-level instructions for the agent. Supports Jinja templating.

  • output_type: Expected output type (default: str). Set to a Pydantic BaseModel for structured output.

  • toolsets: List of pydantic-ai toolsets (SQLToolset, HookToolset, AgentSkillsToolset for Agent Skills: AgentSkillsToolset, etc.).

  • enable_tool_logging: Wrap each toolset in LoggingToolset so that every tool call is logged in real time. Default True.

  • agent_params: Additional keyword arguments passed to the pydantic-ai Agent constructor (e.g. retries, model_settings, capabilities). See Guardrails and capabilities for how to enable pydantic-ai capabilities such as Thinking, WebSearch, and ImageGeneration.

  • usage_limits: Optional pydantic-ai UsageLimits enforced on every agent run (initial run, durable replay, and HITL regeneration), or a dict of the same fields – the dict form is templated via Jinja, then coerced per field type, failing the task with a ValueError naming the field if a rendered value doesn’t parse. Use it to cap requests, tokens, or tool calls per task – agents are particularly prone to runaway tool loops, so tool_calls_limit is a useful guardrail. It also supports a per-run USD cost_limit; see Single prompts: LLMOperator and @task.llm for the caveats (not a hard guarantee; not enforced for models pydantic-ai can’t price, which log a warning instead of failing the run) and an example. Default None.

    Warning

    With durable=True, a task retry replays cached model steps instead of re-calling the model – but pydantic-ai still adds each replayed step’s cost to the retry’s own usage total, since it cannot distinguish a replay from a live call. A cost_limit therefore counts already-paid-for replayed cost against every retry’s fresh budget, leaving less headroom for the new calls the retry actually makes. And if the limit is lowered between attempts – easy to do by accident when usage_limits is templated as a dict – a retry can exceed it with zero new model calls. The LLM run cost line in the task log reports the run’s cumulative cost for the same reason, not what this attempt actually spent.

  • durable: When True, enables step-level caching of model responses and tool results. On retry, cached steps are replayed instead of re-executing expensive LLM calls. On Airflow >= 3.3 the cache uses the task state store (no configuration needed); on older cores it requires the [common.ai] durable_cache_path config option to be set. Default False.

  • code_mode: When True, wraps the agent’s tools in a single run_code tool that the model drives by writing Python, executed in the Monty sandbox. Requires the code-mode extra. Default False. See Code mode.

  • message_history: Prior conversation to seed a multi-turn session, as a list of pydantic-ai ModelMessage objects or their JSON form (str / bytes). When set, the post-run transcript is pushed to XCom under the key message_history for the next run to resume. Default None (single-turn). See Multi-turn sessions and message history.

  • serialize_output: 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.

HITL review parameters: enable_hitl_review, max_hitl_iterations, hitl_timeout and hitl_poll_interval turn on and bound the iterative review loop, which needs the hitl_review plugin. Human-in-the-loop (HITL) review for agents documents each parameter and the review workflow.

Logging

All AI operators automatically log a post-run summary after run_sync() completes. AgentOperator additionally wraps toolsets for real-time per-tool-call logging (controlled by enable_tool_logging).

Real-time tool call logging (AgentOperator only): each tool call is logged as it happens:

INFO - Tool call: list_tables
INFO - Tool list_tables returned in 0.12s
INFO - Tool call: get_schema
INFO - Tool get_schema returned in 0.08s
INFO - Tool call: query
INFO - Tool query returned in 0.34s

Tool arguments are logged at DEBUG level to avoid leaking sensitive data at the default log level.

Post-run summary (all operators): after the LLM run finishes, a summary is logged with model name, token usage, and the full tool call sequence:

INFO - LLM run complete: model=gpt-5, requests=4, tool_calls=3, input_tokens=2847, output_tokens=512, total_tokens=3359
INFO - Tool call sequence: list_tables -> get_schema -> query

At DEBUG level, the LLM output is also logged (truncated to 500 characters).

Both layers use Airflow’s ::group:: / ::endgroup:: log markers, which render as collapsible sections in the Airflow UI task log viewer.

To disable real-time tool logging while keeping the post-run summary:

AgentOperator(
    task_id="my_agent",
    prompt="...",
    llm_conn_id="my_llm",
    toolsets=[SQLToolset(db_conn_id="my_db")],
    enable_tool_logging=False,
)

Security

See also

Securing agent tools for defense layers, allowed_tables limitations, HookToolset guidelines, recommended configurations, and the production checklist.

Was this entry helpful?