Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
de98a9952e
|
@@ -17,10 +17,14 @@ on:
|
|||||||
jobs:
|
jobs:
|
||||||
publish-images:
|
publish-images:
|
||||||
name: publish-images
|
name: publish-images
|
||||||
runs-on: [self-hosted, vm]
|
runs-on: self-hosted
|
||||||
timeout-minutes: 45
|
timeout-minutes: 45
|
||||||
permissions:
|
permissions:
|
||||||
contents: read
|
contents: read
|
||||||
|
container:
|
||||||
|
image: docker.io/gitea/runner-images:ubuntu-latest@sha256:fd911d7417bfbf0f454530e447da95b58001e1df41bbc5e1a8dd35d432575aae
|
||||||
|
volumes:
|
||||||
|
- /run/spire/agent-sockets:/run/spire/agent-sockets:ro
|
||||||
env:
|
env:
|
||||||
PUSH_REGISTRY: zot-push.ad.ddupan.top
|
PUSH_REGISTRY: zot-push.ad.ddupan.top
|
||||||
PULL_REGISTRY: zot.ad.ddupan.top
|
PULL_REGISTRY: zot.ad.ddupan.top
|
||||||
@@ -41,6 +45,19 @@ jobs:
|
|||||||
apt-get install --yes --no-install-recommends shellcheck
|
apt-get install --yes --no-install-recommends shellcheck
|
||||||
shellcheck scripts/*
|
shellcheck scripts/*
|
||||||
|
|
||||||
|
- name: Fetch pinned SPIRE CLI
|
||||||
|
shell: bash
|
||||||
|
run: |
|
||||||
|
set -euo pipefail
|
||||||
|
archive=/tmp/spire.tar.gz
|
||||||
|
curl --fail --location --silent --show-error \
|
||||||
|
--output "$archive" \
|
||||||
|
https://github.com/spiffe/spire/releases/download/v1.15.3/spire-1.15.3-linux-amd64-musl.tar.gz
|
||||||
|
printf '%s %s\n' \
|
||||||
|
ca1a4d1155317bdd2afc7f36663828a10410c7c840e54725b90b4064b0a301c7 \
|
||||||
|
"$archive" | sha256sum --check --status
|
||||||
|
tar -xzf "$archive" -C /tmp spire-1.15.3/bin/spire-agent
|
||||||
|
|
||||||
- name: Build and publish
|
- name: Build and publish
|
||||||
shell: bash
|
shell: bash
|
||||||
run: |
|
run: |
|
||||||
@@ -58,7 +75,7 @@ jobs:
|
|||||||
trap cleanup EXIT
|
trap cleanup EXIT
|
||||||
export DOCKER_CONFIG="$docker_config"
|
export DOCKER_CONFIG="$docker_config"
|
||||||
|
|
||||||
/opt/spire/bin/spire-agent api fetch jwt \
|
/tmp/spire-1.15.3/bin/spire-agent api fetch jwt \
|
||||||
-audience zot \
|
-audience zot \
|
||||||
-socketPath "$SPIRE_AGENT_SOCKET" \
|
-socketPath "$SPIRE_AGENT_SOCKET" \
|
||||||
-output json >"$jwt_file"
|
-output json >"$jwt_file"
|
||||||
|
|||||||
@@ -24,9 +24,6 @@ runs-on: [self-hosted, vm]
|
|||||||
Cloud Hypervisor,退出后完整清理。
|
Cloud Hypervisor,退出后完整清理。
|
||||||
- `guest-runner`:在 guest 中领取一次性 runner registration token,注册 ephemeral
|
- `guest-runner`:在 guest 中领取一次性 runner registration token,注册 ephemeral
|
||||||
runner,执行一个 job 后关机。
|
runner,执行一个 job 后关机。
|
||||||
- `opensandbox-worker`:在正式 sandbox 集群通过 OpenSandbox `ci-vm` Pool 分配
|
|
||||||
Kata guest,按实际 Pod UID 创建临时 SPIFFE entry,并在任务结束后回收两者。身份
|
|
||||||
与 Pool 契约见 [`docs/opensandbox-runner.md`](docs/opensandbox-runner.md)。
|
|
||||||
- `pod-worker`:在 Kubernetes 中创建一次性 privileged Pod;Pod 内的 workflow 使用
|
- `pod-worker`:在 Kubernetes 中创建一次性 privileged Pod;Pod 内的 workflow 使用
|
||||||
host executor,Docker、BuildKit 和 kind 等工具由 pipeline 按需 setup。Runner 固定在
|
host executor,Docker、BuildKit 和 kind 等工具由 pipeline 按需 setup。Runner 固定在
|
||||||
支持原生 job hooks 的 3.x 版本,在 workflow 第一步前等待实际任务对应的 SVID。
|
支持原生 job hooks 的 3.x 版本,在 workflow 第一步前等待实际任务对应的 SVID。
|
||||||
|
|||||||
@@ -14,7 +14,6 @@ COPY --from=runner /usr/local/bin/run.sh /usr/local/bin/run.sh
|
|||||||
COPY --from=spire /opt/spire/bin/spire-agent /opt/spire/bin/spire-agent
|
COPY --from=spire /opt/spire/bin/spire-agent /opt/spire/bin/spire-agent
|
||||||
COPY config/runner.yaml /etc/gitea-runner/config.yaml
|
COPY config/runner.yaml /etc/gitea-runner/config.yaml
|
||||||
COPY --chmod=0755 scripts/gitea-job-started /usr/local/libexec/gitea-job-started
|
COPY --chmod=0755 scripts/gitea-job-started /usr/local/libexec/gitea-job-started
|
||||||
COPY --chmod=0755 scripts/gitea-opensandbox-runner /usr/local/libexec/gitea-opensandbox-runner
|
|
||||||
|
|
||||||
VOLUME ["/data"]
|
VOLUME ["/data"]
|
||||||
WORKDIR /
|
WORKDIR /
|
||||||
|
|||||||
@@ -1,55 +0,0 @@
|
|||||||
# OpenSandbox Kata 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:
|
|
||||||
|
|
||||||
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`.
|
|
||||||
|
|
||||||
## Identity ordering
|
|
||||||
|
|
||||||
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/<pod-uid>`;
|
|
||||||
- workload: `spiffe://ddupan.top/ci/<owner>/<repository>/<job>`;
|
|
||||||
- selector: `unix:uid:2000`.
|
|
||||||
|
|
||||||
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.
|
|
||||||
|
|
||||||
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.
|
|
||||||
|
|
||||||
## Required Pool contract
|
|
||||||
|
|
||||||
The `ci-vm` Pool template owns infrastructure that callers cannot override in
|
|
||||||
Pool mode:
|
|
||||||
|
|
||||||
- `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.
|
|
||||||
|
|
||||||
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`.
|
|
||||||
|
|
||||||
## 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.
|
|
||||||
@@ -16,7 +16,6 @@ test = ["pytest==8.4.2", "pytest-asyncio==1.2.0"]
|
|||||||
gitea-dynamic-runner-controller = "gitea_dynamic_runner.controller:main"
|
gitea-dynamic-runner-controller = "gitea_dynamic_runner.controller:main"
|
||||||
gitea-dynamic-runner-pod-worker = "gitea_dynamic_runner.pod_worker:cli"
|
gitea-dynamic-runner-pod-worker = "gitea_dynamic_runner.pod_worker:cli"
|
||||||
gitea-dynamic-runner-vm-worker = "gitea_dynamic_runner.worker:cli"
|
gitea-dynamic-runner-vm-worker = "gitea_dynamic_runner.worker:cli"
|
||||||
gitea-dynamic-runner-opensandbox-worker = "gitea_dynamic_runner.opensandbox_worker:cli"
|
|
||||||
gitea-spire-jwt-broker = "gitea_dynamic_runner.jwt_broker:main"
|
gitea-spire-jwt-broker = "gitea_dynamic_runner.jwt_broker:main"
|
||||||
|
|
||||||
[tool.pytest.ini_options]
|
[tool.pytest.ini_options]
|
||||||
|
|||||||
@@ -1,21 +0,0 @@
|
|||||||
#!/usr/bin/env bash
|
|
||||||
set -euo pipefail
|
|
||||||
|
|
||||||
: "${CI_SPIFFE_ID:?CI_SPIFFE_ID is required}"
|
|
||||||
: "${SPIFFE_ENDPOINT_SOCKET:?SPIFFE_ENDPOINT_SOCKET is required}"
|
|
||||||
|
|
||||||
socket_path=${SPIFFE_ENDPOINT_SOCKET#unix://}
|
|
||||||
deadline=$((SECONDS + 120))
|
|
||||||
while (( SECONDS < deadline )); do
|
|
||||||
if /opt/spire/bin/spire-agent api fetch x509 \
|
|
||||||
-socketPath "$socket_path" \
|
|
||||||
-output json 2>/dev/null | \
|
|
||||||
jq -e --arg id "$CI_SPIFFE_ID" \
|
|
||||||
'any(.svids[]?; .spiffe_id == $id)' >/dev/null; then
|
|
||||||
exec /usr/local/bin/run.sh
|
|
||||||
fi
|
|
||||||
sleep 1
|
|
||||||
done
|
|
||||||
|
|
||||||
printf 'timed out waiting for OpenSandbox SPIFFE identity %s\n' "$CI_SPIFFE_ID" >&2
|
|
||||||
exit 1
|
|
||||||
@@ -1,79 +0,0 @@
|
|||||||
"""Minimal client for the OpenSandbox lifecycle API."""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
from pathlib import Path
|
|
||||||
from urllib.parse import quote
|
|
||||||
|
|
||||||
from aiohttp import ClientResponseError, ClientSession
|
|
||||||
|
|
||||||
|
|
||||||
class OpenSandboxClient:
|
|
||||||
def __init__(
|
|
||||||
self,
|
|
||||||
*,
|
|
||||||
api_url: str,
|
|
||||||
api_key_file: Path | None = None,
|
|
||||||
) -> None:
|
|
||||||
self.api_url = api_url.rstrip("/")
|
|
||||||
self.api_key_file = api_key_file
|
|
||||||
self.session: ClientSession | None = None
|
|
||||||
|
|
||||||
async def __aenter__(self) -> OpenSandboxClient:
|
|
||||||
headers = {}
|
|
||||||
if self.api_key_file is not None:
|
|
||||||
headers["OPEN-SANDBOX-API-KEY"] = self.api_key_file.read_text().strip()
|
|
||||||
self.session = ClientSession(headers=headers, 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(
|
|
||||||
self,
|
|
||||||
*,
|
|
||||||
pool: str,
|
|
||||||
timeout: int,
|
|
||||||
entrypoint: list[str],
|
|
||||||
env: dict[str, str],
|
|
||||||
metadata: dict[str, str],
|
|
||||||
) -> dict[str, object]:
|
|
||||||
response = await self._request(
|
|
||||||
"POST",
|
|
||||||
"/v1/sandboxes",
|
|
||||||
json={
|
|
||||||
"timeout": timeout,
|
|
||||||
"entrypoint": entrypoint,
|
|
||||||
"env": env,
|
|
||||||
"metadata": metadata,
|
|
||||||
"extensions": {"poolRef": pool},
|
|
||||||
},
|
|
||||||
)
|
|
||||||
return await response.json()
|
|
||||||
|
|
||||||
async def get(self, sandbox_id: str) -> dict[str, object] | None:
|
|
||||||
try:
|
|
||||||
response = await self._request(
|
|
||||||
"GET", f"/v1/sandboxes/{quote(sandbox_id, safe='')}"
|
|
||||||
)
|
|
||||||
except ClientResponseError as error:
|
|
||||||
if error.status == 404:
|
|
||||||
return None
|
|
||||||
raise
|
|
||||||
return await response.json()
|
|
||||||
|
|
||||||
async def delete(self, sandbox_id: str) -> None:
|
|
||||||
try:
|
|
||||||
response = await self._request(
|
|
||||||
"DELETE", f"/v1/sandboxes/{quote(sandbox_id, safe='')}"
|
|
||||||
)
|
|
||||||
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("OpenSandboxClient is not open")
|
|
||||||
return await self.session.request(method, f"{self.api_url}{path}", **kwargs)
|
|
||||||
@@ -1,276 +0,0 @@
|
|||||||
#!/usr/bin/env python3
|
|
||||||
"""JetStream worker that provisions disposable OpenSandbox Kata runners."""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import asyncio
|
|
||||||
import hashlib
|
|
||||||
import json
|
|
||||||
import logging
|
|
||||||
import os
|
|
||||||
from pathlib import Path
|
|
||||||
import ssl
|
|
||||||
import uuid
|
|
||||||
|
|
||||||
import nats
|
|
||||||
from nats.errors import TimeoutError
|
|
||||||
from nats.js.api import AckPolicy, ConsumerConfig
|
|
||||||
|
|
||||||
from .models import RunnerRequest
|
|
||||||
from .opensandbox import OpenSandboxClient
|
|
||||||
from .pod_worker import heartbeat, identity_path
|
|
||||||
from .sandbox_kubernetes import SandboxKubernetesClient, allocated_pod_name
|
|
||||||
|
|
||||||
|
|
||||||
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"))
|
|
||||||
|
|
||||||
|
|
||||||
def sandbox_request(request: RunnerRequest, runner_name: str) -> dict[str, object]:
|
|
||||||
if request.backend != "vm":
|
|
||||||
raise ValueError("OpenSandbox backend only accepts vm requests")
|
|
||||||
path = identity_path(request.repository, request.job_name)
|
|
||||||
return {
|
|
||||||
"pool": OPENSANDBOX_POOL,
|
|
||||||
"timeout": SANDBOX_TIMEOUT,
|
|
||||||
"entrypoint": ["/usr/local/libexec/gitea-opensandbox-runner"],
|
|
||||||
"env": {
|
|
||||||
"GITEA_RUNNER_NAME": runner_name,
|
|
||||||
"GITEA_RUNNER_LABELS": "self-hosted:host,vm:host",
|
|
||||||
"CI_SPIFFE_ID": f"spiffe://{SPIFFE_TRUST_DOMAIN}/ci/{path}",
|
|
||||||
"SPIFFE_ENDPOINT_SOCKET": "unix:///run/spire/agent-sockets/spire-agent.sock",
|
|
||||||
},
|
|
||||||
"metadata": {
|
|
||||||
"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,
|
|
||||||
},
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
def entry_name(sandbox_id: str) -> str:
|
|
||||||
suffix = hashlib.sha256(sandbox_id.encode()).hexdigest()[:12]
|
|
||||||
return f"gitea-ci-{suffix}"
|
|
||||||
|
|
||||||
|
|
||||||
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}"],
|
|
||||||
},
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
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
|
|
||||||
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)
|
|
||||||
)
|
|
||||||
task.add_done_callback(_report_task)
|
|
||||||
active.add(task)
|
|
||||||
|
|
||||||
|
|
||||||
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__),
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
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,
|
|
||||||
):
|
|
||||||
try:
|
|
||||||
await consume_requests(subscription, sandbox_client, kubernetes_client)
|
|
||||||
finally:
|
|
||||||
await nc.drain()
|
|
||||||
|
|
||||||
|
|
||||||
def cli() -> None:
|
|
||||||
logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO"))
|
|
||||||
asyncio.run(main())
|
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
|
||||||
cli()
|
|
||||||
@@ -1,132 +0,0 @@
|
|||||||
"""Kubernetes resources that bind an OpenSandbox Kata guest to SPIRE."""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import json
|
|
||||||
from pathlib import Path
|
|
||||||
import ssl
|
|
||||||
from urllib.parse import quote
|
|
||||||
|
|
||||||
from aiohttp import ClientResponseError, ClientSession, TCPConnector
|
|
||||||
|
|
||||||
|
|
||||||
class SandboxKubernetesClient:
|
|
||||||
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) -> SandboxKubernetesClient:
|
|
||||||
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 get_batchsandbox(self, sandbox_id: str) -> dict[str, object] | None:
|
|
||||||
try:
|
|
||||||
response = await self._request(
|
|
||||||
"GET", f"{self._batchsandboxes_path()}/{quote(sandbox_id, safe='')}"
|
|
||||||
)
|
|
||||||
except ClientResponseError as error:
|
|
||||||
if error.status == 404:
|
|
||||||
return None
|
|
||||||
raise
|
|
||||||
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, safe='')}"
|
|
||||||
)
|
|
||||||
except ClientResponseError as error:
|
|
||||||
if error.status == 404:
|
|
||||||
return None
|
|
||||||
raise
|
|
||||||
return await response.json()
|
|
||||||
|
|
||||||
async def create_entry(self, manifest: dict[str, object]) -> dict[str, object]:
|
|
||||||
response = await self._request(
|
|
||||||
"POST", self._entries_path(), json=manifest
|
|
||||||
)
|
|
||||||
return await response.json()
|
|
||||||
|
|
||||||
async def get_entry(self, name: str) -> dict[str, object] | None:
|
|
||||||
try:
|
|
||||||
response = await self._request(
|
|
||||||
"GET", f"{self._entries_path()}/{quote(name, safe='')}"
|
|
||||||
)
|
|
||||||
except ClientResponseError as error:
|
|
||||||
if error.status == 404:
|
|
||||||
return None
|
|
||||||
raise
|
|
||||||
return await response.json()
|
|
||||||
|
|
||||||
async def delete_entry(self, name: str) -> None:
|
|
||||||
try:
|
|
||||||
response = await self._request(
|
|
||||||
"DELETE",
|
|
||||||
f"{self._entries_path()}/{quote(name, safe='')}",
|
|
||||||
json={"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("SandboxKubernetesClient is not open")
|
|
||||||
return await self.session.request(method, f"{self.api_url}{path}", **kwargs)
|
|
||||||
|
|
||||||
def _batchsandboxes_path(self) -> str:
|
|
||||||
return (
|
|
||||||
"/apis/sandbox.opensandbox.io/v1alpha1/namespaces/"
|
|
||||||
f"{quote(self.namespace, safe='')}/batchsandboxes"
|
|
||||||
)
|
|
||||||
|
|
||||||
def _pods_path(self) -> str:
|
|
||||||
return f"/api/v1/namespaces/{quote(self.namespace, safe='')}/pods"
|
|
||||||
|
|
||||||
@staticmethod
|
|
||||||
def _entries_path() -> str:
|
|
||||||
return "/apis/spire.spiffe.io/v1alpha1/clusterstaticentries"
|
|
||||||
|
|
||||||
|
|
||||||
def allocated_pod_name(batchsandbox: dict[str, object]) -> str | None:
|
|
||||||
metadata = batchsandbox.get("metadata")
|
|
||||||
if not isinstance(metadata, dict):
|
|
||||||
return None
|
|
||||||
annotations = metadata.get("annotations")
|
|
||||||
if not isinstance(annotations, dict):
|
|
||||||
return None
|
|
||||||
raw = annotations.get("sandbox.opensandbox.io/alloc-status")
|
|
||||||
if not isinstance(raw, str):
|
|
||||||
return None
|
|
||||||
try:
|
|
||||||
allocation = json.loads(raw)
|
|
||||||
except json.JSONDecodeError:
|
|
||||||
return None
|
|
||||||
if not isinstance(allocation, dict):
|
|
||||||
return None
|
|
||||||
pods = allocation.get("pods")
|
|
||||||
if not isinstance(pods, list) or len(pods) != 1 or not isinstance(pods[0], str):
|
|
||||||
return None
|
|
||||||
return pods[0]
|
|
||||||
@@ -1,114 +0,0 @@
|
|||||||
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
|
|
||||||
|
|
||||||
|
|
||||||
def request() -> RunnerRequest:
|
|
||||||
return RunnerRequest(
|
|
||||||
job_id=47,
|
|
||||||
run_id=12,
|
|
||||||
backend="vm",
|
|
||||||
repository="panxiao81/example",
|
|
||||||
job_name="publish image",
|
|
||||||
labels=("self-hosted", "vm"),
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
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")
|
|
||||||
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["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 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-"
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
@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:
|
|
||||||
def __init__(self):
|
|
||||||
self.calls = 0
|
|
||||||
|
|
||||||
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 get_pod(self, name):
|
|
||||||
assert name == "pool-pod-1"
|
|
||||||
return {"metadata": {"uid": "pod-uid-1"}}
|
|
||||||
|
|
||||||
|
|
||||||
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",
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
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
|
|
||||||
Reference in New Issue
Block a user