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 10b07b4519b Fix deferrable Kubernetes 401s with a default exec-based 
kubeconfig (#72300)
10b07b4519b is described below

commit 10b07b4519bb513b77af35ca5b81307f3d82c9a5
Author: rjgoyln <[email protected]>
AuthorDate: Mon Sep 21 06:10:59 2026 +0800

    Fix deferrable Kubernetes 401s with a default exec-based kubeconfig (#72300)
    
    * Fix deferrable Kubernetes 401s with a default exec-based kubeconfig
    
    Exec credential plugins such as "aws eks get-token" issue tokens that
    expire after roughly fifteen minutes, so a triggerer holding a cached
    kubeconfig starts getting 401 Unauthorized part-way through a long pod
    or job. Caching was already skipped for exec-based auth reached through
    config_dict, kube_config and kube_config_path, but a kubeconfig picked
    up from the default location was still cached unconditionally and kept
    the stale token for the lifetime of the hook.
    
    When KUBECONFIG genuinely names several files the merge rules that
    decide the active user belong to kubernetes_asyncio, so the config is
    left uncached rather than duplicating them here.
    
    * Say plainly when exec auth is assumed rather than detected
    
    The helpers return True both when a kubeconfig genuinely uses an exec
    credential plugin and when the active user could not be read at all, so
    a reader had no way to tell the two apart. The comment was narrower
    still: it explained only the several-files case, while the same branch
    also covers the default location resolving to no readable file.
    
    * Say what the path helper returns
    
    _resolve_default_kubeconfig_path returns a path and loads nothing.
    
    Generated-by: Claude Opus 5
    
    ---------
    
    Co-authored-by: rjgoyln <[email protected]>
---
 .../providers/cncf/kubernetes/hooks/kubernetes.py  |  63 ++++++++++---
 .../unit/cncf/kubernetes/hooks/test_kubernetes.py  | 101 +++++++++++++++++++++
 2 files changed, 150 insertions(+), 14 deletions(-)

diff --git 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/hooks/kubernetes.py
 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/hooks/kubernetes.py
index bccf6093ea0..8c911363860 100644
--- 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/hooks/kubernetes.py
+++ 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/hooks/kubernetes.py
@@ -19,6 +19,7 @@ from __future__ import annotations
 import asyncio
 import contextlib
 import json
+import os
 import tempfile
 from collections.abc import AsyncGenerator
 from functools import cached_property
@@ -868,6 +869,45 @@ class AsyncKubernetesHook(KubernetesHook):
         # fallback to check all users if active context user cannot be 
resolved; this is a safe fallback since it errs on the side of not caching
         return any("exec" in (u.get("user") or {}) for u in users)
 
+    async def _kubeconfig_file_uses_exec_auth(self, kubeconfig_path: str, 
context: str | None) -> bool:
+        """Detect exec auth in the kubeconfig at ``kubeconfig_path``; an 
unreadable file counts as exec."""
+        try:
+            async with aiofiles.open(kubeconfig_path) as f:
+                kubeconfig_data = yaml.safe_load(await f.read())
+            return self._uses_exec_auth(kubeconfig_data, context=context)
+        except Exception as exc:
+            self.log.warning(
+                "Error while parsing kube_config from %s to detect exec auth; "
+                "continuing without caching the config: %s",
+                kubeconfig_path,
+                exc,
+            )
+            return True
+
+    @staticmethod
+    def _resolve_default_kubeconfig_path() -> str | None:
+        """Return the path to the default-location kubeconfig, if it resolves 
to a single existing file."""
+        paths = [
+            os.path.expanduser(path)
+            for path in 
async_config.KUBE_CONFIG_DEFAULT_LOCATION.split(os.pathsep)
+            if path
+        ]
+        # kubernetes_asyncio skips the paths that do not exist, so only the 
surviving ones
+        # decide whether the active user can be read from a single file.
+        existing = [path for path in paths if os.path.exists(path)]
+        if len(existing) != 1:
+            return None
+        return existing[0]
+
+    async def _default_kubeconfig_uses_exec_auth(self, context: str | None) -> 
bool:
+        """Detect exec auth in the default kubeconfig; anything but one 
readable file counts as exec."""
+        kubeconfig_path = self._resolve_default_kubeconfig_path()
+        if kubeconfig_path is None:
+            # No single file to read: which user is active is 
kubernetes_asyncio's to decide,
+            # so leave the config uncached rather than reimplementing its 
merge rules here.
+            return True
+        return await self._kubeconfig_file_uses_exec_auth(kubeconfig_path, 
context)
+
     async def _load_config(self):
         """
         Load Kubernetes configuration.
@@ -932,19 +972,9 @@ class AsyncKubernetesHook(KubernetesHook):
             )
 
             if self._is_exec_auth is None:
-                try:
-                    async with aiofiles.open(kubeconfig_path) as f:
-                        content = await f.read()
-                        data = yaml.safe_load(content)
-                    self._is_exec_auth = self._uses_exec_auth(data, 
context=cluster_context)
-                except Exception as exc:
-                    self.log.warning(
-                        "Error while parsing kube_config from %s to detect 
exec auth; "
-                        "continuing without caching the config: %s",
-                        kubeconfig_path,
-                        exc,
-                    )
-                    self._is_exec_auth = True
+                self._is_exec_auth = await 
self._kubeconfig_file_uses_exec_auth(
+                    kubeconfig_path, cluster_context
+                )
 
             if not self._is_exec_auth:
                 self._config_loaded = True
@@ -993,7 +1023,12 @@ class AsyncKubernetesHook(KubernetesHook):
             client_configuration=self.client_configuration,
             context=cluster_context,
         )
-        self._config_loaded = True
+
+        if self._is_exec_auth is None:
+            self._is_exec_auth = await 
self._default_kubeconfig_uses_exec_auth(cluster_context)
+
+        if not self._is_exec_auth:
+            self._config_loaded = True
 
     async def get_conn_extras(self) -> dict:
         if self._extras is None:
diff --git 
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/hooks/test_kubernetes.py 
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/hooks/test_kubernetes.py
index 09d686821ca..ff9ec0b8dd3 100644
--- 
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/hooks/test_kubernetes.py
+++ 
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/hooks/test_kubernetes.py
@@ -1347,6 +1347,107 @@ class TestAsyncKubernetesHook:
         await hook._load_config()
         assert hook._config_loaded is expected_cached
 
+    @pytest.mark.asyncio
+    @pytest.mark.parametrize(
+        ("kubeconfig_content", "extra_default_paths", "expected_cached"),
+        [
+            pytest.param(
+                (
+                    "current-context: ctx1\n"
+                    "contexts:\n- name: ctx1\n  context:\n    user: user1\n"
+                    "users:\n- name: user1\n  user:\n    exec:\n      command: 
aws eks get-token\n"
+                ),
+                0,
+                False,
+                id="default_kubeconfig_with_exec",
+            ),
+            pytest.param(
+                (
+                    "current-context: ctx1\n"
+                    "contexts:\n- name: ctx1\n  context:\n    user: user1\n"
+                    "users:\n- name: user1\n  user:\n    token: static-token\n"
+                ),
+                0,
+                True,
+                id="default_kubeconfig_no_exec",
+            ),
+            pytest.param(
+                (
+                    "current-context: ctx1\n"
+                    "contexts:\n- name: ctx1\n  context:\n    user: user1\n"
+                    "users:\n- name: user1\n  user:\n    token: static-token\n"
+                ),
+                1,
+                False,
+                id="default_kubeconfig_merged_from_several_files",
+            ),
+        ],
+    )
+    
@mock.patch("airflow.providers.cncf.kubernetes.hooks.kubernetes.async_config.load_kube_config")
+    async def test_load_config_caching_behavior_default_kubeconfig(
+        self, mock_load_file, tmp_path, kubeconfig_content, 
extra_default_paths, expected_cached
+    ):
+        mock_load_file.return_value = None
+        kubeconfig_file = tmp_path / "config"
+        kubeconfig_file.write_text(kubeconfig_content)
+        extra_files = []
+        for index in range(extra_default_paths):
+            extra_file = tmp_path / f"extra{index}.yaml"
+            extra_file.write_text(kubeconfig_content)
+            extra_files.append(str(extra_file))
+        default_location = os.pathsep.join([str(kubeconfig_file), 
*extra_files])
+
+        hook = AsyncKubernetesHook(conn_id=None, in_cluster=False)
+        hook._get_field = mock.AsyncMock(return_value=None)
+        with 
mock.patch(f"{HOOK_MODULE}.async_config.KUBE_CONFIG_DEFAULT_LOCATION", 
default_location):
+            await hook._load_config()
+
+        mock_load_file.assert_awaited_once()
+        assert hook._config_loaded is expected_cached
+
+    @pytest.mark.asyncio
+    
@mock.patch("airflow.providers.cncf.kubernetes.hooks.kubernetes.async_config.load_kube_config")
+    async def 
test_load_config_default_kubeconfig_ignores_paths_that_do_not_exist(
+        self, mock_load_file, tmp_path
+    ):
+        mock_load_file.return_value = None
+        kubeconfig_file = tmp_path / "config"
+        kubeconfig_file.write_text(
+            "current-context: ctx1\n"
+            "contexts:\n- name: ctx1\n  context:\n    user: user1\n"
+            "users:\n- name: user1\n  user:\n    token: static-token\n"
+        )
+        default_location = os.pathsep.join([str(kubeconfig_file), str(tmp_path 
/ "missing.yaml")])
+
+        hook = AsyncKubernetesHook(conn_id=None, in_cluster=False)
+        hook._get_field = mock.AsyncMock(return_value=None)
+        with 
mock.patch(f"{HOOK_MODULE}.async_config.KUBE_CONFIG_DEFAULT_LOCATION", 
default_location):
+            await hook._load_config()
+
+        assert hook._config_loaded is True
+
+    @pytest.mark.asyncio
+    
@mock.patch("airflow.providers.cncf.kubernetes.hooks.kubernetes.async_config.load_kube_config")
+    async def test_load_config_default_kubeconfig_expands_home_directory(
+        self, mock_load_file, tmp_path, monkeypatch
+    ):
+        mock_load_file.return_value = None
+        kube_dir = tmp_path / ".kube"
+        kube_dir.mkdir()
+        (kube_dir / "config").write_text(
+            "current-context: ctx1\n"
+            "contexts:\n- name: ctx1\n  context:\n    user: user1\n"
+            "users:\n- name: user1\n  user:\n    token: static-token\n"
+        )
+        monkeypatch.setenv("HOME", str(tmp_path))
+
+        hook = AsyncKubernetesHook(conn_id=None, in_cluster=False)
+        hook._get_field = mock.AsyncMock(return_value=None)
+        with 
mock.patch(f"{HOOK_MODULE}.async_config.KUBE_CONFIG_DEFAULT_LOCATION", 
"~/.kube/config"):
+            await hook._load_config()
+
+        assert hook._config_loaded is True
+
     @pytest.mark.asyncio
     @mock.patch(KUBE_API.format("list_namespaced_event"))
     async def test_async_get_pod_events_with_resource_version(

Reply via email to