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..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 @@ -84,6 +84,8 @@ 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, 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 @@ 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 +179,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..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 @@ -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,13 @@ 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 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 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