Skip to content
38 changes: 20 additions & 18 deletions .github/workflows/ci.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -28,31 +28,33 @@ jobs:
fail-fast: false
matrix:
python-version:
# TODO: Requires dbt 1.12 which is in beta at the time of check.
# - '3.14'
- '3.14'
- '3.13'
- '3.12'
- '3.11'
- '3.10'
airflow-version:
# Latest release as of 2026-05-24
# Latest release as of 2026-09-07
# See: https://airflow.apache.org/docs/apache-airflow/stable/release_notes.html
- '3.2.1'
# GCP Cloud Composer latest as of 2026-05-24
- '3.3.1'
# Previous minor, kept for backward-compat coverage.
- '3.2.2'
# GCP Cloud Composer latest and AWS MWAA latest both also match the
# latest release (3.3.1) as of 2026-09-07, so no separate entries needed.
# See: https://docs.cloud.google.com/composer/docs/composer-versions
- '3.1.7'
# AWS MWAA latest as of 2026-05-24
# See: https://docs.aws.amazon.com/mwaa/latest/userguide/airflow-versions.html
# TODO: Uncomment once latest stable no longer matches MWAA latest. At the
# time of check, they are the same.
# - '3.2.1'
dbt-version:
- '1.11.11'
- '1.10.17'
- '1.12.3'
- '1.11.14'
- '1.10.23'
exclude:
# Airflow added 3.14 support in >=3.2
- airflow-version: '3.1.7'
python-version: '3.14'
# dbt-core <1.12 pins mashumaro<3.15, which lacks Python 3.14 support
# (only added in mashumaro 3.17). See mashumaro release notes and
# https://github.com/dbt-labs/dbt-core/issues/12098.
- python-version: '3.14'
dbt-version: '1.11.14'
- python-version: '3.14'
dbt-version: '1.10.23'

runs-on: ubuntu-latest
steps:
Expand Down Expand Up @@ -84,7 +86,7 @@ jobs:

- run: |
sudo apt-get update
sudo apt-get install --yes --no-install-recommends postgresql
sudo apt-get install --yes --no-install-recommends postgresql unixodbc-dev

- name: Checkout
uses: actions/checkout@v6
Expand All @@ -111,8 +113,8 @@ jobs:
run: uv run ruff check .

- name: Static type checking with mypy
# We only run mypy on the latest supported versions of Airflow & dbt,
if: matrix.python-version == '3.13' && matrix.airflow-version == '3.1.1' && matrix.dbt-version == '1.11.12'
# We only run mypy on the latest supported versions of Python, Airflow, and dbt.
if: matrix.python-version == '3.14' && matrix.airflow-version == '3.3.1' && matrix.dbt-version == '1.12.3'
run: uv run mypy .

- name: Code formatting with ruff
Expand Down
4 changes: 3 additions & 1 deletion airflow_dbt_python/hooks/target.py
Original file line number Diff line number Diff line change
Expand Up @@ -211,7 +211,9 @@ def get_db_conn_hook(cls, conn_id: str) -> DbtConnectionHook:
"""Get a dbt hook class depend on Airflow connection type."""
conn = cls.get_connection(conn_id)

