From dc4e9b5c20509225dfe7e75d80f4161ec9bece44 Mon Sep 17 00:00:00 2001 From: mitre88 Date: Fri, 31 Jul 2026 20:24:34 -0600 Subject: [PATCH] Normalize BigQuery DTS sensor expected statuses after rendering Rebased onto current main. Validate expected_statuses after template rendering so templated values are checked at runtime. --- .../google/cloud/sensors/bigquery_dts.py | 4 ++-- .../google/cloud/sensors/test_bigquery_dts.py | 17 +++++++++++++++++ .../prek/validate_operators_init_exemptions.txt | 1 - 3 files changed, 19 insertions(+), 3 deletions(-) diff --git a/providers/google/src/airflow/providers/google/cloud/sensors/bigquery_dts.py b/providers/google/src/airflow/providers/google/cloud/sensors/bigquery_dts.py index 6f12f92f36d16..46bb1e5c64f5f 100644 --- a/providers/google/src/airflow/providers/google/cloud/sensors/bigquery_dts.py +++ b/providers/google/src/airflow/providers/google/cloud/sensors/bigquery_dts.py @@ -99,7 +99,7 @@ def __init__( self.retry = retry self.request_timeout = request_timeout self.metadata = metadata - self.expected_statuses = self._normalize_state_list(expected_statuses) + self.expected_statuses = expected_statuses self.project_id = project_id self.gcp_cloud_conn_id = gcp_conn_id self.impersonation_chain = impersonation_chain @@ -144,4 +144,4 @@ def poke(self, context: Context) -> bool: if run.state in (TransferState.FAILED, TransferState.CANCELLED): message = f"Transfer {self.run_id} did not succeed" raise AirflowException(message) - return run.state in self.expected_statuses + return run.state in self._normalize_state_list(self.expected_statuses) diff --git a/providers/google/tests/unit/google/cloud/sensors/test_bigquery_dts.py b/providers/google/tests/unit/google/cloud/sensors/test_bigquery_dts.py index a0f6581117235..5479712fc4a19 100644 --- a/providers/google/tests/unit/google/cloud/sensors/test_bigquery_dts.py +++ b/providers/google/tests/unit/google/cloud/sensors/test_bigquery_dts.py @@ -89,3 +89,20 @@ def test_poke_returns_true(self, mock_hook): retry=DEFAULT, timeout=None, ) + + @mock.patch( + "airflow.providers.google.cloud.sensors.bigquery_dts.BiqQueryDataTransferServiceHook", + return_value=MM(get_transfer_run=MM(return_value=MM(state=TransferState.SUCCEEDED))), + ) + def test_templated_expected_statuses_normalized_at_poke_time(self, mock_hook): + op = BigQueryDataTransferServiceTransferRunSensor( + transfer_config_id=TRANSFER_CONFIG_ID, + run_id=RUN_ID, + task_id="id", + project_id=PROJECT_ID, + expected_statuses="{{ var.value.expected_status }}", + ) + # Template rendering replaces the Jinja expression with the resolved value before poke. + op.expected_statuses = "succeeded" + + assert op.poke({}) is True diff --git a/scripts/ci/prek/validate_operators_init_exemptions.txt b/scripts/ci/prek/validate_operators_init_exemptions.txt index e6869670a5bea..da57189052933 100644 --- a/scripts/ci/prek/validate_operators_init_exemptions.txt +++ b/scripts/ci/prek/validate_operators_init_exemptions.txt @@ -16,7 +16,6 @@ providers/google/src/airflow/providers/google/cloud/operators/dataproc.py::Datap providers/google/src/airflow/providers/google/cloud/operators/dataproc.py::DataprocSubmitJobOperator providers/google/src/airflow/providers/google/cloud/operators/functions.py::CloudFunctionDeployFunctionOperator providers/google/src/airflow/providers/google/cloud/operators/gcs.py::GCSFileTransformOperator -providers/google/src/airflow/providers/google/cloud/sensors/bigquery_dts.py::BigQueryDataTransferServiceTransferRunSensor providers/google/src/airflow/providers/google/cloud/sensors/cloud_composer.py::CloudComposerExternalTaskSensor providers/google/src/airflow/providers/google/cloud/transfers/azure_fileshare_to_gcs.py::AzureFileShareToGCSOperator providers/google/src/airflow/providers/google/cloud/transfers/gcs_to_bigquery.py::GCSToBigQueryOperator