This is an automated email from the ASF dual-hosted git repository.
josh-fell pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 1ac02ade890 Add validate trigger events in dbt cloud deferrable tasks
(#69739)
1ac02ade890 is described below
commit 1ac02ade890cee17dd8ad6c193c5ca65c0a181da
Author: Takayoshi Makabe <[email protected]>
AuthorDate: Fri Sep 25 00:46:48 2026 +0900
Add validate trigger events in dbt cloud deferrable tasks (#69739)
* Add validate tirgger in deferrable tasks
* Update known_airflow_exceptions.txt for dbt cloud sensor exception change
* Revert unrelated count change in known_airflow_exceptions.txt
---
generated/known_airflow_exceptions.txt | 1 -
.../src/airflow/providers/dbt/cloud/hooks/dbt.py | 26 +++++++++++++++
.../airflow/providers/dbt/cloud/operators/dbt.py | 4 ++-
.../src/airflow/providers/dbt/cloud/sensors/dbt.py | 16 +++++++---
.../cloud/tests/unit/dbt/cloud/hooks/test_dbt.py | 30 ++++++++++++++++++
.../tests/unit/dbt/cloud/operators/test_dbt.py | 27 +++++++++++++++-
.../cloud/tests/unit/dbt/cloud/sensors/test_dbt.py | 37 +++++++++++++++++++---
7 files changed, 129 insertions(+), 12 deletions(-)
diff --git a/generated/known_airflow_exceptions.txt
b/generated/known_airflow_exceptions.txt
index 29512c0c613..cb860b5cd71 100644
--- a/generated/known_airflow_exceptions.txt
+++ b/generated/known_airflow_exceptions.txt
@@ -186,7 +186,6 @@
providers/databricks/tests/unit/databricks/hooks/test_databricks_base.py::1
providers/datadog/src/airflow/providers/datadog/hooks/datadog.py::2
providers/datadog/src/airflow/providers/datadog/sensors/datadog.py::1
providers/dbt/cloud/src/airflow/providers/dbt/cloud/hooks/dbt.py::3
-providers/dbt/cloud/src/airflow/providers/dbt/cloud/sensors/dbt.py::1
providers/dingding/src/airflow/providers/dingding/hooks/dingding.py::2
providers/discord/src/airflow/providers/discord/operators/discord_webhook.py::1
providers/docker/src/airflow/providers/docker/decorators/docker.py::1
diff --git a/providers/dbt/cloud/src/airflow/providers/dbt/cloud/hooks/dbt.py
b/providers/dbt/cloud/src/airflow/providers/dbt/cloud/hooks/dbt.py
index 40129cf41b4..b0d64de7b0c 100644
--- a/providers/dbt/cloud/src/airflow/providers/dbt/cloud/hooks/dbt.py
+++ b/providers/dbt/cloud/src/airflow/providers/dbt/cloud/hooks/dbt.py
@@ -138,6 +138,32 @@ class DbtCloudResourceLookupError(AirflowException):
"""Exception raised when a dbt Cloud resource cannot be uniquely
identified."""
+class DbtCloudTriggerEventException(AirflowException):
+ """Raised when a deferred task resumes with a missing or malformed trigger
event."""
+
+
+#: Statuses the provider's trigger emits in its terminal event.
+TRIGGER_EVENT_STATUSES = frozenset({"success", "cancelled", "error",
"timeout"})
+
+
+def validate_execute_complete_event(event: dict[str, Any] | None = None) ->
dict[str, Any]:
+ """
+ Validate the event a deferred task resumes with, returning it if
well-formed.
+
+ The event crosses the triggerer/worker boundary through the metadata DB,
so a
+ resuming task can receive ``None`` (e.g. a lost payload) or a status its
handlers
+ do not recognize (triggerer/worker version skew, a custom trigger). Both
must fail
+ loudly rather than silently falling through to the success path.
+ """
+ if event is None:
+ raise DbtCloudTriggerEventException("Trigger error: event is None")
+ if event.get("status") not in TRIGGER_EVENT_STATUSES:
+ raise DbtCloudTriggerEventException(
+ f"Unexpected trigger event status {event.get('status')!r}:
{event!r}"
+ )
+ return event
+
+
T = TypeVar("T", bound=Any)
diff --git
a/providers/dbt/cloud/src/airflow/providers/dbt/cloud/operators/dbt.py
b/providers/dbt/cloud/src/airflow/providers/dbt/cloud/operators/dbt.py
index dbd92140b69..0eb178d4d5e 100644
--- a/providers/dbt/cloud/src/airflow/providers/dbt/cloud/operators/dbt.py
+++ b/providers/dbt/cloud/src/airflow/providers/dbt/cloud/operators/dbt.py
@@ -29,6 +29,7 @@ from airflow.providers.dbt.cloud.hooks.dbt import (
DbtCloudJobRunException,
DbtCloudJobRunStatus,
JobRunInfo,
+ validate_execute_complete_event,
)
from airflow.providers.dbt.cloud.triggers.dbt import DbtCloudRunJobTrigger
from airflow.providers.dbt.cloud.utils.openlineage import (
@@ -304,8 +305,9 @@ class DbtCloudRunJobOperator(BaseOperator):
)
return self.run_id
- def execute_complete(self, context: Context, event: dict[str, Any]) -> int:
+ def execute_complete(self, context: Context, event: dict[str, Any] | None
= None) -> int:
"""Execute when the trigger fires - returns immediately."""
+ event = validate_execute_complete_event(event)
self.run_id = event["run_id"]
if event["status"] == "cancelled":
self.hook.log_job_run_failure_details(run_id=int(self.run_id),
account_id=self.account_id)
diff --git a/providers/dbt/cloud/src/airflow/providers/dbt/cloud/sensors/dbt.py
b/providers/dbt/cloud/src/airflow/providers/dbt/cloud/sensors/dbt.py
index 2b40dd8724f..1368727847d 100644
--- a/providers/dbt/cloud/src/airflow/providers/dbt/cloud/sensors/dbt.py
+++ b/providers/dbt/cloud/src/airflow/providers/dbt/cloud/sensors/dbt.py
@@ -20,8 +20,13 @@ import time
from functools import cached_property
from typing import TYPE_CHECKING, Any
-from airflow.providers.common.compat.sdk import AirflowException,
BaseSensorOperator, conf
-from airflow.providers.dbt.cloud.hooks.dbt import DbtCloudHook,
DbtCloudJobRunException, DbtCloudJobRunStatus
+from airflow.providers.common.compat.sdk import BaseSensorOperator, conf
+from airflow.providers.dbt.cloud.hooks.dbt import (
+ DbtCloudHook,
+ DbtCloudJobRunException,
+ DbtCloudJobRunStatus,
+ validate_execute_complete_event,
+)
from airflow.providers.dbt.cloud.triggers.dbt import DbtCloudRunJobTrigger
from airflow.providers.dbt.cloud.utils.openlineage import
generate_openlineage_events_from_dbt_cloud_run
@@ -116,15 +121,16 @@ class DbtCloudJobRunSensor(BaseSensorOperator):
method_name="execute_complete",
)
- def execute_complete(self, context: Context, event: dict[str, Any]) -> int:
+ def execute_complete(self, context: Context, event: dict[str, Any] | None
= None) -> int:
"""
Execute when the trigger fires - returns immediately.
This relies on trigger to throw an exception, otherwise it assumes
execution was successful.
"""
- if event["status"] in ["error", "cancelled"]:
- raise AirflowException()
+ event = validate_execute_complete_event(event)
+ if event["status"] != "success":
+ raise DbtCloudJobRunException(event["message"])
self.log.info(event["message"])
return int(event["run_id"])
diff --git a/providers/dbt/cloud/tests/unit/dbt/cloud/hooks/test_dbt.py
b/providers/dbt/cloud/tests/unit/dbt/cloud/hooks/test_dbt.py
index dc2e501e8a0..59fd1c953c7 100644
--- a/providers/dbt/cloud/tests/unit/dbt/cloud/hooks/test_dbt.py
+++ b/providers/dbt/cloud/tests/unit/dbt/cloud/hooks/test_dbt.py
@@ -36,8 +36,10 @@ from airflow.providers.dbt.cloud.hooks.dbt import (
DbtCloudJobRunException,
DbtCloudJobRunStatus,
DbtCloudResourceLookupError,
+ DbtCloudTriggerEventException,
TokenAuth,
fallback_to_default_account,
+ validate_execute_complete_event,
)
from tests_common.test_utils.compat import timezone
@@ -159,6 +161,34 @@ class TestDbtCloudJobRunStatus:
DbtCloudJobRunStatus.check_is_valid(statuses)
+class TestValidateExecuteCompleteEvent:
+ @pytest.mark.parametrize(
+ ("event", "match"),
+ [
+ pytest.param(None, "event is None", id="none"),
+ pytest.param({}, "Unexpected trigger event status None",
id="missing-status"),
+ pytest.param(
+ {"status": "ended", "run_id": 1234}, "Unexpected trigger event
status", id="unknown-status"
+ ),
+ ],
+ )
+ def test_invalid_event_raises(self, event, match):
+ with pytest.raises(DbtCloudTriggerEventException, match=match):
+ validate_execute_complete_event(event)
+
+ @pytest.mark.parametrize(
+ "event",
+ [
+ pytest.param({"status": "success", "run_id": 1234, "message":
"ok"}, id="success"),
+ pytest.param({"status": "cancelled", "run_id": 1234, "message":
"cancelled"}, id="cancelled"),
+ pytest.param({"status": "error", "run_id": 1234, "message":
"failed"}, id="error"),
+ pytest.param({"status": "timeout", "run_id": 1234, "message":
"timed out"}, id="timeout"),
+ ],
+ )
+ def test_valid_event_is_returned(self, event):
+ assert validate_execute_complete_event(event) is event
+
+
class TestDbtCloudHook:
# TODO: Potential performance issue, converted setup_class to a
setup_connections function level fixture
@pytest.fixture(autouse=True)
diff --git a/providers/dbt/cloud/tests/unit/dbt/cloud/operators/test_dbt.py
b/providers/dbt/cloud/tests/unit/dbt/cloud/operators/test_dbt.py
index 922f3ada714..be1ebec812e 100644
--- a/providers/dbt/cloud/tests/unit/dbt/cloud/operators/test_dbt.py
+++ b/providers/dbt/cloud/tests/unit/dbt/cloud/operators/test_dbt.py
@@ -24,7 +24,12 @@ import pytest
from airflow.models import DAG, Connection
from airflow.providers.common.compat.sdk import TaskDeferred, timezone
-from airflow.providers.dbt.cloud.hooks.dbt import DbtCloudHook,
DbtCloudJobRunException, DbtCloudJobRunStatus
+from airflow.providers.dbt.cloud.hooks.dbt import (
+ DbtCloudHook,
+ DbtCloudJobRunException,
+ DbtCloudJobRunStatus,
+ DbtCloudTriggerEventException,
+)
from airflow.providers.dbt.cloud.operators.dbt import (
DbtCloudGetJobRunArtifactOperator,
DbtCloudListJobRunsOperator,
@@ -376,6 +381,26 @@ class TestDbtCloudRunJobOperator:
run_id=RUN_ID,
)
+ @pytest.mark.parametrize(
+ "event",
+ [
+ pytest.param(None, id="none"),
+ pytest.param({"status": "ended", "run_id": RUN_ID, "message":
"m"}, id="unknown-status"),
+ ],
+ )
+ def test_execute_complete_invalid_event_raises_instead_of_succeeding(self,
event):
+ """Verify a missing or unrecognized trigger event fails loudly instead
of succeeding."""
+ operator = DbtCloudRunJobOperator(
+ task_id=TASK_ID,
+ dbt_cloud_conn_id=ACCOUNT_ID_CONN,
+ job_id=JOB_ID,
+ dag=self.dag,
+ deferrable=True,
+ )
+
+ with pytest.raises(DbtCloudTriggerEventException):
+ operator.execute_complete(context=self.mock_context, event=event)
+
@patch(
"airflow.providers.dbt.cloud.hooks.dbt.DbtCloudHook.get_job_by_name",
return_value=mock_response_json(DEFAULT_ACCOUNT_JOB_RESPONSE),
diff --git a/providers/dbt/cloud/tests/unit/dbt/cloud/sensors/test_dbt.py
b/providers/dbt/cloud/tests/unit/dbt/cloud/sensors/test_dbt.py
index 9435fad5d2b..81e513be659 100644
--- a/providers/dbt/cloud/tests/unit/dbt/cloud/sensors/test_dbt.py
+++ b/providers/dbt/cloud/tests/unit/dbt/cloud/sensors/test_dbt.py
@@ -22,8 +22,13 @@ from unittest.mock import patch
import pytest
from airflow.models.connection import Connection
-from airflow.providers.common.compat.sdk import AirflowException, TaskDeferred
-from airflow.providers.dbt.cloud.hooks.dbt import DbtCloudHook,
DbtCloudJobRunException, DbtCloudJobRunStatus
+from airflow.providers.common.compat.sdk import TaskDeferred
+from airflow.providers.dbt.cloud.hooks.dbt import (
+ DbtCloudHook,
+ DbtCloudJobRunException,
+ DbtCloudJobRunStatus,
+ DbtCloudTriggerEventException,
+)
from airflow.providers.dbt.cloud.sensors.dbt import DbtCloudJobRunSensor
from airflow.providers.dbt.cloud.triggers.dbt import DbtCloudRunJobTrigger
@@ -156,10 +161,11 @@ class TestDbtCloudJobRunSensor:
[
("cancelled", "Job run 1234 has been cancelled."),
("error", "Job run 1234 has failed."),
+ ("timeout", "Job run 1234 has timed out."),
],
)
def test_execute_complete_failure(self, mock_status, mock_message):
- """Assert execute_complete method to raise exception on the cancelled
and error status"""
+ """Assert execute_complete method to raise exception on the cancelled,
error, and timeout status"""
task = DbtCloudJobRunSensor(
dbt_cloud_conn_id=self.CONN_ID,
task_id=self.TASK_ID,
@@ -167,7 +173,30 @@ class TestDbtCloudJobRunSensor:
timeout=self.TIMEOUT,
deferrable=True,
)
- with pytest.raises(AirflowException):
+ with pytest.raises(DbtCloudJobRunException, match=mock_message):
task.execute_complete(
context={}, event={"status": mock_status, "message":
mock_message, "run_id": self.DBT_RUN_ID}
)
+
+ @pytest.mark.parametrize(
+ ("event", "match"),
+ [
+ pytest.param(None, "event is None", id="none"),
+ pytest.param(
+ {"status": "rescheduling", "message": "m", "run_id": 1234},
+ "Unexpected trigger event status",
+ id="unknown-status",
+ ),
+ ],
+ )
+ def test_execute_complete_invalid_event_raises(self, event, match):
+ """Assert execute_complete raises instead of silently succeeding on a
malformed event"""
+ task = DbtCloudJobRunSensor(
+ dbt_cloud_conn_id=self.CONN_ID,
+ task_id=self.TASK_ID,
+ run_id=self.DBT_RUN_ID,
+ timeout=self.TIMEOUT,
+ deferrable=True,
+ )
+ with pytest.raises(DbtCloudTriggerEventException, match=match):
+ task.execute_complete(context={}, event=event)