if hook_cls := cls._dbt_hooks_by_conn_type.get(conn.conn_type):
if conn.conn_type is not None and (
hook_cls := cls._dbt_hooks_by_conn_type.get(conn.conn_type)
):
return hook_cls(conn=conn)
raise KeyError(
f"There are no DbtConnectionHook subclasses with conn_type={conn.conn_type}"
Expand Down
18 changes: 15 additions & 3 deletions airflow_dbt_python/utils/configs.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
from airflow_dbt_python.utils.version import (
DBT_INSTALLED_GTE_1_9,
DBT_INSTALLED_GTE_1_10_7,
DBT_INSTALLED_GTE_1_12,
)


Expand Down Expand Up @@ -192,6 +193,7 @@ class BaseConfig:
require_nested_cumulative_type_params: Optional[bool] = None
require_ref_searches_node_package_before_root: Optional[bool] = None
require_resource_names_without_spaces: Optional[bool] = None
require_source_and_semantic_model_names_without_spaces: Optional[bool] = None
require_unique_project_resource_names: Optional[bool] = None
require_valid_schema_from_generate_schema_name: Optional[bool] = None
require_yaml_configuration_for_mf_time_spines: Optional[bool] = None
Expand Down Expand Up @@ -435,9 +437,19 @@ def create_dbt_task(
# TODO: Support for catalog integrations
active_integrations=[], # type: ignore
)
task = self.dbt_task(
args=local_flags, config=runtime_config, manifest=manifest
)
if DBT_INSTALLED_GTE_1_12 and issubclass(self.dbt_task, FreshnessTask):
# dbt-core>=1.12 made FreshnessTask.__init__ require catalogs, unlike
# its sibling ConfiguredTask subclasses, which default it to None.
task = self.dbt_task(
args=local_flags,
config=runtime_config,
manifest=manifest,
catalogs=[],
)
else:
task = self.dbt_task(
args=local_flags, config=runtime_config, manifest=manifest
)
elif issubclass(self.dbt_task, DepsTask):
task = self.dbt_task(args=local_flags, project=project)
elif issubclass(self.dbt_task, DebugTask):
Expand Down
30 changes: 9 additions & 21 deletions airflow_dbt_python/utils/version.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,31 +3,18 @@
These are only used to ensure backwards compatibility with older versions of dbt.
"""

try:
from dbt.semver import Matchers, VersionSpecifier
except ImportError:
from dbt_common.semver import Matchers, VersionSpecifier
from importlib.metadata import version

from dbt.version import installed
from packaging.version import Version

DBT_1_8 = VersionSpecifier(
major="1", minor="8", patch="0", matcher=Matchers.GREATER_THAN_OR_EQUAL
)
DBT_1_9 = VersionSpecifier(
major="1", minor="9", patch="0", matcher=Matchers.GREATER_THAN_OR_EQUAL
)
DBT_1_10_7 = VersionSpecifier(
major="1", minor="10", patch="7", matcher=Matchers.GREATER_THAN_OR_EQUAL
)
DBT_2_0 = VersionSpecifier(
major="2", minor="0", patch="0", matcher=Matchers.GREATER_THAN_OR_EQUAL
)
DBT_VERSION = Version(version("dbt-core"))

DBT_INSTALLED_GTE_1_9 = installed.compare(DBT_1_9) == 1
DBT_INSTALLED_GTE_1_10_7 = installed.compare(DBT_1_10_7) == 1
DBT_INSTALLED_GTE_1_9 = DBT_VERSION >= Version("1.9.0")
DBT_INSTALLED_GTE_1_10_7 = DBT_VERSION >= Version("1.10.7")
DBT_INSTALLED_GTE_1_12 = DBT_VERSION >= Version("1.12.0")

DBT_INSTALLED_1_8 = DBT_1_8 < installed < DBT_1_9
DBT_INSTALLED_1_9 = DBT_1_9 < installed < DBT_2_0
DBT_INSTALLED_1_8 = Version("1.8.0") <= DBT_VERSION < Version("1.9.0")
DBT_INSTALLED_1_9 = Version("1.9.0") <= DBT_VERSION < Version("2.0.0")


def _get_base_airflow_version_tuple() -> tuple[int, int, int]:
Expand All @@ -40,4 +27,5 @@ def _get_base_airflow_version_tuple() -> tuple[int, int, int]:

AIRFLOW_V_3_0_PLUS = _get_base_airflow_version_tuple() >= (3, 0, 0)
AIRFLOW_V_3_1_PLUS = _get_base_airflow_version_tuple() >= (3, 1, 0)
AIRFLOW_V_3_3_PLUS = _get_base_airflow_version_tuple() >= (3, 3, 0)
AIRFLOW_V_3_0 = AIRFLOW_V_3_0_PLUS and not AIRFLOW_V_3_1_PLUS
5 changes: 2 additions & 3 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ description = "A collection of Airflow operators, hooks, and utilities to execut
authors = [{ name = "Tomás Farías Santana", email = "tomas@tomasfarias.dev" }]
license = "MIT"
readme = "README.md"
requires-python = ">=3.10,<3.14"
requires-python = ">=3.10,<3.15"
classifiers = [
"Development Status :: 5 - Production/Stable",

Expand All @@ -24,6 +24,7 @@ dependencies = [
"apache-airflow>=3.1,<4.0; python_version>='3.13'",
"contextlib-chdir==1.0.2;python_version<'3.11'",
"dbt-core>=1.8.0,<2.0.0",
"packaging>=20.0",
]

[project.urls]
Expand Down Expand Up @@ -90,8 +91,6 @@ dev = [
"pytest-mock>=3.14.0",
"pytest-postgresql>=5",
"ruff>=0.0.254",
# Pinned as Airflow crashes: https://github.com/apache/airflow/issues/57426.
"structlog<=25.4",
"types-PyYAML>=6.0.7",
"types-freezegun>=1.1.6",
]
Expand Down
66 changes: 66 additions & 0 deletions tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,72 @@
"""


@pytest.fixture(scope="session", autouse=True)
def _fix_in_process_execution_api_lifespan_race():
"""Work around a lifespan race in Airflow's ``InProcessExecutionAPI``.

Before apache/airflow#68840 (fixed in Airflow 3.3+), ``InProcessExecutionAPI
.transport`` schedules the FastAPI app's lifespan startup via
``asyncio.run_coroutine_threadsafe`` without waiting for it, so a
``dag.test()`` task run can hit the transport before ``app.state
.svcs_registry`` is set, raising ``AttributeError: 'State' object has no
attribute 'svcs_registry'``.

``InProcessExecutionAPI`` is an ``@attrs.define()`` (slotted) class where
attrs itself manages caching for ``transport`` via a generated
``__getattr__``/setter pair, so we can't simply monkeypatch the
``transport`` property (attrs' setter breaks on a plain replacement). We
instead build our own transport the same way the buggy property does, but
waiting for lifespan startup to actually finish, and hand it out via a
plain (non-attrs) stand-in object patched into
``in_process_api_server()``.
"""
try:
import airflow.sdk.execution_time.supervisor as supervisor_module
from airflow.api_fastapi.execution_api.app import InProcessExecutionAPI
except ImportError:
yield
return

if not hasattr(supervisor_module, "in_process_api_server"):
yield
return

if "transport" not in InProcessExecutionAPI.__dict__:
# Already fixed upstream (owns its own loop/thread; no race to work
# around) or a version this workaround doesn't otherwise apply to.
yield
return

import asyncio
from contextlib import AsyncExitStack
from types import SimpleNamespace

import httpx
from a2wsgi import ASGIMiddleware

app = InProcessExecutionAPI().app
# Same construction as the buggy property (own event loop created by
# ASGIMiddleware internally), just actually waiting for the result.
middleware = ASGIMiddleware(app)

async def start_lifespan(cm: AsyncExitStack) -> None:
await cm.enter_async_context(app.router.lifespan_context(app))

cm = AsyncExitStack()
asyncio.run_coroutine_threadsafe(start_lifespan(cm), middleware.loop).result()

warm_api = SimpleNamespace(app=app, transport=httpx.WSGITransport(app=middleware))

original_in_process_api_server = supervisor_module.in_process_api_server
supervisor_module.in_process_api_server = lambda: warm_api

yield

supervisor_module.in_process_api_server = original_in_process_api_server
asyncio.run_coroutine_threadsafe(cm.aclose(), middleware.loop).result(timeout=5)


@pytest.fixture(scope="session")
def database(postgresql_proc):
"""Initialize a test postgres database."""
Expand Down
10 changes: 8 additions & 2 deletions tests/dags/test_dbt_dags.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
AIRFLOW_V_3_0,
AIRFLOW_V_3_0_PLUS,
AIRFLOW_V_3_1_PLUS,
AIRFLOW_V_3_3_PLUS,
)

if AIRFLOW_V_3_0:
Expand All @@ -39,7 +40,7 @@
try:
from airflow.serialization.serialized_objects import DagSerialization
except ImportError:
DagSerialization = SerializedDAG
DagSerialization = SerializedDAG # type: ignore

DATA_INTERVAL_START = pendulum.datetime(2022, 1, 1, tz="UTC")
DATA_INTERVAL_END = DATA_INTERVAL_START + dt.timedelta(hours=1)
Expand Down Expand Up @@ -130,7 +131,12 @@ def _run_task_instance(ti):
@pytest.fixture(scope="session")
def dagbag():
"""An Airflow DagBag."""
dagbag = DagBag(dag_folder="examples/", include_examples=False)
if AIRFLOW_V_3_3_PLUS:
# DagBag dropped include_examples in Airflow 3.3: it no longer
# auto-discovers example DAGs into a plain folder scan like this.
dagbag = DagBag(dag_folder="examples/")
else:
dagbag = DagBag(dag_folder="examples/", include_examples=False)

return dagbag

Expand Down
Loading
Loading