This is an automated email from the ASF dual-hosted git repository.

shahar1 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 220d2199266 Validate KubernetesResourceBaseOperator yaml_conf after 
rendering (#70337)
220d2199266 is described below

commit 220d219926632b510210ddcfcfc48dfb0424e79c
Author: Stefan Wang <[email protected]>
AuthorDate: Fri Jul 24 00:56:35 2026 -0700

    Validate KubernetesResourceBaseOperator yaml_conf after rendering (#70337)
    
    * Validate KubernetesResourceBaseOperator yaml_conf after rendering
    
    yaml_conf and yaml_conf_file are template fields, rendered after __init__ 
runs.
    The base constructor raised when neither was set, acting on the un-rendered
    values. Move the presence check into a helper called from execute() in both 
the
    create and delete operators, so it runs after rendering.
    
    related: #70296
---
 .../src/airflow/providers/cncf/kubernetes/operators/resource.py   | 4 ++++
 .../tests/unit/cncf/kubernetes/operators/test_resource.py         | 8 ++++++++
 scripts/ci/prek/validate_operators_init_exemptions.txt            | 1 -
 3 files changed, 12 insertions(+), 1 deletion(-)

diff --git 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/resource.py
 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/resource.py
index f7983581715..c16df215d6e 100644
--- 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/resource.py
+++ 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/resource.py
@@ -84,6 +84,8 @@ class KubernetesResourceBaseOperator(BaseOperator):
         self.namespaced = namespaced
         self.config_file = config_file
 
+    def _validate_yaml_conf(self) -> None:
+        # yaml_conf/yaml_conf_file are template fields; validate after 
rendering, called from execute.
         if not any([self.yaml_conf, self.yaml_conf_file]):
             raise AirflowException("One of `yaml_conf` or `yaml_conf_file` 
arguments must be provided")
 
@@ -144,6 +146,7 @@ class 
KubernetesCreateResourceOperator(KubernetesResourceBaseOperator):
             k8s_resource_iterator(self.create_custom_from_yaml_object, objects)
 
     def execute(self, context) -> None:
+        self._validate_yaml_conf()
         if self.yaml_conf:
             self._create_objects(yaml.safe_load_all(self.yaml_conf))
         elif self.yaml_conf_file and os.path.exists(self.yaml_conf_file):
@@ -176,6 +179,7 @@ class 
KubernetesDeleteResourceOperator(KubernetesResourceBaseOperator):
             k8s_resource_iterator(self.delete_custom_from_yaml_object, objects)
 
     def execute(self, context) -> None:
+        self._validate_yaml_conf()
         if self.yaml_conf:
             self._delete_objects(yaml.safe_load_all(self.yaml_conf))
         elif self.yaml_conf_file and os.path.exists(self.yaml_conf_file):
diff --git 
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_resource.py
 
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_resource.py
index 3bea6a051a2..8c5e6d1f7ad 100644
--- 
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_resource.py
+++ 
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_resource.py
@@ -27,6 +27,7 @@ from airflow.providers.cncf.kubernetes.operators.resource 
import (
     KubernetesCreateResourceOperator,
     KubernetesDeleteResourceOperator,
 )
+from airflow.providers.common.compat.sdk import AirflowException
 from airflow.utils import timezone
 
 TEST_VALID_RESOURCE_YAML = """
@@ -93,6 +94,13 @@ class TestKubernetesXResourceOperator:
         args = {"owner": "airflow", "start_date": timezone.datetime(2020, 2, 
1)}
         self.dag = DAG("test_dag_id", schedule=None, default_args=args)
 
+    def test_missing_yaml_conf_rejected_at_execute(self, context):
+        # yaml_conf/yaml_conf_file are template fields; the presence check 
runs at execute.
+        for operator_class in (KubernetesCreateResourceOperator, 
KubernetesDeleteResourceOperator):
+            op = operator_class(task_id="test_task_id")
+            with pytest.raises(AirflowException, match="One of `yaml_conf` or 
`yaml_conf_file`"):
+                op.execute(context={})
+
     @patch("kubernetes.config.load_kube_config")
     
@patch("kubernetes.client.api.CoreV1Api.create_namespaced_persistent_volume_claim")
     def test_create_application_from_yaml(
diff --git a/scripts/ci/prek/validate_operators_init_exemptions.txt 
b/scripts/ci/prek/validate_operators_init_exemptions.txt
index c2e43f97a97..881d2cd90d1 100644
--- a/scripts/ci/prek/validate_operators_init_exemptions.txt
+++ b/scripts/ci/prek/validate_operators_init_exemptions.txt
@@ -27,7 +27,6 @@ 
providers/apache/hive/src/airflow/providers/apache/hive/sensors/named_hive_parti
 
providers/apache/kafka/src/airflow/providers/apache/kafka/operators/produce.py::ProduceToTopicOperator
 
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/kueue.py::KubernetesInstallKueueOperator
 
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py::KubernetesPodOperator
-providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/resource.py::KubernetesResourceBaseOperator
 
providers/cohere/src/airflow/providers/cohere/operators/embedding.py::CohereEmbeddingOperator
 
providers/databricks/src/airflow/providers/databricks/operators/databricks_repos.py::DatabricksReposCreateOperator
 
providers/databricks/src/airflow/providers/databricks/operators/databricks_repos.py::DatabricksReposDeleteOperator

Reply via email to