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.
See also
SQL agent¶
The most common pattern: give an agent access to a database so it can answer questions by writing and executing SQL.
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 |
|---|---|
|
Lists available table names (filtered by |
|
Returns column names and types for a table |
|
Executes a SQL query and returns rows as JSON |
|
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.
@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.
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.
# 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
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.
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.
@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¶
Five features have pages of their own:
Multi-turn sessions and message history: pass
message_historyto carry a conversation across runs.Durable execution: set
durable=Trueto replay completed model and tool steps on retry instead of paying for them again.Capabilities and guardrails: pass pydantic-ai capabilities and
pydantic-ai-shieldsguardrails withcapabilities=.Code mode: pass the
CodeModecapability to collapse the agent’s tools into a singlerun_codetool the model drives by writing Python.Approve an agent’s tool calls: mark tools that need a person’s approval, and the task pauses before a marked call runs.
Durable execution¶
Moved to Durable execution.
Prompt caching¶
Every request an agent makes re-sends its tool definitions, its system prompt and the
conversation so far. An agent that calls three tools makes four requests, and a mapped
@task.agent makes that many per map index, all starting with the same system prompt.
cache_prompt (on by default) asks the provider to keep that repeated prefix, so later
requests read it back instead of paying the full input price for it again.
What it turns on depends on the model the connection resolves to:
Provider |
What |
|---|---|
Anthropic ( |
Marks the end of the tool definitions, the end of the system prompt and the latest message as points to cache up to. |
Bedrock ( |
The same three marks, for models pydantic-ai knows support caching. Nothing for the rest. |
OpenAI, Azure OpenAI, Gemini |
Nothing. These cache long prompts on their own. |
Because each model reads only its own provider’s settings, the same flag covers a fallback chain that spans providers.
On Anthropic, a 5-minute cache write costs 1.25x the normal input price and a read costs 0.1x (less on some newer models), so a prefix read back even once costs less than sending it twice. A prompt shorter than the model’s minimum length for caching (512 to 4,096 tokens, depending on the model) is not cached and costs nothing extra. See Anthropic’s prompt caching guide for the per-model minimums and prices.
Caching costs more than it saves in three cases:
A single long request. An agent with no tools that sends a long prompt once and is not run again within five minutes pays the write and never reads it back.
A large final tool result. The latest message is written to the cache on every request, including the last one, which nothing reads. When the last tool returns much more than the system prompt and tool definitions add up to, as a query returning thousands of rows can, that final write costs more than the earlier reads saved.
Map indexes that start together. A cache entry exists only once the response that wrote it has started, so map indexes that all start at the same moment each write their own copy. Indexes that start later read it back.
Turn it off for such a task:
AgentOperator(
task_id="summarize_quarter",
prompt="Summarize the attached report.",
llm_conn_id="anthropic_default",
system_prompt=long_style_guide,
cache_prompt=False,
)
To choose what is cached or for how long, set the provider’s own settings in
agent_params["model_settings"]. Setting any anthropic_cache* key hands Anthropic
caching back to you, and cache_prompt adds nothing for Anthropic; the same holds for
bedrock_cache* and openrouter_cache*. A CachePoint in the prompt or the message
history hands caching back to you for every provider. For example, a mapped task whose
instances run further apart than five minutes can keep the system prompt for an hour, at 2x
the input price for each write instead of 1.25x:
AgentOperator.partial(
task_id="classify_ticket",
llm_conn_id="anthropic_default",
system_prompt=long_taxonomy,
agent_params={
"model_settings": {
"anthropic_cache_instructions": "1h",
"anthropic_cache_tool_definitions": "1h",
}
},
).expand(prompt=tickets)
When the provider reports cache activity, the task log shows it under the run summary:
LLM run complete: model=claude-sonnet-4-5, requests=2, tool_calls=1, input_tokens=..., ...
LLM prompt cache: cache_read_tokens=..., cache_write_tokens=...
input_tokens includes both counts. With Observability (OpenTelemetry tracing) turned on, each
request’s GenAI span carries them too. The usage XCom carries them as
cache_read_tokens and cache_write_tokens.
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 PydanticBaseModelfor structured output.toolsets: List of pydantic-ai toolsets (SQLToolset,HookToolset,AgentSkillsToolsetfor Agent Skills: AgentSkillsToolset, etc.).capabilities: List of pydantic-ai capabilities (Thinking,WebSearch, guardrails, etc.). See Capabilities and guardrails.enable_tool_logging: Wrap each toolset inLoggingToolsetso that every tool call is logged in real time. DefaultTrue.agent_params: Additional keyword arguments passed to the pydantic-aiAgentconstructor (e.g.retries,model_settings).usage_limits: Optional pydantic-aiUsageLimitsenforced on every agent run (initial run, durable replay, and HITL regeneration), or adictof the same fields – the dict form is templated via Jinja, then coerced per field type, failing the task with aValueErrornaming 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, sotool_calls_limitis a useful guardrail. It also supports a USDcost_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. DefaultNone.On Airflow >= 3.3, setting this counts usage across every attempt of the task instance combined – the initial run, every retry, and every HITL regeneration all add to one running total kept in the AIP-103 task state store under the
__commonai_usage__key – instead of each attempt starting a fresh count. This also applies to the implicitrequest_limit=50default, which can now block a retry that used to pass on its own. A step replayed bydurable=Truedoes not count toward that total – seedurablebelow. To keep the same effective per-attempt headroom this cross-attempt total used to give each attempt on its own, scale each limit byretries + 1, or useusage_limits=Noneto opt back out.Clearing and rerunning a finished (failed or succeeded) task instance gets a fresh budget automatically; clearing a running task instance does not bump
max_tries, so the restarted attempt still sees the prior spend. To reset the budget for a task instance that keeps retrying without a clear of a finished attempt, delete the__commonai_usage__key via the Task State Store UI.Note
A worker killed with SIGKILL – including after
on_kill’s grace period expires, or an OOM kill – cannot persist that attempt’s usage, so the next attempt’s count under-represents actual spend by that amount.On Airflow < 3.3, and whenever
usage_limitsisNone, each attempt is still checked and counted on its own, as before. Within one attempt, a HITL regeneration shares the count of the run before it wheneverusage_limitsis set, on every Airflow version; withusage_limits=Noneeach regeneration starts a fresh count.durable: WhenTrue, 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 Airflow versions it requires the[common.ai] durable_cache_pathconfig option to be set. DefaultFalse. A replayed step adds nothing to the usage counted againstusage_limitsor reported in theusageXCom – not its request, tokens, cost, or tool calls – so each attempt counts only the model and tool calls it actually makes, on every Airflow version. A step that re-runs live because the conversation changed since the previous attempt is counted like any other live call, and a retry whose cross-attempt total already sits at a limit can still start when the steps it needs are cached. Clearing a failed task instance starts a fresh budget but keeps the durable cache its attempts left behind, so what the rerun replays from that cache is free there too.cache_prompt: Ask the provider to cache the tool definitions, system prompt and conversation so later requests read them back at a discount. DefaultTrue; a no-op for providers that cache on their own. See Prompt caching.message_history: Prior conversation to seed a multi-turn session, as a list of pydantic-aiModelMessageobjects or their JSON form (str/bytes). When set, the post-run transcript is pushed to XCom under the keymessage_historyfor the next run to resume. DefaultNone(single-turn). See Multi-turn sessions and message history.serialize_output: IfTrueandoutput_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.
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.