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(