"""Small in-cluster Kubernetes API client used by the Pod backend.""" from __future__ import annotations import json from pathlib import Path from urllib.parse import quote from aiohttp import ClientResponseError, ClientSession, TCPConnector import ssl class KubernetesClient: def __init__( self, *, api_url: str, token_file: Path, ca_file: Path, namespace: str, ) -> None: self.api_url = api_url.rstrip("/") self.token_file = token_file self.ca_file = ca_file self.namespace = namespace self.session: ClientSession | None = None async def __aenter__(self) -> KubernetesClient: context = ssl.create_default_context(cafile=self.ca_file) self.session = ClientSession( connector=TCPConnector(ssl=context), headers={ "Authorization": f"Bearer {self.token_file.read_text().strip()}", }, raise_for_status=True, ) return self async def __aexit__(self, *_: object) -> None: if self.session is not None: await self.session.close() async def create_pod(self, manifest: dict[str, object]) -> dict[str, object]: response = await self._request("POST", self._pods_path(), json=manifest) return await response.json() async def get_pod(self, name: str) -> dict[str, object] | None: try: response = await self._request("GET", f"{self._pods_path()}/{quote(name)}") except ClientResponseError as error: if error.status == 404: return None raise return await response.json() async def bind_identity(self, name: str, identity_path: str) -> None: patch = { "metadata": { "labels": {"ci.ddupan.top/identity-bound": "true"}, "annotations": {"ci.ddupan.top/spiffe-path": identity_path}, } } response = await self._request( "PATCH", f"{self._pods_path()}/{quote(name)}", data=json.dumps(patch), headers={"Content-Type": "application/merge-patch+json"}, ) response.release() async def delete_pod(self, name: str) -> None: try: response = await self._request( "DELETE", f"{self._pods_path()}/{quote(name)}", json={"gracePeriodSeconds": 30, "propagationPolicy": "Background"}, ) response.release() except ClientResponseError as error: if error.status != 404: raise async def _request(self, method: str, path: str, **kwargs: object): if self.session is None: raise RuntimeError("KubernetesClient is not open") return await self.session.request(method, f"{self.api_url}{path}", **kwargs) def _pods_path(self) -> str: return f"/api/v1/namespaces/{quote(self.namespace)}/pods"