airflow.providers.common.ai.toolsets.managed_agent¶
Attributes¶
Classes¶
Base class exposing a vendor-managed agent as a single pydantic-ai tool. |
|
Expose any managed-agent client as one tool. |
Module Contents¶
- 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.AirflowToolsetBase 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
ManagedAgentToolsetover aManagedAgentClient, which every provider hook that adoptsBaseManagedAgentHookproduces viahook.agent(...). Subclass this directly only for an agent that has no hook at all. Subclasses implementagent_refandinvoke_sync()(or override the asyncinvoke()); 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_namerendered as prose, matching howHookToolsethandles 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.
Nonedefers to the platform default. Exposed astimeoutso an implementation can honor it.max_retries (int) – How many times the calling model may rephrase after the remote agent rejects a request.
0turns the first rejection into a hard error.
- property timeout: float | None[source]¶
Seconds to wait for one invocation, or
Nonefor 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
promptto 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 raisesManagedAgentRejectedinstead andManagedAgentToolsettranslates 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 execute_tool(name, tool_args, *, ctx, tool)[source]¶
Run tool
namewith validatedtool_argsand return its result unmasked.This is the method a subclass implements;
call_tool()runs it and masks what it returns.ctxandtoolare 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:
BaseManagedAgentToolsetExpose 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 aBoundManagedAgentfrom a vendor hook’sagent()method, or aFailoverManagedAgentClientover 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 inresponse.rawis 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.Nonemeans 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’sclass_method, for instance). The hook decides which keys are allowed.
- 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.