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

potiuk 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 ee5eefcb9fd Fix KubernetesPodOperator discarding successful XCom when 
sidecar kill fails (#72068)
ee5eefcb9fd is described below

commit ee5eefcb9fd53477b79bd4d6116a53c459d4f11f
Author: Jeremy Schoemaker <[email protected]>
AuthorDate: Mon Sep 21 09:30:07 2026 -0500

    Fix KubernetesPodOperator discarding successful XCom when sidecar kill 
fails (#72068)
    
    closes: #71369
---
 .../providers/cncf/kubernetes/operators/job.py     |  6 ++-
 .../providers/cncf/kubernetes/operators/pod.py     | 12 ++++--
 .../providers/cncf/kubernetes/utils/pod_manager.py | 24 +++++++++--
 .../unit/cncf/kubernetes/operators/test_job.py     | 47 +++++++++++++++++++++-
 .../unit/cncf/kubernetes/utils/test_pod_manager.py | 41 +++++++++++++++++++
 5 files changed, 122 insertions(+), 8 deletions(-)

diff --git 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/job.py
 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/job.py
index e888159c739..21fabd19c01 100644
--- 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/job.py
+++ 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/job.py
@@ -238,7 +238,11 @@ class KubernetesJobOperator(KubernetesPodOperator):
                             pod=pod, container_name=self.base_container_name
                         )
                         
self.pod_manager.await_xcom_sidecar_container_start(pod=pod)
-                        xcom_result.append(self.extract_xcom(pod=pod))
+                        # `wait_until_job_complete` below polls without a 
timeout, and the
+                        # Job can never reach a terminal state while the xcom 
sidecar is
+                        # still looping. So unlike KubernetesPodOperator, this 
path must
+                        # fail loudly when the sidecar cannot be killed rather 
than hang.
+                        xcom_result.append(self.extract_xcom(pod=pod, 
ignore_kill_failure=False))
                 self.job = self.hook.wait_until_job_complete(
                     job_name=self.job.metadata.name,
                     namespace=self.job.metadata.namespace,
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 2436fb7bb5d..5178222c647 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
@@ -810,9 +810,15 @@ class KubernetesPodOperator(BaseOperator):
                 self._read_pod_events(pod, reraise=False)
             raise
 
-    def extract_xcom(self, pod: k8s.V1Pod) -> dict[Any, Any] | None:
-        """Retrieve xcom value and kill xcom sidecar container."""
-        result = self.pod_manager.extract_xcom(pod)
+    def extract_xcom(self, pod: k8s.V1Pod, *, ignore_kill_failure: bool = 
True) -> dict[Any, Any] | None:
+        """
+        Retrieve xcom value and kill xcom sidecar container.
+
+        :param pod: the pod to read the XCom result from.
+        :param ignore_kill_failure: when True (the default), a failure to kill 
the sidecar
+            container does not discard the XCom value that was already read.
+        """
+        result = self.pod_manager.extract_xcom(pod, 
ignore_kill_failure=ignore_kill_failure)
         if isinstance(result, str) and result.rstrip() == EMPTY_XCOM_RESULT:
             self.log.info("xcom result file is empty.")
             return None
diff --git 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py
 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py
index 0914de4e90a..96588fe8334 100644
--- 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py
+++ 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py
@@ -974,8 +974,17 @@ class PodManager(LoggingMixin):
                 )
             time.sleep(1)
 
-    def extract_xcom(self, pod: V1Pod) -> str:
-        """Retrieve XCom value and kill xcom sidecar container."""
+    def extract_xcom(self, pod: V1Pod, *, ignore_kill_failure: bool = True) -> 
str:
+        """
+        Retrieve XCom value and kill xcom sidecar container.
+
+        :param pod: the pod to read the XCom result from.
+        :param ignore_kill_failure: when True (the default), a failure to kill 
the sidecar
+            container is logged as a warning and the successfully read XCom 
value is still
+            returned. Set it to False when the caller cannot tolerate a 
sidecar that keeps
+            running, for example ``KubernetesJobOperator``, whose Job can 
never reach a
+            terminal state while the sidecar is alive.
+        """
         # make sure that xcom sidecar container is still running
         if not self.container_is_running(pod, 
PodDefaults.SIDECAR_CONTAINER_NAME):
             raise XComRetrievalError(
@@ -986,7 +995,16 @@ class PodManager(LoggingMixin):
             result = self.extract_xcom_json(pod)
             return result
         finally:
-            self.extract_xcom_kill(pod)
+            try:
+                self.extract_xcom_kill(pod)
+            except (PodCommandException, ApiException) as e:
+                if not ignore_kill_failure:
+                    raise
+                self.log.warning(
+                    "Failed to kill xcom sidecar container in pod %s, leaving 
it running: %s",
+                    pod.metadata.name,
+                    e,
+                )
 
     @generic_api_retry
     def extract_xcom_json(self, pod: V1Pod) -> str:
diff --git 
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_job.py 
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_job.py
index 60536cbb439..6f2cdd7d425 100644
--- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_job.py
+++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_job.py
@@ -36,7 +36,7 @@ from airflow.providers.cncf.kubernetes.operators.job import (
     KubernetesPatchJobOperator,
 )
 from airflow.providers.cncf.kubernetes.triggers.job import KubernetesJobTrigger
-from airflow.providers.cncf.kubernetes.utils.pod_manager import PodManager
+from airflow.providers.cncf.kubernetes.utils.pod_manager import 
PodCommandException, PodManager
 from airflow.providers.common.compat.sdk import AirflowException, 
TaskDeferred, timezone
 from airflow.utils.session import create_session
 from airflow.utils.types import DagRunType
@@ -867,6 +867,51 @@ class TestKubernetesJobOperator:
         with pytest.raises(AirflowProviderDeprecationWarning):
             assert op.pod == mock_pods_expected[0]
 
+    @patch(f"{POD_MANAGER_CLASS}.extract_xcom_kill")
+    @patch(f"{POD_MANAGER_CLASS}.extract_xcom_json", return_value='{"a": 
"true"}')
+    @patch(f"{POD_MANAGER_CLASS}.container_is_running", return_value=True)
+    @patch(f"{POD_MANAGER_CLASS}.await_xcom_sidecar_container_start")
+    @patch(f"{POD_MANAGER_CLASS}.await_container_completion")
+    @patch(JOB_OPERATORS_PATH.format("KubernetesJobOperator.get_pods"))
+    
@patch(JOB_OPERATORS_PATH.format("KubernetesJobOperator.build_job_request_obj"))
+    @patch(JOB_OPERATORS_PATH.format("KubernetesJobOperator.create_job"))
+    @patch(f"{HOOK_CLASS}.wait_until_job_complete")
+    def test_xcom_sidecar_kill_failure_fails_job_instead_of_hanging(
+        self,
+        mock_wait_until_job_complete,
+        mock_create_job,
+        mock_build_job_request_obj,
+        mock_get_pods,
+        mock_await_container_completion,
+        mock_await_xcom_sidecar_container_start,
+        mock_container_is_running,
+        mock_extract_xcom_json,
+        mock_extract_xcom_kill,
+    ):
+        """
+        The Job path must not swallow a sidecar kill failure.
+
+        ``wait_until_job_complete`` polls without a timeout and the Job can 
never reach a
+        terminal state while the xcom sidecar keeps looping, so the task would 
hang. It
+        has to fail loudly before the polling starts instead.
+        """
+        mock_get_pods.return_value = [mock.MagicMock()]
+        mock_extract_xcom_kill.side_effect = PodCommandException(
+            "Command failed with stderr: Permission denied"
+        )
+
+        op = KubernetesJobOperator(
+            task_id="test_task_id",
+            wait_until_job_complete=True,
+            job_poll_interval=POLL_INTERVAL,
+            do_xcom_push=True,
+        )
+
+        with pytest.raises(PodCommandException, match="Permission denied"):
+            op.execute(context=dict(ti=mock.MagicMock()))
+
+        mock_wait_until_job_complete.assert_not_called()
+
     @pytest.mark.parametrize("do_xcom_push", [True, False])
     @pytest.mark.parametrize("get_logs", [True, False])
     @patch(JOB_OPERATORS_PATH.format("KubernetesJobOperator._write_logs"))
diff --git 
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py
 
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py
index f4c21683762..45e4b20a84f 100644
--- 
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py
+++ 
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py
@@ -32,6 +32,7 @@ from urllib3.exceptions import HTTPError as BaseHTTPError
 from airflow.providers.cncf.kubernetes.exceptions import KubernetesApiError
 from airflow.providers.cncf.kubernetes.utils.pod_manager import (
     AsyncPodManager,
+    PodCommandException,
     PodLogsConsumer,
     PodManager,
     PodPhase,
@@ -1225,6 +1226,46 @@ class TestPodManager:
         assert ret == xcom_json
         assert mock_exec_xcom_kill.call_count == 1
 
+    
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.kubernetes_stream")
+    
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.PodManager.extract_xcom_kill")
+    @mock.patch(
+        
"airflow.providers.cncf.kubernetes.utils.pod_manager.PodManager.container_is_running",
+        return_value=True,
+    )
+    def test_extract_xcom_returns_result_when_sidecar_kill_fails(
+        self, mock_container_is_running, mock_exec_xcom_kill, 
mock_kubernetes_stream
+    ):
+        """A failure to kill the sidecar must not discard the XCom value 
already read."""
+        xcom_json = """{"a": "true"}"""
+        mock_client = MagicMock()
+        mock_client.peek_stderr.return_value = ""
+        mock_client.read_all.return_value = xcom_json
+        mock_kubernetes_stream.return_value = mock_client
+        mock_exec_xcom_kill.side_effect = PodCommandException("Command failed 
with stderr: Permission denied")
+        ret = self.pod_manager.extract_xcom(pod=MagicMock())
+        assert ret == xcom_json
+        assert mock_exec_xcom_kill.call_count == 1
+
+    
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.kubernetes_stream")
+    
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.PodManager.extract_xcom_kill")
+    @mock.patch(
+        
"airflow.providers.cncf.kubernetes.utils.pod_manager.PodManager.container_is_running",
+        return_value=True,
+    )
+    def test_extract_xcom_reraises_kill_failure_when_not_ignored(
+        self, mock_container_is_running, mock_exec_xcom_kill, 
mock_kubernetes_stream
+    ):
+        """With ignore_kill_failure=False the kill failure still propagates to 
the caller."""
+        xcom_json = """{"a": "true"}"""
+        mock_client = MagicMock()
+        mock_client.peek_stderr.return_value = ""
+        mock_client.read_all.return_value = xcom_json
+        mock_kubernetes_stream.return_value = mock_client
+        mock_exec_xcom_kill.side_effect = PodCommandException("Command failed 
with stderr: Permission denied")
+        with pytest.raises(PodCommandException, match="Permission denied"):
+            self.pod_manager.extract_xcom(pod=MagicMock(), 
ignore_kill_failure=False)
+        assert mock_exec_xcom_kill.call_count == 1
+
     
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.kubernetes_stream")
     
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.PodManager.extract_xcom_kill")
     @mock.patch(

Reply via email to