Airflow hooks as tools: HookToolset

Generic adapter that exposes selected methods of any Airflow Hook as pydantic-ai tools via introspection. Requires an explicit allowed_methods list; there is no auto-discovery.

from airflow.providers.http.hooks.http import HttpHook
from airflow.providers.common.ai.toolsets.hook import HookToolset

http_hook = HttpHook(http_conn_id="my_api")

toolset = HookToolset(
    http_hook,
    allowed_methods=["run"],
    tool_name_prefix="http_",
)

For each listed method, the introspection engine:

  1. Builds a JSON Schema from the method signature (inspect.signature + get_type_hints).

  2. Extracts the description from the first paragraph of the docstring.

  3. Enriches parameter descriptions from Sphinx :param: or Google Args: blocks.

Templated connection IDs

The hook’s connection ID is a Jinja template, rendered for each task instance just before it runs, so one toolset can reach a different system depending on the run: PostgresHook(postgres_conn_id="warehouse_{{ var.value.environment }}") switches between staging and production, and a mapped task can give each map index its own connection:

@task.agent(
    llm_conn_id="pydanticai_default",
    toolsets=[
        HookToolset(
            PostgresHook(postgres_conn_id="analytics_{{ task.op_kwargs.customer }}"),
            allowed_methods=["get_records"],
        )
    ],
)
def report(customer: str) -> str:
    return f"Summarize this month's orders for {customer}."


report.expand(customer=customers())

Each task instance runs against a copy of the hook with its rendered connection ID; the hook in the Dag file keeps the template. HookToolset reads the ID from the attribute conn_name_attr names, or from conn_id for hooks such as WasbHook that keep it there; a hook that stores it anywhere else is not templated. This works for hooks that read their connection when a method is called, which is what Airflow hooks are expected to do: a hook that looks the connection up in its constructor fails at Dag parse time, because the template is not a connection ID yet. The same warning as for SQLToolset applies: build the ID from values the Dag controls, not from params or dag_run.conf (see Templated connection IDs).

Fix arguments the model must not choose

Note

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

Exposing a method lets the model pick every argument it takes. When some of them are the Dag author’s decision, such as which bucket a storage hook reads, pin them:

from airflow.providers.amazon.aws.hooks.s3 import S3Hook

from airflow.providers.common.ai.toolsets import HookToolset

reports = HookToolset(
    S3Hook(aws_conn_id="aws_default"),
    allowed_methods=["list_keys", "read_key"],
    pinned_arguments={"bucket_name": "acme-reports"},
)

A pinned argument is left out of the schema the model sees and passed to every allowed method. If the model supplies it anyway, the call is refused while its arguments are validated, before an approval gate or the hook sees it. When the rest of the call is valid, the model is told the argument is fixed (see Restricting the agent).

A pin binds one parameter name, so every allowed method has to take it by that name. When one does not, the toolset raises ValueError when it is created: a method that takes the same thing under another name, such as S3Hook.delete_objects, which takes bucket, or inside a dict or **kwargs, such as S3Hook.generate_presigned_url, would let the model choose it after all. Expose such a method from a second HookToolset, where what it can reach is visible in the Dag. A method that works out the value again from another argument it is given, such as a full URL, is outside what the pin controls.

Pinned values are passed as written: they are not rendered as templates, and they are not part of what AgentOperator(durable=True) fingerprints, so change one only between Dag runs, not between the tries of one.

Restricting the agent

allowed_methods decides which hook methods become tools, and pinned_arguments decides which of their arguments the model cannot set. This agent can list and read one bucket through S3Hook, and nothing else:

