From 9e0f15cb04ecd41e788d191bed4739430cd89038 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Thu, 23 Jul 2026 14:53:10 -0700 Subject: [PATCH 1/2] 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 Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../providers/cncf/kubernetes/operators/resource.py | 5 +++++ .../unit/cncf/kubernetes/operators/test_resource.py | 9 +++++++++ scripts/ci/prek/validate_operators_init_exemptions.txt | 1 - 3 files changed, 14 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 f79835817155e..73b3fa4e9c4e0 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,9 @@ def __init__( 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 (from execute), + # not in __init__ where they are still the un-rendered Jinja expressions. 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 +147,7 @@ def _create_objects(self, objects): 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 +180,7 @@ def _delete_objects(self, objects): 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 3bea6a051a24e..b81252928757b 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 @@ KubernetesCreateResourceOperator, KubernetesDeleteResourceOperator, ) +from airflow.providers.common.compat.sdk import AirflowException from airflow.utils import timezone TEST_VALID_RESOURCE_YAML = """ @@ -93,6 +94,14 @@ def setup_method(self): 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 must run at execute + # (after rendering), so constructing with neither no longer raises in __init__. + 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 01fd5bc56dbb7..6e56b1bee1322 100644 --- a/scripts/ci/prek/validate_operators_init_exemptions.txt +++ b/scripts/ci/prek/validate_operators_init_exemptions.txt @@ -32,7 +32,6 @@ providers/apache/kafka/src/airflow/providers/apache/kafka/operators/produce.py:: providers/apache/spark/src/airflow/providers/apache/spark/operators/spark_submit.py::SparkSubmitOperator 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/common/ai/src/airflow/providers/common/ai/operators/agent.py::AgentOperator providers/common/ai/src/airflow/providers/common/ai/operators/document_loader.py::DocumentLoaderOperator From 9c2be8104b8f82becc9a65cc101cc833a0327f70 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Thu, 23 Jul 2026 23:16:21 -0700 Subject: [PATCH 2/2] Tighten KubernetesResourceBaseOperator validation comments Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../airflow/providers/cncf/kubernetes/operators/resource.py | 3 +-- .../tests/unit/cncf/kubernetes/operators/test_resource.py | 3 +-- 2 files changed, 2 insertions(+), 4 deletions(-) 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 73b3fa4e9c4e0..c16df215d6e25 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 @@ -85,8 +85,7 @@ def __init__( self.config_file = config_file def _validate_yaml_conf(self) -> None: - # yaml_conf/yaml_conf_file are template fields; validate after rendering (from execute), - # not in __init__ where they are still the un-rendered Jinja expressions. + # 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") 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 b81252928757b..8c5e6d1f7ad4b 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 @@ -95,8 +95,7 @@ def setup_method(self): 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 must run at execute - # (after rendering), so constructing with neither no longer raises in __init__. + # 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`"):