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