airflow/providers/common/ai/example_dags/example_hook_toolset.py[source]

    @dag(tags=["example"])
    def example_hook_toolset_restricted():
        AgentOperator(
            task_id="read_reports",
            prompt="Summarize the September finance report.",
            llm_conn_id="pydanticai_default",
            toolsets=[
                HookToolset(
                    # Log in with credentials scoped to the reports bucket.
                    S3Hook(aws_conn_id="aws_reports_reader"),
                    # Offer these two methods as tools; no other S3Hook method is reachable.
                    allowed_methods=["list_keys", "read_key"],
                    # Fix the bucket: the model never sees the argument and cannot set it.
                    pinned_arguments={"bucket_name": "acme-reports"},
                    # Name the tools s3_list_keys and s3_read_key.
                    tool_name_prefix="s3_",
                    # Let the model correct an invalid call up to 2 times.
                    max_retries=2,
                )
            ],
        )

The model is offered two tools: s3_read_key, which takes only key, and s3_list_keys, whose parameters include prefix but not bucket_name. Run against an S3 endpoint that also holds an acme-payroll bucket, a call that names that bucket was refused before it reached S3, with this message to the model:

bucket_name is fixed for this tool: call it again without it.

Fix the errors and try again.

A call to a method that is not listed, such as s3_delete_objects, got Unknown tool name: 's3_delete_objects'. Available tools: 's3_list_keys', 's3_read_key'. The model can correct both kinds of call and carry on.

An exception from the hook is different: it fails the run, and the task with it. Reading a key that does not exist ended the run with ClientError: An error occurred (404) when calling the HeadObject operation: Not Found.

The pin fixes the bucket and leaves every key in it to the model. Give aws_reports_reader credentials that can read only that bucket, so the connection holds the same limit if a method you expose later reaches another bucket some other way.

Parameters

  • hook: An instantiated Airflow Hook. Its connection ID is templated.

  • allowed_methods: Method names to expose as tools. Required. Methods are validated with hasattr + callable at instantiation time.

  • tool_name_prefix: Optional prefix prepended to each tool name (e.g. "s3_" produces "s3_list_keys").

  • pinned_arguments: Arguments fixed by the Dag author rather than chosen by the model. See above.

  • max_retries: How many times the model may correct a call with invalid arguments, or one that supplies a pinned argument. Default None, the agent’s retries. See How often the model may correct a failed call.

When to choose it

Choose it when the target already has an Airflow connection and a hook, and what you want the agent to do is already a method on that hook. This is the cheapest route (no new server, no new credential, no new query dialect) and the only one that reaches any provider hook with synchronous methods without anyone writing an adapter first. HookToolset is a reflection-based adapter, so the work is choosing the method list.

What it cannot do

  • It allow-lists method names, and fixes only the arguments you pin. Once read_key is exposed, the agent picks the key within the pinned bucket; the defense-layer table states this outright. Choose methods whose worst case you accept, not methods you intend to constrain later. To expose a method that changes something and have a person approve the call first, wrap the toolset with .approval_required(). A task instance can pause for approval once per Dag run; see Approve an agent’s tool calls.

  • Its calls act as barriers. The tools are registered with sequential=True and each hook method runs in a worker thread, one blocking hook call at a time in the task process, so a slow call holds up every other tool the model emitted in that step, not only this toolset’s. This is not specific to HookToolset; see Tool calls as barriers.

  • It returns exactly one shape. Every result goes through serialize_for_llm and comes back as a JSON-encoded string; there is no structured error type and no ModelRetry wrapper, so a hook exception fails the agent run, and the task with it, instead of giving the model something it can correct. SQLToolset, by contrast, hands the database’s own error back as a retry.

  • It calls the method and serializes what comes back. The code contains no path that awaits a coroutine result, and none that checks for one, so an async def hook method is not a case this adapter is written to handle. Treat synchronous methods as the supported set.

A real example. The read-only S3 pair from the HookToolset guidance in Securing agent tools:

HookToolset(
    s3_hook,
    allowed_methods=["list_keys", "read_key"],
    tool_name_prefix="s3_",
)

Credentials and where it runs. The hook instance is yours, so the credential is whatever connection that hook resolves; the toolset never looks one up itself. Calls run in the Airflow worker process.

Was this entry helpful?