Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion dev/breeze/src/airflow_breeze/global_constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -862,7 +862,7 @@ def get_airflow_extras():
{
"python-version": "3.10",
"airflow-version": "2.11.1",
"remove-providers": "anthropic common.messaging common.dataquality edge3 fab git keycloak informatica common.ai modal opensearch",
"remove-providers": "anthropic common.messaging common.dataquality edge3 fab git keycloak informatica modal opensearch",
"run-unit-tests": "true",
},
{
Expand Down
3 changes: 2 additions & 1 deletion providers/common/ai/README.rst
Original file line number Diff line number Diff line change
Expand Up @@ -53,10 +53,11 @@ Requirements
========================================== ==================
PIP package Version required
========================================== ==================
``apache-airflow`` ``>=3.0.0``
``apache-airflow`` ``>=2.11.0``
``apache-airflow-providers-common-compat`` ``>=1.15.0``
``apache-airflow-providers-standard`` ``>=1.20.0``
``pydantic-ai-slim`` ``>=2.33.0``
``structlog`` ``>=24.2.0``
========================================== ==================

Optional cross provider package dependencies
Expand Down
5 changes: 3 additions & 2 deletions providers/common/ai/docs/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -172,15 +172,16 @@ For the minimum Airflow version supported, see ``Requirements`` below.
Requirements
------------

The minimum Apache Airflow version supported by this provider distribution is ``3.0.0``.
The minimum Apache Airflow version supported by this provider distribution is ``2.11.0``.

========================================== ==================
PIP package Version required
========================================== ==================
``apache-airflow`` ``>=3.0.0``
``apache-airflow`` ``>=2.11.0``
``apache-airflow-providers-common-compat`` ``>=1.15.0``
``apache-airflow-providers-standard`` ``>=1.20.0``
``pydantic-ai-slim`` ``>=2.33.0``
``structlog`` ``>=24.2.0``
========================================== ==================

Optional cross provider package dependencies
Expand Down
33 changes: 31 additions & 2 deletions providers/common/ai/docs/installation.rst
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
Installation
============

The provider needs Airflow 3.0 or later. Install it with the extra that matches the
The provider needs Airflow 2.11 or later. Install it with the extra that matches the
model vendor your connection will point at:

.. code-block:: bash
Expand Down Expand Up @@ -62,21 +62,50 @@ package each extra installs.
Features gated on the Airflow version
-------------------------------------

The provider runs on Airflow 3.0, but some features need a newer core:
The provider runs on Airflow 2.11, but some features need a newer core:

.. list-table::
:header-rows: 1
:widths: 60 40

* - Feature
- Needs
* - The ``skills`` and ``git`` extras (``apache-airflow-providers-git`` needs Airflow 3)
- Airflow 3.0
* - :doc:`Approval gates <approval_gates>` and :doc:`HITL review <hitl_review>`
- Airflow 3.1
* - The **Model** field in the connection form; on older cores put the model in
**Extra**, for example ``{"model": "openai:gpt-5"}``
- Airflow 3.2
* - :doc:`Retry policies <retry_policies>`
- Airflow 3.3
* - :doc:`Durable execution <durable_execution>` without configuring
``[common.ai] durable_cache_path`` (the task state store)
- Airflow 3.3
* - :doc:`Tool approval <tool_approval>` that pauses the task; on older cores a tool
marked for approval fails the task
- Airflow 3.3
* - A :doc:`structured output <structured_output>` reaching downstream tasks as the
Pydantic model; on older cores it arrives as a ``dict``
- Airflow 3.3

Airflow 2.11
------------

On Airflow 2.11 the operators, decorators, hooks and toolsets run as they do on Airflow
3.0, apart from the table above. Three things differ from an Airflow 3 install:

* The examples in these docs import ``dag``, ``task`` and ``Param`` from ``airflow.sdk``.
On Airflow 2 import ``dag`` and ``task`` from ``airflow.decorators`` and ``Param`` from
``airflow.models.param``; the provider's own imports stay the same.
* Install Airflow with its constraints file as usual, then add the provider without it. The
Airflow 2.11 constraints pin ``apache-airflow-providers-common-compat`` and
``apache-airflow-providers-common-sql`` to releases older than this provider needs.
Installing Airflow 2.11.0 without its constraints can also pull in a ``universal-pathlib``
0.3 release, which Airflow 2's ``ObjectStoragePath`` rejects; 2.11.1 and later cap it.
Leave out the ``skills`` and ``git`` extras: they need Airflow 3, and without
constraints ``pip`` upgrades Airflow to satisfy them.
* Python 3.10 to 3.12: the provider needs 3.10 or later, and Airflow 2.11 supports up to 3.12.

Next steps
----------
Expand Down
3 changes: 3 additions & 0 deletions providers/common/ai/docs/observability.rst
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,9 @@ How it works
tool approval (see :doc:`tool_approval`) continues as
``<task-instance id>-resumed``, which is the ``run_id`` the operator pushes;
``usage`` covers both parts.
Airflow 2 has no task-instance id, so there the key is
``<dag_id>/<run_id>/<task_id>/<map_index>/<try_number>``, and spans carry the five
identity keys without ``airflow.task_instance.id``.
* **Scope.** The ``run_id`` / ``usage`` XComs come only from ``AgentOperator`` and
``@task.agent``, and so do the ``airflow.*`` identity attributes, apart from a Strands or
ADK agent run inside ``agent_framework_tracing`` (see below). The other LLM
Expand Down
4 changes: 2 additions & 2 deletions providers/common/ai/docs/operators/llm_batch.rst
Original file line number Diff line number Diff line change
Expand Up @@ -282,8 +282,8 @@ provider's own batch listing first.
``cancel_on_kill`` cancels the batch if the task is killed. In deferrable mode this runs from the
trigger's ``on_kill``, which only **Airflow 3.3+** calls; on those versions clearing, marking
success or marking failed on a deferred task from the UI counts as a kill, so the batch is
cancelled and the next attempt submits a fresh one rather than re-attaching. On Airflow 3.0 to
3.2 a killed deferred task's batch keeps running and a clear re-attaches to it. Set
cancelled and the next attempt submits a fresh one rather than re-attaching. Before
Airflow 3.3 a killed deferred task's batch keeps running and a clear re-attaches to it. Set
``cancel_on_kill=False`` if you want clear-to-re-attach on 3.3+ as well.

``cancel_on_timeout=False`` lets a batch keep running (and billing) past this task's own
Expand Down
3 changes: 2 additions & 1 deletion providers/common/ai/docs/quickstart.rst
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,8 @@ which one task asks a model to summarize release notes and a second task uses th
At the end you know where the model's output lands and what a successful run looks like.

You need a working :doc:`Airflow installation <apache-airflow:installation/index>` on
Airflow 3.0 or later and an API key for the model vendor you plan to use. Step 4 makes one
Airflow 2.11 or later and an API key for the model vendor you plan to use. On Airflow 2,
see :ref:`howto/installation` for what differs. Step 4 makes one
real, billed API call.

1. Install the provider
Expand Down
3 changes: 2 additions & 1 deletion providers/common/ai/docs/self_hosted_models.rst
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,8 @@ Before you start
------------------

This guide assumes a working :doc:`apache-airflow:installation/index`
(Airflow 3.0+) already exists. Its job stops at wiring Airflow to a server
(Airflow 2.11+; on Airflow 2 see :ref:`howto/installation` for
what differs) already exists. Its job stops at wiring Airflow to a server
that's already running -- it doesn't cover installing or operating the
model-serving stack itself.

Expand Down
7 changes: 5 additions & 2 deletions providers/common/ai/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -67,14 +67,17 @@ requires-python = ">=3.10"
# Make sure to run ``prek update-providers-dependencies --all-files``
# After you modify the dependencies, and rebuild your Breeze CI image with ``breeze ci-image build``
dependencies = [
"apache-airflow>=3.0.0",
"apache-airflow-providers-common-compat>=1.15.0",
"apache-airflow>=2.11.0",
"apache-airflow-providers-common-compat>=1.15.0", # use next version
Comment thread
kaxil marked this conversation as resolved.
"apache-airflow-providers-standard>=1.20.0",
# 2.33.0 is the first release that works with anthropic>=1: it moved to httpx2 alongside
# the SDK and stopped passing temperature/top_p/top_k as messages.create() kwargs, both
# of which raise TypeError on 2.31 and earlier. The cost API this provider relies on
# (RunUsage.cost, UsageLimits.cost_limit) landed earlier, in 2.23.0.
"pydantic-ai-slim>=2.33.0",
# Airflow 3 brings structlog in through the Task SDK; Airflow 2 does not. 24.2.0 is the first
# release whose render_to_log_kwargs hands ``stacklevel`` to stdlib logging under that name.
"structlog>=24.2.0",
]

# The optional dependencies should be modified in place in the generated file
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,8 @@
__version__ = "0.10.0"

if packaging.version.parse(packaging.version.parse(airflow_version).base_version) < packaging.version.parse(
"3.0.0"
"2.11.0"
):
raise RuntimeError(
f"The package `apache-airflow-providers-common-ai:{__version__}` needs Apache Airflow 3.0.0+"
f"The package `apache-airflow-providers-common-ai:{__version__}` needs Apache Airflow 2.11.0+"
)
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,6 @@
from collections.abc import Iterator
from typing import TYPE_CHECKING, Any

import structlog

from airflow.providers.common.ai.batch.base import (
BatchAdapter,
BatchState,
Expand All @@ -43,9 +41,10 @@
SubmitResult,
)
from airflow.providers.common.ai.exceptions import LLMBatchLimitExceededError
from airflow.providers.common.ai.utils.task_logger import get_task_logger
from airflow.providers.common.compat.sdk import AirflowOptionalProviderFeatureException

log = structlog.get_logger(logger_name="task")
log = get_task_logger()

if TYPE_CHECKING:
from airflow.providers.common.ai.batch.base import BatchRequest
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,6 @@
from datetime import datetime
from typing import TYPE_CHECKING, Any

import structlog

from airflow.providers.common.ai.batch.base import (
BatchAdapter,
BatchState,
Expand All @@ -43,9 +41,10 @@
SubmitResult,
)
from airflow.providers.common.ai.exceptions import LLMBatchLimitExceededError, LLMBatchModelMismatchError
from airflow.providers.common.ai.utils.task_logger import get_task_logger
from airflow.providers.common.compat.sdk import AirflowOptionalProviderFeatureException

log = structlog.get_logger(logger_name="task")
log = get_task_logger()

if TYPE_CHECKING:
from airflow.providers.common.ai.batch.base import BatchRequest
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,17 +34,16 @@
from dataclasses import dataclass
from typing import TYPE_CHECKING, Any

import structlog

from airflow.providers.common.ai.batch.base import evaluate_batch_counts
from airflow.providers.common.ai.batch.output_schema import validate_extracted_output
from airflow.providers.common.ai.utils.task_logger import get_task_logger

if TYPE_CHECKING:
from airflow.providers.common.ai.batch.base import BatchAdapter, RawResultItem
from airflow.providers.common.ai.batch.output_schema import OutputSpec
from airflow.sdk import ObjectStoragePath

log = structlog.get_logger(logger_name="task")
log = get_task_logger()

STATUS_SUCCESS = "success"
STATUS_ERROR = "error"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,13 +33,13 @@
validate_prompt,
)
from airflow.providers.common.compat.sdk import (
SET_DURING_EXECUTION,
DecoratedOperator,
TaskDecorator,
context_merge,
determine_kwargs,
task_decorator_factory,
)
from airflow.sdk.definitions._internal.types import SET_DURING_EXECUTION

