This is an automated email from the ASF dual-hosted git repository.
potiuk 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 00bc44480c7 Apply impersonation_chain to deferred BigQuery existence
checks (#71648)
00bc44480c7 is described below
commit 00bc44480c7f00e781e5bac8227509c98578483a
Author: Sepuri Sai Krishna <[email protected]>
AuthorDate: Tue Oct 6 04:05:29 2026 +0530
Apply impersonation_chain to deferred BigQuery existence checks (#71648)
* Apply impersonation_chain to deferred BigQuery existence checks
The sensors handed the impersonation chain to the trigger inside
hook_params,
which nothing reads, while the trigger authenticated from an attribute those
call sites never set. A deferred existence check therefore ran as the
connection's service account, and only in deferrable mode, so disabling
deferral appeared to fix the resulting permission error.
* Move the BigQuery impersonation tests beside their sensors
A reviewer asked for the coverage to live with each sensor's own tests
rather than in a standalone parametrized test, so the file keeps one place to
look per sensor.
---
.../providers/google/cloud/sensors/bigquery.py | 2 ++
.../unit/google/cloud/sensors/test_bigquery.py | 37 ++++++++++++++++++++++
2 files changed, 39 insertions(+)
diff --git
a/providers/google/src/airflow/providers/google/cloud/sensors/bigquery.py
b/providers/google/src/airflow/providers/google/cloud/sensors/bigquery.py
index 29f77ffd27f..65531b857bb 100644
--- a/providers/google/src/airflow/providers/google/cloud/sensors/bigquery.py
+++ b/providers/google/src/airflow/providers/google/cloud/sensors/bigquery.py
@@ -129,6 +129,7 @@ class BigQueryTableExistenceSensor(BaseSensorOperator):
project_id=self.project_id,
poll_interval=self.poke_interval,
gcp_conn_id=self.gcp_conn_id,
+ impersonation_chain=self.impersonation_chain,
hook_params={
"impersonation_chain": self.impersonation_chain,
},
@@ -301,6 +302,7 @@ class
BigQueryTablePartitionExistenceSensor(BaseSensorOperator):
partition_id=self.partition_id,
poll_interval=self.poke_interval,
gcp_conn_id=self.gcp_conn_id,
+ impersonation_chain=self.impersonation_chain,
hook_params={
"impersonation_chain": self.impersonation_chain,
},
diff --git a/providers/google/tests/unit/google/cloud/sensors/test_bigquery.py
b/providers/google/tests/unit/google/cloud/sensors/test_bigquery.py
index 1881c360e5e..eb953d6af03 100644
--- a/providers/google/tests/unit/google/cloud/sensors/test_bigquery.py
+++ b/providers/google/tests/unit/google/cloud/sensors/test_bigquery.py
@@ -103,6 +103,24 @@ class TestBigqueryTableExistenceSensor:
"Trigger is not a BigQueryTableExistenceTrigger"
)
+ @mock.patch("airflow.providers.google.cloud.sensors.bigquery.BigQueryHook")
+ def test_deferred_trigger_receives_impersonation_chain(self, mock_hook):
+ task = BigQueryTableExistenceSensor(
+ task_id="check_table_exists",
+ project_id=TEST_PROJECT_ID,
+ dataset_id=TEST_DATASET_ID,
+ table_id=TEST_TABLE_ID,
+ gcp_conn_id=TEST_GCP_CONN_ID,
+ impersonation_chain=TEST_IMPERSONATION_CHAIN,
+ deferrable=True,
+ )
+ mock_hook.return_value.table_exists.return_value = False
+
+ with pytest.raises(TaskDeferred) as exc:
+ task.execute(mock.MagicMock())
+
+ assert exc.value.trigger.impersonation_chain ==
TEST_IMPERSONATION_CHAIN
+
def test_execute_deferred_failure(self):
"""Tests that an expected exception is raised in case of error event"""
task = BigQueryTableExistenceSensor(
@@ -206,6 +224,25 @@ class TestBigqueryTablePartitionExistenceSensor:
"Trigger is not a BigQueryTablePartitionExistenceTrigger"
)
+ @mock.patch("airflow.providers.google.cloud.sensors.bigquery.BigQueryHook")
+ def test_deferred_trigger_receives_impersonation_chain(self, mock_hook):
+ task = BigQueryTablePartitionExistenceSensor(
+ task_id="test_task_id",
+ project_id=TEST_PROJECT_ID,
+ dataset_id=TEST_DATASET_ID,
+ table_id=TEST_TABLE_ID,
+ partition_id=TEST_PARTITION_ID,
+ gcp_conn_id=TEST_GCP_CONN_ID,
+ impersonation_chain=TEST_IMPERSONATION_CHAIN,
+ deferrable=True,
+ )
+ mock_hook.return_value.table_partition_exists.return_value = False
+
+ with pytest.raises(TaskDeferred) as exc:
+ task.execute(context={})
+
+ assert exc.value.trigger.impersonation_chain ==
TEST_IMPERSONATION_CHAIN
+
def test_execute_with_deferrable_mode_execute_failure(self):
"""Tests that an AirflowException is raised in case of error event"""
task = BigQueryTablePartitionExistenceSensor(