From 08f503ff6725daead7f2c21882a789a797e90c5f Mon Sep 17 00:00:00 2001 From: Amogh Atreya Date: Fri, 24 Jul 2026 11:11:06 +0530 Subject: [PATCH 1/3] fix(providers): move S3 and GCE operator template validation out of __init__ (#70296) --- generated/known_airflow_exceptions.txt | 2 +- .../providers/amazon/aws/operators/s3.py | 5 --- .../unit/amazon/aws/operators/test_s3.py | 40 +++++++++---------- .../google/cloud/operators/compute.py | 10 ++--- .../google/cloud/operators/test_compute.py | 24 +++++++++++ .../validate_operators_init_exemptions.txt | 2 - 6 files changed, 50 insertions(+), 33 deletions(-) diff --git a/generated/known_airflow_exceptions.txt b/generated/known_airflow_exceptions.txt index 1fb461854905d..66ee6590bc84e 100644 --- a/generated/known_airflow_exceptions.txt +++ b/generated/known_airflow_exceptions.txt @@ -85,7 +85,7 @@ providers/amazon/src/airflow/providers/amazon/aws/operators/neptune.py::4 providers/amazon/src/airflow/providers/amazon/aws/operators/rds.py::4 providers/amazon/src/airflow/providers/amazon/aws/operators/redshift_cluster.py::9 providers/amazon/src/airflow/providers/amazon/aws/operators/redshift_data.py::2 -providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py::5 +providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py::4 providers/amazon/src/airflow/providers/amazon/aws/operators/sagemaker.py::22 providers/amazon/src/airflow/providers/amazon/aws/operators/sagemaker_unified_studio.py::2 providers/amazon/src/airflow/providers/amazon/aws/operators/step_function.py::2 diff --git a/providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py b/providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py index b74fb60a1b763..381473de9435a 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py @@ -697,11 +697,6 @@ def __init__( self._keys: str | list[str] = "" - if not exactly_one(keys is None, all(var is None for var in [prefix, from_datetime, to_datetime])): - raise AirflowException( - "Either keys or at least one of prefix, from_datetime, to_datetime should be set." - ) - def execute(self, context: Context): if not exactly_one( self.keys is None, all(var is None for var in [self.prefix, self.from_datetime, self.to_datetime]) diff --git a/providers/amazon/tests/unit/amazon/aws/operators/test_s3.py b/providers/amazon/tests/unit/amazon/aws/operators/test_s3.py index cd907171927e4..b3d10fa23705c 100644 --- a/providers/amazon/tests/unit/amazon/aws/operators/test_s3.py +++ b/providers/amazon/tests/unit/amazon/aws/operators/test_s3.py @@ -1118,19 +1118,22 @@ def test_s3_delete_empty_string(self): pytest.param(None, None, None, None, id="all-none"), ], ) - def test_validate_keys_and_filters_in_constructor(self, keys, prefix, from_datetime, to_datetime): - with pytest.raises( - AirflowException, - match=r"Either keys or at least one of prefix, from_datetime, to_datetime should be set.", - ): - S3DeleteObjectsOperator( - task_id="test_validate_keys_and_prefix_in_constructor", - bucket="foo-bar-bucket", - keys=keys, - prefix=prefix, - from_datetime=from_datetime, - to_datetime=to_datetime, - ) + def test_no_validation_in_constructor(self, keys, prefix, from_datetime, to_datetime): + # Template fields are rendered after __init__, so keys/prefix/from_datetime/to_datetime + # must not be validated in the constructor — construction must succeed even for combinations + # that are invalid once rendered (validation happens in execute()). + op = S3DeleteObjectsOperator( + task_id="test_no_validation_in_constructor", + bucket="foo-bar-bucket", + keys=keys, + prefix=prefix, + from_datetime=from_datetime, + to_datetime=to_datetime, + ) + assert op.keys == keys + assert op.prefix == prefix + assert op.from_datetime == from_datetime + assert op.to_datetime == to_datetime @pytest.mark.parametrize( ("keys", "prefix", "from_datetime", "to_datetime"), @@ -1162,17 +1165,14 @@ def test_validate_keys_and_prefix_in_execute(self, keys, prefix, from_datetime, conn.create_bucket(Bucket=bucket) conn.upload_fileobj(Bucket=bucket, Key=key_of_test, Fileobj=BytesIO(b"input")) - # Set valid values for constructor, and change them later for emulate rendering template op = S3DeleteObjectsOperator( task_id="test_validate_keys_and_prefix_in_execute", bucket=bucket, - keys="keys-exists", - prefix=None, + keys=keys, + prefix=prefix, + from_datetime=from_datetime, + to_datetime=to_datetime, ) - op.keys = keys - op.prefix = prefix - op.from_datetime = from_datetime - op.to_datetime = to_datetime # The object should be detected before the DELETE action is tested objects_in_dest_bucket = conn.list_objects(Bucket=bucket, Prefix=key_of_test) diff --git a/providers/google/src/airflow/providers/google/cloud/operators/compute.py b/providers/google/src/airflow/providers/google/cloud/operators/compute.py index 46d448213774e..706eccfa2305c 100644 --- a/providers/google/src/airflow/providers/google/cloud/operators/compute.py +++ b/providers/google/src/airflow/providers/google/cloud/operators/compute.py @@ -565,11 +565,7 @@ def __init__( self.retry = retry self.timeout = timeout self.metadata = metadata - - if validate_body: - self._field_validator = GcpBodyFieldValidator( - GCE_INSTANCE_TEMPLATE_VALIDATION_PATCH_SPECIFICATION, api_version=api_version - ) + self._validate_body = validate_body self._field_sanitizer = GcpBodyFieldSanitizer(GCE_INSTANCE_FIELDS_TO_SANITIZE) super().__init__( project_id=project_id, @@ -587,6 +583,10 @@ def _validate_inputs(self) -> None: raise AirflowException("The required parameter 'resource_id' is missing. ") def execute(self, context: Context) -> None: + if self._validate_body: + self._field_validator = GcpBodyFieldValidator( + GCE_INSTANCE_TEMPLATE_VALIDATION_PATCH_SPECIFICATION, api_version=self.api_version + ) hook = ComputeEngineHook( gcp_conn_id=self.gcp_conn_id, api_version=self.api_version, diff --git a/providers/google/tests/unit/google/cloud/operators/test_compute.py b/providers/google/tests/unit/google/cloud/operators/test_compute.py index 5d8bc187b07fd..5644ee5c1ff74 100644 --- a/providers/google/tests/unit/google/cloud/operators/test_compute.py +++ b/providers/google/tests/unit/google/cloud/operators/test_compute.py @@ -516,6 +516,30 @@ def test_delete_instance_should_execute_successfully(self, mock_hook): zone=GCE_ZONE, ) + @mock.patch(COMPUTE_ENGINE_HOOK_PATH) + def test_delete_instance_builds_field_validator_in_execute(self, mock_hook): + # api_version is a template field, so the body validator (which reads it) must be built in + # execute() — after rendering — not in the constructor. + op = ComputeEngineDeleteInstanceOperator( + resource_id=GCE_RESOURCE_ID, + zone=GCE_ZONE, + task_id=TASK_ID, + ) + assert op._field_validator is None + op.execute(context=mock.MagicMock()) + assert op._field_validator is not None + + @mock.patch(COMPUTE_ENGINE_HOOK_PATH) + def test_delete_instance_skips_field_validator_when_validate_body_false(self, mock_hook): + op = ComputeEngineDeleteInstanceOperator( + resource_id=GCE_RESOURCE_ID, + zone=GCE_ZONE, + task_id=TASK_ID, + validate_body=False, + ) + op.execute(context=mock.MagicMock()) + assert op._field_validator is None + def test_delete_instance_should_throw_ex_when_missing_zone(self): with pytest.raises(AirflowException, match=r"The required parameter 'zone' is missing"): ComputeEngineDeleteInstanceOperator( diff --git a/scripts/ci/prek/validate_operators_init_exemptions.txt b/scripts/ci/prek/validate_operators_init_exemptions.txt index 7ff6497cf73f0..f3cacfeed3809 100644 --- a/scripts/ci/prek/validate_operators_init_exemptions.txt +++ b/scripts/ci/prek/validate_operators_init_exemptions.txt @@ -14,7 +14,6 @@ providers/amazon/src/airflow/providers/amazon/aws/operators/emr.py::EmrAddStepsO providers/amazon/src/airflow/providers/amazon/aws/operators/glue.py::GlueDataQualityOperator providers/amazon/src/airflow/providers/amazon/aws/operators/neptune.py::NeptuneStartDbClusterOperator providers/amazon/src/airflow/providers/amazon/aws/operators/neptune.py::NeptuneStopDbClusterOperator -providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py::S3DeleteObjectsOperator providers/amazon/src/airflow/providers/amazon/aws/operators/sagemaker.py::SageMakerCreateNotebookOperator providers/amazon/src/airflow/providers/amazon/aws/operators/sagemaker.py::SageMakerProcessingOperator providers/amazon/src/airflow/providers/amazon/aws/operators/step_function.py::StepFunctionStartExecutionOperator @@ -46,7 +45,6 @@ providers/google/src/airflow/providers/google/cloud/operators/cloud_build.py::Cl providers/google/src/airflow/providers/google/cloud/operators/cloud_storage_transfer_service.py::CloudDataTransferServiceCreateJobOperator providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineCopyInstanceTemplateOperator providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineDeleteInstanceGroupManagerOperator -providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineDeleteInstanceOperator providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineDeleteInstanceTemplateOperator providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineInsertInstanceFromTemplateOperator providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineInsertInstanceGroupManagerOperator From f180f243b517aa84c17801a7eb105a1a96a0535f Mon Sep 17 00:00:00 2001 From: Amogh Atreya Date: Mon, 27 Jul 2026 11:07:34 +0530 Subject: [PATCH 2/3] Revert Compute Engine delete operator changes to focus PR on S3 A reviewer asked to keep this PR scoped to a single provider, and the Compute Engine operators are being addressed in a separate PR. This restores ComputeEngineDeleteInstanceOperator (and its exemption entry) to their original state, leaving only the S3DeleteObjectsOperator fix. Co-authored-by: Cursor --- .../google/cloud/operators/compute.py | 10 ++++---- .../google/cloud/operators/test_compute.py | 24 ------------------- .../validate_operators_init_exemptions.txt | 1 + 3 files changed, 6 insertions(+), 29 deletions(-) diff --git a/providers/google/src/airflow/providers/google/cloud/operators/compute.py b/providers/google/src/airflow/providers/google/cloud/operators/compute.py index 706eccfa2305c..46d448213774e 100644 --- a/providers/google/src/airflow/providers/google/cloud/operators/compute.py +++ b/providers/google/src/airflow/providers/google/cloud/operators/compute.py @@ -565,7 +565,11 @@ def __init__( self.retry = retry self.timeout = timeout self.metadata = metadata - self._validate_body = validate_body + + if validate_body: + self._field_validator = GcpBodyFieldValidator( + GCE_INSTANCE_TEMPLATE_VALIDATION_PATCH_SPECIFICATION, api_version=api_version + ) self._field_sanitizer = GcpBodyFieldSanitizer(GCE_INSTANCE_FIELDS_TO_SANITIZE) super().__init__( project_id=project_id, @@ -583,10 +587,6 @@ def _validate_inputs(self) -> None: raise AirflowException("The required parameter 'resource_id' is missing. ") def execute(self, context: Context) -> None: - if self._validate_body: - self._field_validator = GcpBodyFieldValidator( - GCE_INSTANCE_TEMPLATE_VALIDATION_PATCH_SPECIFICATION, api_version=self.api_version - ) hook = ComputeEngineHook( gcp_conn_id=self.gcp_conn_id, api_version=self.api_version, diff --git a/providers/google/tests/unit/google/cloud/operators/test_compute.py b/providers/google/tests/unit/google/cloud/operators/test_compute.py index 5644ee5c1ff74..5d8bc187b07fd 100644 --- a/providers/google/tests/unit/google/cloud/operators/test_compute.py +++ b/providers/google/tests/unit/google/cloud/operators/test_compute.py @@ -516,30 +516,6 @@ def test_delete_instance_should_execute_successfully(self, mock_hook): zone=GCE_ZONE, ) - @mock.patch(COMPUTE_ENGINE_HOOK_PATH) - def test_delete_instance_builds_field_validator_in_execute(self, mock_hook): - # api_version is a template field, so the body validator (which reads it) must be built in - # execute() — after rendering — not in the constructor. - op = ComputeEngineDeleteInstanceOperator( - resource_id=GCE_RESOURCE_ID, - zone=GCE_ZONE, - task_id=TASK_ID, - ) - assert op._field_validator is None - op.execute(context=mock.MagicMock()) - assert op._field_validator is not None - - @mock.patch(COMPUTE_ENGINE_HOOK_PATH) - def test_delete_instance_skips_field_validator_when_validate_body_false(self, mock_hook): - op = ComputeEngineDeleteInstanceOperator( - resource_id=GCE_RESOURCE_ID, - zone=GCE_ZONE, - task_id=TASK_ID, - validate_body=False, - ) - op.execute(context=mock.MagicMock()) - assert op._field_validator is None - def test_delete_instance_should_throw_ex_when_missing_zone(self): with pytest.raises(AirflowException, match=r"The required parameter 'zone' is missing"): ComputeEngineDeleteInstanceOperator( diff --git a/scripts/ci/prek/validate_operators_init_exemptions.txt b/scripts/ci/prek/validate_operators_init_exemptions.txt index f3cacfeed3809..fd83af9d572c3 100644 --- a/scripts/ci/prek/validate_operators_init_exemptions.txt +++ b/scripts/ci/prek/validate_operators_init_exemptions.txt @@ -45,6 +45,7 @@ providers/google/src/airflow/providers/google/cloud/operators/cloud_build.py::Cl providers/google/src/airflow/providers/google/cloud/operators/cloud_storage_transfer_service.py::CloudDataTransferServiceCreateJobOperator providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineCopyInstanceTemplateOperator providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineDeleteInstanceGroupManagerOperator +providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineDeleteInstanceOperator providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineDeleteInstanceTemplateOperator providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineInsertInstanceFromTemplateOperator providers/google/src/airflow/providers/google/cloud/operators/compute.py::ComputeEngineInsertInstanceGroupManagerOperator From 43d87651158dae4e6be33566789b22338bac0d2f Mon Sep 17 00:00:00 2001 From: Amogh Atreya Date: Mon, 27 Jul 2026 15:21:54 +0530 Subject: [PATCH 3/3] Raises ValueError for conflicting S3DeleteObjectsOperator arguments --- .../providers/amazon/aws/operators/s3.py | 9 +++++ .../unit/amazon/aws/operators/test_s3.py | 40 +++++++++---------- 2 files changed, 29 insertions(+), 20 deletions(-) diff --git a/providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py b/providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py index 381473de9435a..164777c78c136 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/operators/s3.py @@ -697,6 +697,15 @@ def __init__( self._keys: str | list[str] = "" + # Checked here and again in execute(): this guard keeps a plain authoring mistake a + # parse-time error, while execute() catches a templated `keys` that renders to None + # (which would otherwise list — and delete from — the whole bucket). + by_scan = prefix is not None or from_datetime is not None or to_datetime is not None + if not exactly_one(keys is not None, by_scan): + raise ValueError( + "Either keys or at least one of prefix, from_datetime, to_datetime should be set." + ) + def execute(self, context: Context): if not exactly_one( self.keys is None, all(var is None for var in [self.prefix, self.from_datetime, self.to_datetime]) diff --git a/providers/amazon/tests/unit/amazon/aws/operators/test_s3.py b/providers/amazon/tests/unit/amazon/aws/operators/test_s3.py index b3d10fa23705c..9d062caa4ace8 100644 --- a/providers/amazon/tests/unit/amazon/aws/operators/test_s3.py +++ b/providers/amazon/tests/unit/amazon/aws/operators/test_s3.py @@ -1118,22 +1118,19 @@ def test_s3_delete_empty_string(self): pytest.param(None, None, None, None, id="all-none"), ], ) - def test_no_validation_in_constructor(self, keys, prefix, from_datetime, to_datetime): - # Template fields are rendered after __init__, so keys/prefix/from_datetime/to_datetime - # must not be validated in the constructor — construction must succeed even for combinations - # that are invalid once rendered (validation happens in execute()). - op = S3DeleteObjectsOperator( - task_id="test_no_validation_in_constructor", - bucket="foo-bar-bucket", - keys=keys, - prefix=prefix, - from_datetime=from_datetime, - to_datetime=to_datetime, - ) - assert op.keys == keys - assert op.prefix == prefix - assert op.from_datetime == from_datetime - assert op.to_datetime == to_datetime + def test_validate_keys_and_filters_in_constructor(self, keys, prefix, from_datetime, to_datetime): + with pytest.raises( + ValueError, + match=r"Either keys or at least one of prefix, from_datetime, to_datetime should be set.", + ): + S3DeleteObjectsOperator( + task_id="test_validate_keys_and_prefix_in_constructor", + bucket="foo-bar-bucket", + keys=keys, + prefix=prefix, + from_datetime=from_datetime, + to_datetime=to_datetime, + ) @pytest.mark.parametrize( ("keys", "prefix", "from_datetime", "to_datetime"), @@ -1165,14 +1162,17 @@ def test_validate_keys_and_prefix_in_execute(self, keys, prefix, from_datetime, conn.create_bucket(Bucket=bucket) conn.upload_fileobj(Bucket=bucket, Key=key_of_test, Fileobj=BytesIO(b"input")) + # Set valid values for constructor, and change them later for emulate rendering template op = S3DeleteObjectsOperator( task_id="test_validate_keys_and_prefix_in_execute", bucket=bucket, - keys=keys, - prefix=prefix, - from_datetime=from_datetime, - to_datetime=to_datetime, + keys="keys-exists", + prefix=None, ) + op.keys = keys + op.prefix = prefix + op.from_datetime = from_datetime + op.to_datetime = to_datetime # The object should be detected before the DELETE action is tested objects_in_dest_bucket = conn.list_objects(Bucket=bucket, Prefix=key_of_test)