if TYPE_CHECKING:
from airflow.sdk import Context
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,13 +35,13 @@
validate_prompt,
)
from airflow.providers.common.compat.sdk import (
SET_DURING_EXECUTION,
DecoratedOperator,
TaskDecorator,
context_merge,
determine_kwargs,
task_decorator_factory,
)
from airflow.sdk.definitions._internal.types import SET_DURING_EXECUTION

if TYPE_CHECKING:
from airflow.sdk import Context
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,13 +29,13 @@

from airflow.providers.common.ai.operators.llm_batch import LLMBatchOperator
from airflow.providers.common.compat.sdk import (
SET_DURING_EXECUTION,
DecoratedOperator,
TaskDecorator,
context_merge,
determine_kwargs,
task_decorator_factory,
)
from airflow.sdk.definitions._internal.types import SET_DURING_EXECUTION

if TYPE_CHECKING:
from airflow.sdk import Context
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,13 +33,13 @@
validate_prompt,
)
from airflow.providers.common.compat.sdk import (
SET_DURING_EXECUTION,
DecoratedOperator,
TaskDecorator,
context_merge,
determine_kwargs,
task_decorator_factory,
)
from airflow.sdk.definitions._internal.types import SET_DURING_EXECUTION

if TYPE_CHECKING:
from airflow.sdk import Context
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,13 +23,13 @@

