airflow.providers.common.ai.toolsets.managed_agent

Attributes

log

Classes

BaseManagedAgentToolset

Base class exposing a vendor-managed agent as a single pydantic-ai tool.

ManagedAgentToolset

Expose any managed-agent client as one tool.

Module Contents

airflow.providers.common.ai.toolsets.managed_agent.log[source]
class airflow.providers.common.ai.toolsets.managed_agent.BaseManagedAgentToolset(*, tool_name, description=None, timeout=None, max_retries=1)[source]

Bases: airflow.providers.common.ai.utils.toolset_base.AirflowToolset

Base class exposing a vendor-managed agent as a single pydantic-ai tool.

Note

Experimental: this can change or be removed in a minor release of this provider. See Stable and experimental features.

A managed agent runs its own reasoning loop on the vendor’s infrastructure. Airflow submits one request and reads one answer, so the Airflow-side agent features – toolsets, human-in-the-loop review, durable step replay – apply to the calling agent and never reach inside the managed agent.

Most code should not subclass this. Use ManagedAgentToolset over a ManagedAgentClient, which every provider hook that adopts BaseManagedAgentHook produces via hook.agent(...). Subclass this directly only for an agent that has no hook at all. Subclasses implement agent_ref and invoke_sync() (or override the async invoke()); tool naming, argument validation, result serialization, logging and metrics are handled here so every implementation presents the same surface to the model.

Parameters:
  • tool_name (str) – Name the calling model sees, and the identifier it emits when calling the tool. A verb phrase naming the specialist reads best, e.g. ask_bookings_analyst.

  • description (str | None) – What this agent knows and when to consult it. Optional – it falls back to tool_name rendered as prose, matching how HookToolset handles a method with no docstring. Worth writing anyway: it is what tells the model to consult the agent rather than answer from its own knowledge, and it is the only place to state a scope limit the name cannot carry (“cannot see revenue figures”). Since the argument schema is always a bare prompt, the name and this string are the whole of what the model knows about the agent.

  • timeout (float | None) – Seconds to wait for a single invocation. None defers to the platform default. Exposed as timeout so an implementation can honor it.

  • max_retries (int) – How many times the calling model may rephrase after the remote agent rejects a request. 0 turns the first rejection into a hard error.

replayable: bool = False[source]
property timeout: float | None[source]

Seconds to wait for one invocation, or None for the platform default.

property agent_ref: airflow.providers.common.ai.managed_agents.base.ManagedAgentRef[source]
Abstractmethod:

Normalized identity of the remote agent.

Logged after every successful call, so the resolved remote identity behind a task appears in that task’s log even though the Dag only names a connection. Resolution is never on the call’s critical path: a failure here is logged, not raised.

async invoke(prompt)[source]

Send prompt to the remote agent and return the agent’s answer.

Override this when the vendor call is already asynchronous. When it blocks, implement invoke_sync() instead and let the default implementation here run it in a worker thread, which keeps it off the event loop that the whole agent run shares.

Return the answer, not the transport envelope. Failures sort into three classes: ModelRetry (the model can fix it by rephrasing; a hook-backed client raises ManagedAgentRejected instead and ManagedAgentToolset translates it), ManagedAgentInvocationError (terminal), and anything transient, which should propagate unchanged so Airflow’s task-level retry handles it.

Release anything you allocate, on every path. Platforms that require a session bill for its lifetime, so an implementation that opens one here must close it in a finally. A tool call has no post-task cleanup hook to fall back on.

Parameters:

prompt (str) – The question or instruction to send to the remote agent.

abstractmethod invoke_sync(prompt)[source]

Blocking variant of invoke(), run in a worker thread.

A thread cannot be cancelled, so set a timeout on the underlying request: a caller that stops waiting does not stop this call.

Parameters:

prompt (str) – The question or instruction to send to the remote agent.

property id: str[source]

An ID for the toolset that is unique among all toolsets registered with the same agent.

If you’re implementing a concrete implementation that users can instantiate more than once, you should let them optionally pass a custom ID to the constructor and return that here.

A toolset needs to have an ID in order to be used in a durable execution environment like Temporal, in which case the ID will be used to identify the toolset’s activities within the workflow.

IDs wrapped in angle brackets (‘<agent>’ for an agent’s own function toolset, ‘<output>’ for its output tools) name a role the framework fills on the user’s behalf rather than a registered toolset. Don’t return one from your own toolset.

async get_tools(ctx)[source]

The tools that are available in this toolset.

async execute_tool(name, tool_args, *, ctx, tool)[source]

Run tool name with validated tool_args and return its result unmasked.

This is the method a subclass implements; call_tool() runs it and masks what it returns. ctx and tool are keyword-only so that arguments can be added here later without breaking subclasses.

class airflow.providers.common.ai.toolsets.managed_agent.ManagedAgentToolset(client, *, tool_name, description=None, timeout=None, max_retries=1, replayable=False, vendor_options=None)[source]

Bases: BaseManagedAgentToolset

Expose any managed-agent client as one tool.

Note

Experimental: this can change or be removed in a minor release of this provider. See Stable and experimental features.

This is the toolset to use. It accepts any ManagedAgentClient. The client is usually a BoundManagedAgent from a vendor hook’s agent() method, or a FailoverManagedAgentClient over several of them:

from airflow.providers.amazon.aws.hooks.bedrock import (
    BedrockAgentCoreHook,
)
from airflow.providers.common.ai.toolsets import ManagedAgentToolset

claims = BedrockAgentCoreHook(aws_conn_id="aws_prod").agent(RUNTIME_ARN)
toolset = ManagedAgentToolset(
    claims,
    tool_name="ask_claims_agent",
    description="Reviews an insurance claim and returns a coverage determination.",
)

The model receives response.text. The vendor envelope in response.raw is for Python callers of the client and never reaches the model.

Parameters:
  • client (airflow.providers.common.ai.managed_agents.base.ManagedAgentClient) – The agent to consult.

  • tool_name (str) – See BaseManagedAgentToolset.

  • description (str | None) – See BaseManagedAgentToolset.

  • timeout (float | None) – Passed to the client on every request as ManagedAgentRequest.timeout. None means whatever the vendor client defaults to, which for Agent Engine is no deadline at all.

  • max_retries (int) – See BaseManagedAgentToolset.

  • replayable (bool) – Whether the durable cache may replay a completed call. Only set it for an agent that is read-only.

  • vendor_options (dict[str, Any] | None) – Sent with every request as ManagedAgentRequest.vendor_options, for per-agent settings the vendor hook accepts there (Agent Engine’s class_method, for instance). The hook decides which keys are allowed.

replayable = False[source]
property client: airflow.providers.common.ai.managed_agents.base.ManagedAgentClient[source]
property agent_ref: airflow.providers.common.ai.managed_agents.base.ManagedAgentRef[source]

Normalized identity of the remote agent.

Logged after every successful call, so the resolved remote identity behind a task appears in that task’s log even though the Dag only names a connection. Resolution is never on the call’s critical path: a failure here is logged, not raised.

invoke_sync(prompt)[source]

Blocking variant of invoke(), run in a worker thread.

A thread cannot be cancelled, so set a timeout on the underlying request: a caller that stops waiting does not stop this call.

Parameters:

prompt (str) – The question or instruction to send to the remote agent.

Was this entry helpful?