Files
gitea-dynamic-runner/src/gitea_dynamic_runner/kubernetes.py
T
panxiao81 83eb87bcec
test / python (pull_request) Successful in 9s
test / shell (pull_request) Successful in 15s
重命名为动态 Runner 并记录协议调度路线
2026-09-16 14:51:06 +00:00

90 lines
2.9 KiB
Python

"""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"