from airflow.providers.common.ai.operators.llm_file_analysis import LLMFileAnalysisOperator
from airflow.providers.common.compat.sdk import (
SET_DURING_EXECUTION,
DecoratedOperator,
TaskDecorator,
context_merge,
determine_kwargs,
task_decorator_factory,
)
from airflow.sdk.definitions._internal.types import SET_DURING_EXECUTION

if TYPE_CHECKING:
from airflow.sdk import Context
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,13 +33,13 @@
validate_prompt,
)
from airflow.providers.common.compat.sdk import (
SET_DURING_EXECUTION,
DecoratedOperator,
TaskDecorator,
context_merge,
determine_kwargs,
task_decorator_factory,
)
from airflow.sdk.definitions._internal.types import SET_DURING_EXECUTION

if TYPE_CHECKING:
from airflow.sdk import Context
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,13 +33,13 @@
validate_prompt,
)
from airflow.providers.common.compat.sdk import (
SET_DURING_EXECUTION,
DecoratedOperator,
TaskDecorator,
context_merge,
determine_kwargs,
task_decorator_factory,
)
from airflow.sdk.definitions._internal.types import SET_DURING_EXECUTION

if TYPE_CHECKING:
from airflow.sdk import Context
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,14 +21,14 @@
from dataclasses import dataclass, field
from typing import TYPE_CHECKING, Any

