airflow.providers.common.ai.managed_agents.base

The contract a provider hook implements to expose a vendor-managed agent.

A managed agent runs its own reasoning loop on the vendor’s infrastructure – Snowflake Cortex Agents, Amazon Bedrock AgentCore, Azure AI Foundry hosted agents, Vertex AI Agent Engine. Airflow submits a request and reads an answer. This module defines the shape of that exchange once, so that every consumer in common.ai (the toolset, the failover group) is written against one interface rather than one per cloud.

The design follows DbApiHook in common.sql: a small base mixed into each vendor’s own hook, with the agent as an argument to every method, because a hook is scoped to a connection and one connection reaches many agents. Vendor providers adopt it the way they adopt BaseMessageQueueProvider from common.messaging: behind an optional extra, in a module whose import of this contract is guarded, so a provider that floors Airflow 2 never has to raise its floor.

This module imports nothing from pydantic-ai on purpose. A vendor hook’s guarded import of the contract must stay cheap.

Classes

ManagedAgentRef

Normalized identity of a remote agent.

ManagedAgentCapabilities

What a (hook, agent) pair can do.

ManagedAgentRequest

One request to a managed agent.

ManagedAgentUsage

Usage the vendor reported for one invocation. Every field is optional because vendors differ.

ManagedAgentResponse

One answer from a managed agent.

BaseManagedAgentHook

Mixin a vendor hook adopts to expose its managed agents through the common contract.

ManagedAgentClient

What common.ai's consumers are typed against.

BoundManagedAgent

A (hook, agent) pair. Forwards to the hook and resolves identity lazily.

Functions

describe(client)

Render a client's identity for a log line without letting identity resolution fail the caller.

Module Contents

class airflow.providers.common.ai.managed_agents.base.ManagedAgentRef[source]

Normalized identity of a remote agent.

Parameters:
  • platform – A stable, dotted platform id such as aws.bedrock_agentcore or gcp.vertex_agent_engine. Used as a metric tag, so keep the set small.

  • name – The vendor’s canonical identifier for the agent: an ARN, a full resource name, DATABASE.SCHEMA.NAME.

  • version – The resolved version or revision when the platform exposes one. Recorded so a behaviour change can be attributed to a deployment rather than to Airflow.

platform: str[source]
name: str[source]
version: str | None = None[source]
class airflow.providers.common.ai.managed_agents.base.ManagedAgentCapabilities[source]

What a (hook, agent) pair can do.

Consumers check these and refuse rather than degrade: a bound agent rejects a request that carries a session_id when sessions is False, and a failover group never offers sessions at all, because failing over discards the conversation the primary was holding.

sessions: bool = False[source]

Whether ManagedAgentRequest.session_id continues a conversation.

structured_output: bool = False[source]

Whether ManagedAgentResponse.structured can carry a typed value.

usage: bool = False[source]

Whether ManagedAgentResponse.usage is populated.

trace: bool = False[source]

Whether ManagedAgentResponse.trace_ref is populated.

class airflow.providers.common.ai.managed_agents.base.ManagedAgentRequest[source]

One request to a managed agent.

Exactly one of prompt and messages must be set. Everything the contract does not type travels in vendor_options. A hook may read some of them itself and passes the rest through to the vendor call. Hooks reject options that would re-target the call (the agent identity, the connection), because a model-facing caller must not be able to change what it is talking to.

prompt: str | None = None[source]
messages: collections.abc.Sequence[dict[str, Any]] | None = None[source]
session_id: str | None = None[source]
timeout: float | None = None[source]

Seconds to wait for the vendor call. The hook must enforce it on the request itself.

vendor_options: dict[str, Any][source]
__post_init__()[source]
as_messages()[source]

Return the request as a message list, for vendors that only accept messages.

class airflow.providers.common.ai.managed_agents.base.ManagedAgentUsage[source]

Usage the vendor reported for one invocation. Every field is optional because vendors differ.

input_tokens: int | None = None[source]
output_tokens: int | None = None[source]
class airflow.providers.common.ai.managed_agents.base.ManagedAgentResponse[source]

One answer from a managed agent.

text is what a calling model should read; the hook unwraps the vendor envelope to produce it. raw is that envelope, always populated and never handed to a model, so a Python caller loses nothing.

text: str[source]
raw: Any[source]
structured: Any | None = None[source]
session_id: str | None = None[source]
usage: ManagedAgentUsage | None = None[source]
trace_ref: str | None = None[source]

A vendor request, invocation or trace id, for joining Airflow’s record to the vendor’s.

class airflow.providers.common.ai.managed_agents.base.BaseManagedAgentHook[source]

Bases: abc.ABC

Mixin a vendor hook adopts to expose its managed agents through the common contract.

Mixed in beside the vendor’s own base and never replacing it:

class BedrockAgentCoreHook(AwsBaseHook, BaseManagedAgentHook): ...

It therefore has no __init__ and makes no assumption about get_conn. The agent is an argument to every method, the way a statement is an argument to DbApiHook.run.

Method names are chosen to collide with nothing on the shipped vendor hooks. That matters more than it looks: SnowflakeCortexAgentHook already defines run_agent and describe_agent, and an abstract method that a vendor base happens to define is silently satisfied with the wrong signature.

Implementations sort failures into three classes, and conflating them is the most common way an adoption goes wrong:

  • ManagedAgentRejected – the agent rejected the request in a way rephrasing could fix. The toolset turns it into a pydantic-ai ModelRetry so the calling model tries again.

  • ManagedAgentInvocationError – terminal: bad credentials, missing agent, revoked quota. Nothing on the agent side recovers it; whether the task retries is the task’s retry policy.

  • Anything transient (429, 5xx, connection reset, read timeout) – propagate unchanged. Airflow’s task-level retry is the right layer; a rephrase does nothing for a 503.

agent_platform: ClassVar[str][source]

The platform every ManagedAgentRef from this hook carries.

abstractmethod resolve_agent(agent)[source]

Normalize agent into a platform-qualified reference. Must not make a network call.

abstractmethod get_agent_capabilities(agent)[source]

Report what agent on this connection can do. Must not make a network call.

abstractmethod invoke_agent(agent, request)[source]

Send request to agent and return its answer. Blocking.

agent(agent)[source]

Bind one agent on this connection. This is what the common.ai toolsets consume.

class airflow.providers.common.ai.managed_agents.base.ManagedAgentClient[source]

Bases: Protocol

What common.ai’s consumers are typed against.

A BoundManagedAgent satisfies it. So does FailoverManagedAgentClient, and so can anything that needs no Airflow connection at all.

property ref: ManagedAgentRef[source]
property capabilities: ManagedAgentCapabilities[source]
invoke(request)[source]
class airflow.providers.common.ai.managed_agents.base.BoundManagedAgent[source]

A (hook, agent) pair. Forwards to the hook and resolves identity lazily.

This is also where the contract’s “refuse rather than degrade” rule is enforced for every adopter: a request that asks for something the pair’s capabilities do not include is rejected before the hook is called.

hook: BaseManagedAgentHook[source]
agent: str[source]
property ref: ManagedAgentRef[source]
property capabilities: ManagedAgentCapabilities[source]
invoke(request)[source]
airflow.providers.common.ai.managed_agents.base.describe(client)[source]

Render a client’s identity for a log line without letting identity resolution fail the caller.

A standby whose connection is misconfigured must not fail a call the primary served.

Was this entry helpful?