From f3a199e7bad5d2240a0d6c0b6efc507713bb3dd7 Mon Sep 17 00:00:00 2001 From: panxiao81 Date: Fri, 18 Sep 2026 18:12:10 +0000 Subject: [PATCH] =?UTF-8?q?=E9=87=8D=E6=9E=84=E4=B8=BA=E7=9B=B4=E6=8E=A5?= =?UTF-8?q?=20OpenSandbox=20=E7=94=9F=E5=91=BD=E5=91=A8=E6=9C=9F=E8=B0=83?= =?UTF-8?q?=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 21 +- docs/opensandbox-runner.md | 81 ++-- pyproject.toml | 2 +- scripts/gitea-opensandbox-runner | 11 + src/gitea_dynamic_runner/controller.py | 165 ++++----- .../opensandbox_identity.py | 183 +++++++++ .../opensandbox_worker.py | 349 ++++++------------ .../sandbox_kubernetes.py | 16 + tests/test_opensandbox_identity.py | 53 +++ tests/test_opensandbox_worker.py | 140 +++---- 10 files changed, 556 insertions(+), 465 deletions(-) create mode 100644 src/gitea_dynamic_runner/opensandbox_identity.py create mode 100644 tests/test_opensandbox_identity.py diff --git a/README.md b/README.md index 0f10256..ec508f6 100644 --- a/README.md +++ b/README.md @@ -17,28 +17,23 @@ runs-on: [self-hosted, vm] 组件: -- `controller`:接收 Gitea `workflow_job` webhook,将指定 label 的 queued job - 发布到 NATS JetStream。 -- `worker`:领取任务、限制并发,并通过 Pod 或 microVM backend 创建一次性环境。 +- `controller`:接收 Gitea `workflow_job` webhook,直接调用 OpenSandbox Lifecycle + API,从 `ci-pod` 或 `ci-vm` Pool 创建一次性环境。 - `microvm-runner-launch`:为每个任务以 direct I/O 转换出 flat qcow2 root disk、创建 NoCloud seed 和 TAP,运行 Cloud Hypervisor,退出后完整清理。 - `guest-runner`:在 guest 中领取一次性 runner registration token,注册 ephemeral runner,执行一个 job 后关机。 -- `opensandbox-worker`:在正式 sandbox 集群通过 OpenSandbox `ci-vm` Pool 分配 - Kata guest,按实际 Pod UID 创建临时 SPIFFE entry,并在任务结束后回收两者。身份 - 与 Pool 契约见 [`docs/opensandbox-runner.md`](docs/opensandbox-runner.md)。 +- `opensandbox-identity`:在 sandbox 集群按实际 Pod UID 创建并清理临时 SPIFFE + entry;不持有 OpenSandbox API key、Gitea token 或 Bao 凭据。身份与 Pool 契约见 + [`docs/opensandbox-runner.md`](docs/opensandbox-runner.md)。 - `pod-worker`:在 Kubernetes 中创建一次性 privileged Pod;Pod 内的 workflow 使用 host executor,Docker、BuildKit 和 kind 等工具由 pipeline 按需 setup。Runner 固定在 支持原生 job hooks 的 3.x 版本,在 workflow 第一步前等待实际任务对应的 SVID。 - `jwt-broker`:早期共享 Kubernetes runner 的过渡实验;目标架构不部署它,每个 动态 Pod 或 VM 直接取得自己的 SPIFFE 身份。 -消息流使用一个 `WorkQueuePolicy` stream。相同 runner label 的所有 worker 共享同一 -durable consumer;扩容只需要增加 worker 或提高单机 capacity。 - -当前 webhook → NATS 流程是用于尽快验证 Pod/VM 生命周期的 bootstrap 实现,不是 -长期调度接口。长期目标是让 controller 作为兼容 Gitea Runner 协议的调度器直接注册、 -声明 labels、领取 task,并把已领取 task 交给 Pod/VM executor;路线与迁移边界见 +OpenSandbox 路径不使用 NATS。旧 Pod/microVM worker 只作为迁移期代码保留,不应重新 +部署。长期 RunnerService 协议路线见 [`docs/runner-protocol-roadmap.md`](docs/runner-protocol-roadmap.md)。 ## 开发 @@ -52,7 +47,7 @@ pytest ## 安全边界 -- NATS 密码、webhook secret 和 Gitea registration token 只从文件读取。 +- OpenSandbox API key、webhook secret 和 Gitea registration token 只从文件读取。 - registration token 不写入 seed image;worker 通过单次 nonce endpoint 交给 guest。 - guest 启动时从仅监听 microVM bridge 的 worker endpoint 获取固定版本 Runner 和配置 资产;基础镜像无需为 Runner 发布而重做。 diff --git a/docs/opensandbox-runner.md b/docs/opensandbox-runner.md index 783e968..b95c28d 100644 --- a/docs/opensandbox-runner.md +++ b/docs/opensandbox-runner.md @@ -1,55 +1,52 @@ -# OpenSandbox Kata runner +# OpenSandbox runner -The OpenSandbox worker replaces the legacy direct Cloud Hypervisor launcher for -jobs labelled `self-hosted, vm`. It runs in the sandbox Kubernetes cluster and -uses two local control planes: +Gitea 的 `workflow_job` webhook 只负责发现带 `self-hosted,pod` 或 +`self-hosted,vm` label 的 queued job。controller 不再把请求写入 NATS,而是携带由 +homelab ExternalSecret 挂载的 API key 直接调用 OpenSandbox Lifecycle API: -1. OpenSandbox Lifecycle API creates a sandbox from the `ci-vm` Pool. -2. Kubernetes API exposes the concrete BatchSandbox allocation and manages its - short-lived `ClusterStaticEntry`. +```text +http://10.60.0.13:8080/v1/sandboxes +``` -## Identity ordering +请求通过 `extensions.poolRef` 选择 `ci-pod` 或 `ci-vm`。地址是 VyOS HAProxy 的 +内网 TCP frontend,backend 为 sandbox 两个节点上的固定 NodePort;不通过公网或 +Cloudflare Tunnel。OpenSandbox key 只从 OpenBao `kv/k8s/opensandbox-api.api_key` +进入 homelab Secret,进程通过 `OPENSANDBOX_API_KEY_FILE` 读取。 -The worker sends the stable repository/job identity and a generated runner name -as task environment. Once OpenSandbox has allocated a Pool Pod, the worker reads -the Pod UID and creates an entry with: +## 身份顺序 -- parent: `spiffe://ddupan.top/spire/agent/k8s_psat/sandbox-kata/pod/`; -- workload: `spiffe://ddupan.top/ci///`; -- selector: `unix:uid:2000`. +Lifecycle 请求把稳定的 repository/task SPIFFE ID 放入 task environment。sandbox +集群内的 `opensandbox-identity` controller 读取 BatchSandbox allocation 得到实际 +Pod UID,然后创建: -The Pool must run the runner task as UID 2000 and set -`shareProcessNamespace: true`. Its guest-local SPIRE Agent uses a Pod-bound PSAT -and exposes the Workload API through the shared `spire-agent-socket` emptyDir. -The runner image starts through `gitea-opensandbox-runner`, which waits until the -exact expected SVID is available before it registers with Gitea. This prevents a -job from starting between Pod allocation and entry reconciliation. +- parent:`spiffe://ddupan.top/spire/agent/k8s_psat/sandbox-kata/pod/`; +- workload:`spiffe://ddupan.top/ci///`; +- selector:`unix:uid:2000`。 -The UID selector is the boundary between containers in the same Kata Pod. The -SPIRE Agent and privileged Docker daemon must not run as UID 2000. The runner may -access Docker only through a group-owned Unix socket. +job ID 只进入诊断 label,不进入业务身份。Pool 内 runner 进程固定使用 UID 2000, +Pod 设置 `shareProcessNamespace: true`;guest-local SPIRE Agent 使用 Pod-bound PSAT, +通过内存 emptyDir 暴露 Workload API。`gitea-opensandbox-runner` 等待精确 SVID,随后 +通过一次性 nonce URL 领取 Gitea registration token。token 不进入 Lifecycle 请求、 +BatchSandbox、镜像或 sandbox Secret。 -## Required Pool contract +identity controller 只能读取 opensandbox namespace 的 BatchSandbox/Pod,并维护带 +自身 label 的 ClusterStaticEntry。它不持有 OpenSandbox API key、Gitea token 或 Bao +凭据。BatchSandbox 消失后,对应 entry 在下一次 reconcile 删除;Pod 删除后,按 Pod +UID attestation 的临时 Agent 失去父级。 -The `ci-vm` Pool template owns infrastructure that callers cannot override in -Pool mode: +## Pool 契约 -- `runtimeClassName: kata-clh-runtime-rs` with block-backed emptyDir storage; -- runner image containing `gitea-opensandbox-runner` and SPIRE CLI; -- guest-local SPIRE Agent sidecar and projected audience `spire-server` token; -- `shareProcessNamespace: true`; -- runner UID 2000 and a distinct UID for every sidecar; -- ephemeral Gitea registration token delivery; -- Docker/BuildKit storage and socket entirely inside the Kata guest. +`ci-vm` 使用 `kata-clh-runtime-rs`;`ci-pod` 使用默认 runc。两者都要求: -The worker ServiceAccount needs read access to BatchSandboxes and Pods and -create/get/delete access to ClusterStaticEntries. OpenSandbox API credentials, -when enabled, are mounted from a Secret and read from -`OPENSANDBOX_API_KEY_FILE`. +- runner 镜像包含 Gitea Runner、Node.js action userspace、SPIRE CLI 和 identity gate; +- runner UID 2000,SPIRE Agent 与 privileged dockerd 使用不同 UID; +- Docker socket 通过 group 2000 共享,Docker 数据仅存在于 sandbox emptyDir; +- `self-hosted` 必须是所有 runner labels 的前缀; +- ephemeral/once runner 完成一项任务后退出。 -## Cleanup +## 清理与恢复 -On success, failure, timeout, or cancellation the worker deletes the -ClusterStaticEntry before deleting the sandbox. Both deletes accept `404`, so a -JetStream redelivery can safely repeat cleanup. The SPIRE Agent registration is -bound to the Pod UID and is removed by SPIRE after the Pod disappears. +controller 监控 Lifecycle 状态,在终止、失败、超时或取消时调用 DELETE。API delete、 +identity entry delete 均接受对象已不存在。controller 重启时,OpenSandbox timeout +仍是最终回收边界;后续可基于 metadata list 恢复主动监控,但不得为此重新引入消息 +队列。 diff --git a/pyproject.toml b/pyproject.toml index 40b8091..000dd86 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -16,7 +16,7 @@ test = ["pytest==8.4.2", "pytest-asyncio==1.2.0"] gitea-dynamic-runner-controller = "gitea_dynamic_runner.controller:main" gitea-dynamic-runner-pod-worker = "gitea_dynamic_runner.pod_worker:cli" gitea-dynamic-runner-vm-worker = "gitea_dynamic_runner.worker:cli" -gitea-dynamic-runner-opensandbox-worker = "gitea_dynamic_runner.opensandbox_worker:cli" +gitea-dynamic-runner-opensandbox-identity = "gitea_dynamic_runner.opensandbox_identity:cli" gitea-spire-jwt-broker = "gitea_dynamic_runner.jwt_broker:main" [tool.pytest.ini_options] diff --git a/scripts/gitea-opensandbox-runner b/scripts/gitea-opensandbox-runner index 9d40499..cf54f2a 100644 --- a/scripts/gitea-opensandbox-runner +++ b/scripts/gitea-opensandbox-runner @@ -3,8 +3,15 @@ set -euo pipefail : "${CI_SPIFFE_ID:?CI_SPIFFE_ID is required}" : "${SPIFFE_ENDPOINT_SOCKET:?SPIFFE_ENDPOINT_SOCKET is required}" +: "${GITEA_RUNNER_REGISTRATION_TOKEN_URL:?GITEA_RUNNER_REGISTRATION_TOKEN_URL is required}" socket_path=${SPIFFE_ENDPOINT_SOCKET#unix://} +token_file=$(mktemp) +cleanup() { + rm -f -- "$token_file" +} +trap cleanup EXIT +export GITEA_RUNNER_REGISTRATION_TOKEN_FILE="$token_file" deadline=$((SECONDS + 120)) while (( SECONDS < deadline )); do if /opt/spire/bin/spire-agent api fetch x509 \ @@ -12,6 +19,10 @@ while (( SECONDS < deadline )); do -output json 2>/dev/null | \ jq -e --arg id "$CI_SPIFFE_ID" \ 'any(.svids[]?; .spiffe_id == $id)' >/dev/null; then + umask 077 + curl --fail --silent --show-error --retry 10 --retry-all-errors \ + --connect-timeout 2 --max-time 30 \ + "$GITEA_RUNNER_REGISTRATION_TOKEN_URL" >"$token_file" exec /usr/local/bin/run.sh fi sleep 1 diff --git a/src/gitea_dynamic_runner/controller.py b/src/gitea_dynamic_runner/controller.py index 582c3e5..d464d04 100644 --- a/src/gitea_dynamic_runner/controller.py +++ b/src/gitea_dynamic_runner/controller.py @@ -1,63 +1,54 @@ #!/usr/bin/env python3 -"""Gitea workflow_job webhook to NATS JetStream producer.""" +"""Gitea workflow_job webhook to the OpenSandbox Lifecycle API.""" + +from __future__ import annotations import hashlib import hmac import json import logging import os -import ssl from pathlib import Path -import nats from aiohttp import web -from nats.js.api import DiscardPolicy, RetentionPolicy, StorageType, StreamConfig -from nats.js.errors import NotFoundError -from .models import IdentityBinding, RunnerRequest +from .models import RunnerRequest +from .opensandbox import OpenSandboxClient +from .opensandbox_worker import OpenSandboxScheduler, RegistrationTokens LOG = logging.getLogger(__name__) -SUBJECT_PREFIX = os.environ.get("NATS_SUBJECT_PREFIX", "ci.runner") -STREAM = os.environ.get("NATS_STREAM", "CI_RUNNER") -NATS_URL = os.environ.get("NATS_URL", "tls://nats.ad.ddupan.top:4222") -NATS_USER = os.environ.get("NATS_USER", "ci-producer") -NATS_PASSWORD_FILE = Path(os.environ.get("NATS_PASSWORD_FILE", "/run/secrets/nats/password")) -NATS_CA_FILE = os.environ.get("NATS_CA_FILE", "/etc/ssl/certs/ca-certificates.crt") -WEBHOOK_SECRET_FILE = Path(os.environ.get("WEBHOOK_SECRET_FILE", "/run/secrets/gitea/webhook-secret")) +WEBHOOK_SECRET_FILE = Path( + os.environ.get("WEBHOOK_SECRET_FILE", "/run/secrets/gitea/webhook-secret") +) +REGISTRATION_TOKEN_FILE = Path( + os.environ.get( + "REGISTRATION_TOKEN_FILE", "/run/secrets/gitea/registration-token" + ) +) +OPENSANDBOX_API = os.environ.get( + "OPENSANDBOX_API", "http://10.60.0.13:8080" +) +OPENSANDBOX_API_KEY_FILE = Path( + os.environ.get( + "OPENSANDBOX_API_KEY_FILE", + "/run/secrets/opensandbox/api-key", + ) +) def accepts(payload: object) -> tuple[bool, str | None]: - """Compatibility helper for callers that only need acceptance and identity.""" request = RunnerRequest.from_webhook(payload) return (request is not None, str(request.job_id) if request else None) def valid_signature(body: bytes, signature: str) -> bool: - expected = hmac.new(WEBHOOK_SECRET_FILE.read_bytes().strip(), body, hashlib.sha256).hexdigest() + expected = hmac.new( + WEBHOOK_SECRET_FILE.read_bytes().strip(), body, hashlib.sha256 + ).hexdigest() return hmac.compare_digest(signature.removeprefix("sha256="), expected) -async def ensure_stream(js: object) -> None: - config = StreamConfig( - name=STREAM, - subjects=[f"{SUBJECT_PREFIX}.>"], - retention=RetentionPolicy.WORK_QUEUE, - storage=StorageType.FILE, - discard=DiscardPolicy.OLD, - max_age=24 * 60 * 60, - max_msgs=10_000, - max_bytes=256 * 1024 * 1024, - duplicate_window=24 * 60 * 60, - ) - try: - await js.stream_info(STREAM) - except NotFoundError: - await js.add_stream(config=config) - else: - await js.update_stream(config=config) - - async def webhook(request: web.Request) -> web.Response: body = await request.read() if not valid_signature(body, request.headers.get("X-Gitea-Signature", "")): @@ -70,49 +61,28 @@ async def webhook(request: web.Request) -> web.Response: LOG.info("workflow_job webhook received %s", context) runner_request = RunnerRequest.from_webhook(payload) - if runner_request is not None: - subject = f"{SUBJECT_PREFIX}.{runner_request.backend}" - await request.app["js"].publish( - subject, - runner_request.to_json(), - headers={"Nats-Msg-Id": f"gitea-workflow-job-{runner_request.job_id}-queued"}, - ) - LOG.info( - "runner request published subject=%s job_id=%d run_id=%d backend=%s " - "repository=%s job_name=%r labels=%s", - subject, - runner_request.job_id, - runner_request.run_id, - runner_request.backend, - runner_request.repository, - runner_request.job_name, - runner_request.labels, - ) - return web.Response(status=202, text="queued\n") + if runner_request is None: + LOG.info("workflow_job webhook ignored %s", context) + return web.Response(status=204) - binding = IdentityBinding.from_webhook(payload) - if binding is not None: - subject = f"{SUBJECT_PREFIX}.{binding.backend}.binding" - await request.app["js"].publish( - subject, - binding.to_json(), - headers={"Nats-Msg-Id": f"gitea-workflow-job-{binding.job_id}-in-progress"}, - ) - LOG.info( - "identity binding published subject=%s job_id=%d run_id=%d backend=%s " - "runner_name=%s repository=%s job_name=%r", - subject, - binding.job_id, - binding.run_id, - binding.backend, - binding.runner_name, - binding.repository, - binding.job_name, - ) - return web.Response(status=202, text="binding queued\n") + scheduler: OpenSandboxScheduler = request.app["scheduler"] + try: + sandbox_id = await scheduler.create(runner_request) + except ValueError as error: + LOG.info("duplicate workflow_job webhook ignored: %s", error) + return web.Response(status=202, text="already scheduled\n") + except Exception: + LOG.exception("failed to create OpenSandbox runner %s", context) + raise web.HTTPServiceUnavailable(text="sandbox unavailable\n") + return web.Response(status=202, text=f"sandbox={sandbox_id}\n") - LOG.warning("workflow_job webhook ignored %s", context) - return web.Response(status=204) + +async def registration_token(request: web.Request) -> web.Response: + scheduler: OpenSandboxScheduler = request.app["scheduler"] + value = await scheduler.tokens.consume(request.match_info["nonce"]) + if value is None: + raise web.HTTPNotFound() + return web.Response(body=value, headers={"Cache-Control": "no-store"}) def _webhook_context(payload: object) -> str: @@ -131,37 +101,42 @@ def _webhook_context(payload: object) -> str: async def health(request: web.Request) -> web.Response: - connected = request.app["nc"].is_connected - return web.Response(text="ok\n" if connected else "disconnected\n", status=200 if connected else 503) - - -async def nats_context(app: web.Application): - tls = ssl.create_default_context(cafile=NATS_CA_FILE) - nc = await nats.connect( - NATS_URL, - user=NATS_USER, - password=NATS_PASSWORD_FILE.read_text().strip(), - tls=tls, - name="microvm-runner-controller", + client: OpenSandboxClient = request.app["opensandbox_client"] + return web.Response( + text="ok\n" if client.session is not None else "disconnected\n", + status=200 if client.session is not None else 503, ) - app["nc"] = nc - app["js"] = nc.jetstream() - await ensure_stream(app["js"]) - yield - await nc.drain() + + +async def opensandbox_context(app: web.Application): + tokens = RegistrationTokens(REGISTRATION_TOKEN_FILE.read_bytes()) + async with OpenSandboxClient( + api_url=OPENSANDBOX_API, + api_key_file=OPENSANDBOX_API_KEY_FILE, + ) as client: + scheduler = OpenSandboxScheduler(client, tokens) + app["opensandbox_client"] = client + app["scheduler"] = scheduler + yield + await scheduler.close() def create_app() -> web.Application: app = web.Application(client_max_size=1024 * 1024) - app.cleanup_ctx.append(nats_context) + app.cleanup_ctx.append(opensandbox_context) app.router.add_post("/webhook", webhook) + app.router.add_get("/token/{nonce}", registration_token) app.router.add_get("/healthz", health) return app def main() -> None: logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO")) - web.run_app(create_app(), host=os.environ.get("LISTEN", "0.0.0.0"), port=int(os.environ.get("PORT", "8787"))) + web.run_app( + create_app(), + host=os.environ.get("LISTEN", "0.0.0.0"), + port=int(os.environ.get("PORT", "8787")), + ) if __name__ == "__main__": diff --git a/src/gitea_dynamic_runner/opensandbox_identity.py b/src/gitea_dynamic_runner/opensandbox_identity.py new file mode 100644 index 0000000..a8b7af6 --- /dev/null +++ b/src/gitea_dynamic_runner/opensandbox_identity.py @@ -0,0 +1,183 @@ +"""Reconcile OpenSandbox Pod UIDs to narrowly scoped SPIRE entries.""" + +from __future__ import annotations + +import asyncio +import hashlib +import logging +import os +from pathlib import Path + +from aiohttp import ClientResponseError + +from .sandbox_kubernetes import SandboxKubernetesClient, allocated_pod_name + + +LOG = logging.getLogger(__name__) +KUBERNETES_API = os.environ.get("KUBERNETES_API", "https://kubernetes.default.svc") +KUBERNETES_TOKEN_FILE = Path( + os.environ.get( + "KUBERNETES_TOKEN_FILE", + "/var/run/secrets/kubernetes.io/serviceaccount/token", + ) +) +KUBERNETES_CA_FILE = Path( + os.environ.get( + "KUBERNETES_CA_FILE", + "/var/run/secrets/kubernetes.io/serviceaccount/ca.crt", + ) +) +NAMESPACE = os.environ.get("OPENSANDBOX_NAMESPACE", "opensandbox") +SPIFFE_TRUST_DOMAIN = os.environ.get("SPIFFE_TRUST_DOMAIN", "ddupan.top") +SPIRE_CLUSTER_NAME = os.environ.get("SPIRE_CLUSTER_NAME", "sandbox-kata") +SPIRE_CLASS_NAME = os.environ.get("SPIRE_CLASS_NAME", "spire-mgmt-spire") +RUNNER_UID = int(os.environ.get("RUNNER_UID", "2000")) +RECONCILE_INTERVAL = int(os.environ.get("RECONCILE_INTERVAL", "2")) + + +def entry_name(sandbox_id: str) -> str: + suffix = hashlib.sha256(sandbox_id.encode()).hexdigest()[:12] + return f"gitea-ci-{suffix}" + + +def sandbox_name(document: dict[str, object]) -> str | None: + metadata = document.get("metadata") + value = metadata.get("name") if isinstance(metadata, dict) else None + return value if isinstance(value, str) and value else None + + +def is_runner_sandbox(document: dict[str, object]) -> bool: + metadata = document.get("metadata") + labels = metadata.get("labels") if isinstance(metadata, dict) else None + return isinstance(labels, dict) and labels.get("ci.ddupan.top/runner") == "true" + + +def task_environment(document: dict[str, object]) -> dict[str, str]: + spec = document.get("spec") + task = spec.get("taskTemplate") if isinstance(spec, dict) else None + task_spec = task.get("spec") if isinstance(task, dict) else None + process = task_spec.get("process") if isinstance(task_spec, dict) else None + values = process.get("env") if isinstance(process, dict) else None + result: dict[str, str] = {} + if not isinstance(values, list): + return result + for item in values: + if not isinstance(item, dict): + continue + name, value = item.get("name"), item.get("value") + if isinstance(name, str) and isinstance(value, str): + result[name] = value + return result + + +def identity_entry( + *, sandbox_id: str, pod_uid: str, spiffe_id: str +) -> dict[str, object]: + expected_prefix = f"spiffe://{SPIFFE_TRUST_DOMAIN}/ci/" + if not spiffe_id.startswith(expected_prefix): + raise ValueError("runner SPIFFE ID is outside the CI namespace") + return { + "apiVersion": "spire.spiffe.io/v1alpha1", + "kind": "ClusterStaticEntry", + "metadata": { + "name": entry_name(sandbox_id), + "labels": { + "app.kubernetes.io/name": "gitea-dynamic-runner", + "app.kubernetes.io/component": "opensandbox-identity", + "ci.ddupan.top/sandbox-id": sandbox_id, + }, + }, + "spec": { + "className": SPIRE_CLASS_NAME, + "parentID": ( + f"spiffe://{SPIFFE_TRUST_DOMAIN}/spire/agent/k8s_psat/" + f"{SPIRE_CLUSTER_NAME}/pod/{pod_uid}" + ), + "spiffeID": spiffe_id, + "selectors": [f"unix:uid:{RUNNER_UID}"], + }, + } + + +async def reconcile(client: SandboxKubernetesClient) -> None: + sandboxes = await client.list_batchsandboxes() + live_names = { + name + for document in sandboxes + if is_runner_sandbox(document) and (name := sandbox_name(document)) + } + entries = await client.list_entries() + existing = { + name + for document in entries + if (name := sandbox_name(document)) is not None + } + + for document in sandboxes: + if not is_runner_sandbox(document): + continue + name = sandbox_name(document) + pod_name = allocated_pod_name(document) + spiffe_id = task_environment(document).get("CI_SPIFFE_ID") + if not name or not pod_name or not spiffe_id or entry_name(name) in existing: + continue + pod = await client.get_pod(pod_name) + metadata = pod.get("metadata") if isinstance(pod, dict) else None + pod_uid = metadata.get("uid") if isinstance(metadata, dict) else None + if not isinstance(pod_uid, str) or not pod_uid: + continue + try: + await client.create_entry( + identity_entry( + sandbox_id=name, + pod_uid=pod_uid, + spiffe_id=spiffe_id, + ) + ) + except ClientResponseError as error: + if error.status != 409: + raise + LOG.info( + "SPIRE entry ready sandbox=%s pod=%s pod_uid=%s spiffe_id=%s", + name, + pod_name, + pod_uid, + spiffe_id, + ) + + for document in entries: + metadata = document.get("metadata") + labels = metadata.get("labels") if isinstance(metadata, dict) else None + sandbox_id = labels.get("ci.ddupan.top/sandbox-id") if isinstance(labels, dict) else None + name = metadata.get("name") if isinstance(metadata, dict) else None + if ( + isinstance(name, str) + and isinstance(sandbox_id, str) + and sandbox_id not in live_names + ): + await client.delete_entry(name) + LOG.info("removed stale SPIRE entry=%s sandbox=%s", name, sandbox_id) + + +async def main() -> None: + async with SandboxKubernetesClient( + api_url=KUBERNETES_API, + token_file=KUBERNETES_TOKEN_FILE, + ca_file=KUBERNETES_CA_FILE, + namespace=NAMESPACE, + ) as client: + while True: + try: + await reconcile(client) + except Exception: + LOG.exception("OpenSandbox identity reconcile failed") + await asyncio.sleep(RECONCILE_INTERVAL) + + +def cli() -> None: + logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO")) + asyncio.run(main()) + + +if __name__ == "__main__": + cli() diff --git a/src/gitea_dynamic_runner/opensandbox_worker.py b/src/gitea_dynamic_runner/opensandbox_worker.py index deb407d..88baa55 100644 --- a/src/gitea_dynamic_runner/opensandbox_worker.py +++ b/src/gitea_dynamic_runner/opensandbox_worker.py @@ -1,68 +1,51 @@ -#!/usr/bin/env python3 -"""JetStream worker that provisions disposable OpenSandbox Kata runners.""" +"""Direct OpenSandbox lifecycle scheduler used by the webhook controller.""" from __future__ import annotations import asyncio -import hashlib -import json import logging import os -from pathlib import Path -import ssl +import secrets import uuid - -import nats -from nats.errors import TimeoutError -from nats.js.api import AckPolicy, ConsumerConfig +from collections.abc import Callable from .models import RunnerRequest from .opensandbox import OpenSandboxClient -from .pod_worker import heartbeat, identity_path -from .sandbox_kubernetes import SandboxKubernetesClient, allocated_pod_name +from .pod_worker import identity_path LOG = logging.getLogger(__name__) -STREAM = os.environ.get("NATS_STREAM", "CI_RUNNER") -REQUEST_SUBJECT = os.environ.get("NATS_SUBJECT", "ci.runner.vm") -NATS_URL = os.environ.get("NATS_URL", "tls://nats.ad.ddupan.top:4222") -NATS_USER = os.environ.get("NATS_USER", "ci-worker") -NATS_PASSWORD_FILE = Path(os.environ.get("NATS_PASSWORD_FILE", "/run/secrets/nats/password")) -NATS_CA_FILE = os.environ.get("NATS_CA_FILE", "/etc/ssl/certs/ca-certificates.crt") -KUBERNETES_API = os.environ.get("KUBERNETES_API", "https://kubernetes.default.svc") -KUBERNETES_TOKEN_FILE = Path(os.environ.get("KUBERNETES_TOKEN_FILE", "/var/run/secrets/kubernetes.io/serviceaccount/token")) -KUBERNETES_CA_FILE = Path(os.environ.get("KUBERNETES_CA_FILE", "/var/run/secrets/kubernetes.io/serviceaccount/ca.crt")) -NAMESPACE = os.environ.get("OPENSANDBOX_NAMESPACE", "opensandbox") -OPENSANDBOX_API = os.environ.get( - "OPENSANDBOX_API", "http://opensandbox-server.opensandbox-system.svc" -) -OPENSANDBOX_API_KEY_FILE_VALUE = os.environ.get("OPENSANDBOX_API_KEY_FILE") -OPENSANDBOX_API_KEY_FILE = Path(OPENSANDBOX_API_KEY_FILE_VALUE) if OPENSANDBOX_API_KEY_FILE_VALUE else None -OPENSANDBOX_POOL = os.environ.get("OPENSANDBOX_POOL", "ci-vm") SANDBOX_TIMEOUT = int(os.environ.get("RUNNER_SANDBOX_TIMEOUT", str(4 * 60 * 60))) -ALLOCATION_TIMEOUT = int(os.environ.get("RUNNER_ALLOCATION_TIMEOUT", "120")) -CAPACITY = int(os.environ.get("RUNNER_CAPACITY", "2")) SPIFFE_TRUST_DOMAIN = os.environ.get("SPIFFE_TRUST_DOMAIN", "ddupan.top") -SPIRE_CLUSTER_NAME = os.environ.get("SPIRE_CLUSTER_NAME", "sandbox-kata") -SPIRE_CLASS_NAME = os.environ.get("SPIRE_CLASS_NAME", "spire-mgmt-spire") -RUNNER_UID = int(os.environ.get("RUNNER_UID", "2000")) +TOKEN_BASE_URL = os.environ.get( + "RUNNER_TOKEN_BASE_URL", "http://192.168.10.127:8787/token" +).rstrip("/") -def sandbox_request(request: RunnerRequest, runner_name: str) -> dict[str, object]: - if request.backend != "vm": - raise ValueError("OpenSandbox backend only accepts vm requests") +def sandbox_request( + request: RunnerRequest, + runner_name: str, + registration_token_url: str, +) -> dict[str, object]: path = identity_path(request.repository, request.job_name) return { - "pool": OPENSANDBOX_POOL, + "pool": f"ci-{request.backend}", "timeout": SANDBOX_TIMEOUT, "entrypoint": ["/usr/local/libexec/gitea-opensandbox-runner"], "env": { "GITEA_RUNNER_NAME": runner_name, - "GITEA_RUNNER_LABELS": "self-hosted:host,vm:host", + "GITEA_RUNNER_LABELS": f"self-hosted:host,{request.backend}:host", + "GITEA_RUNNER_REGISTRATION_TOKEN_URL": registration_token_url, + "GITEA_RUNNER_EPHEMERAL": "1", + "GITEA_RUNNER_ONCE": "1", + "CONFIG_FILE": "/etc/gitea-runner/config.yaml", "CI_SPIFFE_ID": f"spiffe://{SPIFFE_TRUST_DOMAIN}/ci/{path}", - "SPIFFE_ENDPOINT_SOCKET": "unix:///run/spire/agent-sockets/spire-agent.sock", + "SPIFFE_ENDPOINT_SOCKET": ( + "unix:///run/spire/agent-sockets/spire-agent.sock" + ), }, "metadata": { + "ci.ddupan.top/runner": "true", "ci.ddupan.top/job-id": str(request.job_id), "ci.ddupan.top/run-id": str(request.run_id), "ci.ddupan.top/runner-name": runner_name, @@ -70,207 +53,111 @@ def sandbox_request(request: RunnerRequest, runner_name: str) -> dict[str, objec } -def entry_name(sandbox_id: str) -> str: - suffix = hashlib.sha256(sandbox_id.encode()).hexdigest()[:12] - return f"gitea-ci-{suffix}" +class RegistrationTokens: + """Single-use registration-token URLs; values never enter Sandbox CRs.""" + + def __init__(self, token: bytes) -> None: + self._token = token.strip() + self._values: dict[str, bytes] = {} + self._lock = asyncio.Lock() + + async def issue(self) -> tuple[str, str]: + nonce = secrets.token_urlsafe(32) + async with self._lock: + self._values[nonce] = self._token + return nonce, f"{TOKEN_BASE_URL}/{nonce}" + + async def consume(self, nonce: str) -> bytes | None: + async with self._lock: + return self._values.pop(nonce, None) + + async def revoke(self, nonce: str) -> None: + async with self._lock: + self._values.pop(nonce, None) -def identity_entry( - request: RunnerRequest, - *, - sandbox_id: str, - pod_uid: str, -) -> dict[str, object]: - path = identity_path(request.repository, request.job_name) - return { - "apiVersion": "spire.spiffe.io/v1alpha1", - "kind": "ClusterStaticEntry", - "metadata": { - "name": entry_name(sandbox_id), - "labels": { - "app.kubernetes.io/name": "gitea-dynamic-runner", - "app.kubernetes.io/component": "opensandbox-identity", - "ci.ddupan.top/sandbox-id": sandbox_id, - }, - }, - "spec": { - "className": SPIRE_CLASS_NAME, - "parentID": ( - f"spiffe://{SPIFFE_TRUST_DOMAIN}/spire/agent/k8s_psat/" - f"{SPIRE_CLUSTER_NAME}/pod/{pod_uid}" - ), - "spiffeID": f"spiffe://{SPIFFE_TRUST_DOMAIN}/ci/{path}", - "selectors": [f"unix:uid:{RUNNER_UID}"], - }, - } +class OpenSandboxScheduler: + def __init__( + self, + client: OpenSandboxClient, + tokens: RegistrationTokens, + *, + on_finished: Callable[[int], None] | None = None, + ) -> None: + self.client = client + self.tokens = tokens + self.on_finished = on_finished + self.active: dict[int, asyncio.Task[None]] = {} - -async def wait_for_allocation( - client: SandboxKubernetesClient, sandbox_id: str -) -> tuple[str, str]: - deadline = asyncio.get_running_loop().time() + ALLOCATION_TIMEOUT - while asyncio.get_running_loop().time() < deadline: - batchsandbox = await client.get_batchsandbox(sandbox_id) - if batchsandbox is not None and (pod_name := allocated_pod_name(batchsandbox)): - pod = await client.get_pod(pod_name) - if pod is not None: - metadata = pod.get("metadata") - pod_uid = metadata.get("uid") if isinstance(metadata, dict) else None - if isinstance(pod_uid, str) and pod_uid: - return pod_name, pod_uid - await asyncio.sleep(1) - raise asyncio.TimeoutError(f"OpenSandbox {sandbox_id} allocation timed out") - - -async def wait_for_sandbox(client: OpenSandboxClient, sandbox_id: str) -> bool: - deadline = asyncio.get_running_loop().time() + SANDBOX_TIMEOUT - while asyncio.get_running_loop().time() < deadline: - sandbox = await client.get(sandbox_id) - if sandbox is None: - raise RuntimeError(f"OpenSandbox {sandbox_id} disappeared") - status = sandbox.get("status") - state = status.get("state") if isinstance(status, dict) else None - if state == "Terminated": - return True - if state == "Failed": - return False - await asyncio.sleep(2) - raise asyncio.TimeoutError(f"OpenSandbox {sandbox_id} timed out") - - -async def run_request( - message: object, - sandbox_client: OpenSandboxClient, - kubernetes_client: SandboxKubernetesClient, -) -> None: - try: - request = RunnerRequest.from_json(message.data) - except (json.JSONDecodeError, UnicodeDecodeError, ValueError) as error: - LOG.error("discarding invalid OpenSandbox request: %s", error) - await message.ack() - return - if request.backend != "vm": - LOG.error("discarding %s request received by OpenSandbox worker", request.backend) - await message.ack() - return - - runner_name = f"gitea-vm-{uuid.uuid4().hex[:12]}" - sandbox_id: str | None = None - static_entry: str | None = None - stop = asyncio.Event() - pulse = asyncio.create_task(heartbeat(message, stop)) - try: - create = await sandbox_client.create(**sandbox_request(request, runner_name)) - candidate = create.get("id") - if not isinstance(candidate, str) or not candidate: - raise RuntimeError("OpenSandbox create response has no id") - sandbox_id = candidate - pod_name, pod_uid = await wait_for_allocation(kubernetes_client, sandbox_id) - manifest = identity_entry(request, sandbox_id=sandbox_id, pod_uid=pod_uid) - static_entry = entry_name(sandbox_id) - await kubernetes_client.create_entry(manifest) - LOG.info( - "OpenSandbox identity created sandbox=%s pod=%s pod_uid=%s entry=%s runner=%s job_id=%d", - sandbox_id, pod_name, pod_uid, static_entry, runner_name, request.job_id, - ) - if await wait_for_sandbox(sandbox_client, sandbox_id): - await message.ack() - else: - await message.nak(delay=30) - except Exception: - LOG.exception("OpenSandbox runner failed sandbox=%s job_id=%d", sandbox_id, request.job_id) - await message.nak(delay=30) - raise - finally: - stop.set() - await pulse + async def create(self, request: RunnerRequest) -> str: + if request.job_id in self.active: + raise ValueError(f"job {request.job_id} already has an active sandbox") + runner_name = f"gitea-{request.backend}-{uuid.uuid4().hex[:12]}" + nonce, token_url = await self.tokens.issue() try: - if static_entry is not None: - await kubernetes_client.delete_entry(static_entry) - finally: - if sandbox_id is not None: - await sandbox_client.delete(sandbox_id) - - -async def consume_requests( - subscription: object, - sandbox_client: OpenSandboxClient, - kubernetes_client: SandboxKubernetesClient, -) -> None: - active: set[asyncio.Task[None]] = set() - while True: - active = {task for task in active if not task.done()} - free = CAPACITY - len(active) - if free < 1: - await asyncio.wait(active, return_when=asyncio.FIRST_COMPLETED) - continue - try: - messages = await subscription.fetch(batch=free, timeout=5) - except TimeoutError: - continue - for message in messages: - task = asyncio.create_task( - run_request(message, sandbox_client, kubernetes_client) + response = await self.client.create( + **sandbox_request(request, runner_name, token_url) ) - task.add_done_callback(_report_task) - active.add(task) + sandbox_id = response.get("id") + if not isinstance(sandbox_id, str) or not sandbox_id: + raise RuntimeError("OpenSandbox create response has no id") + except Exception: + await self.tokens.revoke(nonce) + raise - -def _report_task(task: asyncio.Task[None]) -> None: - if not task.cancelled() and (error := task.exception()) is not None: - LOG.error( - "OpenSandbox runner task failed", - exc_info=(type(error), error, error.__traceback__), + task = asyncio.create_task( + self._monitor(request, sandbox_id, nonce), + name=f"opensandbox-{sandbox_id}", ) + task.add_done_callback(self._report) + self.active[request.job_id] = task + LOG.info( + "OpenSandbox runner created sandbox=%s runner=%s job_id=%d " + "repository=%s job_name=%r pool=ci-%s", + sandbox_id, + runner_name, + request.job_id, + request.repository, + request.job_name, + request.backend, + ) + return sandbox_id - -async def main() -> None: - if CAPACITY < 1: - raise ValueError("RUNNER_CAPACITY must be at least 1") - context = ssl.create_default_context(cafile=NATS_CA_FILE) - nc = await nats.connect( - NATS_URL, - user=NATS_USER, - password=NATS_PASSWORD_FILE.read_text().strip(), - tls=context, - name="opensandbox-runner-worker", - ) - js = nc.jetstream() - subscription = await js.pull_subscribe( - REQUEST_SUBJECT, - durable="vm", - stream=STREAM, - config=ConsumerConfig( - durable_name="vm", - filter_subject=REQUEST_SUBJECT, - ack_policy=AckPolicy.EXPLICIT, - ack_wait=5 * 60, - max_ack_pending=max(CAPACITY, 1), - max_deliver=5, - ), - ) - async with ( - OpenSandboxClient( - api_url=OPENSANDBOX_API, - api_key_file=OPENSANDBOX_API_KEY_FILE, - ) as sandbox_client, - SandboxKubernetesClient( - api_url=KUBERNETES_API, - token_file=KUBERNETES_TOKEN_FILE, - ca_file=KUBERNETES_CA_FILE, - namespace=NAMESPACE, - ) as kubernetes_client, - ): + async def _monitor( + self, request: RunnerRequest, sandbox_id: str, nonce: str + ) -> None: try: - await consume_requests(subscription, sandbox_client, kubernetes_client) + deadline = asyncio.get_running_loop().time() + SANDBOX_TIMEOUT + while asyncio.get_running_loop().time() < deadline: + sandbox = await self.client.get(sandbox_id) + if sandbox is None: + return + status = sandbox.get("status") + state = status.get("state") if isinstance(status, dict) else None + if state in {"Terminated", "Failed"}: + return + await asyncio.sleep(2) + raise asyncio.TimeoutError(f"OpenSandbox {sandbox_id} timed out") finally: - await nc.drain() + await self.tokens.revoke(nonce) + try: + await self.client.delete(sandbox_id) + finally: + self.active.pop(request.job_id, None) + if self.on_finished is not None: + self.on_finished(request.job_id) + @staticmethod + def _report(task: asyncio.Task[None]) -> None: + if not task.cancelled() and (error := task.exception()) is not None: + LOG.error( + "OpenSandbox lifecycle task failed", + exc_info=(type(error), error, error.__traceback__), + ) -def cli() -> None: - logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO")) - asyncio.run(main()) - - -if __name__ == "__main__": - cli() + async def close(self) -> None: + tasks = list(self.active.values()) + for task in tasks: + task.cancel() + if tasks: + await asyncio.gather(*tasks, return_exceptions=True) diff --git a/src/gitea_dynamic_runner/sandbox_kubernetes.py b/src/gitea_dynamic_runner/sandbox_kubernetes.py index 8a6551b..bfbf1c1 100644 --- a/src/gitea_dynamic_runner/sandbox_kubernetes.py +++ b/src/gitea_dynamic_runner/sandbox_kubernetes.py @@ -51,6 +51,12 @@ class SandboxKubernetesClient: raise return await response.json() + async def list_batchsandboxes(self) -> list[dict[str, object]]: + response = await self._request("GET", self._batchsandboxes_path()) + document = await response.json() + items = document.get("items") if isinstance(document, dict) else None + return [item for item in items if isinstance(item, dict)] if isinstance(items, list) else [] + async def get_pod(self, name: str) -> dict[str, object] | None: try: response = await self._request( @@ -68,6 +74,16 @@ class SandboxKubernetesClient: ) return await response.json() + async def list_entries(self) -> list[dict[str, object]]: + response = await self._request( + "GET", + f"{self._entries_path()}?labelSelector=" + "app.kubernetes.io%2Fcomponent%3Dopensandbox-identity", + ) + document = await response.json() + items = document.get("items") if isinstance(document, dict) else None + return [item for item in items if isinstance(item, dict)] if isinstance(items, list) else [] + async def get_entry(self, name: str) -> dict[str, object] | None: try: response = await self._request( diff --git a/tests/test_opensandbox_identity.py b/tests/test_opensandbox_identity.py new file mode 100644 index 0000000..34ee28b --- /dev/null +++ b/tests/test_opensandbox_identity.py @@ -0,0 +1,53 @@ +import json + +from gitea_dynamic_runner.opensandbox_identity import ( + identity_entry, + is_runner_sandbox, + task_environment, +) +from gitea_dynamic_runner.sandbox_kubernetes import allocated_pod_name + + +def batchsandbox(): + return { + "metadata": { + "name": "sandbox-1", + "labels": {"ci.ddupan.top/runner": "true"}, + "annotations": { + "sandbox.opensandbox.io/alloc-status": json.dumps( + {"pods": ["ci-vm-pod-1"]} + ) + }, + }, + "spec": { + "taskTemplate": { + "spec": { + "process": { + "env": [ + { + "name": "CI_SPIFFE_ID", + "value": "spiffe://ddupan.top/ci/org/repo/test", + } + ] + } + } + } + }, + } + + +def test_extracts_allocation_and_identity_environment(): + document = batchsandbox() + assert is_runner_sandbox(document) + assert allocated_pod_name(document) == "ci-vm-pod-1" + assert task_environment(document)["CI_SPIFFE_ID"].endswith("/org/repo/test") + + +def test_entry_binds_exact_pod_agent_and_runner_uid(): + entry = identity_entry( + sandbox_id="sandbox-1", + pod_uid="pod-uid", + spiffe_id="spiffe://ddupan.top/ci/org/repo/test", + ) + assert entry["spec"]["parentID"].endswith("/sandbox-kata/pod/pod-uid") + assert entry["spec"]["selectors"] == ["unix:uid:2000"] diff --git a/tests/test_opensandbox_worker.py b/tests/test_opensandbox_worker.py index 9fef46e..f623df2 100644 --- a/tests/test_opensandbox_worker.py +++ b/tests/test_opensandbox_worker.py @@ -1,114 +1,88 @@ import asyncio -import json import pytest from gitea_dynamic_runner.models import RunnerRequest -from gitea_dynamic_runner import opensandbox_worker -from gitea_dynamic_runner.sandbox_kubernetes import allocated_pod_name +from gitea_dynamic_runner.opensandbox_worker import ( + OpenSandboxScheduler, + RegistrationTokens, + sandbox_request, +) -def request() -> RunnerRequest: +def request(backend: str = "vm") -> RunnerRequest: return RunnerRequest( - job_id=47, - run_id=12, - backend="vm", + job_id=42, + run_id=7, + backend=backend, repository="panxiao81/example", job_name="publish image", - labels=("self-hosted", "vm"), + labels=("self-hosted", backend), ) -def test_sandbox_request_uses_pool_and_identity_gate(monkeypatch): - monkeypatch.setattr(opensandbox_worker, "OPENSANDBOX_POOL", "ci-vm") - document = opensandbox_worker.sandbox_request(request(), "gitea-vm-abcd") +def test_sandbox_request_uses_pool_identity_and_one_time_token_url(): + document = sandbox_request( + request(), "gitea-vm-abcd", "http://scheduler/token/nonce" + ) assert document["pool"] == "ci-vm" assert document["entrypoint"] == [ "/usr/local/libexec/gitea-opensandbox-runner" ] - assert document["env"]["GITEA_RUNNER_NAME"] == "gitea-vm-abcd" assert document["env"]["GITEA_RUNNER_LABELS"] == "self-hosted:host,vm:host" - assert document["env"]["CI_SPIFFE_ID"].startswith( - "spiffe://ddupan.top/ci/panxiao81/example/publish-image-" + assert document["env"]["GITEA_RUNNER_REGISTRATION_TOKEN_URL"].endswith( + "/nonce" ) - assert document["metadata"]["ci.ddupan.top/job-id"] == "47" - - -def test_identity_entry_is_bound_to_kata_pod_uid(monkeypatch): - monkeypatch.setattr(opensandbox_worker, "SPIRE_CLUSTER_NAME", "sandbox-kata") - monkeypatch.setattr(opensandbox_worker, "RUNNER_UID", 2000) - manifest = opensandbox_worker.identity_entry( - request(), sandbox_id="sandbox-123", pod_uid="pod-uid-456" + assert document["env"]["CI_SPIFFE_ID"] == ( + "spiffe://ddupan.top/ci/panxiao81/example/publish-image-3e72cdc4a97e" ) - assert manifest["metadata"]["name"] == opensandbox_worker.entry_name( - "sandbox-123" - ) - spec = manifest["spec"] - assert spec["parentID"] == ( - "spiffe://ddupan.top/spire/agent/k8s_psat/" - "sandbox-kata/pod/pod-uid-456" - ) - assert spec["selectors"] == ["unix:uid:2000"] - assert spec["spiffeID"].startswith( - "spiffe://ddupan.top/ci/panxiao81/example/publish-image-" + assert document["metadata"]["ci.ddupan.top/runner"] == "true" + + +async def test_registration_token_is_single_use(monkeypatch): + monkeypatch.setattr( + "gitea_dynamic_runner.opensandbox_worker.TOKEN_BASE_URL", + "http://scheduler/token", ) + tokens = RegistrationTokens(b"secret\n") + nonce, url = await tokens.issue() + assert url == f"http://scheduler/token/{nonce}" + assert await tokens.consume(nonce) == b"secret" + assert await tokens.consume(nonce) is None -@pytest.mark.parametrize( - ("annotation", "expected"), - [ - (json.dumps({"pods": ["pool-pod-1"], "poolRef": "ci-vm"}), "pool-pod-1"), - (json.dumps({"pods": []}), None), - (json.dumps({"pods": ["one", "two"]}), None), - ("not-json", None), - ], -) -def test_allocated_pod_name(annotation, expected): - batchsandbox = { - "metadata": { - "annotations": {"sandbox.opensandbox.io/alloc-status": annotation} - } - } - assert allocated_pod_name(batchsandbox) == expected - - -class FakeKubernetesClient: +class FakeOpenSandbox: def __init__(self): - self.calls = 0 + self.created = None + self.deleted = [] - async def get_batchsandbox(self, sandbox_id): - self.calls += 1 - if self.calls == 1: - return {"metadata": {"annotations": {}}} - return { - "metadata": { - "annotations": { - "sandbox.opensandbox.io/alloc-status": json.dumps( - {"pods": ["pool-pod-1"], "poolRef": "ci-vm"} - ) - } - } - } + async def create(self, **document): + self.created = document + return {"id": "sandbox-1"} - async def get_pod(self, name): - assert name == "pool-pod-1" - return {"metadata": {"uid": "pod-uid-1"}} + async def get(self, sandbox_id): + return {"status": {"state": "Terminated"}} + + async def delete(self, sandbox_id): + self.deleted.append(sandbox_id) -async def test_wait_for_allocation_returns_real_pod_uid(monkeypatch): - async def no_sleep(_): - return None - - monkeypatch.setattr(asyncio, "sleep", no_sleep) - client = FakeKubernetesClient() - assert await opensandbox_worker.wait_for_allocation(client, "sandbox-1") == ( - "pool-pod-1", - "pod-uid-1", +async def test_scheduler_creates_monitors_and_deletes(monkeypatch): + monkeypatch.setattr( + "gitea_dynamic_runner.opensandbox_worker.TOKEN_BASE_URL", + "http://scheduler/token", ) + client = FakeOpenSandbox() + scheduler = OpenSandboxScheduler(client, RegistrationTokens(b"secret")) + assert await scheduler.create(request()) == "sandbox-1" + await scheduler.active[42] + assert client.created["pool"] == "ci-vm" + assert client.deleted == ["sandbox-1"] -def test_entry_name_is_stable_and_dns_safe(): - name = opensandbox_worker.entry_name("sandbox/with unsafe characters") - assert name == opensandbox_worker.entry_name("sandbox/with unsafe characters") - assert name.startswith("gitea-ci-") - assert "/" not in name +async def test_scheduler_rejects_duplicate_active_job(): + client = FakeOpenSandbox() + scheduler = OpenSandboxScheduler(client, RegistrationTokens(b"secret")) + scheduler.active[42] = asyncio.Future() + with pytest.raises(ValueError, match="already"): + await scheduler.create(request())