import structlog
from pydantic_ai.messages import ModelResponse, ToolCallPart
from pydantic_ai.models.wrapper import WrapperModel

from airflow.providers.common.ai.durable.base import build_model_step_key, build_tool_step_key
from airflow.providers.common.ai.durable.fingerprint import fingerprint_model_request
from airflow.providers.common.ai.utils.task_logger import get_task_logger

log = structlog.get_logger(logger_name="task")
log = get_task_logger()

if TYPE_CHECKING:
from pydantic_ai.messages import ModelMessage
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,11 @@
from dataclasses import dataclass, field
from typing import TYPE_CHECKING, Any

import structlog
from pydantic_ai.toolsets.wrapper import WrapperToolset

from airflow.providers.common.ai.durable.base import build_tool_step_key
from airflow.providers.common.ai.durable.fingerprint import fingerprint_tool_call
from airflow.providers.common.ai.utils.task_logger import get_task_logger
from airflow.providers.common.ai.utils.tool_metrics import record_tool_call
from airflow.providers.common.ai.utils.toolset_base import AirflowToolset

Expand All @@ -36,7 +36,7 @@
from airflow.providers.common.ai.durable.replay_usage import ReplayUsageLedger
from airflow.providers.common.ai.durable.step_counter import DurableStepCounter

log = structlog.get_logger(logger_name="task")
log = get_task_logger()


@dataclass
Expand Down
Loading
Loading