This is an automated email from the ASF dual-hosted git repository. kaxil pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/airflow.git
commit 2eb11c14b853c0f73bc5f24d34a78d28d123a951 Author: Kaxil Naik <[email protected]> AuthorDate: Wed Sep 30 07:04:25 2026 +0100 Document running your own Pydantic AI agent, and sandboxes in other frameworks (#73902) Add a Pydantic AI guide for teams with an agent of their own: a model from PydanticAIHook and the provider's toolsets passed straight to it, and what AgentOperator adds on top. The sandbox guide now shows a Strands agent owning its sandbox's life with a with block. --- providers/common/ai/docs/frameworks/index.rst | 33 +++++++-- .../common/ai/docs/frameworks/pydantic_ai.rst | 78 ++++++++++++++++++++++ providers/common/ai/docs/sandbox/index.rst | 34 ++++++++++ .../ai/example_dags/example_pydantic_ai_agent.py | 72 ++++++++++++++++++++ 4 files changed, 210 insertions(+), 7 deletions(-) diff --git a/providers/common/ai/docs/frameworks/index.rst b/providers/common/ai/docs/frameworks/index.rst index 9ed5632091b..e99537b90a7 100644 --- a/providers/common/ai/docs/frameworks/index.rst +++ b/providers/common/ai/docs/frameworks/index.rst @@ -41,17 +41,22 @@ is Airflow's and which part stays yours. - What this provider gives you - Who runs the agent - Install - * - Pydantic AI + * - Pydantic AI, through ``AgentOperator`` - :class:`~airflow.providers.common.ai.operators.agent.AgentOperator`, ``@task.agent``, the ``@task.llm`` family and every toolset, plus durable replay, human review, approval gates, retry policies and tracing. - The operator - Included + * - Pydantic AI, your own agent + - A model from a connection through ``PydanticAIHook``, and every toolset, since they + are Pydantic AI toolsets. Guide: :doc:`pydantic_ai`. + - You + - Included * - Strands Agents - The :class:`~airflow.providers.common.ai.tools.strands.AirflowTools` plugin gives - a Strands agent the tools of ``SQLToolset``, ``HookToolset`` and the other - toolsets, with Airflow's secret masker applied to every result. Guide: - :doc:`strands`. + a Strands agent the tools of ``SQLToolset``, ``HookToolset``, + ``ObjectStorageToolset`` and the other toolsets, with Airflow's secret masker + applied to every result. Guide: :doc:`strands`. - You - ``strands-agents`` * - Google ADK @@ -77,7 +82,7 @@ is Airflow's and which part stays yours. The Strands and ADK integrations, the framework-neutral tool interface under them, and the tracing helper are experimental: they can change or be removed in a minor release of -this provider. +this provider. See :ref:`howto/stability`. Tested versions --------------- @@ -107,6 +112,8 @@ Choosing a route Start from the agent you already have: +- **An existing Pydantic AI agent**: keep it, give it a model with ``PydanticAIHook`` and + pass Airflow's toolsets to it. See :doc:`pydantic_ai`. - **An existing Strands agent**: keep it, and add Airflow's toolsets with the ``AirflowTools`` plugin. See :doc:`strands`. - **An existing ADK agent**: keep it, and add Airflow's toolsets with the ADK @@ -141,15 +148,27 @@ what the Strands and ADK integrations are written against: - :class:`~airflow.providers.common.ai.tools.ToolProvider` is anything with an ``airflow_tools()`` method returning ``AirflowTool`` objects. Every toolset this provider ships that reads from a connection implements it, except - ``MCPToolset``, which works in Pydantic AI agents and through the LangChain bridge. + ``MCPToolset`` and the Agent Skills toolset. ``MCPToolset`` works in Pydantic AI agents + and through the LangChain bridge; for Strands, ADK and other frameworks, connect the + framework's own MCP client to the server. An adapter maps these onto the framework's own tool type and error status, and makes sure a ``ToolCallError`` ends the run rather than reaching the model, as the Strands -plugin does. +plugin does. :func:`~airflow.providers.common.ai.tools.collect_tools` turns the toolsets and +tools an adapter is given into one list, as the Strands and ADK adapters do. + +Blocking hook calls, such as a SQL query, run one at a time in the task's process, so an +agent that calls two database tools at once gets its answers one after the other. + +Run a Strands or ADK agent inside +:func:`~airflow.providers.common.ai.tools.tracing.agent_framework_tracing` to keep its +spans tied to the task, and free of prompt text unless ``[common.ai] capture_content`` is +on; see :doc:`../observability`. .. toctree:: :hidden: :titlesonly: + Pydantic AI <pydantic_ai> Strands Agents <strands> Google ADK <adk> diff --git a/providers/common/ai/docs/frameworks/pydantic_ai.rst b/providers/common/ai/docs/frameworks/pydantic_ai.rst new file mode 100644 index 00000000000..f88e3d30887 --- /dev/null +++ b/providers/common/ai/docs/frameworks/pydantic_ai.rst @@ -0,0 +1,78 @@ + .. Licensed to the Apache Software Foundation (ASF) under one + or more contributor license agreements. See the NOTICE file + distributed with this work for additional information + regarding copyright ownership. The ASF licenses this file + to you under the Apache License, Version 2.0 (the + "License"); you may not use this file except in compliance + with the License. You may obtain a copy of the License at + + .. http://www.apache.org/licenses/LICENSE-2.0 + + .. Unless required by applicable law or agreed to in writing, + software distributed under the License is distributed on an + "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + KIND, either express or implied. See the License for the + specific language governing permissions and limitations + under the License. + +.. _howto/frameworks:pydantic_ai: + +Pydantic AI +=========== + +This provider's operators run `Pydantic AI <https://ai.pydantic.dev/>`__ agents, and +:class:`~airflow.providers.common.ai.operators.agent.AgentOperator` builds one for you +from a connection, a prompt and a list of toolsets. If you already have a Pydantic AI +agent, with its own instructions, output types, capabilities or history processing, you +do not have to rebuild it as an ``AgentOperator``. Run it in a ``@task`` and give it what +Airflow has: a model from a connection, and toolsets bound to your connections. + +Run your own agent in a task +---------------------------- + +.. exampleinclude:: /../../ai/src/airflow/providers/common/ai/example_dags/example_pydantic_ai_agent.py + :language: python + :start-after: [START example_pydantic_ai_agent] + :end-before: [END example_pydantic_ai_agent] + +``PydanticAIHook.get_hook`` returns the hook for the connection's type, and its +:meth:`~airflow.providers.common.ai.hooks.pydantic_ai.PydanticAIHook.get_conn` returns +the Pydantic AI model the connection describes, with the same vendor prefixes, fallback +connections and self-hosted endpoints as ``AgentOperator``; see :doc:`../model_providers`. +Every toolset this provider ships is a Pydantic AI toolset, so it goes into +``toolsets=`` as it is, next to any toolsets and tools of your own. + +For the OpenTelemetry spans ``AgentOperator`` emits, build the agent with the hook's +:meth:`~airflow.providers.common.ai.hooks.pydantic_ai.PydanticAIHook.create_agent` +instead of ``Agent(...)``. It takes the same arguments, and under ``[common.ai] +otel_export_enabled`` it sends the agent's spans through Airflow's tracing, without prompt +text unless ``capture_content`` is on; see :doc:`../observability`. Pydantic AI's own +instrumentation records prompt and completion text by default. + +What you get, and what you do not +--------------------------------- + +The SQL, hook, object storage, DataFusion, MCP, sandbox and managed-agent toolsets behave +as they do inside ``AgentOperator``: SQL validation, ``allowed_tables``, result bounds, +object-storage path checks, and the secret masker on everything a tool returns, including +the text of a failure the model is asked to correct. Their calls are counted in the +``common_ai.tool_calls`` metric described in :doc:`../observability`. The Agent Skills +toolset is the exception: its results are masked only inside ``AgentOperator``, and its +calls are not counted. + +``AgentOperator`` adds things on top of the agent that a task running your own agent +does not get: + +- Durable replay of model and tool steps across task retries (``durable=True``). +- Human review of the output (``enable_hitl_review``), and a pause for a person to + approve marked tool calls before they run (:doc:`../tool_approval`). +- Masking of the results of the Agent Skills toolset and of toolsets you wrote yourself. + There is no public masking wrapper, so run the agent through ``AgentOperator`` if their + results can carry a secret. +- Rendering of templated connection IDs in toolsets, such as ``SQLToolset("{{ ... }}")``. + In your own task, pass the connection ID itself. +- Tool call logging, and the task's identity (``airflow.dag_id``, ``airflow.task_id`` + and the rest) on the agent's spans. + +If you find yourself rebuilding one of these, that is a sign ``AgentOperator`` fits: +its ``agent_params`` passes any other argument through to the Pydantic AI ``Agent``. diff --git a/providers/common/ai/docs/sandbox/index.rst b/providers/common/ai/docs/sandbox/index.rst index 6d201ee09fc..e21f8112897 100644 --- a/providers/common/ai/docs/sandbox/index.rst +++ b/providers/common/ai/docs/sandbox/index.rst @@ -438,6 +438,40 @@ worker for Modal. Work runs in a per-run microVM on the worker host with act as barriers, as they do for the other routes that build their own tools; see :ref:`toolset-call-barriers`. +.. _sandbox-other-frameworks: + +With another agent framework +---------------------------- + +A Strands or Google ADK agent can use the same sandbox through ``AirflowTools`` +(see :doc:`../frameworks/index`). Outside a Pydantic AI run nothing ends the run for the +toolset, so the task owns the sandbox's life: open the toolset with ``with`` (or +``async with``) around the agent, and the sandbox it provisions is destroyed when the +block ends, however the agent finishes: + +.. code-block:: python + + from strands import Agent + + from airflow.providers.common.ai.sandbox.modal import ModalSandboxBackend + from airflow.providers.common.ai.tools.strands import AirflowTools + from airflow.providers.common.ai.toolsets import SandboxToolset, SQLToolset + from airflow.sdk import task + + warehouse = SQLToolset("warehouse", allowed_tables=["ledger"]) + + + @task + def reconcile() -> str: + with SandboxToolset(ModalSandboxBackend()) as sandbox: + agent = Agent(plugins=[AirflowTools(warehouse, sandbox)]) + return str(agent("Reconcile the September ledger against the warehouse.")) + +A tool call made before the block or after it is refused rather than provisioning a +sandbox nothing would destroy. The toolset's own error rules hold as well: a command +that fails is output the model reads, and a sandbox that cannot be provisioned ends the +agent run so the task fails. + .. _sandbox-limitations: Limitations diff --git a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_pydantic_ai_agent.py b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_pydantic_ai_agent.py new file mode 100644 index 00000000000..bcef2356980 --- /dev/null +++ b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_pydantic_ai_agent.py @@ -0,0 +1,72 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +""" +Run a Pydantic AI agent you build yourself, with Airflow's toolsets as its tools. + +``AgentOperator`` builds the agent for you. When you already have a Pydantic AI agent, +keep it and run it in an ordinary ``@task``: ``PydanticAIHook`` gives it a model from an +Airflow connection, and Airflow's toolsets are Pydantic AI toolsets, so they go straight +into ``toolsets=``. +""" + +from __future__ import annotations + +import os + +from pydantic_ai import Agent + +from airflow.providers.common.ai.hooks.pydantic_ai import PydanticAIHook +from airflow.providers.common.ai.toolsets import ObjectStorageToolset +from airflow.providers.common.compat.sdk import dag, task + +LLM_CONN_ID = os.environ.get("LLM_CONN_ID", "pydanticai_default") +DB_CONN_ID = os.environ.get("DB_CONN_ID", "sql_default") +FILES = os.environ.get("FILES", "s3://acme-reports/finance/") +FILES_CONN_ID = os.environ.get("FILES_CONN_ID", "aws_default") + +DEFAULT_QUESTION = "Does the September report's revenue match the orders table?" + + +# [START example_pydantic_ai_agent] +@dag(tags=["example"]) +def example_pydantic_ai_agent(): + """Answer a question across a database and a reports bucket with your own agent.""" + + @task + def run_pydantic_ai_agent(question: str = DEFAULT_QUESTION) -> str: + from airflow.providers.common.ai.toolsets.sql import SQLToolset + + agent = Agent( + PydanticAIHook.get_hook(LLM_CONN_ID).get_conn(), + instructions=( + "You check reports against the warehouse. Read the report files, query the " + "orders table, and say whether the numbers agree." + ), + toolsets=[ + SQLToolset(db_conn_id=DB_CONN_ID, allowed_tables=["orders"]), + ObjectStorageToolset(FILES, conn_id=FILES_CONN_ID), + ], + ) + return agent.run_sync(question).output + + run_pydantic_ai_agent() + + +# [END example_pydantic_ai_agent] + + +example_pydantic_ai_agent()
