From a7ed135bff4b4d0aff5ffa0c899d43b1047cb8a2 Mon Sep 17 00:00:00 2001 From: Dr Alex Mitre <30060514+mitre88@users.noreply.github.com> Date: Mon, 27 Jul 2026 09:40:39 -0600 Subject: [PATCH] Validate Google operator templated parameters after rendering CloudSpeechToTextRecognizeSpeechOperator, CloudTextToSpeechSynthesizeOperator, and CloudFirestoreExportDatabaseOperator validate template fields inside _validate_inputs() called from __init__, so the checks ran against un-rendered Jinja expressions and an empty rendered value was never caught. The validation now runs at execute time, and the raises are narrowed from AirflowException to ValueError per the ongoing exception clean-up. --- generated/known_airflow_exceptions.txt | 3 --- .../google/cloud/operators/speech_to_text.py | 7 +++---- .../google/cloud/operators/text_to_speech.py | 5 ++--- .../google/firebase/operators/firestore.py | 5 ++--- .../cloud/operators/test_speech_to_text.py | 16 ++++++++++++++++ .../cloud/operators/test_text_to_speech.py | 3 +-- .../google/firebase/operators/test_firestore.py | 17 +++++++++++++++++ 7 files changed, 41 insertions(+), 15 deletions(-) diff --git a/generated/known_airflow_exceptions.txt b/generated/known_airflow_exceptions.txt index 1fb461854905d..912cf5e4842e6 100644 --- a/generated/known_airflow_exceptions.txt +++ b/generated/known_airflow_exceptions.txt @@ -270,8 +270,6 @@ providers/google/src/airflow/providers/google/cloud/operators/looker.py::1 providers/google/src/airflow/providers/google/cloud/operators/managed_kafka.py::15 providers/google/src/airflow/providers/google/cloud/operators/pubsub.py::1 providers/google/src/airflow/providers/google/cloud/operators/spanner.py::19 -providers/google/src/airflow/providers/google/cloud/operators/speech_to_text.py::2 -providers/google/src/airflow/providers/google/cloud/operators/text_to_speech.py::1 providers/google/src/airflow/providers/google/cloud/operators/translate.py::9 providers/google/src/airflow/providers/google/cloud/operators/translate_speech.py::2 providers/google/src/airflow/providers/google/cloud/operators/vertex_ai/batch_prediction_job.py::1 @@ -317,7 +315,6 @@ providers/google/src/airflow/providers/google/cloud/utils/credentials_provider.p providers/google/src/airflow/providers/google/common/hooks/base_google.py::7 providers/google/src/airflow/providers/google/common/hooks/operation_helpers.py::2 providers/google/src/airflow/providers/google/firebase/hooks/firestore.py::1 -providers/google/src/airflow/providers/google/firebase/operators/firestore.py::1 providers/google/src/airflow/providers/google/leveldb/hooks/leveldb.py::4 providers/google/src/airflow/providers/google/marketing_platform/hooks/campaign_manager.py::2 providers/google/src/airflow/providers/google/marketing_platform/hooks/search_ads.py::1 diff --git a/providers/google/src/airflow/providers/google/cloud/operators/speech_to_text.py b/providers/google/src/airflow/providers/google/cloud/operators/speech_to_text.py index 5424fba7191c6..3ac73c54abfac 100644 --- a/providers/google/src/airflow/providers/google/cloud/operators/speech_to_text.py +++ b/providers/google/src/airflow/providers/google/cloud/operators/speech_to_text.py @@ -25,7 +25,6 @@ from google.api_core.gapic_v1.method import DEFAULT, _MethodDefault from google.protobuf.json_format import MessageToDict -from airflow.providers.common.compat.sdk import AirflowException from airflow.providers.google.cloud.hooks.speech_to_text import CloudSpeechToTextHook, RecognitionAudio from airflow.providers.google.cloud.operators.cloud_base import GoogleCloudBaseOperator from airflow.providers.google.common.hooks.base_google import PROVIDE_PROJECT_ID @@ -99,17 +98,17 @@ def __init__( self.gcp_conn_id = gcp_conn_id self.retry = retry self.timeout = timeout - self._validate_inputs() self.impersonation_chain = impersonation_chain super().__init__(**kwargs) def _validate_inputs(self) -> None: if self.audio == "": - raise AirflowException("The required parameter 'audio' is empty") + raise ValueError("The required parameter 'audio' is empty") if self.config == "": - raise AirflowException("The required parameter 'config' is empty") + raise ValueError("The required parameter 'config' is empty") def execute(self, context: Context): + self._validate_inputs() hook = CloudSpeechToTextHook( gcp_conn_id=self.gcp_conn_id, impersonation_chain=self.impersonation_chain, diff --git a/providers/google/src/airflow/providers/google/cloud/operators/text_to_speech.py b/providers/google/src/airflow/providers/google/cloud/operators/text_to_speech.py index 5deded3ef8b7d..bfd014090829e 100644 --- a/providers/google/src/airflow/providers/google/cloud/operators/text_to_speech.py +++ b/providers/google/src/airflow/providers/google/cloud/operators/text_to_speech.py @@ -25,7 +25,6 @@ from google.api_core.gapic_v1.method import DEFAULT, _MethodDefault -from airflow.providers.common.compat.sdk import AirflowException from airflow.providers.google.cloud.hooks.gcs import GCSHook from airflow.providers.google.cloud.hooks.text_to_speech import CloudTextToSpeechHook from airflow.providers.google.cloud.operators.cloud_base import GoogleCloudBaseOperator @@ -112,7 +111,6 @@ def __init__( self.gcp_conn_id = gcp_conn_id self.retry = retry self.timeout = timeout - self._validate_inputs() self.impersonation_chain = impersonation_chain super().__init__(**kwargs) @@ -125,9 +123,10 @@ def _validate_inputs(self) -> None: "target_filename", ]: if getattr(self, parameter) == "": - raise AirflowException(f"The required parameter '{parameter}' is empty") + raise ValueError(f"The required parameter '{parameter}' is empty") def execute(self, context: Context) -> None: + self._validate_inputs() hook = CloudTextToSpeechHook( gcp_conn_id=self.gcp_conn_id, impersonation_chain=self.impersonation_chain, diff --git a/providers/google/src/airflow/providers/google/firebase/operators/firestore.py b/providers/google/src/airflow/providers/google/firebase/operators/firestore.py index 0c35f5729239f..d39f4bf879bc5 100644 --- a/providers/google/src/airflow/providers/google/firebase/operators/firestore.py +++ b/providers/google/src/airflow/providers/google/firebase/operators/firestore.py @@ -19,7 +19,6 @@ from collections.abc import Sequence from typing import TYPE_CHECKING -from airflow.providers.common.compat.sdk import AirflowException from airflow.providers.google.common.hooks.base_google import PROVIDE_PROJECT_ID from airflow.providers.google.firebase.hooks.firestore import CloudFirestoreHook from airflow.providers.google.version_compat import BaseOperator @@ -78,14 +77,14 @@ def __init__( self.project_id = project_id self.gcp_conn_id = gcp_conn_id self.api_version = api_version - self._validate_inputs() self.impersonation_chain = impersonation_chain def _validate_inputs(self) -> None: if not self.body: - raise AirflowException("The required parameter 'body' is missing") + raise ValueError("The required parameter 'body' is missing") def execute(self, context: Context): + self._validate_inputs() hook = CloudFirestoreHook( gcp_conn_id=self.gcp_conn_id, api_version=self.api_version, diff --git a/providers/google/tests/unit/google/cloud/operators/test_speech_to_text.py b/providers/google/tests/unit/google/cloud/operators/test_speech_to_text.py index 6e515a3d45d51..c2e0ad737cd0a 100644 --- a/providers/google/tests/unit/google/cloud/operators/test_speech_to_text.py +++ b/providers/google/tests/unit/google/cloud/operators/test_speech_to_text.py @@ -55,6 +55,22 @@ def test_recognize_speech_green_path(self, mock_hook): config=CONFIG, audio=AUDIO, retry=DEFAULT, timeout=None ) + @patch("airflow.providers.google.cloud.operators.speech_to_text.CloudSpeechToTextHook") + def test_empty_audio_fails_at_execute_time(self, mock_hook): + op = CloudSpeechToTextRecognizeSpeechOperator( + project_id=PROJECT_ID, + gcp_conn_id=GCP_CONN_ID, + audio="{{ var.value.audio }}", + config=CONFIG, + task_id="id", + ) + # Template rendering replaces the Jinja expression with the resolved value before execute. + op.audio = "" + + with pytest.raises(ValueError, match="The required parameter 'audio' is empty"): + op.execute(context={"task_instance": Mock()}) + mock_hook.assert_not_called() + @patch("airflow.providers.google.cloud.operators.speech_to_text.CloudSpeechToTextHook") def test_missing_config(self, mock_hook): mock_hook.return_value.recognize_speech.return_value = True diff --git a/providers/google/tests/unit/google/cloud/operators/test_text_to_speech.py b/providers/google/tests/unit/google/cloud/operators/test_text_to_speech.py index d4c2a59787049..76c53c5c3f4c2 100644 --- a/providers/google/tests/unit/google/cloud/operators/test_text_to_speech.py +++ b/providers/google/tests/unit/google/cloud/operators/test_text_to_speech.py @@ -22,7 +22,6 @@ import pytest from google.api_core.gapic_v1.method import DEFAULT -from airflow.providers.common.compat.sdk import AirflowException from airflow.providers.google.cloud.operators.text_to_speech import CloudTextToSpeechSynthesizeOperator PROJECT_ID = "project-id" @@ -98,7 +97,7 @@ def test_missing_arguments( ): mocked_context = Mock() - with pytest.raises(AirflowException) as ctx: + with pytest.raises(ValueError, match="is empty") as ctx: CloudTextToSpeechSynthesizeOperator( project_id="project-id", input_data=input_data, diff --git a/providers/google/tests/unit/google/firebase/operators/test_firestore.py b/providers/google/tests/unit/google/firebase/operators/test_firestore.py index 340058ef6eb42..e9fd938b01777 100644 --- a/providers/google/tests/unit/google/firebase/operators/test_firestore.py +++ b/providers/google/tests/unit/google/firebase/operators/test_firestore.py @@ -18,6 +18,8 @@ from unittest import mock +import pytest + from airflow.providers.google.firebase.operators.firestore import CloudFirestoreExportDatabaseOperator TEST_OUTPUT_URI_PREFIX: str = "gs://example-bucket/path" @@ -42,3 +44,18 @@ def test_execute(self, mock_firestore_hook): mock_firestore_hook.return_value.export_documents.assert_called_once_with( body=EXPORT_DOCUMENT_BODY, database_id="(default)", project_id=TEST_PROJECT_ID ) + + @mock.patch("airflow.providers.google.firebase.operators.firestore.CloudFirestoreHook") + def test_empty_body_fails_at_execute_time(self, mock_firestore_hook): + op = CloudFirestoreExportDatabaseOperator( + task_id="test-task", + body="{{ var.value.export_body }}", + gcp_conn_id="google_cloud_default", + project_id=TEST_PROJECT_ID, + ) + # Template rendering replaces the Jinja expression with the resolved value before execute. + op.body = None + + with pytest.raises(ValueError, match="The required parameter 'body' is missing"): + op.execute(mock.MagicMock()) + mock_firestore_hook.return_value.export_documents.assert_not_called()