Compare commits

..
Author SHA1 Message Date
panxiao81 2679f81cdd fix: 使用 OpenSandbox chart 服务地址
test / python (pull_request) Successful in 9s
test / shell (pull_request) Successful in 17s
2026-09-18 01:09:37 +00:00
panxiao81 d7ba64d6f3 接入 OpenSandbox Kata Runner 控制面
test / python (pull_request) Successful in 9s
test / shell (pull_request) Successful in 21s
2026-09-18 00:56:10 +00:00
panxiao81 8d70df2128 Merge pull request '在 VM Runner 中构建发布镜像' (#17) from refactor/publish-on-vm into main
test / python (push) Successful in 10s
test / shell (push) Successful in 15s
publish images / publish-images (push) Failing after 0s
Reviewed-on: #17
2026-09-16 18:29:26 +00:00
panxiao81 3e69d102dd 在 VM Runner 中发布镜像
test / python (pull_request) Successful in 14s
test / shell (pull_request) Successful in 22s
2026-09-16 18:28:33 +00:00
panxiao81 0387b2592a Merge pull request '使用 native snapshotter 发布镜像证明' (#16) from fix/publish-native-snapshotter into main
test / python (push) Successful in 10s
test / shell (push) Successful in 16s
publish images / publish-images (push) Canceled after 9m5s
Reviewed-on: #16
2026-09-16 18:20:10 +00:00
panxiao81 057731ae62 使用 native snapshotter 生成镜像证明
test / python (pull_request) Successful in 8s
test / shell (pull_request) Successful in 14s
2026-09-16 18:19:27 +00:00
panxiao81 529daa4c75 Merge pull request '修复 Pod 内嵌套 Docker overlay 构建' (#15) from fix/publish-vfs-builder into main
test / python (push) Successful in 9s
publish images / publish-images (push) Failing after 35s
test / shell (push) Successful in 16s
Reviewed-on: #15
2026-09-16 18:17:24 +00:00
panxiao81 38946e8def 使用 vfs 构建嵌套容器镜像
test / shell (pull_request) Successful in 20s
test / python (pull_request) Successful in 9s
2026-09-16 18:16:41 +00:00
panxiao81 893323a89b Merge pull request '添加 Zot 镜像构建发布流水线' (#14) from feat/image-publish-pipeline into main
test / python (push) Successful in 9s
test / shell (push) Successful in 16s
publish images / publish-images (push) Failing after 57s
Reviewed-on: #14
2026-09-16 18:12:25 +00:00
panxiao81 6f149c85e2 添加 Zot 镜像发布流水线
test / python (pull_request) Successful in 10s
test / shell (pull_request) Successful in 16s
2026-09-16 18:10:02 +00:00
panxiao81 852afbf02e Merge pull request '添加动态 Pod Runner 冒烟测试' (#9) from test/pod-smoke into main
test / python (push) Successful in 10s
test / shell (push) Failing after 17s
Reviewed-on: #9
2026-09-16 18:03:01 +00:00
panxiao81 6a29a10244 Merge pull request '修复 LXC 内 microVM Runner 启动链路' (#11) from fix/microvm-lxc-bootstrap into main
test / shell (push) Failing after 19s
test / python (push) Successful in 9s
Reviewed-on: #11
2026-09-16 18:02:51 +00:00
panxiao81 94a962459c Merge pull request '为 Pod Runner 提供 Ubuntu workflow 环境' (#12) from fix/ubuntu-runner-image into main
test / python (push) Successful in 8s
test / shell (push) Successful in 18s
Reviewed-on: #12
2026-09-16 18:02:16 +00:00
panxiao81 d03927d73c Merge pull request '修复 Runner SPIFFE 身份路径并补充关联日志' (#13) from fix/identity-binding-observability into main
test / shell (push) Successful in 15s
test / python (push) Successful in 8s
Reviewed-on: #13
2026-09-16 18:01:57 +00:00
panxiao81 68bb02b20a 修复 Runner 身份路径并补充关联日志
test / python (pull_request) Successful in 8s
test / shell (pull_request) Successful in 15s
2026-09-16 18:00:50 +00:00
panxiao81 c4aa6ee0af 记录 LXC guest 内存共享要求
test / python (pull_request) Successful in 26s
test / shell (pull_request) Failing after 1m2s
2026-09-16 17:42:52 +00:00
panxiao81 401c1f9a00 共享 guest memfd 避免 LXC 内存双重计费
test / shell (pull_request) Failing after 22s
test / python (pull_request) Failing after 1m3s
2026-09-16 17:38:08 +00:00
panxiao81 6dac9897fd 稳定 kind smoke 与 microVM 内存配置
test / python (pull_request) Successful in 14s
test / shell (pull_request) Failing after 19s
2026-09-16 17:28:29 +00:00
panxiao81 072a5bad77 为 microVM 磁盘启用 direct I/O
test / python (pull_request) Successful in 18s
test / shell (pull_request) Failing after 21s
2026-09-16 17:19:32 +00:00
panxiao81 b25b1fbf62 避免 microVM 镜像复制占满 LXC page cache
test / python (pull_request) Successful in 9s
test / shell (pull_request) Failing after 18s
2026-09-16 17:15:08 +00:00
panxiao81 8578bee895 修正 kind 二进制校验
test / python (pull_request) Successful in 9s
test / shell (pull_request) Failing after 16s
2026-09-16 16:50:52 +00:00
panxiao81 02e0b698c5 为 cloud-init Runner 设置 HOME
test / python (pull_request) Successful in 11s
test / shell (pull_request) Failing after 21s
2026-09-16 16:48:44 +00:00
panxiao81 ac85d5fe58 校验 microVM Runner 任务完成状态
test / python (pull_request) Successful in 10s
test / shell (pull_request) Successful in 16s
2026-09-16 16:44:05 +00:00
panxiao81 bb73a48695 将 Pod smoke 改为手动触发
test / python (pull_request) Successful in 14s
test / shell (pull_request) Successful in 20s
2026-09-16 16:42:49 +00:00
panxiao81 70a8aa68c6 将 VM smoke 改为手动触发
test / python (pull_request) Successful in 10s
test / shell (pull_request) Successful in 15s
2026-09-16 16:42:20 +00:00
panxiao81 ce4c9f13b5 添加 microVM kind 冒烟测试
test / python (pull_request) Successful in 10s
test / shell (pull_request) Successful in 17s
VM kind smoke / kind (pull_request) Failing after 2s
2026-09-16 16:39:06 +00:00
panxiao81 3641e6ffb3 修复 LXC 内 microVM Runner 启动链路
test / python (pull_request) Successful in 10s
test / shell (pull_request) Successful in 16s
2026-09-16 16:38:03 +00:00
panxiao81 99ef62bfa0 Merge remote-tracking branch 'origin/main' into test/pod-smoke
dynamic Pod smoke test / pod-smoke (pull_request) Successful in 5s
test / python (pull_request) Successful in 10s
test / shell (pull_request) Successful in 18s
2026-09-16 16:06:55 +00:00
panxiao81 21aabe162f 添加动态 Pod Runner 冒烟测试
test / python (pull_request) Successful in 9s
test / shell (pull_request) Successful in 16s
dynamic Pod smoke test / pod-smoke (pull_request) Failing after 1s
2026-09-16 15:32:38 +00:00
20 changed files with 1057 additions and 22 deletions
+19
View File
@@ -0,0 +1,19 @@
name: dynamic Pod smoke test
on:
workflow_dispatch:
jobs:
pod-smoke:
name: pod-smoke
runs-on: [self-hosted, pod]
timeout-minutes: 5
steps:
- name: Verify disposable host environment
run: |
test -S /run/spire/agent-sockets/spire-agent.sock
/opt/spire/bin/spire-agent api fetch jwt \
-audience ci-smoke \
-socketPath /run/spire/agent-sockets/spire-agent.sock \
>/dev/null
test "$(id -u)" = 0
+122
View File
@@ -0,0 +1,122 @@
---
name: publish images
on:
push:
branches: [main]
paths:
- '.gitea/workflows/publish-images.yml'
- 'config/**'
- 'container/**'
- 'scripts/**'
- 'src/**'
- 'pyproject.toml'
- 'README.md'
workflow_dispatch:
jobs:
publish-images:
name: publish-images
runs-on: [self-hosted, vm]
timeout-minutes: 45
permissions:
contents: read
env:
PUSH_REGISTRY: zot-push.ad.ddupan.top
PULL_REGISTRY: zot.ad.ddupan.top
CONTROLLER_REPOSITORY: panxiao81/gitea-dynamic-runner-controller
RUNNER_REPOSITORY: panxiao81/gitea-dynamic-runner-runner
SPIRE_AGENT_SOCKET: /run/spire/agent-sockets/spire-agent.sock
steps:
- uses: actions/checkout@v4
- name: Test source
shell: bash
run: |
set -euo pipefail
python3 -m pip install --break-system-packages -e '.[test]'
pytest -q
python3 -m compileall -q src tests
apt-get update
apt-get install --yes --no-install-recommends shellcheck
shellcheck scripts/*
- name: Build and publish
shell: bash
run: |
set -euo pipefail
set +x
: "${GITHUB_SHA:?GITHUB_SHA is required}"
image_tag="sha-${GITHUB_SHA}"
docker_config=$(mktemp -d)
jwt_file=$(mktemp)
cleanup() {
docker buildx rm ci-builder >/dev/null 2>&1 || true
rm -rf -- "$docker_config" "$jwt_file"
}
trap cleanup EXIT
export DOCKER_CONFIG="$docker_config"
/opt/spire/bin/spire-agent api fetch jwt \
-audience zot \
-socketPath "$SPIRE_AGENT_SOCKET" \
-output json >"$jwt_file"
# shellcheck disable=SC2016
jq -er '.[0].svids[0].svid' "$jwt_file" | \
docker login "$PUSH_REGISTRY" --username zot --password-stdin
docker buildx create \
--name ci-builder \
--driver docker-container \
--use
publish() {
local repository=$1
local dockerfile=$2
local metadata=$3
docker buildx build \
--builder ci-builder \
--platform linux/amd64 \
--file "$dockerfile" \
--tag "${PUSH_REGISTRY}/${repository}:${image_tag}" \
--tag "${PUSH_REGISTRY}/${repository}:main" \
--provenance=mode=max \
--sbom=true \
--metadata-file "$metadata" \
--push \
.
}
publish \
"$CONTROLLER_REPOSITORY" \
container/controller.Dockerfile \
controller-metadata.json
publish \
"$RUNNER_REPOSITORY" \
container/runner.Dockerfile \
runner-metadata.json
controller_digest=$(
# shellcheck disable=SC2016
jq -er '."containerimage.digest"' controller-metadata.json
)
runner_digest=$(
# shellcheck disable=SC2016
jq -er '."containerimage.digest"' runner-metadata.json
)
controller_ref="${PULL_REGISTRY}/${CONTROLLER_REPOSITORY}@${controller_digest}"
runner_ref="${PULL_REGISTRY}/${RUNNER_REPOSITORY}@${runner_digest}"
printf 'controller=%s\nrunner=%s\n' "$controller_ref" "$runner_ref"
if [[ -n "${GITHUB_STEP_SUMMARY:-}" ]]; then
{
printf '## Published images\n\n'
# shellcheck disable=SC2016
printf -- '- Controller: `%s`\n' "$controller_ref"
# shellcheck disable=SC2016
printf -- '- Runner: `%s`\n' "$runner_ref"
# shellcheck disable=SC2016
printf -- '- Source: `%s`\n' "$GITHUB_SHA"
} >>"$GITHUB_STEP_SUMMARY"
fi
+31
View File
@@ -0,0 +1,31 @@
name: VM kind smoke
on:
workflow_dispatch:
jobs:
kind:
runs-on: [self-hosted, vm]
steps:
- name: Verify Docker
run: docker info
- name: Install kind
shell: bash
run: |
set -euo pipefail
version=v0.33.0
curl --fail --location --silent --show-error \
--output /tmp/kind "https://kind.sigs.k8s.io/dl/${version}/kind-linux-amd64"
curl --fail --location --silent --show-error \
--output /tmp/kind.sha256sum "https://kind.sigs.k8s.io/dl/${version}/kind-linux-amd64.sha256sum"
printf '%s %s\n' "$(cut -d ' ' -f1 /tmp/kind.sha256sum)" /tmp/kind | sha256sum --check
chmod 0755 /tmp/kind
- name: Create and delete kind cluster
shell: bash
run: |
set -euo pipefail
trap '/tmp/kind delete cluster --name smoke' EXIT
/tmp/kind create cluster --name smoke --wait 180s
/tmp/kind get clusters | grep -Fx smoke
+5
View File
@@ -0,0 +1,5 @@
---
self-hosted-runner:
labels:
- pod
- vm
+11 -1
View File
@@ -20,10 +20,13 @@ runs-on: [self-hosted, vm]
- `controller`:接收 Gitea `workflow_job` webhook,将指定 label 的 queued job
发布到 NATS JetStream。
- `worker`:领取任务、限制并发,并通过 Pod 或 microVM backend 创建一次性环境。
- `microvm-runner-launch`:为每个任务创建 COW disk、NoCloud seed 和 TAP,运行
- `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)。
- `pod-worker`:在 Kubernetes 中创建一次性 privileged Pod;Pod 内的 workflow 使用
host executor,Docker、BuildKit 和 kind 等工具由 pipeline 按需 setup。Runner 固定在
支持原生 job hooks 的 3.x 版本,在 workflow 第一步前等待实际任务对应的 SVID。
@@ -51,9 +54,16 @@ pytest
- NATS 密码、webhook secret 和 Gitea registration token 只从文件读取。
- registration token 不写入 seed image;worker 通过单次 nonce endpoint 交给 guest。
- guest 启动时从仅监听 microVM bridge 的 worker endpoint 获取固定版本 Runner 和配置
资产;基础镜像无需为 Runner 发布而重做。
- `runner-vm-bootstrap.yaml` 暂时只验证 VM 调度和生命周期,不提供 SPIFFE
identity;VM agent attestation 完成前不得将它当作身份链路验证结果。
- guest runner 使用 `--ephemeral`,每台 VM 只执行一个 job。
- launcher 只接受 UUID instance ID 和 URL-safe nonce,所有临时文件都位于独立目录。
- base image 不得包含 runner identity、registration token、SSH 密码或 host key。
- LXC 内运行 Cloud Hypervisor 必须为 guest memory 启用 `shared=on`。默认的 private
memfd 映射会在 guest 写入后同时产生 shmem 与 anonymous CoW charge,使 LXC cgroup
对 guest RAM 接近双倍计费。
homelab 的 Kubernetes、OpenBao、LXC、bridge 和容量配置保留在
`panxiao81/homelab-infra`。
+11
View File
@@ -0,0 +1,11 @@
runner:
capacity: 1
timeout: 3h
shutdown_timeout: 1m
host:
workdir_parent: /workspace
container:
require_docker: false
valid_volumes: []
+1
View File
@@ -14,6 +14,7 @@ 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 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-opensandbox-runner /usr/local/libexec/gitea-opensandbox-runner
VOLUME ["/data"]
WORKDIR /
+55
View File
@@ -0,0 +1,55 @@
# 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.
+1
View File
@@ -16,6 +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-spire-jwt-broker = "gitea_dynamic_runner.jwt_broker:main"
[tool.pytest.ini_options]
+4 -1
View File
@@ -6,6 +6,8 @@ instance=${2:?Gitea instance is required}
runner_name=${3:?runner name is required}
runner_labels=${4:?runner labels are required}
token_file=/run/gitea-runner-registration-token
config_file=${GITEA_RUNNER_CONFIG_FILE:-/etc/gitea-runner/config-vm-bootstrap.yaml}
export HOME="${HOME:-/root}"
cleanup() {
rm -f -- "$token_file"
@@ -25,4 +27,5 @@ gitea-runner register \
--labels "$runner_labels" \
--token-file "$token_file"
rm -f -- "$token_file"
gitea-runner daemon
gitea-runner daemon --config "$config_file" --once
echo gitea-runner-job-complete >/dev/ttyS0
+21
View File
@@ -0,0 +1,21 @@
#!/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
+20 -7
View File
@@ -29,8 +29,13 @@ tap="mvr${instance_id%%-*}"
overlay="$vm_dir/root.qcow2"
seed="$vm_dir/seed.img"
serial="$vm_dir/serial.log"
log_dir="$state_root/logs"
cleanup() {
if test -f "$serial"; then
install -d -m 0700 "$log_dir"
cp "$serial" "$log_dir/$instance_id.log"
fi
ip link delete "$tap" 2>/dev/null || true
rm -rf -- "$vm_dir"
}
@@ -39,7 +44,10 @@ trap cleanup EXIT INT TERM
test -r "$base_image"
test -r "$firmware"
install -d -m 0700 "$state_root/instances" "$vm_dir"
qemu-img create -q -f qcow2 -F qcow2 -b "$base_image" "$overlay"
# Cloud Hypervisor cannot open qcow2 backing chains. Keep each disposable
# root disk flat and bypass the LXC page cache: caching both a 5 GiB copy and
# guest RAM can otherwise trigger the container memory limit.
qemu-img convert -q -T none -t none -f qcow2 -O qcow2 "$base_image" "$overlay"
cat >"$vm_dir/meta-data" <<EOF
instance-id: $instance_id
@@ -48,7 +56,8 @@ EOF
cat >"$vm_dir/user-data" <<EOF
#cloud-config
runcmd:
- [ /usr/local/libexec/gitea-microvm-guest-runner, "$token_url/token/$nonce", "$gitea_instance", "gitea-${instance_id%%-*}", "$runner_labels" ]
- [ sh, -c, "curl --fail --silent --show-error --retry 10 --retry-all-errors $token_url/assets/guest-assets.tar.gz | tar -xz -C /" ]
- [ /usr/local/libexec/gitea-microvm-guest-runner, "$token_url/token/$nonce", "$gitea_instance", "gitea-vm-${instance_id%%-*}", "$runner_labels" ]
EOF
cloud-localds "$seed" "$vm_dir/user-data" "$vm_dir/meta-data"
@@ -63,12 +72,16 @@ ip tuntap add dev "$tap" mode tap
ip link set "$tap" master "$bridge"
ip link set "$tap" up
timeout --signal=TERM --kill-after=30s "$vm_timeout" "$cloud_hypervisor" \
if timeout --signal=TERM --kill-after=30s "$vm_timeout" "$cloud_hypervisor" \
--firmware "$firmware" \
--cpus "boot=$cpus" \
--memory "size=$memory" \
--disk "path=$overlay" \
--disk "path=$seed,readonly=on" \
--memory "size=$memory,shared=on" \
--disk "path=$overlay,image_type=qcow2,direct=on,sparse=off" \
--disk "path=$seed,readonly=on,image_type=raw,direct=on,sparse=off" \
--net "tap=$tap,mac=$mac" \
--serial "file=$serial" \
--console off
--console off; then
grep -Fq gitea-runner-job-complete "$serial"
else
exit $?
fi
+48 -2
View File
@@ -4,6 +4,7 @@
import hashlib
import hmac
import json
import logging
import os
import ssl
from pathlib import Path
@@ -16,6 +17,7 @@ from nats.js.errors import NotFoundError
from .models import IdentityBinding, RunnerRequest
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")
@@ -64,27 +66,70 @@ async def webhook(request: web.Request) -> web.Response:
payload = json.loads(body)
except json.JSONDecodeError as error:
raise web.HTTPBadRequest(text="invalid JSON\n") from error
context = _webhook_context(payload)
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(
f"{SUBJECT_PREFIX}.{runner_request.backend}",
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")
binding = IdentityBinding.from_webhook(payload)
if binding is not None:
subject = f"{SUBJECT_PREFIX}.{binding.backend}.binding"
await request.app["js"].publish(
f"{SUBJECT_PREFIX}.{binding.backend}.binding",
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")
LOG.warning("workflow_job webhook ignored %s", context)
return web.Response(status=204)
def _webhook_context(payload: object) -> str:
if not isinstance(payload, dict):
return f"payload_type={type(payload).__name__}"
job = payload.get("workflow_job")
repository = payload.get("repository")
job = job if isinstance(job, dict) else {}
repository = repository if isinstance(repository, dict) else {}
return (
f"action={payload.get('action')!r} job_id={job.get('id')!r} "
f"run_id={job.get('run_id')!r} runner_name={job.get('runner_name')!r} "
f"repository={repository.get('full_name')!r} job_name={job.get('name')!r} "
f"labels={job.get('labels')!r}"
)
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)
@@ -115,6 +160,7 @@ def create_app() -> web.Application:
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")))
+79
View File
@@ -0,0 +1,79 @@
"""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)
@@ -0,0 +1,276 @@
#!/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()
+76 -5
View File
@@ -4,12 +4,13 @@
from __future__ import annotations
import asyncio
import hashlib
import json
import logging
import os
from pathlib import Path
import re
import ssl
from urllib.parse import quote
import uuid
import nats
@@ -49,16 +50,26 @@ RUNNER_TOKEN_SECRET = os.environ.get("RUNNER_TOKEN_SECRET", "gitea-dynamic-runne
GITEA_INSTANCE = os.environ.get("GITEA_INSTANCE", "https://git.ddupan.top")
CAPACITY = int(os.environ.get("RUNNER_CAPACITY", "4"))
POD_TIMEOUT = int(os.environ.get("RUNNER_POD_TIMEOUT", str(4 * 60 * 60)))
SPIFFE_PATH_SEGMENT = re.compile(r"^[A-Za-z0-9._-]+$")
def identity_path(repository: str, job_name: str) -> str:
parts = repository.split("/")
if len(parts) != 2 or not all(parts):
raise ValueError("repository must be owner/name")
encoded = [quote(part, safe="-._~") for part in (*parts, job_name)]
if not job_name.strip() or any(not part for part in encoded):
if not job_name.strip():
raise ValueError("identity components must not be empty")
return "/".join(encoded)
return "/".join(safe_identity_segment(part) for part in (*parts, job_name))
def safe_identity_segment(value: str) -> str:
"""Map arbitrary Gitea names to stable SPIFFE Operator path segments."""
if SPIFFE_PATH_SEGMENT.fullmatch(value):
return value
slug = re.sub(r"[^A-Za-z0-9._-]+", "-", value).strip("-._")
slug = slug[:48].rstrip("-._") or "segment"
digest = hashlib.sha256(value.encode()).hexdigest()[:12]
return f"{slug}-{digest}"
def pod_manifest(request: RunnerRequest, pod_name: str) -> dict[str, object]:
@@ -174,36 +185,90 @@ async def run_request(message: object, client: KubernetesClient) -> None:
return
pod_name = f"gitea-pod-{uuid.uuid4().hex[:12]}"
LOG.info(
"runner request received job_id=%d run_id=%d repository=%s job_name=%r "
"backend=%s pod=%s",
request.job_id,
request.run_id,
request.repository,
request.job_name,
request.backend,
pod_name,
)
stop = asyncio.Event()
pulse = asyncio.create_task(heartbeat(message, stop))
try:
await client.create_pod(pod_manifest(request, pod_name))
LOG.info(
"runner Pod created pod=%s job_id=%d run_id=%d image=%s",
pod_name,
request.job_id,
request.run_id,
RUNNER_IMAGE,
)
if await wait_for_pod(client, pod_name):
await message.ack()
LOG.info(
"runner request acknowledged pod=%s job_id=%d result=succeeded",
pod_name,
request.job_id,
)
else:
LOG.error("runner Pod %s failed", pod_name)
LOG.error(
"runner Pod failed pod=%s job_id=%d; request will be retried",
pod_name,
request.job_id,
)
await message.nak(delay=30)
except Exception:
LOG.exception(
"runner request failed pod=%s job_id=%d; request will be retried",
pod_name,
request.job_id,
)
await message.nak(delay=30)
raise
finally:
stop.set()
await pulse
LOG.info("deleting runner Pod pod=%s job_id=%d", pod_name, request.job_id)
await client.delete_pod(pod_name)
async def bind_request(message: object, client: KubernetesClient) -> None:
binding: IdentityBinding | None = None
try:
document = json.loads(message.data)
binding = IdentityBinding(**document)
if binding.backend != "pod" or not binding.runner_name.startswith("gitea-pod-"):
raise ValueError("invalid Pod identity binding")
path = identity_path(binding.repository, binding.job_name)
LOG.info(
"identity binding received job_id=%d run_id=%d runner_name=%s "
"repository=%s job_name=%r identity_path=%s",
binding.job_id,
binding.run_id,
binding.runner_name,
binding.repository,
binding.job_name,
path,
)
await client.bind_identity(binding.runner_name, path)
except ClientResponseError as error:
if error.status == 404:
LOG.warning(
"identity binding Pod not found runner_name=%s job_id=%s; binding will be retried",
binding.runner_name if binding else None,
binding.job_id if binding else None,
)
await message.nak(delay=2)
return
LOG.exception(
"identity binding Kubernetes request failed runner_name=%s job_id=%s status=%d",
binding.runner_name if binding else None,
binding.job_id if binding else None,
error.status,
)
await message.nak(delay=30)
raise
except (json.JSONDecodeError, TypeError, ValueError) as error:
@@ -211,6 +276,12 @@ async def bind_request(message: object, client: KubernetesClient) -> None:
await message.ack()
return
await message.ack()
LOG.info(
"identity binding applied and acknowledged runner_name=%s job_id=%d identity_path=%s",
binding.runner_name,
binding.job_id,
path,
)
async def consume_requests(subscription: object, client: KubernetesClient) -> None:
@@ -0,0 +1,132 @@
"""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]
+11
View File
@@ -31,6 +31,9 @@ REGISTRATION_TOKEN_FILE = Path(os.environ.get("REGISTRATION_TOKEN_FILE", "/etc/m
LAUNCHER = os.environ.get("LAUNCHER", "/usr/local/libexec/microvm-runner-launch")
TOKEN_LISTEN = os.environ.get("TOKEN_LISTEN", "172.30.0.1")
TOKEN_PORT = int(os.environ.get("TOKEN_PORT", "8787"))
GUEST_ASSETS = Path(
os.environ.get("GUEST_ASSETS", "/opt/gitea-dynamic-runner/guest-assets.tar.gz")
)
tokens: dict[str, bytes] = {}
token_lock = asyncio.Lock()
@@ -45,6 +48,13 @@ async def token(request: web.Request) -> web.Response:
return web.Response(body=value, headers={"Cache-Control": "no-store"})
async def guest_assets(_: web.Request) -> web.FileResponse:
return web.FileResponse(
GUEST_ASSETS,
headers={"Cache-Control": "public, immutable"},
)
async def heartbeat(message: object, stop: asyncio.Event) -> None:
while True:
try:
@@ -162,6 +172,7 @@ async def consume() -> None:
async def main() -> None:
app = web.Application()
app.router.add_get("/token/{nonce}", token)
app.router.add_get("/assets/guest-assets.tar.gz", guest_assets)
runner = web.AppRunner(app)
await runner.setup()
await web.TCPSite(runner, TOKEN_LISTEN, TOKEN_PORT).start()
+114
View File
@@ -0,0 +1,114 @@
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
+16 -2
View File
@@ -16,8 +16,22 @@ def request() -> RunnerRequest:
def test_identity_path_is_meaningful_and_uri_safe():
assert pod_worker.identity_path("panxiao81/example", "publish image") == (
"panxiao81/example/publish%20image"
path = pod_worker.identity_path("panxiao81/example", "publish image")
assert path.startswith("panxiao81/example/publish-image-")
assert "%" not in path
assert all(pod_worker.SPIFFE_PATH_SEGMENT.fullmatch(part) for part in path.split("/"))
assert pod_worker.identity_path("panxiao81/example", "lint") == "panxiao81/example/lint"
assert path == pod_worker.identity_path("panxiao81/example", "publish image")
assert path != pod_worker.identity_path("panxiao81/example", "publish-image")
real_path = pod_worker.identity_path(
"panxiao81/postgresql-tenant-operator", "Run on Ubuntu"
)
assert real_path.startswith(
"panxiao81/postgresql-tenant-operator/Run-on-Ubuntu-"
)
assert all(
pod_worker.SPIFFE_PATH_SEGMENT.fullmatch(part)
for part in real_path.split("/")
)
with pytest.raises(ValueError):
pod_worker.identity_path("invalid", "test")