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 81f09b0bee426..cf7ef8101376f 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 @@ -818,6 +818,19 @@ def __init__( self._event_polling_fallback = False self._config_loaded = False + def _uses_exec_auth(self, kubeconfig_data: dict) -> bool: + """ + Detect if kubeconfig uses exec-based authentication. + + Exec plugins return short-lived tokens (EKS, GKE, etc). + """ + users = kubeconfig_data.get("users", []) + for user in users: + user_auth = user.get("user", {}) + if "exec" in user_auth: + return True + return False + async def _load_config(self): """Load Kubernetes configuration once per hook instance.""" if self._config_loaded: @@ -846,39 +859,54 @@ async def _load_config(self): self._config_loaded = True return - # If above block does not return, we are not in a cluster. self._is_in_cluster = False - if self.config_dict: self.log.debug(LOADING_KUBE_CONFIG_FILE_RESOURCE.format("config dictionary")) - await async_config.load_kube_config_from_dict(self.config_dict, context=cluster_context) - self._config_loaded = True - return + await async_config.load_kube_config_from_dict( + self.config_dict, + context=cluster_context, + ) + + if not self._uses_exec_auth(self.config_dict): + self._config_loaded = True + + return if kubeconfig_path is not None: self.log.debug("loading kube_config from: %s", kubeconfig_path) + await async_config.load_kube_config( config_file=kubeconfig_path, client_configuration=self.client_configuration, context=cluster_context, ) - self._config_loaded = True - return + try: + async with aiofiles.open(kubeconfig_path) as f: + content = await f.read() + data = yaml.safe_load(content) + + if not self._uses_exec_auth(data): + self._config_loaded = True + 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 if kubeconfig is not None: async with aiofiles.tempfile.NamedTemporaryFile() as temp_config: - self.log.debug( - "Reading kubernetes configuration file from connection " - "object and writing temporary config file with its content", - ) if isinstance(kubeconfig, dict): - self.log.debug( - LOADING_KUBE_CONFIG_FILE_RESOURCE.format( - "connection kube_config dictionary (serializing)" - ) - ) - kubeconfig = json.dumps(kubeconfig) - await temp_config.write(kubeconfig.encode()) + kubeconfig_data = kubeconfig + kubeconfig_str = json.dumps(kubeconfig) + else: + kubeconfig_data = None + kubeconfig_str = kubeconfig + + await temp_config.write(kubeconfig_str.encode()) await temp_config.flush() await async_config.load_kube_config( @@ -886,10 +914,13 @@ async def _load_config(self): client_configuration=self.client_configuration, context=cluster_context, ) - self._config_loaded = True - return + if kubeconfig_data and not self._uses_exec_auth(kubeconfig_data): + self._config_loaded = True + + return self.log.debug(LOADING_KUBE_CONFIG_FILE_RESOURCE.format("default configuration file")) + await async_config.load_kube_config( client_configuration=self.client_configuration, context=cluster_context,