Strands Agents¶
You have an agent built with Strands Agents and
you want it to run as an Airflow task: after an upstream task, on a schedule, or
when an asset updates, reading data through connections your deployment already
manages. Keep the agent as it is. Build the model and the Agent yourself, and
give it Airflow’s toolsets with the
AirflowTools plugin.
Note
AirflowTools and the framework-neutral tool interface it builds on are
experimental. They may change in a minor release of this provider.
See Stable and experimental features.
Run a Strands agent in a task¶
The agent below answers a question about a database. Its model’s API key comes
from an Airflow connection, and its tools come from a
SQLToolset, so the agent can
list tables, read schemas and run read-only queries, with results bounded to
what fits in the model’s context:
@dag(tags=["example"])
def example_strands_agent():
"""Answer a question about a database with a Strands agent."""
@task
def run_strands_agent(question: str = DEFAULT_QUESTION) -> str:
from strands import Agent
from strands.models.anthropic import AnthropicModel
from airflow.providers.common.ai.tools.strands import AirflowTools
from airflow.providers.common.ai.tools.tracing import agent_framework_tracing
from airflow.providers.common.ai.toolsets.sql import SQLToolset
llm = BaseHook.get_connection(LLM_CONN_ID)
model = AnthropicModel(
client_args={"api_key": llm.password, "base_url": llm.host or None},
model_id=LLM_MODEL,
max_tokens=2048,
)
# Spans carry the task's identity and no prompt text; see the tracing section of the guide.
with agent_framework_tracing():
agent = Agent(
model=model,
plugins=[AirflowTools(SQLToolset(db_conn_id=DB_CONN_ID))],
system_prompt=(
"You are a SQL analyst. Use list_tables and get_schema to explore "
"the database, then run read-only queries to answer the question."
),
# Strands streams the reply to stdout by default; the task returns it instead.
callback_handler=None,
)
return str(agent(question))
run_strands_agent()
Everything except AirflowTools is plain Strands. The agent’s loop, model
features, hooks and session handling work as the Strands documentation describes,
and tools you wrote with Strands’ own @tool decorator go in tools= as
usual, next to the plugin.
The model can come from any connection that holds its credentials. For Strands’
default Bedrock model, build BedrockModel(boto_session=...) from
AwsBaseHook(aws_conn_id=...).get_session() in the Amazon provider, so the
model uses the same AWS connection as the rest of your Dags.
What the plugin gives the agent¶
AirflowTools accepts any toolset that implements
ToolProvider, such as
SQLToolset and
HookToolset, and adds one
Strands tool per toolset tool. MCPToolset is the exception: connect Strands’ own
MCPClient to the server instead. Each tool:
Keeps the toolset tool’s name, description and argument schema.
Runs through the toolset itself, so connection resolution, SQL validation,
allowed_tables,allowed_methodsand result bounds behave as they do inAgentOperator.Passes every result through Airflow’s secret masker before it goes back to the model.
Writes a
Tool call: <name>group to the task log with the call’s duration.
Failures follow the same rule as in AgentOperator:
A failure the model can correct, such as a rejected SQL statement or an invalid argument, reaches it as a Strands result with
status="error", and the model can try again. The tool’s retry limit bounds how many times in a row it can do so.SQLToolsettreats every database error this way, so a query that fails because the database rejects the connection’s credentials ends the run once that limit is reached.A refusal the tool reports as final, by raising pydantic-ai’s
ToolFailed, also reaches the model as an error, and does not count against the retry limit. Bound the task withexecution_timeoutif a model could keep asking for it.Any other failure, such as a hook method raising, ends the agent run. Strands raises
EventLoopExceptionwith the maskedToolCallErroras itsoriginal_exception, the task fails, and Airflow’s own retry takes over. Strands on its own would hand such an error to the model and let the run carry on.
Outside AgentOperator, a toolset’s connection ID is used as written: it is not
rendered as a template, so pass the value itself rather than a Jinja expression.
The plugin is built on the framework-neutral tool interface described in Agent frameworks, which is also the starting point for an adapter to another framework.
Why masking tool results matters¶
Airflow masks connection passwords in task logs. Tool results are not logs:
they go to the model, to the model provider’s API, into any traces the framework
records, and often into the agent’s final answer, which your task may return as
an XCom. A tool whose client library puts a credentialed URL in its error message
would send the raw password to all of those places, even though the task log
shows ***.
Tools added by AirflowTools apply the masker to what they return, so a
registered secret is replaced with *** before it leaves the tool. The masker
knows the secrets Airflow has registered, such as connection passwords and
sensitive connection extras. It does not recognize a credential that exists only
in the data a tool returns, so the permissions of the connection’s database role
or cloud identity remain the real boundary on what the agent can reach.
Tools you write yourself with Strands’ @tool decorator do not pass through
the masker. To give one of your own functions the same treatment, wrap it as an
AirflowTool and pass it to
AirflowTools alongside the toolsets:
import asyncio
from airflow.providers.common.ai.tools import AirflowTool, ToolResult
from airflow.providers.common.ai.tools.strands import AirflowTools
async def lookup_customer(arguments: dict) -> ToolResult:
# crm_client.get blocks, so run it off the agent's event loop.
customer = await asyncio.to_thread(crm_client.get, arguments["customer_id"])
return ToolResult(content=customer)
lookup = AirflowTool(
name="lookup_customer",
description="Look up one customer in the CRM by id.",
parameters={
"type": "object",
"properties": {"customer_id": {"type": "integer"}},
"required": ["customer_id"],
},
function=lookup_customer,
)
agent = Agent(model=model, plugins=[AirflowTools(warehouse, lookup)])
Tracing¶
The example runs the agent inside
agent_framework_tracing(), so Strands’
OpenTelemetry spans carry the task’s identity and leave out prompts, completions and tool
inputs and outputs unless [common.ai] capture_content is on. Create the Agent
inside the block: Strands reads the switch once per process, when it creates its one
tracer, so an Agent created earlier in the process, such as at module level, keeps
content capture on. See
Observability (OpenTelemetry tracing).
Differences from AgentOperator¶
The agent is Strands’, so features that AgentOperator implements on top of
Pydantic AI do not apply:
A task retry runs the agent again from the start. There is no step-level replay as with
AgentOperator(durable=True), so tools that change data must be safe to repeat.There is no human review step like
AgentOperator(enable_hitl_review=True).
Installation¶
pip install apache-airflow-providers-common-ai "strands-agents>=1.56.0"
Pin the strands-agents lower bound as above; 1.56 is the release this
integration is tested with. Without a bound, a resolver that cannot satisfy Strands’
dependencies may pick the 0.0.1 placeholder release instead of reporting a
conflict. For AnthropicModel, install anthropic through this provider’s
anthropic extra rather than through Strands’ own anthropic extra, which
requires an anthropic release older than 1.0. Strands also pins mcp below 2.2, so
installing it downgrades a newer mcp, which MCPToolset uses.