From 910515b7b3dd8873d909877aee5cdd7aaaeb6b47 Mon Sep 17 00:00:00 2001 From: Aaryan Mahajan Date: Wed, 22 Jul 2026 16:17:15 +0400 Subject: [PATCH] Fix KubernetesPodOperator dry_run requiring live Kubernetes API client airflow dags test and dry_run() failed outside a cluster because build_pod_request_obj() unconditionally read hook.is_in_cluster to set the airflow_kpo_in_cluster label, which instantiates a Kubernetes API client and requires kube credentials/config to be available. A dry run should not need cluster access at all. closes: #45812 --- .../cncf/kubernetes/operators/pod.py | 19 +++++++------ .../cncf/kubernetes/operators/test_pod.py | 28 +++++++++++++++++++ 2 files changed, 39 insertions(+), 8 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py index afbe07b827a08..f52333d605c38 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py @@ -1557,12 +1557,16 @@ def _delete_with_retry(): except kubernetes.client.exceptions.ApiException: self.log.exception("Unable to delete pod %s", self.pod.metadata.name) - def build_pod_request_obj(self, context: Context | None = None) -> k8s.V1Pod: + def build_pod_request_obj(self, context: Context | None = None, *, dry_run: bool = False) -> k8s.V1Pod: """ Return V1Pod object based on pod template file, full pod spec, and other operator parameters. The V1Pod attributes are derived (in order of precedence) from operator params, full pod spec, pod template file. + + :param dry_run: if True, skip anything that requires a live Kubernetes API client + (e.g. determining whether the hook is running in-cluster), since dry runs must not + require kube credentials or config to be available. """ self.log.debug("Creating pod for KubernetesPodOperator task %s", self.task_id) @@ -1673,12 +1677,11 @@ def build_pod_request_obj(self, context: Context | None = None) -> k8s.V1Pod: pod.metadata.labels.update(labels) # Add Airflow Version to the label # And a label to identify that pod is launched by KubernetesPodOperator - pod.metadata.labels.update( - { - "airflow_version": airflow_version.replace("+", "-"), - "airflow_kpo_in_cluster": str(self.hook.is_in_cluster), - } - ) + pod.metadata.labels.update({"airflow_version": airflow_version.replace("+", "-")}) + if not dry_run: + # self.hook.is_in_cluster instantiates a Kubernetes API client, which requires + # kube credentials/config to be available and must be skipped during a dry run. + pod.metadata.labels["airflow_kpo_in_cluster"] = str(self.hook.is_in_cluster) pod_mutation_hook(pod) return pod @@ -1689,7 +1692,7 @@ def dry_run(self) -> None: Does not include labels specific to the task instance (since there isn't one in a dry_run) and excludes all empty elements. """ - pod = self.build_pod_request_obj() + pod = self.build_pod_request_obj(dry_run=True) print(yaml.dump(prune_dict(pod.to_dict(), mode="strict"))) def process_duplicate_label_pods(self, pod_list: list[k8s.V1Pod]) -> k8s.V1Pod: diff --git a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_pod.py b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_pod.py index 0e903ec644b1a..640d1a199c75f 100644 --- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_pod.py +++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_pod.py @@ -472,6 +472,34 @@ def test_labels_mapped(self): "airflow_kpo_in_cluster": str(k.hook.is_in_cluster), } + @patch(HOOK_CLASS) + def test_build_pod_request_obj_dry_run_skips_live_kube_client(self, hook_mock): + """dry_run must not require a live Kubernetes API client (e.g. no kube config available).""" + type(hook_mock.return_value).is_in_cluster = mock.PropertyMock( + side_effect=RuntimeError("kube config not available") + ) + k = KubernetesPodOperator( + name="test", + task_id="task", + ) + pod = k.build_pod_request_obj(dry_run=True) + assert "airflow_kpo_in_cluster" not in pod.metadata.labels + + with pytest.raises(RuntimeError, match="kube config not available"): + k.build_pod_request_obj() + + @patch(HOOK_CLASS) + def test_dry_run_method_does_not_require_live_kube_client(self, hook_mock): + type(hook_mock.return_value).is_in_cluster = mock.PropertyMock( + side_effect=RuntimeError("kube config not available") + ) + hook_mock.return_value.get_namespace.return_value = "default" + k = KubernetesPodOperator( + name="test", + task_id="task", + ) + k.dry_run() + def test_find_custom_pod_labels(self): k = KubernetesPodOperator( labels={"foo": "bar", "hello": "airflow"},