Skip to content

Support Airflow 2.11 in the Common AI provider - #73991

Merged
kaxil merged 3 commits into
apache:mainfrom
astronomer:common-ai-airflow-2
Oct 1, 2026
Merged

kaxil merged 3 commits into
apache:mainfrom
astronomer:common-ai-airflow-2

Conversation

@kaxil

@kaxil kaxil commented Sep 30, 2026 •

Copy link
Copy Markdown
Member

Summary

The Common AI provider required Airflow 3.0, so an Airflow 2.11 deployment could not install it at all. This lowers the floor to apache-airflow>=2.11.0, the floor common.compat, standard and common.sql already carry. When the community provider floor moves to 3.1, this provider moves with it.

On 2.11 the operators, decorators, hooks and toolsets run as they do on Airflow 3.0. Features that need a newer core already have version gates, which this PR leaves in place: approval gates and HITL review (3.1), and retry policies, tool-approval pause and resume, and the task state store (3.3). Approval gates and HITL review fail at Dag parse time with a message naming the version. LLMRetryPolicy fails at import. An approval-required tool fails the task, as it already does on 3.0 to 3.2. Durable execution falls back to durable_cache_path as on 3.0 to 3.2.

How this was checked

The branch ran on a real Airflow 2.11.2 in breeze, with the provider wheels built from this branch and a real model (Anthropic claude-sonnet-5). One Dag was triggered through the REST API and run by the scheduler with the LocalExecutor:

breeze release-management prepare-provider-distributions common.ai common.compat standard common.sql \
    --distribution-format wheel --skip-tag-check
breeze shell --python 3.12 --backend postgres --use-airflow-version 2.11.2 \
    --airflow-constraints-reference constraints-2.11.2 --install-airflow-with-constraints \
    --providers-skip-constraints --use-distributions-from-dist --distribution-format wheel
# inside the container, with the Dag in /files/dags:
AIRFLOW__SCHEDULER__STANDALONE_DAG_PROCESSOR=False airflow standalone

Breeze turns STANDALONE_DAG_PROCESSOR on for Airflow 3. Airflow 2's airflow standalone does not start a separate Dag processor, so without that override it never parses a Dag.

The LLM endpoint host is redacted in the HTTP request lines of the log screenshots.

Core 2.11.2 with this branch's providers (the list filtered to those four):

Providers page on Airflow 2.11.2

The run: every task succeeded; path_a was skipped because the model picked path_b.

Grid view of the run

AgentOperator with SQLToolset: the model lists the tables, writes a query, and the toolset runs it against SQLite.

AgentOperator with SQLToolset log

AgentOperator(durable=True) with a function toolset: the tool call and the durable cache on the durable_cache_path backend.

Durable AgentOperator log

LLMBranchOperator:

LLMBranchOperator log

@task.agent(output_type=...) consumed downstream, which arrives as a dict before Airflow 3.3:

Downstream task receiving the structured output

The Dag
"""Common AI on Airflow 2.11: the core path plus the neighbouring operators."""

from __future__ import annotations

import datetime

from pydantic import BaseModel
from pydantic_ai.toolsets import FunctionToolset

from airflow.decorators import dag, task
from airflow.operators.empty import EmptyOperator
from airflow.providers.common.ai.operators.agent import AgentOperator
from airflow.providers.common.ai.operators.llm import LLMOperator
from airflow.providers.common.ai.operators.llm_branch import LLMBranchOperator
from airflow.providers.common.ai.operators.llm_sql import LLMSQLQueryOperator
from airflow.providers.common.ai.toolsets.sql import SQLToolset

CONN = "pydanticai_e2e"


class Summary(BaseModel):
    title: str
    score: int


def get_weather(city: str) -> str:
    """Return the weather for a city."""
    print(f"TOOL_CALLED get_weather city={city!r}")
    return f"It is sunny in {city}."


