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
The following commit(s) were added to refs/heads/main by this push:
new 35633455ed7 Add first-class `capabilities` to Common AI
`AgentOperator` (#73984)
35633455ed7 is described below
commit 35633455ed746655242ad1adcc59fbd6d097c35f
Author: Kaxil Naik <[email protected]>
AuthorDate: Thu Oct 1 10:39:18 2026 +0100
Add first-class `capabilities` to Common AI `AgentOperator` (#73984)
AgentOperator and @task.agent take a capabilities= list of pydantic-ai
capabilities, kept out of the serialized Dag and templated like toolsets=.
capabilities inside agent_params still works; passing both fails the task.
A CodeMode capability is now refused with durable=True and turns off
per-tool
approval, matching code_mode=True. The guardrails docs page becomes a
capabilities page.
---
providers/common/ai/docs/capabilities.rst | 137 ++++++++++++++++
providers/common/ai/docs/code_mode.rst | 7 +-
providers/common/ai/docs/features.rst | 6 +-
providers/common/ai/docs/guardrails.rst | 66 --------
providers/common/ai/docs/operators/agent.rst | 10 +-
providers/common/ai/docs/redirects.txt | 1 +
providers/common/ai/docs/stability.rst | 2 +-
providers/common/ai/docs/tool_approval.rst | 5 +-
.../ai/example_dags/example_agent_capabilities.py | 27 ++--
.../src/airflow/providers/common/ai/exceptions.py | 2 +-
.../airflow/providers/common/ai/operators/agent.py | 139 ++++++++++++----
.../tests/unit/common/ai/decorators/test_agent.py | 19 +++
.../tests/unit/common/ai/operators/test_agent.py | 177 +++++++++++++++++++--
13 files changed, 457 insertions(+), 141 deletions(-)
diff --git a/providers/common/ai/docs/capabilities.rst
b/providers/common/ai/docs/capabilities.rst
new file mode 100644
index 00000000000..eeca7d716fd
--- /dev/null
+++ b/providers/common/ai/docs/capabilities.rst
@@ -0,0 +1,137 @@
+ .. 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.
+
+.. _capabilities:
+
+Capabilities and guardrails
+===========================
+
+A pydantic-ai `capability <https://ai.pydantic.dev/capabilities/>`__ adds a
behavior to an agent
+in one declaration: tools, instructions, model settings, and hooks that run
around each model
+request or tool call. ``Thinking`` turns on the model's reasoning at a chosen
effort level,
+``WebSearch`` and ``WebFetch`` give the model the web through its provider's
native tool, and
+guardrail packages such as ``pydantic-ai-shields`` check inputs and outputs.
For the full
+catalog, see the pydantic-ai documentation and the
+`pydantic-ai-harness capability matrix
<https://github.com/pydantic/pydantic-ai-harness#capability-matrix>`__.
+
+Pass capabilities to ``AgentOperator`` or ``@task.agent`` with
``capabilities=``:
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_agent_capabilities.py
+ :language: python
+ :start-after: [START howto_operator_agent_capabilities_thinking]
+ :end-before: [END howto_operator_agent_capabilities_thinking]
+
+Capabilities and toolsets work together: the agent gets the tools from both.
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_agent_capabilities.py
+ :language: python
+ :start-after: [START howto_operator_agent_capabilities_composed]
+ :end-before: [END howto_operator_agent_capabilities_composed]
+
+pydantic-ai wraps hooks in list order, the first capability outermost, so a
guard listed first
+sees a request before the capabilities after it. A capability can declare its
own position
+(for example, always outermost), which takes precedence over the list.
+
+Guardrails
+----------
+
+A guardrail is a capability that checks what goes into or comes out of the
agent and stops the
+run when a check fails. This example uses ``InputGuard`` from
``pydantic-ai-shields`` to reject a
+prompt before the agent run starts.
+
+.. note::
+
+ Experimental: the ``shields`` extra can change or be removed in a minor
release of this
+ provider.
+ See :ref:`howto/stability`.
+
+.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_agent_capabilities.py
+ :language: python
+ :start-after: [START howto_operator_agent_capabilities_input_guard]
+ :end-before: [END howto_operator_agent_capabilities_input_guard]
+
+Toolsets as capabilities
+------------------------
+
+pydantic-ai's ``Toolset`` capability holds a toolset, so any toolset from this
provider can be
+passed that way. The connection IDs of ``SQLToolset``, ``MCPToolset`` and
``HookToolset`` are
+templated inside a ``Toolset`` capability the same way as in ``toolsets=``:
+
+.. code-block:: python
+
+ from pydantic_ai.capabilities import Toolset
+
+ AgentOperator(
+ task_id="analyst",
+ prompt="How many orders shipped yesterday?",
+ llm_conn_id="pydanticai_default",
+ capabilities=[Toolset(SQLToolset(db_conn_id="warehouse_{{
var.value.environment }}"))],
+ )
+
+A ``Toolset`` capability built from a function is resolved when the run
starts, so its
+connection IDs are not templated.
+
+Tool results from a ``Toolset`` capability are masked like those from
``toolsets=``, but
+``enable_tool_logging`` only logs calls to ``toolsets=``. Pass a toolset in
``toolsets=`` unless
+you need it inside the capability list, for example to order it against a
guardrail.
+
+With durable execution
+----------------------
+
+With ``durable=True``, a retry replays completed steps from the cache instead
of running them
+again. Whether a capability's work is replayed depends on where it runs:
+
+.. list-table::
+ :header-rows: 1
+
+ * - Capability
+ - On retry
+ * - ``Thinking``, and ``WebSearch``, ``WebFetch`` or ``ImageGeneration``
when the model's
+ provider runs the tool natively
+ - Replayed with the cached model response.
+ * - ``WebSearch``, ``WebFetch`` or ``ImageGeneration`` falling back to a
local tool, for a
+ provider without the native one
+ - The local tool runs again.
+ * - ``Toolset`` holding a toolset
+ - Tool results are replayed.
+ * - ``MCP``, ``PrefixTools``, ``CombinedCapability``, a ``Toolset`` built
from a function,
+ and capabilities from an agent spec file
+ - Tools run again. Pass tools you need replayed in ``toolsets=`` instead.
+ * - pydantic-ai-harness ``CodeMode``
+ - Not allowed: the operator raises ``ValueError``, as it does for
``code_mode=True``. This
+ includes a ``CodeMode`` inside a ``CombinedCapability`` or a wrapper
such as
+ ``PrefixTools``, but not one a capability function builds when the run
starts.
+
+See :doc:`durable_execution` for how the cache works.
+
+Serialization
+-------------
+
+Capabilities passed with ``capabilities=`` are not stored in the serialized
Dag. The worker
+builds them from the Dag file when the task runs, so a capability can hold
functions and
+clients that do not serialize.
+
+``capabilities`` inside ``agent_params`` is still accepted and reaches the
agent the same way.
+``agent_params`` is a template field, though, so Airflow stores each
capability's repr in the
+serialized Dag. For a capability holding a function, such as the
``InputGuard`` above, that
+repr includes the function's memory address, which can differ from one parse
to the next and
+change the serialized Dag with it. Prefer ``capabilities=``. Passing both
fails the task with a
+``ValueError``, since there would be no single order to run the hooks in.
+
+A mapped task (``AgentOperator.partial(...).expand(...)`` or a mapped
``@task.agent``) is the
+exception: Airflow stores every argument given to ``partial``, so there
``capabilities`` is
+serialized as a repr too, as ``toolsets`` is.
diff --git a/providers/common/ai/docs/code_mode.rst
b/providers/common/ai/docs/code_mode.rst
index c32fa15700c..90069cfd225 100644
--- a/providers/common/ai/docs/code_mode.rst
+++ b/providers/common/ai/docs/code_mode.rst
@@ -81,10 +81,9 @@ Requires the ``code-mode`` extra::
:start-after: [START howto_operator_agent_code_mode]
:end-before: [END howto_operator_agent_code_mode]
-Unlike passing a capability through ``agent_params`` (see
-:ref:`capabilities-passthrough`), ``code_mode`` is a plain boolean and is
-serialization-safe: the ``CodeMode`` capability is built at execution time, not
-stored on the serialized operator.
+``code_mode=True`` is the same as adding ``CodeMode()`` at the end of
``capabilities=`` (see
+:ref:`capabilities`), except that the operator builds the capability when the
task runs, so the
+Dag file does not import ``pydantic_ai_harness``.
.. note::
diff --git a/providers/common/ai/docs/features.rst
b/providers/common/ai/docs/features.rst
index f8c7a7cfbf0..3664b80fb84 100644
--- a/providers/common/ai/docs/features.rst
+++ b/providers/common/ai/docs/features.rst
@@ -26,8 +26,8 @@ you use. Each is a parameter on the operator or decorator.
- :doc:`structured_output`: ``output_type`` returns a typed Pydantic object
through XCom
instead of a string.
- :doc:`message_history`: ``message_history`` carries a conversation across
agent runs.
-- :doc:`guardrails`: pydantic-ai capabilities and ``pydantic-ai-shields``
guards pass through
- ``agent_params``.
+- :doc:`capabilities`: ``capabilities=`` adds pydantic-ai capabilities such as
``Thinking`` and
+ ``WebSearch``, and ``pydantic-ai-shields`` guardrails, to an agent.
- :doc:`code_mode`: ``code_mode=True`` lets the model call several tools from
one Python
snippet instead of one round trip per call.
- :doc:`approval_gates`: ``require_approval=True`` pauses an LLM operator
until a person
@@ -50,7 +50,7 @@ Making retries cheap with ``durable=True`` is a reliability
feature and lives un
Structured output <structured_output>
Message history <message_history>
- Guardrails <guardrails>
+ Capabilities and guardrails <capabilities>
Code mode <code_mode>
Approve outputs <approval_gates>
Review agent sessions <hitl_review>
diff --git a/providers/common/ai/docs/guardrails.rst
b/providers/common/ai/docs/guardrails.rst
deleted file mode 100644
index dff3cf6093e..00000000000
--- a/providers/common/ai/docs/guardrails.rst
+++ /dev/null
@@ -1,66 +0,0 @@
- .. 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.
-
-.. _capabilities-passthrough:
-
-Guardrails and capabilities
-===========================
-
-pydantic-ai `capabilities <https://ai.pydantic.dev/capabilities/>`__ bundle
-tools, lifecycle hooks, instructions, and model settings into composable units.
-Common ones include ``Thinking`` (reasoning at a configurable effort level),
-``WebSearch``, ``WebFetch``, ``ImageGeneration``, and ``MCP``.
-For the current capability catalog and package-specific installation notes, see
-the pydantic-ai documentation and the
-`pydantic-ai-harness capability matrix
<https://github.com/pydantic/pydantic-ai-harness#capability-matrix>`__.
-
-``AgentOperator`` does not yet expose a first-class ``capabilities=`` kwarg,
-but anything passed through ``agent_params`` is forwarded to the underlying
-``Agent(...)`` constructor.
-
-.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_agent_capabilities.py
- :language: python
- :start-after: [START howto_operator_agent_capabilities_thinking]
- :end-before: [END howto_operator_agent_capabilities_thinking]
-
-Capabilities compose with toolsets -- pydantic-ai merges tools from both.
-
-.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_agent_capabilities.py
- :language: python
- :start-after: [START howto_operator_agent_capabilities_composed]
- :end-before: [END howto_operator_agent_capabilities_composed]
-
-Guardrail capabilities use the same passthrough pattern. This example uses
-``InputGuard`` from ``pydantic-ai-shields`` to reject a prompt before the agent
-run starts.
-
-.. note::
-
- Experimental: the ``shields`` extra can change or be removed in a minor
release of this
- provider.
- See :ref:`howto/stability`.
-
-.. exampleinclude::
/../../ai/src/airflow/providers/common/ai/example_dags/example_agent_capabilities.py
- :language: python
- :start-after: [START howto_operator_agent_capabilities_input_guard]
- :end-before: [END howto_operator_agent_capabilities_input_guard]
-
-.. warning::
-
- ``agent_params`` is a templated field, which Airflow serializes by calling
- ``str()`` on values it doesn't natively understand. Capability instances
- are not yet round-trip-safe through Dag serialization, so the examples
above construct them inside the ``@dag`` function -- not at module level.
diff --git a/providers/common/ai/docs/operators/agent.rst
b/providers/common/ai/docs/operators/agent.rst
index e8e64b5d180..2d6473e1966 100644
--- a/providers/common/ai/docs/operators/agent.rst
+++ b/providers/common/ai/docs/operators/agent.rst
@@ -249,8 +249,8 @@ Five features have pages of their own:
- :doc:`../message_history`: pass ``message_history`` to carry a conversation
across runs.
- :doc:`../durable_execution`: set ``durable=True`` to replay completed model
and tool steps
on retry instead of paying for them again.
-- :doc:`../guardrails`: pass pydantic-ai capabilities and
``pydantic-ai-shields`` guardrails
- through ``agent_params``.
+- :doc:`../capabilities`: pass pydantic-ai capabilities and
``pydantic-ai-shields`` guardrails
+ with ``capabilities=``.
- :doc:`../code_mode`: set ``code_mode=True`` to collapse the agent's tools
into a single
``run_code`` tool the model drives by writing Python.
- :doc:`../tool_approval`: mark tools that need a person's approval, and the
task pauses before
@@ -277,13 +277,13 @@ Parameters
``BaseModel`` for structured output.
- ``toolsets``: List of pydantic-ai toolsets (``SQLToolset``, ``HookToolset``,
``AgentSkillsToolset`` for :ref:`agent-skills`, etc.).
+- ``capabilities``: List of pydantic-ai capabilities (``Thinking``,
``WebSearch``, guardrails,
+ etc.). See :ref:`capabilities`.
- ``enable_tool_logging``: Wrap each toolset in
:class:`~airflow.providers.common.ai.toolsets.logging.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 :ref:`capabilities-passthrough` for how to enable pydantic-ai
capabilities
- such as ``Thinking``, ``WebSearch``, and ``ImageGeneration``.
+ ``Agent`` constructor (e.g. ``retries``, ``model_settings``).
- .. _agent-usage-budget:
``usage_limits``: Optional pydantic-ai ``UsageLimits`` enforced on every
diff --git a/providers/common/ai/docs/redirects.txt
b/providers/common/ai/docs/redirects.txt
index 68309579391..295150c2603 100644
--- a/providers/common/ai/docs/redirects.txt
+++ b/providers/common/ai/docs/redirects.txt
@@ -21,3 +21,4 @@ choosing_a_toolset.rst toolsets/index.rst
sandbox.rst sandbox/index.rst
hooks/index.rst concepts.rst
end_to_end_pipelines.rst use_cases/index.rst
+guardrails.rst capabilities.rst
diff --git a/providers/common/ai/docs/stability.rst
b/providers/common/ai/docs/stability.rst
index e47ed969296..100ba376555 100644
--- a/providers/common/ai/docs/stability.rst
+++ b/providers/common/ai/docs/stability.rst
@@ -132,7 +132,7 @@ Everything this provider ships that is not in the table
above is experimental.
:doc:`tool_approval`)
- New, and needs Airflow 3.3. How a paused run resumes may change.
* - ``code_mode`` (:doc:`code_mode`), the Agent Skills toolset
- (:doc:`toolsets/skills`) and the ``shields`` extra (used in
:doc:`guardrails`)
+ (:doc:`toolsets/skills`) and the ``shields`` extra (used in
:doc:`capabilities`)
- Thin integrations of packages outside this provider whose APIs are still
changing: ``pydantic-ai-harness``, ``pydantic-ai-skills`` and
``pydantic-ai-shields``.
diff --git a/providers/common/ai/docs/tool_approval.rst
b/providers/common/ai/docs/tool_approval.rst
index a5ede5270e0..c568bd00a4c 100644
--- a/providers/common/ai/docs/tool_approval.rst
+++ b/providers/common/ai/docs/tool_approval.rst
@@ -127,8 +127,9 @@ Requirements and limits
- Airflow 3.3 or later. On older versions a tool marked for approval fails the
task, as it did before.
-- Not together with ``durable=True``, ``enable_hitl_review=True``,
- ``code_mode=True``, or a ``SandboxToolset`` that provisions its own sandbox.
Each
+- Not together with ``durable=True``, ``enable_hitl_review=True``, code mode
+ (``code_mode=True`` or a ``CodeMode`` capability), or a ``SandboxToolset``
that
+ provisions its own sandbox. Each
assumes the run finishes in one go; that sandbox, for one, is destroyed when
the
run pauses. With any of them, a marked tool fails the task. A
``SandboxToolset``
attached to a sandbox another task owns keeps its files through the pause,
so it
diff --git
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_agent_capabilities.py
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_agent_capabilities.py
index 6aa01d3410d..6abfb741b55 100644
---
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_agent_capabilities.py
+++
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_agent_capabilities.py
@@ -14,13 +14,12 @@
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
-"""Example DAGs demonstrating pydantic-ai capabilities via ``agent_params``.
+"""Example DAGs demonstrating pydantic-ai capabilities on ``AgentOperator``.
Capabilities (https://ai.pydantic.dev/capabilities/) are pydantic-ai's
composable units for thinking, web search, image generation, MCP, and more.
-``AgentOperator`` forwards anything in ``agent_params`` to the underlying
-``Agent(...)`` constructor, so capabilities work today without operator-level
-support. A first-class ``capabilities=`` kwarg is on the roadmap.
+``AgentOperator`` passes its ``capabilities`` list to the underlying
+``Agent(...)`` constructor.
"""
from __future__ import annotations
@@ -52,9 +51,7 @@ def example_agent_capabilities_thinking():
prompt="Walk through the steps to compute the 10th Fibonacci number,
then give the answer.",
llm_conn_id="pydanticai_default",
system_prompt="You are a careful mathematician. Think before
answering.",
- agent_params={
- "capabilities": [Thinking(effort="high")],
- },
+ capabilities=[Thinking(effort="high")],
)
@@ -76,9 +73,7 @@ def example_agent_capabilities_web_search():
prompt="Summarize the latest Apache Airflow 3.x release notes from
airflow.apache.org.",
llm_conn_id="pydanticai_default",
system_prompt="You are a release-notes summarizer. Cite the source
URL.",
- agent_params={
- "capabilities": [WebSearch()],
- },
+ capabilities=[WebSearch()],
)
@@ -113,9 +108,7 @@ if SQLToolset is not None:
max_rows=20,
),
],
- agent_params={
- "capabilities": [Thinking(effort="medium"), WebSearch()],
- },
+ capabilities=[Thinking(effort="medium"), WebSearch()],
)
# [END howto_operator_agent_capabilities_composed]
@@ -142,11 +135,9 @@ if InputGuard is not None:
),
llm_conn_id="pydanticai_default",
system_prompt="You summarize customer support requests safely.",
- agent_params={
- "capabilities": [
- InputGuard(guard=lambda prompt: "ignore previous
instructions" not in prompt.lower())
- ],
- },
+ capabilities=[
+ InputGuard(guard=lambda prompt: "ignore previous instructions"
not in prompt.lower())
+ ],
)
example_agent_capabilities_input_guard()
diff --git a/providers/common/ai/src/airflow/providers/common/ai/exceptions.py
b/providers/common/ai/src/airflow/providers/common/ai/exceptions.py
index 38ddab555a5..9b1a3e1d929 100644
--- a/providers/common/ai/src/airflow/providers/common/ai/exceptions.py
+++ b/providers/common/ai/src/airflow/providers/common/ai/exceptions.py
@@ -48,7 +48,7 @@ class UnsupportedToolDeferralError(AirflowFailException):
Either the tool hands its work to an external system, or it needs approval
where
approval is not available (before Airflow 3.3, or with ``durable``,
- ``enable_hitl_review``, ``code_mode`` or a ``SandboxToolset``). A retry
would repeat
+ ``enable_hitl_review``, code mode or a ``SandboxToolset``). A retry would
repeat
the same call, so the task fails without retrying.
"""
diff --git
a/providers/common/ai/src/airflow/providers/common/ai/operators/agent.py
b/providers/common/ai/src/airflow/providers/common/ai/operators/agent.py
index 75a077c16f8..967ab412cff 100644
--- a/providers/common/ai/src/airflow/providers/common/ai/operators/agent.py
+++ b/providers/common/ai/src/airflow/providers/common/ai/operators/agent.py
@@ -22,7 +22,8 @@ import collections
import copy
import hashlib
import json
-from collections.abc import Iterable, Sequence
+import sys
+from collections.abc import Callable, Iterable, Sequence
from dataclasses import replace
from datetime import timedelta
from functools import cached_property
@@ -30,7 +31,7 @@ from typing import TYPE_CHECKING, Any, ClassVar, Literal,
NoReturn
from pydantic import BaseModel, TypeAdapter
from pydantic_ai import DeferredToolRequests, DeferredToolResults, ToolDenied
-from pydantic_ai.capabilities import Toolset
+from pydantic_ai.capabilities import AbstractCapability, Toolset,
WrapperCapability
from pydantic_ai.messages import ModelMessagesTypeAdapter
from pydantic_ai.toolsets.abstract import AbstractToolset
from pydantic_ai.usage import RunUsage
@@ -91,6 +92,7 @@ except ImportError: # pragma: no cover - cores before the
worker-side registrat
if TYPE_CHECKING:
import jinja2
from pydantic_ai import Agent
+ from pydantic_ai.capabilities import AgentCapability
from pydantic_ai.messages import ModelMessage
from pydantic_ai.usage import UsageLimits
@@ -146,9 +148,49 @@ class HITLReviewLink(BaseOperatorLink):
)
-def _is_concrete_toolset_capability(capability: Any) -> bool:
- """Whether *capability* is a ``Toolset`` holding a toolset, not a callable
factory resolved per run."""
- return isinstance(capability, Toolset) and isinstance(capability.toolset,
AbstractToolset)
+def _resolve_capability_toolset(capability: object) -> AbstractToolset[Any] |
None:
+ """Return the toolset a ``Toolset`` capability holds; ``None`` for a
factory resolved per run or any other capability."""
+ if isinstance(capability, Toolset) and isinstance(capability.toolset,
AbstractToolset):
+ return capability.toolset
+ return None
+
+
+def _replace_capability_toolset(
+ capability: AgentCapability[Any], wrap: Callable[[AbstractToolset[Any]],
AbstractToolset[Any]]
+) -> AgentCapability[Any]:
+ """Return a ``Toolset`` capability holding ``wrap(toolset)``; any other
capability comes back as is."""
+ if isinstance(capability, Toolset) and isinstance(capability.toolset,
AbstractToolset):
+ return replace(capability, toolset=wrap(capability.toolset))
+ return capability
+
+
+def _contains_code_mode(capabilities: Iterable[AgentCapability[Any]]) -> bool:
+ """
+ Whether any capability, or one nested inside a combined or wrapper
capability, is ``CodeMode``.
+
+ A capability function, or a ``DynamicCapability``, builds its capability
when the run
+ starts, so there is nothing to inspect here.
+ """
+ # CodeMode's own module is in sys.modules once CodeMode has been imported,
and only then.
+ # The pydantic_ai_harness package root is not a safe place to look: it
exports CodeMode
+ # through a module __getattr__ that imports that module, which fails
without the
+ # ``code-mode`` extra.
+ code_mode_cls = getattr(sys.modules.get("pydantic_ai_harness.code_mode"),
"CodeMode", None)
+ if code_mode_cls is None:
+ return False
+ pending = [capability for capability in capabilities if
isinstance(capability, AbstractCapability)]
+ while pending:
+ capability = pending.pop()
+ if isinstance(capability, code_mode_cls):
+ return True
+ if isinstance(capability, WrapperCapability):
+ # apply() does not visit a wrapper's single wrapped capability,
only a combined one's children.
+ pending.append(capability.wrapped)
+ else:
+ children: list[AbstractCapability[Any]] = []
+ capability.apply(children.append)
+ pending.extend(child for child in children if child is not
capability)
+ return False
def _declares_agent_template_fields(toolset: Any) -> bool:
@@ -236,6 +278,18 @@ class AgentOperator(CancellableAgentRunMixin,
BaseOperator, HITLReviewMixin):
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.
+ :param capabilities: 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).
:param enable_tool_logging: 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.
@@ -309,7 +363,8 @@ class AgentOperator(CancellableAgentRunMixin, BaseOperator,
HITLReviewMixin):
how the model invokes them. Requires the ``code-mode`` extra
(``pip install "apache-airflow-providers-common-ai[code-mode]"``).
Cannot be combined with ``durable=True`` (durable replay assumes a
- stable per-step call order that code mode does not guarantee).
+ stable per-step call order that code mode does not guarantee), whether
+ code mode comes from this flag or from a ``CodeMode`` capability.
Default ``False``.
:param message_history: Prior conversation to seed the run with, for
multi-turn sessions that span task runs. Accepts a ``list`` of
@@ -364,9 +419,10 @@ class AgentOperator(CancellableAgentRunMixin,
BaseOperator, HITLReviewMixin):
(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``,
``code_mode``,
- or a ``SandboxToolset`` that provisions its own sandbox; there, a tool that
- requires approval fails the task as before. A ``SandboxToolset`` attached
to a
+ Not available together with ``durable``, ``enable_hitl_review``, code mode
+ (``code_mode=True`` or a ``CodeMode`` capability), or a ``SandboxToolset``
+ that provisions its own sandbox; there, a tool that requires approval fails
+ the task as before. A ``SandboxToolset`` attached to a
sandbox another task owns is fine: the sandbox outlives the pause.
:param tool_approval_timeout: Experimental. How long the pause waits for a
decision.
@@ -414,6 +470,7 @@ class AgentOperator(CancellableAgentRunMixin, BaseOperator,
HITLReviewMixin):
system_prompt: str = "",
output_type: type = str,
toolsets: list[AbstractToolset] | None = None,
+ capabilities: list[AgentCapability[Any]] | None = None,
enable_tool_logging: bool = True,
agent_params: dict[str, Any] | None = None,
usage_limits: UsageLimits | dict[str, Any] | None = None,
@@ -446,6 +503,7 @@ class AgentOperator(CancellableAgentRunMixin, BaseOperator,
HITLReviewMixin):
self.toolsets = toolsets
self.enable_tool_logging = enable_tool_logging
self.agent_params = agent_params or {}
+ self.capabilities = capabilities
# No validation here -- see coerce_usage_limits() docstring for why.
self.usage_limits = usage_limits
self.message_history = message_history
@@ -487,6 +545,14 @@ class AgentOperator(CancellableAgentRunMixin,
BaseOperator, HITLReviewMixin):
# replay. Reject the combination rather than silently
mis-replaying.
raise ValueError("durable=True and code_mode=True cannot be used
together.")
+ if (durable or code_mode) and
_contains_code_mode(self._declared_capabilities):
+ if durable:
+ # The same conflict as code_mode=True, reached through the
capability itself.
+ raise ValueError("durable=True cannot be used with a CodeMode
capability.")
+ # code_mode=True adds a second CodeMode, and pydantic-ai then
fails the run on a
+ # duplicate ``run_code`` tool without saying where the second one
came from.
+ raise ValueError("code_mode=True adds a CodeMode capability; pass
one or the other, not both.")
+
if message_history is not None and enable_hitl_review:
# The post-review transcript is not recoverable today
(run_hitl_review
# returns only the final string), so emitting the pre-review
transcript
@@ -625,19 +691,23 @@ class AgentOperator(CancellableAgentRunMixin,
BaseOperator, HITLReviewMixin):
for toolset in toolsets
]
+ def render_capabilities(capabilities: list[Any]) -> list[Any]:
+ return [
+ _replace_capability_toolset(capability, lambda toolset:
toolset.visit_and_replace(render))
+ if
_declares_agent_template_fields(_resolve_capability_toolset(capability))
+ else capability
+ for capability in capabilities
+ ]
+
if self.toolsets:
self.toolsets = render_all(self.toolsets)
+ if self.capabilities:
+ self.capabilities = render_capabilities(self.capabilities)
agent_params = dict(self.agent_params)
if agent_params.get("toolsets"):
agent_params["toolsets"] = render_all(agent_params["toolsets"])
if agent_params.get("capabilities"):
- agent_params["capabilities"] = [
- replace(capability,
toolset=capability.toolset.visit_and_replace(render))
- if _is_concrete_toolset_capability(capability)
- and _declares_agent_template_fields(capability.toolset)
- else capability
- for capability in agent_params["capabilities"]
- ]
+ agent_params["capabilities"] =
render_capabilities(agent_params["capabilities"])
self.agent_params = agent_params
@cached_property
@@ -652,6 +722,11 @@ class AgentOperator(CancellableAgentRunMixin,
BaseOperator, HITLReviewMixin):
def _build_agent(self) -> Agent[object, Any]:
"""Build and return a pydantic-ai Agent from the operator's config."""
extra_kwargs = dict(self.agent_params)
+ passed_through = extra_kwargs.pop("capabilities", None)
+ if passed_through is not None and self.capabilities is not None:
+ # pydantic-ai wraps capability hooks in list order, so merging the
two lists
+ # would pick an order the Dag author never wrote down.
+ raise ValueError("Pass capabilities either as capabilities=... or
in agent_params, not both.")
storage = self._durable_storage
counter = self._durable_counter
if self.toolsets:
@@ -665,10 +740,8 @@ class AgentOperator(CancellableAgentRunMixin,
BaseOperator, HITLReviewMixin):
elif extra_kwargs.get("toolsets"):
extra_kwargs["toolsets"] = [ensure_masked(ts) for ts in
extra_kwargs["toolsets"]]
capabilities = [
- replace(capability, toolset=ensure_masked(capability.toolset))
- if _is_concrete_toolset_capability(capability)
- else capability
- for capability in extra_kwargs.get("capabilities") or []
+ _replace_capability_toolset(capability, ensure_masked)
+ for capability in self.capabilities or passed_through or []
]
if self.durable and storage is not None and counter is not None:
# Tools supplied through a ``Toolset`` capability bypass the
@@ -695,7 +768,13 @@ class AgentOperator(CancellableAgentRunMixin,
BaseOperator, HITLReviewMixin):
files would be gone on resume. A sandbox another task owns
(``attach_to``)
outlives the pause, and the resumed run attaches to it again.
"""
- if not AIRFLOW_V_3_3_PLUS or self.durable or self.enable_hitl_review
or self.code_mode:
+ if (
+ not AIRFLOW_V_3_3_PLUS
+ or self.durable
+ or self.enable_hitl_review
+ or self.code_mode
+ or _contains_code_mode(self._declared_capabilities)
+ ):
return False
return all(sandbox.attach_to is not None for sandbox in
self._sandbox_toolsets())
@@ -719,11 +798,16 @@ class AgentOperator(CancellableAgentRunMixin,
BaseOperator, HITLReviewMixin):
for toolset in (*(self.toolsets or []),
*(self.agent_params.get("toolsets") or []))
if isinstance(toolset, AbstractToolset)
]
- for capability in self.agent_params.get("capabilities") or ():
- if _is_concrete_toolset_capability(capability):
- candidates.append(capability.toolset)
+ for capability in self._declared_capabilities:
+ if (toolset := _resolve_capability_toolset(capability)) is not
None:
+ candidates.append(toolset)
return candidates
+ @property
+ def _declared_capabilities(self) -> list[AgentCapability[Any]]:
+ """Capabilities passed via ``capabilities=`` and
``agent_params["capabilities"]``."""
+ return [*(self.capabilities or ()),
*(self.agent_params.get("capabilities") or ())]
+
def _toolset_ids(self) -> list[str]:
"""Ids of every leaf toolset, which for SQL and MCP toolsets name the
connection."""
# Declared order, not sorted: two toolsets that swapped connections
must not compare equal.
@@ -766,9 +850,10 @@ class AgentOperator(CancellableAgentRunMixin,
BaseOperator, HITLReviewMixin):
for capability in capabilities:
# ``Toolset.toolset`` can be a concrete toolset or a callable
factory
# resolved per run; only a concrete toolset can be wrapped here.
- if _is_concrete_toolset_capability(capability):
+ toolset = _resolve_capability_toolset(capability)
+ if toolset is not None:
cached = CachingToolset(
- wrapped=capability.toolset,
+ wrapped=toolset,
storage=storage,
counter=counter,
replay_usage=self._replay_usage,
@@ -1121,7 +1206,7 @@ class AgentOperator(CancellableAgentRunMixin,
BaseOperator, HITLReviewMixin):
raise UnsupportedToolDeferralError(
f"The agent called tools that need approval ({pending_names}),
but tool approval "
"needs Airflow 3.3+ and is not available with durable,
enable_hitl_review, "
- "code_mode or a SandboxToolset."
+ "code mode (code_mode=True or a CodeMode capability) or a
SandboxToolset."
)
store = context["task_state_store"]
if store.get(_TOOL_APPROVAL_REQUESTED_KEY):
diff --git a/providers/common/ai/tests/unit/common/ai/decorators/test_agent.py
b/providers/common/ai/tests/unit/common/ai/decorators/test_agent.py
index 0ed8b5cd8d2..65190ab6105 100644
--- a/providers/common/ai/tests/unit/common/ai/decorators/test_agent.py
+++ b/providers/common/ai/tests/unit/common/ai/decorators/test_agent.py
@@ -20,6 +20,7 @@ from unittest.mock import ANY, MagicMock, patch
import pytest
from pydantic import BaseModel
+from pydantic_ai.capabilities import Thinking
from pydantic_ai.messages import ImageUrl
from pydantic_ai.toolsets.function import FunctionToolset
@@ -179,6 +180,24 @@ class TestAgentDecoratedOperator:
assert isinstance(passed_toolsets[0], LoggingToolset)
assert passed_toolsets[0].wrapped == MaskingToolset(wrapped=toolset)
+ @patch("airflow.providers.common.ai.operators.agent.PydanticAIHook",
autospec=True)
+ def test_execute_passes_capabilities_through(self, mock_hook_cls,
make_mock_run_result):
+ mock_agent = MagicMock(spec=["run_sync", "instrument"])
+ mock_agent.run_sync.return_value = make_mock_run_result("result")
+ mock_hook_cls.get_hook.return_value.create_agent.return_value =
mock_agent
+ thinking = Thinking(effort="high")
+
+ op = _AgentDecoratedOperator(
+ task_id="test",
+ python_callable=lambda: "Do something",
+ llm_conn_id="my_llm",
+ capabilities=[thinking],
+ )
+ op.execute(context=_make_context())
+
+ create_call =
mock_hook_cls.get_hook.return_value.create_agent.call_args
+ assert create_call.kwargs["capabilities"] == [thinking]
+
@requires_typed_xcom
@patch("airflow.providers.common.ai.operators.agent.PydanticAIHook",
autospec=True)
def test_execute_structured_output(self, mock_hook_cls,
make_mock_run_result):
diff --git a/providers/common/ai/tests/unit/common/ai/operators/test_agent.py
b/providers/common/ai/tests/unit/common/ai/operators/test_agent.py
index 1f8575be2ce..b49bcaf1a30 100644
--- a/providers/common/ai/tests/unit/common/ai/operators/test_agent.py
+++ b/providers/common/ai/tests/unit/common/ai/operators/test_agent.py
@@ -23,13 +23,20 @@ import sys
from contextlib import nullcontext
from datetime import timedelta
from decimal import Decimal
-from types import SimpleNamespace
+from types import ModuleType, SimpleNamespace
from unittest.mock import ANY, MagicMock, PropertyMock, call, patch
import pytest
from pydantic import BaseModel
from pydantic_ai import Agent, DeferredToolRequests, Tool
-from pydantic_ai.capabilities import Toolset
+from pydantic_ai.capabilities import (
+ AbstractCapability,
+ CombinedCapability,
+ PrefixTools,
+ Thinking,
+ Toolset,
+ WebSearch,
+)
from pydantic_ai.exceptions import UsageLimitExceeded
from pydantic_ai.messages import (
ModelMessage,
@@ -80,6 +87,7 @@ from airflow.providers.common.ai.utils.usage_budget import (
from airflow.providers.common.compat.sdk import
AirflowOptionalProviderFeatureException, BaseHook
from airflow.sdk import DAG, task
+from tests_common.test_utils.compat import OperatorSerialization
from tests_common.test_utils.version_compat import AIRFLOW_V_3_1_PLUS,
AIRFLOW_V_3_3_PLUS
from unit.common.ai.sandbox.fake_tags import TaggedBackend
@@ -381,28 +389,30 @@ class TestAgentOperatorToolsetTemplating:
assert found._db_conn_id == "tenant_acme"
assert inner._db_conn_id == "tenant_{{ params.customer }}"
- def test_toolset_capability_is_rendered(self):
+ @pytest.mark.parametrize("passed_as", ["capabilities", "agent_params"])
+ def test_toolset_capability_is_rendered(self, passed_as):
capability = Toolset(SQLToolset(db_conn_id="tenant_{{ params.customer
}}"))
- op = AgentOperator(
- task_id="t", prompt="p", llm_conn_id="llm",
agent_params={"capabilities": [capability]}
+ kwargs = (
+ {"capabilities": [capability]}
+ if passed_as == "capabilities"
+ else {"agent_params": {"capabilities": [capability]}}
)
+ op = AgentOperator(task_id="t", prompt="p", llm_conn_id="llm",
**kwargs)
op.render_template_fields(self.CONTEXT)
- (rendered,) = op.agent_params["capabilities"]
+ (rendered,) = op._declared_capabilities
assert rendered.toolset._db_conn_id == "tenant_acme"
assert capability.toolset._db_conn_id == "tenant_{{ params.customer }}"
def test_callable_toolset_capability_is_left_as_is(self):
"""A factory resolved per run has no toolset to render until the run
starts."""
capability = Toolset(lambda ctx: SQLToolset(db_conn_id="tenant_{{
params.customer }}"))
- op = AgentOperator(
- task_id="t", prompt="p", llm_conn_id="llm",
agent_params={"capabilities": [capability]}
- )
+ op = AgentOperator(task_id="t", prompt="p", llm_conn_id="llm",
capabilities=[capability])
op.render_template_fields(self.CONTEXT)
- assert op.agent_params["capabilities"][0] is capability
+ assert op.capabilities[0] is capability
def test_only_connection_ids_are_templated(self):
"""allowed_tables is validated and canonicalised in __init__, so
rendering it later
@@ -433,13 +443,11 @@ class TestAgentOperatorToolsetTemplating:
def test_wrapped_toolset_inside_a_capability_is_rendered(self):
capability = Toolset(SQLToolset(db_conn_id="tenant_{{ params.customer
}}").prefixed("crm"))
- op = AgentOperator(
- task_id="t", prompt="p", llm_conn_id="llm",
agent_params={"capabilities": [capability]}
- )
+ op = AgentOperator(task_id="t", prompt="p", llm_conn_id="llm",
capabilities=[capability])
op.render_template_fields(self.CONTEXT)
- found = find_toolset([op.agent_params["capabilities"][0].toolset],
SQLToolset)
+ found = find_toolset([op.capabilities[0].toolset], SQLToolset)
assert found is not None
assert found.id == "sql-tenant_acme"
@@ -1132,6 +1140,147 @@ class TestAgentOperatorExecute:
op.execute(context=context)
[email protected]
+class _FakeCodeMode(AbstractCapability):
+ """Stands in for pydantic-ai-harness ``CodeMode``, which the ``code-mode``
extra installs."""
+
+
+def _harness_root_without_code_mode_extra() -> ModuleType:
+ """
+ The pydantic_ai_harness package as installed without the ``code-mode``
extra.
+
+ Like the real package, it exports CodeMode through a module
``__getattr__`` that imports
+ ``pydantic_ai_harness.code_mode``, and that import fails because
pydantic-monty is missing.
+ """
+ root = ModuleType("pydantic_ai_harness")
+
+ def __getattr__(name: str) -> object:
+ raise ImportError("pydantic-monty is required for CodeMode")
+
+ root.__getattr__ = __getattr__ # type: ignore[method-assign]
+ return root
+
+
[email protected]
+def fake_harness(monkeypatch):
+ """The harness with CodeMode imported: its ``code_mode`` module is
loaded."""
+ monkeypatch.setitem(sys.modules, "pydantic_ai_harness",
_harness_root_without_code_mode_extra())
+ monkeypatch.setitem(sys.modules, "pydantic_ai_harness.code_mode",
SimpleNamespace(CodeMode=_FakeCodeMode))
+
+
[email protected]
+def harness_without_code_mode_extra(monkeypatch):
+ """The harness installed for other capabilities, with CodeMode never
imported."""
+ monkeypatch.setitem(sys.modules, "pydantic_ai_harness",
_harness_root_without_code_mode_extra())
+ monkeypatch.delitem(sys.modules, "pydantic_ai_harness.code_mode",
raising=False)
+
+
+class TestAgentOperatorCapabilities:
+ @patch("airflow.providers.common.ai.operators.agent.PydanticAIHook",
autospec=True)
+ def test_capabilities_are_passed_to_the_agent_in_order(self,
mock_hook_cls, make_mock_run_result):
+ mock_hook_cls.get_hook.return_value.create_agent.return_value =
_make_mock_agent(
+ "ok", make_mock_run_result
+ )
+ thinking, search = Thinking(effort="high"), WebSearch()
+
+ op = AgentOperator(task_id="t", prompt="p", llm_conn_id="llm",
capabilities=[thinking, search])
+ op.execute(context=_make_context())
+
+ create_call =
mock_hook_cls.get_hook.return_value.create_agent.call_args
+ assert create_call.kwargs["capabilities"] == [thinking, search]
+
+ @patch("airflow.providers.common.ai.operators.agent.PydanticAIHook",
autospec=True)
+ def test_capabilities_in_both_places_are_refused(self, mock_hook_cls):
+ op = AgentOperator(
+ task_id="t",
+ prompt="p",
+ llm_conn_id="llm",
+ capabilities=[Thinking()],
+ agent_params={"capabilities": [WebSearch()]},
+ )
+
+ with pytest.raises(ValueError, match="not both"):
+ op.execute(context=_make_context())
+ mock_hook_cls.get_hook.return_value.create_agent.assert_not_called()
+
+ def test_capabilities_are_left_out_of_the_serialized_dag(self):
+ with DAG(dag_id="d", schedule=None):
+ op = AgentOperator(task_id="t", prompt="p", llm_conn_id="llm",
capabilities=[Thinking()])
+
+ serialized = OperatorSerialization.serialize_operator(op)
+
+ assert "Thinking" not in json.dumps(serialized, default=str)
+
+ @pytest.mark.parametrize(
+ "wrap",
+ [
+ pytest.param(lambda code_mode: code_mode, id="direct"),
+ pytest.param(lambda code_mode: CombinedCapability([Thinking(),
code_mode]), id="combined"),
+ pytest.param(lambda code_mode:
code_mode.prefix_tools("sandboxed"), id="prefixed"),
+ pytest.param(
+ lambda code_mode:
CombinedCapability([PrefixTools(wrapped=code_mode, prefix="sandboxed")]),
+ id="prefixed-in-combined",
+ ),
+ ],
+ )
+ @pytest.mark.parametrize("passed_as", ["capabilities", "agent_params"])
+ def test_durable_refuses_a_code_mode_capability(self, fake_harness, wrap,
passed_as):
+ capabilities = [wrap(_FakeCodeMode())]
+ kwargs = (
+ {"capabilities": capabilities}
+ if passed_as == "capabilities"
+ else {"agent_params": {"capabilities": capabilities}}
+ )
+
+ with pytest.raises(ValueError, match="CodeMode capability"):
+ AgentOperator(task_id="t", prompt="p", llm_conn_id="llm",
durable=True, **kwargs)
+
+ def test_durable_accepts_capabilities_without_code_mode(self,
fake_harness):
+ thinking = Thinking()
+
+ op = AgentOperator(task_id="t", prompt="p", llm_conn_id="llm",
durable=True, capabilities=[thinking])
+
+ assert op.capabilities == [thinking]
+
+ @pytest.mark.skipif(not AIRFLOW_V_3_3_PLUS, reason="Per-tool approval
needs Airflow >= 3.3")
+ @pytest.mark.parametrize(
+ ("capability", "supported"),
+ [
+ pytest.param(Thinking(), True, id="thinking"),
+ pytest.param(_FakeCodeMode(), False, id="code-mode"),
+ ],
+ )
+ def test_code_mode_capability_turns_off_tool_approval(self, fake_harness,
capability, supported):
+ op = AgentOperator(task_id="t", prompt="p", llm_conn_id="llm",
capabilities=[capability])
+
+ assert op._supports_tool_approval() is supported
+
+ @pytest.mark.skipif(not AIRFLOW_V_3_3_PLUS, reason="Per-tool approval
needs Airflow >= 3.3")
+ def test_harness_without_the_code_mode_extra_does_not_break_other_agents(
+ self, harness_without_code_mode_extra
+ ):
+ # Neither the durable check at construction nor the approval check at
run time may
+ # touch the harness root, whose CodeMode export would raise
ImportError here.
+ AgentOperator(task_id="t", prompt="p", llm_conn_id="llm",
durable=True, capabilities=[Thinking()])
+ op = AgentOperator(task_id="u", prompt="p", llm_conn_id="llm",
capabilities=[Thinking()])
+
+ assert op._supports_tool_approval() is True
+
+ def
test_code_mode_flag_and_code_mode_capability_are_refused_together(self,
fake_harness):
+ with pytest.raises(ValueError, match="one or the other"):
+ AgentOperator(
+ task_id="t", prompt="p", llm_conn_id="llm", code_mode=True,
capabilities=[_FakeCodeMode()]
+ )
+
+ def test_capability_function_is_left_for_the_run_to_resolve(self,
fake_harness):
+ def build(ctx):
+ return Thinking()
+
+ op = AgentOperator(task_id="t", prompt="p", llm_conn_id="llm",
durable=True, capabilities=[build])
+
+ assert op.capabilities == [build]
+
+
@pytest.mark.skipif(
not AIRFLOW_V_3_1_PLUS, reason="Human in the loop is only compatible with
Airflow >= 3.1.0"
)