@dag(schedule=None, start_date=datetime.datetime(2026, 1, 1), catchup=False)
def commonai_af2_e2e():
    llm = LLMOperator(task_id="llm_operator", prompt="Say hello in one word.", llm_conn_id=CONN)

    @task.llm(llm_conn_id=CONN, system_prompt="Be terse.")
    def llm_decorator(topic: str):
        return f"One sentence about {topic}."

    agent_fn = AgentOperator(
        task_id="agent_function_toolset",
        prompt="What is the weather in Paris?",
        llm_conn_id=CONN,
        toolsets=[FunctionToolset([get_weather])],
    )

    agent_sql = AgentOperator(
        task_id="agent_sql_toolset",
        prompt="How many rows are in the orders table?",
        llm_conn_id=CONN,
        toolsets=[SQLToolset(db_conn_id="sqlite_e2e", allowed_tables=["orders"], max_rows=5)],
    )

    @task.agent(llm_conn_id=CONN, output_type=Summary)
    def agent_structured():
        return "Summarise Airflow in a title and a score out of 10."

    @task
    def consume(summary):
        print(f"CONSUMED type={type(summary).__name__} value={summary!r}")
        if not isinstance(summary, (dict, Summary)):
            raise TypeError(f"unexpected XCom type {type(summary)}")
        return summary

    branch = LLMBranchOperator(
        task_id="llm_branch",
        prompt="Pick the right path.",
        llm_conn_id=CONN,
        branches={"path_a": "Choose for greetings", "path_b": "Choose for anything else"},
    )
    path_a = EmptyOperator(task_id="path_a")
    path_b = EmptyOperator(task_id="path_b")

    sql_gen = LLMSQLQueryOperator(
        task_id="llm_sql",
        prompt="Count orders per customer.",
        llm_conn_id=CONN,
        db_conn_id="sqlite_e2e",
        table_names=["orders"],
    )

    durable = AgentOperator(
        task_id="agent_durable",
        prompt="What is the weather in Rome?",
        llm_conn_id=CONN,
        toolsets=[FunctionToolset([get_weather])],
        durable=True,
    )

    llm >> llm_decorator("Airflow")
    agent_fn >> agent_sql
    consume(agent_structured())
    branch >> [path_a, path_b]
    mapped = LLMOperator.partial(task_id="llm_mapped", llm_conn_id=CONN).expand(
        prompt=["one", "two", "three"]
    )


commonai_af2_e2e()

The provider's unit tests now run in the existing Compat 2.11.1 job. Locally they pass on 2.11.1 (2313 passed) and on Airflow 3 (2890 passed); common.compat passes on both.

Design rationale

  • Most of the provider already went through common.compat.sdk. The gaps were a few direct Airflow 3 imports and APIs:
    • SET_DURING_EXECUTION in the decorators
    • BaseHook.get_hook(hook_params=...)
    • TaskInstance.id
    • ObjectStoragePath from airflow.sdk
    • structlog, which the provider imported but never declared; Airflow 3 brings it in through the Task SDK and Airflow 2 does not.
  • common.compat changes behaviour on Airflow 2, deliberately.
    • get_current_context now resolves to the standard provider's Airflow 2 fallback before core's. Outside a task it raises RuntimeError, as Airflow 3 does, instead of core's AirflowException. Without the standard provider it still falls back to core. No caller in this repository catches AirflowException from the compat name; SandboxToolset catches RuntimeError, so on Airflow 2 its attach_to outside a task failed instead of falling back to owner=.
    • SET_DURING_EXECUTION is new in compat. On Airflow 2 it is an ArgNotSet whose repr matches Airflow 3's sentinel. With a bare NOTSET, Airflow 2 stores the decorator's prompt template field as an object address, which differs per process and changes the serialized Dag's hash.
    • The common-compat line in the Common AI pyproject.toml is marked # use next version, since the decorators need the new export.
  • PydanticAIHook.get_hook mirrors Airflow 3's BaseHook.get_hook body (get_connection(conn_id).get_hook(hook_params=...)). Overriding the classmethod keeps the operators, and the tests that patch get_hook, unchanged. Calling get_connection().get_hook() from the operators instead would have bypassed those patches.
  • The agent's per-attempt run key (pydantic-ai run_id, the run_id XCom, gen_ai.agent.call.id) falls back to dag_id/run_id/task_id/map_index/try_number where the task instance has no id. Spans leave out airflow.task_instance.id there, because a composite is not a task-instance id.
  • Logging on Airflow 2. Nothing configures structlog on Airflow 2, so the provider's module loggers printed every level, debug included, into task logs. get_task_logger() wraps the airflow.task stdlib logger on Airflow 2, so the task log's level and handlers apply and the record names the caller's line. It does not change the process-wide structlog configuration. On Airflow 3 it returns the same logger as before.
  • AgentOperator declares the HITL review extra link only on 3.1+. On Airflow 2 the webserver logs an error for the unregistered link class on every Dag load, and the link can never render there.
  • CI: common.ai is removed from remove-providers in the Compat 2.11.1 row, so its tests run in the existing job. Test files that imported airflow.sdk directly now import through common.compat.sdk. Two tests are gated to Airflow 3:
    • Airflow 2 reports a missing constructor argument as AirflowException.
    • Airflow 2's MappedOperator resolves an expansion through the metadata database and a task session, which a unit test with a dict context cannot drive.

Gotchas

  • Install Airflow with its constraints file, then add the provider without it: the 2.11 constraints pin common-compat and common-sql below this provider's floors. Airflow 2.11.0 installed without constraints can pick up universal-pathlib 0.3, which its ObjectStoragePath rejects; 2.11.1 and later cap it. The installation page now says this.
  • The skills and git extras need apache-airflow-providers-git, which requires Airflow 3.
  • Before Airflow 3.3, a structured output reaches downstream tasks as a dict, the same as on 3.0 to 3.2. The connection form's Model field needs 3.2; on older cores the model goes in Extra.

  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Comment thread providers/common/compat/tests/unit/common/compat/test__set_during_execution.py Outdated
Comment thread providers/common/ai/docs/installation.rst Outdated
Comment thread providers/common/ai/src/airflow/providers/common/ai/observability.py Outdated
Comment thread providers/common/ai/pyproject.toml
@Lee-W

Lee-W commented Oct 1, 2026

Copy link
Copy Markdown
Member

would be nice if we could drop the co-author in the latest commit. thanks!

@github-actions

github-actions Bot commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

uv.lock on main just moved via #73532 ("Add a vendor-neutral managed-agent hook contract to Common AI"), commit e9c552c and this PR currently conflicts.

Quickest fix:

git fetch upstream main && git rebase upstream/main
rm uv.lock && uv lock
git add uv.lock && git rebase --continue
git push --force-with-lease

Automated nudge — ignore if you're not ready to rebase. This comment is updated in place on future uv.lock bumps.

@vatsrahul1001 vatsrahul1001 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Need to fix conflicts.

@kaxil
kaxil force-pushed the common-ai-airflow-2 branch from b55be6a to 954c24f Compare October 1, 2026 11:37
kaxil added 2 commits October 1, 2026 14:06
Lower the provider's floor from Airflow 3.0 to 2.11, the floor the compat,
standard and common-sql providers carry. The operators, decorators, hooks and
toolsets run on 2.11 as they do on 3.0; features that need Airflow 3.1 or 3.3
(approval gates, HITL review, retry policies, tool approval, the task state
store) keep their existing gates.

- common.compat: export SET_DURING_EXECUTION, backed on Airflow 2 by an
  ArgNotSet that renders as Airflow 3's sentinel does, and resolve
  get_current_context through the standard provider first on Airflow 2, so it
  raises RuntimeError outside a task as Airflow 3 does.
- PydanticAIHook.get_hook accepts hook_params on Airflow 2, mirroring Airflow 3.
- The agent's per-attempt run key falls back to dag/run/task/map/try where the
  task instance has no id; spans then omit airflow.task_instance.id.
- Declare structlog, which the provider imports but Airflow 2 does not ship,
  and route its output through the airflow.task logger on Airflow 2.
- Declare the HITL review extra link only on Airflow 3.1+.
- Run the provider's tests in the Airflow 2.11 compatibility job.
- Rename task_instance_run_key to make_task_instance_run_key.
- Title the installation section "Airflow 2.11", so it does not read as
  covering every Airflow 2 release.
- Silence the Airflow 2-only ArgNotSet import for mypy on Airflow 3.
@kaxil
kaxil force-pushed the common-ai-airflow-2 branch from 954c24f to de79909 Compare October 1, 2026 13:13
@kaxil
kaxil merged commit 8f3e846 into apache:main Oct 1, 2026
156 of 157 checks passed
@kaxil
kaxil deleted the common-ai-airflow-2 branch October 1, 2026 14:10
@kaxil

kaxil commented Oct 1, 2026

Copy link
Copy Markdown
Member Author

Static check failure is unrelated

@github-actions

github-actions Bot commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

Backport failed to create: v3-3-test. View the failure log Run details

Note: As of Merging PRs targeted for Airflow 3.X
the committer who merges the PR is responsible for backporting the PRs that are bug fixes (generally speaking) to the maintenance branches.

In matter of doubt please ask in #release-management Slack channel.

Status Branch Result
❌ v3-3-test Commit Link

You can attempt to backport this manually by running:

cherry_picker 8f3e846 v3-3-test

This should apply the commit to the v3-3-test branch and leave the commit in conflict state marking
the files that need manual conflict resolution.

After you have resolved the conflicts, you can continue the backport process by running:

cherry_picker --continue

If you don't have cherry-picker installed, see the installation guide.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants