Compare commits

...
Author SHA1 Message Date
panxiao81 cc94438bad fix: 为 Pod Docker 数据配置独立卷
test / shell (pull_request) Successful in 28s
test / python (pull_request) Successful in 1m6s
test / go (pull_request) Successful in 3m8s
2026-09-21 15:34:54 +00:00
panxiao81 7d90f28b73 Merge 放宽 Kata kind 的 kubeadm API 超时
test / go (push) Successful in 2m56s
test / shell (push) Successful in 34s
test / python (push) Successful in 51s
2026-09-21 12:57:19 +00:00
panxiao81 8f8ec04b18 ci: 放宽 Kata kind 的 kubeadm API 超时
test / go (pull_request) Successful in 2m53s
test / python (pull_request) Successful in 45s
test / shell (pull_request) Successful in 27s
2026-09-21 12:53:02 +00:00
panxiao81 74ffd49d08 Merge Kata nested kubelet kmsg 修复
test / shell (push) Successful in 50s
test / python (push) Successful in 1m3s
test / go (push) Successful in 2m59s
2026-09-21 12:26:25 +00:00
panxiao81 e3ff308772 ci: 为 Kata nested kubelet 补充 kmsg 设备
test / shell (pull_request) Successful in 45s
test / go (pull_request) Successful in 3m0s
test / python (pull_request) Successful in 1m3s
2026-09-21 12:21:33 +00:00
panxiao81 feb0b84b7c Merge Kata guest loop 设备节点修复
test / shell (push) Successful in 50s
test / python (push) Successful in 1m10s
test / go (push) Successful in 3m10s
2026-09-21 12:15:14 +00:00
panxiao81 66b90146f8 ci: 在 Kata guest 中创建设备节点挂载 loop 盘
test / python (pull_request) Successful in 1m1s
test / shell (pull_request) Successful in 55s
test / go (pull_request) Successful in 3m2s
2026-09-21 12:11:05 +00:00
panxiao81 3d8a04e4f7 Merge 统一 Pod 与 VM executor 镜像和 Docker 启动模型
test / shell (push) Successful in 24s
test / python (push) Successful in 1m0s
test / go (push) Successful in 3m36s
2026-09-21 12:05:21 +00:00
panxiao81 fef7e5a214 ci: 由 VM 任务启动本地 Docker daemon
test / shell (pull_request) Successful in 41s
test / python (pull_request) Successful in 1m29s
test / go (pull_request) Successful in 3m7s
2026-09-21 12:00:55 +00:00
panxiao81 51b940468e ci: 收集 nested containerd 阻塞现场 2026-09-21 11:41:16 +00:00
panxiao81 b614ef4c2f ci: 收集 kind CRI 容器日志 2026-09-21 11:24:40 +00:00
panxiao81 8fa8e46320 ci: 输出 DinD cgroup 层级 2026-09-21 11:09:55 +00:00
panxiao81 cc4405f788 ci: 保留 kind 失败现场 2026-09-21 11:05:03 +00:00
panxiao81 4a428c4384 Merge pull request '修正 VM smoke 的 SPIFFE socket 检查' (#39) from fix/vm-runtime-socket-check into main
test / shell (push) Successful in 25s
test / python (push) Successful in 57s
test / go (push) Successful in 3m3s
2026-09-21 10:49:53 +00:00
panxiao81 9a28a1573e ci: 修正 SPIFFE socket URI 检查
test / python (pull_request) Successful in 40s
test / shell (pull_request) Successful in 25s
test / go (pull_request) Successful in 2m49s
2026-09-21 10:44:36 +00:00
panxiao81 01995bf084 Merge pull request '增加 VM runtime 最小集成测试' (#38) from test/vm-runtime-smoke into main
test / python (push) Successful in 59s
test / shell (push) Successful in 29s
test / go (push) Successful in 3m6s
2026-09-21 10:41:42 +00:00
panxiao81 3d787b2dd3 ci: 增加 VM runtime 最小集成测试
test / python (pull_request) Successful in 1m12s
test / shell (pull_request) Successful in 37s
test / go (pull_request) Successful in 3m20s
2026-09-21 10:37:14 +00:00
panxiao81 54661411e3 Merge pull request '增加 Runner 生命周期结构化日志' (#37) from fix/structured-runner-lifecycle-logs into main
test / python (push) Successful in 53s
test / shell (push) Successful in 33s
test / go (push) Successful in 3m14s
publish images / publish-images (push) Failing after 1m40s
2026-09-21 09:38:45 +00:00
panxiao81 53a080b310 ci: 为测试显式选择 Pod runner
test / shell (pull_request) Successful in 1m34s
test / python (pull_request) Successful in 1m13s
test / go (pull_request) Successful in 3m14s
2026-09-21 09:34:30 +00:00
panxiao81 38e8d59541 feat: 增加 Runner 生命周期结构化日志
test / go (pull_request) Successful in 3m10s
test / shell (pull_request) Failing after 10m4s
test / python (pull_request) Failing after 10m4s
2026-09-21 09:31:46 +00:00
panxiao81 661b5e9218 Merge vm-dev 集成测试隔离标签
publish images / publish-images (push) Failing after 11m12s
test / python (push) Successful in 24s
test / shell (push) Successful in 28s
test / go (push) Successful in 6m44s
2026-09-21 08:51:17 +00:00
panxiao81 48b6b8038e feat: 隔离 VM 集成测试标签
test / python (pull_request) Failing after 11m44s
test / shell (pull_request) Successful in 21s
test / go (pull_request) Successful in 5m12s
2026-09-21 08:50:18 +00:00
panxiao81 70c5ff422f Merge assignment 重投预算修复
publish images / publish-images (push) Failing after 12m48s
test / python (push) Successful in 31s
test / shell (push) Successful in 34s
test / go (push) Successful in 5m40s
2026-09-21 08:44:42 +00:00
panxiao81 906e6a2e18 fix: 保持后端 assignment 可重试
test / shell (pull_request) Failing after 13m28s
test / python (pull_request) Failing after 13m29s
test / go (pull_request) Successful in 6m15s
2026-09-21 08:44:00 +00:00
panxiao81 ace84373f6 Merge OpenSandbox Runner workload 标记
test / python (push) Successful in 25s
test / shell (push) Successful in 30s
test / go (push) Successful in 4m30s
publish images / publish-images (push) Failing after 13m58s
2026-09-21 08:28:32 +00:00
panxiao81 87cf359ef6 fix: 标记 OpenSandbox Runner workload
test / shell (pull_request) Successful in 26s
test / go (pull_request) Successful in 2m58s
test / python (pull_request) Failing after 14m15s
2026-09-21 08:22:56 +00:00
panxiao81 ad3934e7e1 Merge OpenSandbox metadata 编码修复
test / python (push) Successful in 21s
test / shell (push) Successful in 24s
test / go (push) Successful in 6m52s
publish images / publish-images (push) Failing after 10m40s
2026-09-21 07:56:50 +00:00
panxiao81 d82687382a fix: 编码 OpenSandbox assignment metadata
test / go (pull_request) Successful in 4m45s
test / shell (pull_request) Failing after 13m27s
test / python (pull_request) Failing after 13m27s
2026-09-21 07:48:10 +00:00
panxiao81 7f52cf393f Merge Pod 与 VM 独立执行容量池
publish images / publish-images (push) Failing after 11m8s
test / python (push) Successful in 10s
test / shell (push) Successful in 26s
test / go (push) Successful in 5m51s
2026-09-21 07:26:10 +00:00
panxiao81 8a186dcd86 Merge Runner 无中断恢复与容量隔离
test / python (push) Successful in 25s
test / shell (push) Successful in 29s
test / go (push) Successful in 2m28s
publish images / publish-images (push) Failing after 11m13s
2026-09-21 06:48:11 +00:00
panxiao81 327c73e744 Merge Go Runner 与 OpenSandbox VM 集成
test / python (push) Successful in 15s
test / shell (push) Successful in 21s
publish images / publish-images (push) Failing after 13m58s
test / go (push) Failing after 15m46s
2026-09-21 06:11:40 +00:00
19 changed files with 480 additions and 44 deletions
+40
View File
@@ -1,5 +1,7 @@
---
name: dynamic Pod smoke test name: dynamic Pod smoke test
# yamllint disable-line rule:truthy
on: on:
workflow_dispatch: workflow_dispatch:
@@ -17,3 +19,41 @@ jobs:
-socketPath /run/spire/agent-sockets/spire-agent.sock \ -socketPath /run/spire/agent-sockets/spire-agent.sock \
>/dev/null >/dev/null
test "$(id -u)" = 2000 test "$(id -u)" = 2000
- name: Start job-local Docker
shell: bash
run: |
set -euo pipefail
findmnt /var/lib/docker
sudo nohup dockerd \
--host=unix:///var/run/docker.sock \
--storage-driver=overlay2 \
>/tmp/dockerd.log 2>&1 &
for _ in {1..60}; do
if docker info >/dev/null 2>&1; then
exit 0
fi
sleep 1
done
cat /tmp/dockerd.log
exit 1
- name: Build and run image
shell: bash
run: |
set -euo pipefail
context=$(mktemp -d)
cleanup() {
docker image rm --force pod-docker-smoke:test \
>/dev/null 2>&1 || true
rm -rf -- "$context"
}
trap cleanup EXIT
printf '%s\n' \
'FROM alpine:3.22' \
'RUN printf pod-docker-ok >/result' \
>"$context/Dockerfile"
docker build --tag pod-docker-smoke:test "$context"
output=$(docker run --rm pod-docker-smoke:test cat /result)
test "$output" = pod-docker-ok
test "$(docker info --format '{{.Driver}}')" = overlay2
+2 -2
View File
@@ -17,7 +17,7 @@ jobs:
- run: go vet ./... - run: go vet ./...
python: python:
runs-on: self-hosted runs-on: [self-hosted, pod]
steps: steps:
- uses: actions/checkout@v4 - uses: actions/checkout@v4
- uses: actions/setup-python@v5 - uses: actions/setup-python@v5
@@ -28,7 +28,7 @@ jobs:
- run: python -m compileall -q src - run: python -m compileall -q src
shell: shell:
runs-on: self-hosted runs-on: [self-hosted, pod]
steps: steps:
- uses: actions/checkout@v4 - uses: actions/checkout@v4
- run: | - run: |
+150 -7
View File
@@ -1,31 +1,174 @@
---
name: VM kind smoke name: VM kind smoke
# yamllint disable-line rule:truthy
on: on:
workflow_dispatch: workflow_dispatch:
jobs: jobs:
kind: kind:
runs-on: [self-hosted, vm] runs-on: [self-hosted, vm-dev]
steps: steps:
- name: Start job-local Docker
shell: bash
run: |
set -euo pipefail
if [[ ! -e /dev/kmsg ]]; then
sudo mknod /dev/kmsg c 1 11
fi
sudo install -d /var/lib/docker
sudo truncate -s 20G /tmp/docker-data.img
sudo mkfs.ext4 -F /tmp/docker-data.img
if [[ ! -e /dev/loop-control ]]; then
sudo mknod /dev/loop-control c 10 237
fi
for minor in {0..7}; do
if [[ ! -e "/dev/loop${minor}" ]]; then
sudo mknod "/dev/loop${minor}" b 7 "$minor"
fi
done
loop_device=$(sudo losetup --find --show /tmp/docker-data.img)
sudo mount "$loop_device" /var/lib/docker
findmnt /var/lib/docker
# A nested systemd needs a domain cgroup namespace. This is the
# cgroup v2 nesting initialization performed by the official DinD
# entrypoint, kept here because the shared runner image deliberately
# does not carry a second DinD-specific entrypoint.
if [[ -f /sys/fs/cgroup/cgroup.controllers ]]; then
sudo sh -c '
mkdir -p /sys/fs/cgroup/init
while read -r pid; do
printf "%s\n" "$pid" \
>/sys/fs/cgroup/init/cgroup.procs 2>/dev/null || true
done </sys/fs/cgroup/cgroup.procs
sed -e "s/ / +/g" -e "s/^/+/" \
/sys/fs/cgroup/cgroup.controllers \
>/sys/fs/cgroup/cgroup.subtree_control
'
fi
sudo nohup dockerd \
--host=unix:///var/run/docker.sock \
--storage-driver=overlay2 \
>/tmp/dockerd.log 2>&1 &
for _ in {1..60}; do
if docker info >/dev/null 2>&1; then
exit 0
fi
sleep 1
done
cat /tmp/dockerd.log
exit 1
- name: Verify Docker - name: Verify Docker
run: docker info shell: bash
run: |
set -euo pipefail
docker info
echo '### runner cgroup'
cat /proc/self/cgroup
cat /sys/fs/cgroup/cgroup.type
cat /sys/fs/cgroup/cgroup.controllers
cat /sys/fs/cgroup/cgroup.subtree_control
echo '### nested private cgroup namespace'
docker run --rm --privileged --cgroupns=private alpine:3.22 \
sh -c 'cat /proc/self/cgroup; cat /sys/fs/cgroup/cgroup.type'
echo '### nested host cgroup namespace'
docker run --rm --privileged --cgroupns=host alpine:3.22 \
sh -c 'cat /proc/self/cgroup; cat /sys/fs/cgroup/cgroup.type'
- name: Install kind - name: Install kind
shell: bash shell: bash
run: | run: |
set -euo pipefail set -euo pipefail
version=v0.33.0 version=v0.33.0
base_url="https://kind.sigs.k8s.io/dl/${version}"
curl --fail --location --silent --show-error \ curl --fail --location --silent --show-error \
--output /tmp/kind "https://kind.sigs.k8s.io/dl/${version}/kind-linux-amd64" --output /tmp/kind "${base_url}/kind-linux-amd64"
curl --fail --location --silent --show-error \ curl --fail --location --silent --show-error \
--output /tmp/kind.sha256sum "https://kind.sigs.k8s.io/dl/${version}/kind-linux-amd64.sha256sum" --output /tmp/kind.sha256sum \
printf '%s %s\n' "$(cut -d ' ' -f1 /tmp/kind.sha256sum)" /tmp/kind | sha256sum --check "${base_url}/kind-linux-amd64.sha256sum"
checksum=$(cut -d ' ' -f1 /tmp/kind.sha256sum)
printf '%s %s\n' "$checksum" /tmp/kind | sha256sum --check
chmod 0755 /tmp/kind chmod 0755 /tmp/kind
- name: Create and delete kind cluster - name: Create and delete kind cluster
shell: bash shell: bash
run: | run: |
set -euo pipefail set -euo pipefail
trap '/tmp/kind delete cluster --name smoke' EXIT diagnose_and_cleanup() {
/tmp/kind create cluster --name smoke --wait 180s status=$?
node=smoke-control-plane
if (( status != 0 )) && docker inspect "$node" >/dev/null 2>&1; then
echo '::group::kind node inspect'
docker inspect "$node"
echo '::endgroup::'
echo '::group::kind node logs'
docker logs "$node" 2>&1 || true
echo '::endgroup::'
echo '::group::kind node guest state'
docker exec "$node" bash -c '
set +e
echo "### pid 1"
ps -p 1 -o pid,ppid,user,stat,comm,args
cat /proc/1/status
echo "### cgroup"
cat /proc/1/cgroup
findmnt -R /sys/fs/cgroup
stat -fc "%T %a" /sys/fs/cgroup
echo "### systemd"
systemctl --no-pager --failed
systemctl --no-pager status \
multi-user.target containerd.service kubelet.service
echo "### CRI containers"
endpoint=unix:///run/containerd/containerd.sock
crictl --runtime-endpoint "$endpoint" ps --all
for id in $(
crictl --runtime-endpoint "$endpoint" ps --all --quiet
); do
echo "### CRI container $id"
crictl --runtime-endpoint "$endpoint" inspect "$id"
crictl --runtime-endpoint "$endpoint" logs "$id"
done
echo "### containerd metadata"
timeout 10 ctr --namespace k8s.io containers list
timeout 10 ctr --namespace k8s.io snapshots list
echo "### runtime process stacks"
ps -e -o pid,ppid,stat,wchan:32,comm,args
for pid in $(pidof containerd containerd-shim-runc-v2); do
echo "### kernel stack $pid"
cat "/proc/$pid/stack"
done
kill -USR1 "$(pidof containerd)"
sleep 2
journalctl --no-pager -b -n 500
' 2>&1 || true
echo '::endgroup::'
fi
/tmp/kind delete cluster --name smoke || true
exit "$status"
}
trap diagnose_and_cleanup EXIT
# Nested Kata + Docker + kind cold starts can take longer than
# kubeadm's one-minute API-call default even after the static pods
# have been accepted. Give the API server enough time to become
# responsive before kubeadm creates its initial RBAC objects.
cat >/tmp/kind-config.yml <<'EOF'
kind: Cluster
apiVersion: kind.x-k8s.io/v1alpha4
nodes:
- role: control-plane
kubeadmConfigPatches:
- |
apiVersion: kubeadm.k8s.io/v1beta4
kind: InitConfiguration
timeouts:
kubernetesAPICall: 5m0s
EOF
/tmp/kind create cluster \
--name smoke \
--config /tmp/kind-config.yml \
--wait 300s \
--retain
/tmp/kind get clusters | grep -Fx smoke /tmp/kind get clusters | grep -Fx smoke
+36
View File
@@ -0,0 +1,36 @@
---
name: VM runtime smoke
# yamllint disable-line rule:truthy
on:
workflow_dispatch:
jobs:
runtime:
runs-on: [self-hosted, vm-dev]
steps:
- name: Verify workload identity socket
shell: bash
run: |
set -euo pipefail
: "${SPIFFE_ENDPOINT_SOCKET:?SPIFFE_ENDPOINT_SOCKET is required}"
socket_path=${SPIFFE_ENDPOINT_SOCKET#unix://}
test -S "$socket_path"
- name: Verify Docker daemon
shell: bash
run: |
set -euo pipefail
format='{{json .ServerVersion}} {{json .Driver}}'
format="$format {{json .CgroupVersion}}"
docker info --format "$format"
- name: Run and clean nested container
shell: bash
run: |
set -euo pipefail
name=vm-runtime-smoke
trap 'docker rm --force "$name" >/dev/null 2>&1 || true' EXIT
output=$(docker run --name "$name" alpine:3.22 /bin/sh -c \
'test "$(uname -m)" = x86_64 && printf vm-runtime-ok')
test "$output" = vm-runtime-ok
+4
View File
@@ -15,6 +15,10 @@ runs-on: [self-hosted, vm]
只执行一个 job,并在 job 结束后连同本地状态一起销毁。完整的设计约束见 只执行一个 job,并在 job 结束后连同本地状态一起销毁。完整的设计约束见
[`docs/design-principles.md`](docs/design-principles.md)。 [`docs/design-principles.md`](docs/design-principles.md)。
集成期间可将 `VM_RUNNER_LABEL=vm-dev`,只接取显式使用
`runs-on: [self-hosted, vm-dev]` 的测试任务;生产 `vm` job 将保持在 Gitea pending,
不会在 backend 修复过程中继续涌入。
目标 Go controller 组件: 目标 Go controller 组件:
- `scheduler`:以常驻 Gitea RunnerService 身份直接领取 task,并把完整 assignment - `scheduler`:以常驻 Gitea RunnerService 身份直接领取 task,并把完整 assignment
+28 -11
View File
@@ -4,7 +4,7 @@ import (
"context" "context"
"errors" "errors"
"fmt" "fmt"
"log" "log/slog"
"net/http" "net/http"
"os" "os"
"slices" "slices"
@@ -49,7 +49,7 @@ type controllerConfig struct {
PodNamespace, PodImage, PodServiceAccount, SPIRECluster, SPIREClass string PodNamespace, PodImage, PodServiceAccount, SPIRECluster, SPIREClass string
SPIREAgentID string SPIREAgentID string
PodExecutorUID, PodCapacity int PodExecutorUID, PodCapacity int
OpenSandboxURL, OpenSandboxAPIKey, OpenSandboxPool string OpenSandboxURL, OpenSandboxAPIKey, OpenSandboxPool, VMRunnerLabel string
VMTimeout, VMCapacity int VMTimeout, VMCapacity int
} }
@@ -131,7 +131,7 @@ func runController(ctx context.Context) error {
labels = append(labels, string(taskassignment.BackendPod)) labels = append(labels, string(taskassignment.BackendPod))
} }
if slices.Contains(config.Components, controller.VMWorker) { if slices.Contains(config.Components, controller.VMWorker) {
labels = append(labels, string(taskassignment.BackendVM)) labels = append(labels, config.VMRunnerLabel)
} }
poller := taskscheduler.Poller{ poller := taskscheduler.Poller{
Client: giteaClient, Client: giteaClient,
@@ -139,7 +139,7 @@ func runController(ctx context.Context) error {
JetStream: producerJS, SubjectBase: config.SubjectBase, JetStream: producerJS, SubjectBase: config.SubjectBase,
}}, }},
Config: taskscheduler.PollerConfig{Version: "gitea-dynamic-runner/0.4", Labels: labels, Capacity: config.PodCapacity + config.VMCapacity}, Config: taskscheduler.PollerConfig{Version: "gitea-dynamic-runner/0.4", Labels: labels, Capacity: config.PodCapacity + config.VMCapacity},
OnError: func(err error) { log.Printf("scheduler: %v", err) }, OnError: func(err error) { slog.Error("scheduler error", "component", "scheduler", "error", err) },
} }
kubernetesConfig, err := rest.InClusterConfig() kubernetesConfig, err := rest.InClusterConfig()
if err != nil { if err != nil {
@@ -190,11 +190,13 @@ func runController(ctx context.Context) error {
return fmt.Errorf("recover Pod facade claim %s: %w", assignment.ID, err) return fmt.Errorf("recover Pod facade claim %s: %w", assignment.ID, err)
} }
} }
component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendPod, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry, podPool) component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendPod, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap, OnEvent: workerEventLogger(taskassignment.BackendPod)}, registry, podPool)
if err != nil { if err != nil {
return err return err
} }
lifecycle := podbackend.Lifecycle{Backend: backend, OnError: func(err error) { log.Printf("pod lifecycle: %v", err) }} lifecycle := podbackend.Lifecycle{Backend: backend, OnError: func(err error) {
slog.Error("backend lifecycle error", "component", "lifecycle", "backend", taskassignment.BackendPod, "error", err)
}}
components[controller.PodWorker] = runComponent(func(ctx context.Context) error { components[controller.PodWorker] = runComponent(func(ctx context.Context) error {
group, groupContext := errgroup.WithContext(ctx) group, groupContext := errgroup.WithContext(ctx)
group.Go(func() error { return component.Run(groupContext) }) group.Go(func() error { return component.Run(groupContext) })
@@ -220,11 +222,13 @@ func runController(ctx context.Context) error {
return fmt.Errorf("recover VM facade claim %s: %w", assignment.ID, err) return fmt.Errorf("recover VM facade claim %s: %w", assignment.ID, err)
} }
} }
component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendVM, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry, vmPool) component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendVM, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap, OnEvent: workerEventLogger(taskassignment.BackendVM)}, registry, vmPool)
if err != nil { if err != nil {
return err return err
} }
lifecycleReconciler := opensandboxbackend.LifecycleReconciler{Backend: backend, OnError: func(err error) { log.Printf("VM lifecycle: %v", err) }} lifecycleReconciler := opensandboxbackend.LifecycleReconciler{Backend: backend, OnError: func(err error) {
slog.Error("backend lifecycle error", "component", "lifecycle", "backend", taskassignment.BackendVM, "error", err)
}}
components[controller.VMWorker] = runComponent(func(ctx context.Context) error { components[controller.VMWorker] = runComponent(func(ctx context.Context) error {
group, groupContext := errgroup.WithContext(ctx) group, groupContext := errgroup.WithContext(ctx)
group.Go(func() error { return component.Run(groupContext) }) group.Go(func() error { return component.Run(groupContext) })
@@ -294,11 +298,21 @@ func workerComponent(ctx context.Context, js jetstream.JetStream, config control
} }
return assignmentqueue.ConsumerComponent{ return assignmentqueue.ConsumerComponent{
Consumer: consumer, Capacity: capacity, Consumer: consumer, Capacity: capacity,
Processor: assignmentqueue.Processor{TrustDomain: config.TrustDomain, Accepter: accepter, Claims: claims, Admission: admission}, Processor: assignmentqueue.Processor{TrustDomain: config.TrustDomain, Accepter: accepter, Claims: claims, Admission: admission, OnEvent: func(event assignmentqueue.Event) {
OnError: func(err error) { log.Printf("%s worker: %v", backend, err) }, slog.Info("assignment transition", "component", "worker", "event", event.Name, "backend", event.Backend, "assignment", event.AssignmentID, "retry_delay", event.RetryDelay)
}},
OnError: func(err error) {
slog.Error("assignment processing error", "component", "worker", "backend", backend, "error", err)
},
}, nil }, nil
} }
func workerEventLogger(backend taskassignment.Backend) func(taskworker.Event) {
return func(event taskworker.Event) {
slog.Info("executor transition", "component", "worker", "event", event.Name, "backend", backend, "assignment", event.AssignmentID, "executor", event.Executor, "phase", event.Phase)
}
}
func loadControllerConfig() (controllerConfig, error) { func loadControllerConfig() (controllerConfig, error) {
selection, err := controller.ParseSelection(os.Getenv("COMPONENTS")) selection, err := controller.ParseSelection(os.Getenv("COMPONENTS"))
if err != nil { if err != nil {
@@ -345,7 +359,7 @@ func loadControllerConfig() (controllerConfig, error) {
FacadeListen: env("RUNNER_FACADE_LISTEN", ":8443"), FacadeURL: os.Getenv("RUNNER_FACADE_URL"), FacadeSPIFFEID: os.Getenv("RUNNER_FACADE_SPIFFE_ID"), CapabilityKey: []byte(capabilityKey), FacadeListen: env("RUNNER_FACADE_LISTEN", ":8443"), FacadeURL: os.Getenv("RUNNER_FACADE_URL"), FacadeSPIFFEID: os.Getenv("RUNNER_FACADE_SPIFFE_ID"), CapabilityKey: []byte(capabilityKey),
PodNamespace: env("POD_NAMESPACE", "gitea-actions"), PodImage: os.Getenv("POD_EXECUTOR_IMAGE"), PodServiceAccount: env("POD_SERVICE_ACCOUNT", "gitea-task-executor"), PodNamespace: env("POD_NAMESPACE", "gitea-actions"), PodImage: os.Getenv("POD_EXECUTOR_IMAGE"), PodServiceAccount: env("POD_SERVICE_ACCOUNT", "gitea-task-executor"),
SPIRECluster: env("SPIRE_CLUSTER", "homelab"), SPIREClass: env("SPIRE_CLASS", "spire-mgmt-spire"), SPIREAgentID: os.Getenv("SPIRE_AGENT_ID"), PodExecutorUID: envInt("POD_EXECUTOR_UID", 2000), PodCapacity: envInt("POD_CAPACITY", 4), SPIRECluster: env("SPIRE_CLUSTER", "homelab"), SPIREClass: env("SPIRE_CLASS", "spire-mgmt-spire"), SPIREAgentID: os.Getenv("SPIRE_AGENT_ID"), PodExecutorUID: envInt("POD_EXECUTOR_UID", 2000), PodCapacity: envInt("POD_CAPACITY", 4),
OpenSandboxURL: os.Getenv("OPENSANDBOX_API"), OpenSandboxPool: env("OPENSANDBOX_POOL", "ci-vm"), VMTimeout: envInt("VM_TIMEOUT_SECONDS", 14400), VMCapacity: envInt("VM_CAPACITY", 1), OpenSandboxURL: os.Getenv("OPENSANDBOX_API"), OpenSandboxPool: env("OPENSANDBOX_POOL", "ci-vm"), VMRunnerLabel: env("VM_RUNNER_LABEL", "vm"), VMTimeout: envInt("VM_TIMEOUT_SECONDS", 14400), VMCapacity: envInt("VM_CAPACITY", 1),
} }
if config.WorkloadAPIAddr == "" || config.FacadeURL == "" || config.FacadeSPIFFEID == "" { if config.WorkloadAPIAddr == "" || config.FacadeURL == "" || config.FacadeSPIFFEID == "" {
return controllerConfig{}, errors.New("SPIFFE_ENDPOINT_SOCKET, RUNNER_FACADE_URL, and RUNNER_FACADE_SPIFFE_ID are required") return controllerConfig{}, errors.New("SPIFFE_ENDPOINT_SOCKET, RUNNER_FACADE_URL, and RUNNER_FACADE_SPIFFE_ID are required")
@@ -357,6 +371,9 @@ func loadControllerConfig() (controllerConfig, error) {
return controllerConfig{}, errors.New("SPIRE_AGENT_ID is required for pod-worker") return controllerConfig{}, errors.New("SPIRE_AGENT_ID is required for pod-worker")
} }
if slices.Contains(selection, controller.VMWorker) { if slices.Contains(selection, controller.VMWorker) {
if config.VMRunnerLabel != "vm" && config.VMRunnerLabel != "vm-dev" {
return controllerConfig{}, errors.New("VM_RUNNER_LABEL must be vm or vm-dev")
}
if config.OpenSandboxURL == "" { if config.OpenSandboxURL == "" {
return controllerConfig{}, errors.New("OPENSANDBOX_API is required for vm-worker") return controllerConfig{}, errors.New("OPENSANDBOX_API is required for vm-worker")
} }
@@ -56,6 +56,7 @@ func TestLoadControllerConfigRequiresOpenSandboxSecretOnlyForVM(t *testing.T) {
t.Setenv("RUNNER_FACADE_URL", "https://facade:8443") t.Setenv("RUNNER_FACADE_URL", "https://facade:8443")
t.Setenv("RUNNER_FACADE_SPIFFE_ID", "spiffe://ddupan.top/controller") t.Setenv("RUNNER_FACADE_SPIFFE_ID", "spiffe://ddupan.top/controller")
t.Setenv("OPENSANDBOX_API", "http://opensandbox.internal") t.Setenv("OPENSANDBOX_API", "http://opensandbox.internal")
t.Setenv("VM_RUNNER_LABEL", "vm-dev")
if _, err := loadControllerConfig(); err == nil { if _, err := loadControllerConfig(); err == nil {
t.Fatal("expected missing OpenSandbox API key file error") t.Fatal("expected missing OpenSandbox API key file error")
} }
+3 -1
View File
@@ -4,6 +4,7 @@ import (
"context" "context"
"errors" "errors"
"fmt" "fmt"
"log/slog"
"os" "os"
"os/signal" "os/signal"
"syscall" "syscall"
@@ -12,8 +13,9 @@ import (
) )
func main() { func main() {
slog.SetDefault(slog.New(slog.NewJSONHandler(os.Stderr, nil)))
if err := run(); err != nil { if err := run(); err != nil {
fmt.Fprintln(os.Stderr, err) slog.Error("runner stopped", "error", err)
os.Exit(1) os.Exit(1)
} }
} }
+5 -3
View File
@@ -40,9 +40,11 @@ UID attestation 的临时 Agent 失去父级。
`ci-vm` 使用 `kata-clh-runtime-rs`;`ci-pod` 使用默认 runc。两者都要求: `ci-vm` 使用 `kata-clh-runtime-rs`;`ci-pod` 使用默认 runc。两者都要求:
- runner 镜像包含 Gitea Runner、Node.js action userspace、SPIRE CLI 和 identity gate; - runner 镜像包含 Gitea Runner、Node.js action userspace、SPIRE CLI、Docker 工具和
- runner UID 2000,SPIRE Agent 与 privileged dockerd 使用不同 UID; identity gate;Pod 与 VM backend 使用同一个镜像;
- Docker socket 通过 group 2000 共享,Docker 数据仅存在于 sandbox emptyDir; - runner UID 2000;SPIRE Agent 独立运行,需要 Docker 的 workflow 通过 sudo 在
privileged executor 内启动 job-local daemon;
- Kata VM 中 Docker 数据使用 guest 内的 loop-backed ext4,并随 sandbox 一起删除;
- `self-hosted` 必须是所有 runner labels 的前缀; - `self-hosted` 必须是所有 runner labels 的前缀;
- ephemeral/once runner 完成一项任务后退出。 - ephemeral/once runner 完成一项任务后退出。
+26 -1
View File
@@ -78,6 +78,22 @@ type Processor struct {
Admission Admission Admission Admission
RetryDelay time.Duration RetryDelay time.Duration
ClaimTimeout time.Duration ClaimTimeout time.Duration
OnEvent func(Event)
}
// Event describes a non-sensitive assignment handoff transition. It never
// contains task payloads, credentials, capabilities, or workload identities.
type Event struct {
Name string
AssignmentID string
Backend taskassignment.Backend
RetryDelay time.Duration
}
func (p Processor) event(name string, assignment taskassignment.Assignment, retryDelay time.Duration) {
if p.OnEvent != nil {
p.OnEvent(Event{Name: name, AssignmentID: assignment.ID, Backend: assignment.Backend, RetryDelay: retryDelay})
}
} }
func (p Processor) Process(ctx context.Context, message Message) error { func (p Processor) Process(ctx context.Context, message Message) error {
@@ -88,6 +104,7 @@ func (p Processor) Process(ctx context.Context, message Message) error {
if err != nil { if err != nil {
return errors.Join(err, message.TermWithReason("invalid assignment")) return errors.Join(err, message.TermWithReason("invalid assignment"))
} }
p.event("received", assignment, 0)
if _, err := p.Claims.Offer(assignment); err != nil { if _, err := p.Claims.Offer(assignment); err != nil {
return errors.Join(err, message.TermWithReason("conflicting assignment")) return errors.Join(err, message.TermWithReason("conflicting assignment"))
} }
@@ -96,8 +113,10 @@ func (p Processor) Process(ctx context.Context, message Message) error {
if delay <= 0 { if delay <= 0 {
delay = 2 * time.Second delay = 2 * time.Second
} }
p.event("capacity_wait", assignment, delay)
return message.NakWithDelay(delay) return message.NakWithDelay(delay)
} }
p.event("capacity_acquired", assignment, 0)
accepted, err := p.Accepter.Accept(ctx, assignment) accepted, err := p.Accepter.Accept(ctx, assignment)
if err != nil { if err != nil {
p.Admission.Release(assignment.ID) p.Admission.Release(assignment.ID)
@@ -105,9 +124,11 @@ func (p Processor) Process(ctx context.Context, message Message) error {
if delay <= 0 { if delay <= 0 {
delay = 15 * time.Second delay = 15 * time.Second
} }
p.event("backend_retry", assignment, delay)
return errors.Join(err, message.NakWithDelay(delay)) return errors.Join(err, message.NakWithDelay(delay))
} }
if accepted { if accepted {
p.event("backend_ready", assignment, 0)
timeout := p.ClaimTimeout timeout := p.ClaimTimeout
if timeout <= 0 { if timeout <= 0 {
timeout = 4 * time.Minute timeout = 4 * time.Minute
@@ -120,17 +141,21 @@ func (p Processor) Process(ctx context.Context, message Message) error {
if delay <= 0 { if delay <= 0 {
delay = 2 * time.Second delay = 2 * time.Second
} }
p.event("claim_timeout", assignment, delay)
return errors.Join(err, message.NakWithDelay(delay)) return errors.Join(err, message.NakWithDelay(delay))
} }
p.event("runner_claimed", assignment, 0)
if err := message.DoubleAck(ctx); err != nil { if err := message.DoubleAck(ctx); err != nil {
return fmt.Errorf("ack assignment %s: %w", assignment.ID, err) return fmt.Errorf("ack assignment %s: %w", assignment.ID, err)
} }
p.event("acked", assignment, 0)
return nil return nil
} }
delay := p.RetryDelay delay := p.RetryDelay
if delay <= 0 { if delay <= 0 {
delay = 2 * time.Second delay = 2 * time.Second
} }
p.event("backend_pending", assignment, delay)
return message.NakWithDelay(delay) return message.NakWithDelay(delay)
} }
@@ -167,7 +192,7 @@ func OpenConsumer(ctx context.Context, manager consumerManager, stream, subjectB
AckPolicy: jetstream.AckExplicitPolicy, AckPolicy: jetstream.AckExplicitPolicy,
AckWait: 5 * time.Minute, AckWait: 5 * time.Minute,
MaxAckPending: capacity, MaxAckPending: capacity,
MaxDeliver: 20, MaxDeliver: 1000,
}) })
if err != nil { if err != nil {
return nil, fmt.Errorf("open %s assignment consumer: %w", backend, err) return nil, fmt.Errorf("open %s assignment consumer: %w", backend, err)
+12 -2
View File
@@ -126,13 +126,23 @@ func encodedAssignment(t *testing.T) []byte {
func TestProcessorAcknowledgesPersistedHandoff(t *testing.T) { func TestProcessorAcknowledgesPersistedHandoff(t *testing.T) {
message := &fakeMessage{data: encodedAssignment(t)} message := &fakeMessage{data: encodedAssignment(t)}
processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{accepted: true}, Claims: &fakeClaims{claimed: true}, Admission: &fakeAdmission{allowed: true}} var events []Event
processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{accepted: true}, Claims: &fakeClaims{claimed: true}, Admission: &fakeAdmission{allowed: true}, OnEvent: func(event Event) { events = append(events, event) }}
if err := processor.Process(context.Background(), message); err != nil { if err := processor.Process(context.Background(), message); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if message.acked != 1 || message.nacked != 0 { if message.acked != 1 || message.nacked != 0 {
t.Fatalf("message = %#v", message) t.Fatalf("message = %#v", message)
} }
want := []string{"received", "capacity_acquired", "backend_ready", "runner_claimed", "acked"}
if len(events) != len(want) {
t.Fatalf("events = %#v", events)
}
for index := range want {
if events[index].Name != want[index] || events[index].AssignmentID != "gitea-task-42" || events[index].Backend != taskassignment.BackendPod {
t.Fatalf("event[%d] = %#v", index, events[index])
}
}
} }
func TestProcessorRetriesUntilBackendHandoffIsDurable(t *testing.T) { func TestProcessorRetriesUntilBackendHandoffIsDurable(t *testing.T) {
@@ -205,7 +215,7 @@ func TestOpenConsumerUsesIndependentDurablePerBackend(t *testing.T) {
if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.BackendPod, 4); err != nil { if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.BackendPod, 4); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if manager.config.Durable != "pod" || manager.config.FilterSubject != "ci.assignment.pod" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 { if manager.config.Durable != "pod" || manager.config.FilterSubject != "ci.assignment.pod" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 || manager.config.MaxDeliver != 1000 {
t.Fatalf("config = %#v", manager.config) t.Fatalf("config = %#v", manager.config)
} }
} }
+69 -4
View File
@@ -3,9 +3,13 @@ package opensandboxbackend
import ( import (
"context" "context"
"encoding/base64"
"encoding/json"
"errors" "errors"
"fmt" "fmt"
"net/http" "net/http"
"sort"
"strings"
"time" "time"
opensandbox "github.com/alibaba/OpenSandbox/sdks/sandbox/go" opensandbox "github.com/alibaba/OpenSandbox/sdks/sandbox/go"
@@ -17,6 +21,8 @@ import (
const assignmentMetadata = "ci.ddupan.top/assignment-id" const assignmentMetadata = "ci.ddupan.top/assignment-id"
const terminalMetadata = "ci.ddupan.top/terminal" const terminalMetadata = "ci.ddupan.top/terminal"
const annotationsMetadataPrefix = "ci.ddupan.top/annotations-"
const metadataValueLimit = 63
type Lifecycle interface { type Lifecycle interface {
ListSandboxes(context.Context, opensandbox.ListOptions) (*opensandbox.ListSandboxesResponse, error) ListSandboxes(context.Context, opensandbox.ListOptions) (*opensandbox.ListSandboxesResponse, error)
@@ -123,7 +129,11 @@ func (b Backend) RecoverAssignments(ctx context.Context, trustDomain string) ([]
if sandbox.Metadata[terminalMetadata] == "true" { if sandbox.Metadata[terminalMetadata] == "true" {
continue continue
} }
assignment, err := taskassignment.FromMetadata(sandbox.Metadata, sandbox.Metadata, trustDomain) annotations, err := decodeAnnotations(sandbox.Metadata)
if err != nil {
return nil, fmt.Errorf("decode sandbox %s annotations: %w", sandbox.ID, err)
}
assignment, err := taskassignment.FromMetadata(sandbox.Metadata, annotations, trustDomain)
if err != nil { if err != nil {
return nil, fmt.Errorf("recover sandbox %s: %w", sandbox.ID, err) return nil, fmt.Errorf("recover sandbox %s: %w", sandbox.ID, err)
} }
@@ -169,8 +179,12 @@ func (b Backend) Create(ctx context.Context, assignment taskassignment.Assignmen
for key, value := range launch.Environment { for key, value := range launch.Environment {
environment[key] = value environment[key] = value
} }
sandboxMetadata := clone(launch.Metadata.Annotations) sandboxMetadata := clone(launch.Metadata.Labels)
for key, value := range launch.Metadata.Labels { annotations, err := encodeAnnotations(launch.Metadata.Annotations)
if err != nil {
return nil, fmt.Errorf("encode sandbox annotations: %w", err)
}
for key, value := range annotations {
sandboxMetadata[key] = value sandboxMetadata[key] = value
} }
request := opensandbox.CreateSandboxRequest{ request := opensandbox.CreateSandboxRequest{
@@ -194,7 +208,11 @@ func (b Backend) BindIdentity(ctx context.Context, executor *taskworker.Executor
if err != nil { if err != nil {
return fmt.Errorf("verify sandbox identity metadata: %w", err) return fmt.Errorf("verify sandbox identity metadata: %w", err)
} }
if sandbox.Metadata["ci.ddupan.top/spiffe-id"] != identity.SPIFFEID { annotations, err := decodeAnnotations(sandbox.Metadata)
if err != nil {
return fmt.Errorf("decode sandbox identity metadata: %w", err)
}
if annotations["ci.ddupan.top/spiffe-id"] != identity.SPIFFEID {
return fmt.Errorf("sandbox %s has inconsistent SPIFFE identity metadata", executor.Name) return fmt.Errorf("sandbox %s has inconsistent SPIFFE identity metadata", executor.Name)
} }
return nil return nil
@@ -243,3 +261,50 @@ func clone(source map[string]string) map[string]string {
} }
return result return result
} }
// OpenSandbox metadata follows Kubernetes label-value constraints, unlike Pod
// annotations. Store the annotation map as deterministic URL-safe base64
// chunks so repository paths and SPIFFE IDs remain lossless and recoverable.
func encodeAnnotations(annotations map[string]string) (map[string]string, error) {
data, err := json.Marshal(annotations)
if err != nil {
return nil, err
}
encoded := base64.RawURLEncoding.EncodeToString(data)
result := make(map[string]string, (len(encoded)+metadataValueLimit-1)/metadataValueLimit)
for index := 0; len(encoded) > 0; index++ {
length := min(metadataValueLimit, len(encoded))
result[fmt.Sprintf("%s%03d", annotationsMetadataPrefix, index)] = encoded[:length]
encoded = encoded[length:]
}
return result, nil
}
func decodeAnnotations(metadata map[string]string) (map[string]string, error) {
keys := make([]string, 0)
for key := range metadata {
if strings.HasPrefix(key, annotationsMetadataPrefix) {
keys = append(keys, key)
}
}
if len(keys) == 0 {
return nil, errors.New("sandbox annotation metadata is missing")
}
sort.Strings(keys)
var encoded strings.Builder
for index, key := range keys {
if key != fmt.Sprintf("%s%03d", annotationsMetadataPrefix, index) {
return nil, errors.New("sandbox annotation metadata chunks are incomplete")
}
encoded.WriteString(metadata[key])
}
data, err := base64.RawURLEncoding.DecodeString(encoded.String())
if err != nil {
return nil, err
}
var annotations map[string]string
if err := json.Unmarshal(data, &annotations); err != nil {
return nil, err
}
return annotations, nil
}
+29 -2
View File
@@ -31,7 +31,8 @@ func (f *fakeLifecycle) GetSandbox(_ context.Context, id string) (*opensandbox.S
return &item, nil return &item, nil
} }
} }
return &opensandbox.SandboxInfo{ID: id, Metadata: map[string]string{"ci.ddupan.top/spiffe-id": assignment().Identity.SPIFFEID}}, nil metadata, _ := encodeAnnotations(map[string]string{"ci.ddupan.top/spiffe-id": assignment().Identity.SPIFFEID})
return &opensandbox.SandboxInfo{ID: id, Metadata: metadata}, nil
} }
func (f *fakeLifecycle) PatchSandboxMetadata(_ context.Context, id string, patch opensandbox.MetadataPatch) (*opensandbox.SandboxInfo, error) { func (f *fakeLifecycle) PatchSandboxMetadata(_ context.Context, id string, patch opensandbox.MetadataPatch) (*opensandbox.SandboxInfo, error) {
for index := range f.items { for index := range f.items {
@@ -89,12 +90,38 @@ func TestCreateUsesPoolAndPersistsRecoveryMetadata(t *testing.T) {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if lifecycle.created.Extensions["poolRef"] != "ci-vm" || lifecycle.created.Metadata[assignmentMetadata] != assignment().ID || lifecycle.created.Env["CI_SPIFFE_ID"] != assignment().Identity.SPIFFEID || executor.Name != "sandbox-42" { if lifecycle.created.Extensions["poolRef"] != "ci-vm" || lifecycle.created.Metadata[assignmentMetadata] != assignment().ID || lifecycle.created.Metadata["ci.ddupan.top/runner"] != "true" || lifecycle.created.Env["CI_SPIFFE_ID"] != assignment().Identity.SPIFFEID || executor.Name != "sandbox-42" {
t.Fatalf("request=%#v executor=%#v", lifecycle.created, executor) t.Fatalf("request=%#v executor=%#v", lifecycle.created, executor)
} }
if lifecycle.created.Env["CI_RUNNER_CAPABILITY"] != "capability" { if lifecycle.created.Env["CI_RUNNER_CAPABILITY"] != "capability" {
t.Fatalf("environment = %#v", lifecycle.created.Env) t.Fatalf("environment = %#v", lifecycle.created.Env)
} }
annotations, err := decodeAnnotations(lifecycle.created.Metadata)
if err != nil || annotations["ci.ddupan.top/repository"] != "owner/repo" || annotations["ci.ddupan.top/spiffe-id"] != assignment().Identity.SPIFFEID {
t.Fatalf("annotations=%#v err=%v", annotations, err)
}
for _, value := range lifecycle.created.Metadata {
if len(value) > metadataValueLimit {
t.Fatalf("metadata value exceeds %d characters: %q", metadataValueLimit, value)
}
}
}
func TestAnnotationMetadataRoundTripPreservesSlashValues(t *testing.T) {
want := taskworker.BackendMetadata(assignment()).Annotations
encoded, err := encodeAnnotations(want)
if err != nil {
t.Fatal(err)
}
got, err := decodeAnnotations(encoded)
if err != nil {
t.Fatal(err)
}
for key, value := range want {
if got[key] != value {
t.Fatalf("%s=%q, want %q", key, got[key], value)
}
}
} }
func TestBindIdentityVerifiesPersistedMetadata(t *testing.T) { func TestBindIdentityVerifiesPersistedMetadata(t *testing.T) {
+24 -9
View File
@@ -8,6 +8,7 @@ import (
corev1 "k8s.io/api/core/v1" corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors" apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/runtime/schema"
@@ -76,16 +77,25 @@ func (c *Client) CreatePod(ctx context.Context, manifest PodManifest) (Pod, erro
Containers: []corev1.Container{{ Containers: []corev1.Container{{
Name: "executor", Image: manifest.Image, Args: manifest.Args, Env: environment, Name: "executor", Image: manifest.Image, Args: manifest.Args, Env: environment,
SecurityContext: &corev1.SecurityContext{Privileged: boolPointer(true)}, SecurityContext: &corev1.SecurityContext{Privileged: boolPointer(true)},
VolumeMounts: []corev1.VolumeMount{{ VolumeMounts: []corev1.VolumeMount{
Name: "spire-agent-socket", MountPath: "/run/spire/agent-sockets", ReadOnly: true, {Name: "spire-agent-socket", MountPath: "/run/spire/agent-sockets", ReadOnly: true},
}}, {Name: "docker-data", MountPath: "/var/lib/docker"},
}}, },
Volumes: []corev1.Volume{{
Name: "spire-agent-socket",
VolumeSource: corev1.VolumeSource{CSI: &corev1.CSIVolumeSource{
Driver: "csi.spiffe.io", ReadOnly: boolPointer(true),
}},
}}, }},
Volumes: []corev1.Volume{
{
Name: "spire-agent-socket",
VolumeSource: corev1.VolumeSource{CSI: &corev1.CSIVolumeSource{
Driver: "csi.spiffe.io", ReadOnly: boolPointer(true),
}},
},
{
Name: "docker-data",
VolumeSource: corev1.VolumeSource{EmptyDir: &corev1.EmptyDirVolumeSource{
SizeLimit: resourceQuantity("20Gi"),
}},
},
},
}, },
} }
created, err := c.Kubernetes.CoreV1().Pods(manifest.Namespace).Create(ctx, document, metav1.CreateOptions{}) created, err := c.Kubernetes.CoreV1().Pods(manifest.Namespace).Create(ctx, document, metav1.CreateOptions{})
@@ -169,6 +179,11 @@ func podFromKubernetes(pod corev1.Pod) Pod {
func boolPointer(value bool) *bool { return &value } func boolPointer(value bool) *bool { return &value }
func resourceQuantity(value string) *resource.Quantity {
quantity := resource.MustParse(value)
return &quantity
}
func stringMap(values map[string]string) map[string]any { func stringMap(values map[string]string) map[string]any {
result := make(map[string]any, len(values)) result := make(map[string]any, len(values))
for key, value := range values { for key, value := range values {
+6
View File
@@ -37,6 +37,12 @@ func TestClientPodLifecycleUsesTypedClient(t *testing.T) {
if got := pod.Spec.Containers[0].Env; len(got) != 1 || got[0].Name != "CI_RUNNER_CAPABILITY" || got[0].Value != "capability" { if got := pod.Spec.Containers[0].Env; len(got) != 1 || got[0].Name != "CI_RUNNER_CAPABILITY" || got[0].Value != "capability" {
t.Fatalf("environment = %#v", got) t.Fatalf("environment = %#v", got)
} }
if got := pod.Spec.Containers[0].VolumeMounts; len(got) != 2 || got[1].Name != "docker-data" || got[1].MountPath != "/var/lib/docker" {
t.Fatalf("volume mounts = %#v", got)
}
if got := pod.Spec.Volumes; len(got) != 2 || got[1].EmptyDir == nil || got[1].EmptyDir.SizeLimit == nil || got[1].EmptyDir.SizeLimit.String() != "20Gi" {
t.Fatalf("volumes = %#v", got)
}
if err := client.LabelPod(context.Background(), "gitea-actions", created.Name, map[string]string{terminalLabel: "true"}); err != nil { if err := client.LabelPod(context.Background(), "gitea-actions", created.Name, map[string]string{terminalLabel: "true"}); err != nil {
t.Fatal(err) t.Fatal(err)
} }
+1 -1
View File
@@ -102,7 +102,7 @@ func backendFromTask(task *runnerv1.Task) (Backend, error) {
return "", fmt.Errorf("task runs-on labels must include self-hosted: %v", labels) return "", fmt.Errorf("task runs-on labels must include self-hosted: %v", labels)
} }
hasPod := slices.Contains(labels, string(BackendPod)) hasPod := slices.Contains(labels, string(BackendPod))
hasVM := slices.Contains(labels, string(BackendVM)) hasVM := slices.Contains(labels, string(BackendVM)) || slices.Contains(labels, "vm-dev")
if hasPod == hasVM { if hasPod == hasVM {
return "", fmt.Errorf("task runs-on labels must select exactly one of pod or vm: %v", labels) return "", fmt.Errorf("task runs-on labels must select exactly one of pod or vm: %v", labels)
} }
@@ -28,6 +28,7 @@ func TestNewSelectsBackendFromRunsOn(t *testing.T) {
}{ }{
{"[self-hosted, pod]", BackendPod}, {"[self-hosted, pod]", BackendPod},
{"[self-hosted, vm]", BackendVM}, {"[self-hosted, vm]", BackendVM},
{"[self-hosted, vm-dev]", BackendVM},
} { } {
assignment, err := New(task(t, test.labels), "ddupan.top") assignment, err := New(task(t, test.labels), "ddupan.top")
if err != nil { if err != nil {
+29
View File
@@ -71,6 +71,28 @@ type Worker struct {
Backend Backend Backend Backend
Tasks TaskState Tasks TaskState
Bootstrap Bootstrap Bootstrap Bootstrap
OnEvent func(Event)
}
// Event describes a backend lifecycle transition without exposing launch
// environment values or other credentials.
type Event struct {
Name string
AssignmentID string
Executor string
Phase Phase
}
func (w Worker) event(name string, assignment taskassignment.Assignment, executor *Executor) {
if w.OnEvent == nil {
return
}
event := Event{Name: name, AssignmentID: assignment.ID}
if executor != nil {
event.Executor = executor.Name
event.Phase = executor.Phase
}
w.OnEvent(event)
} }
// Accept completes the durable handoff from JetStream to the backend. Once it // Accept completes the durable handoff from JetStream to the backend. Once it
@@ -88,6 +110,7 @@ func (w Worker) Accept(ctx context.Context, assignment taskassignment.Assignment
return false, err return false, err
} }
if executor == nil { if executor == nil {
w.event("executor_absent", assignment, nil)
launch, launchErr := w.launchSpec(assignment) launch, launchErr := w.launchSpec(assignment)
if launchErr != nil { if launchErr != nil {
return false, launchErr return false, launchErr
@@ -96,13 +119,18 @@ func (w Worker) Accept(ctx context.Context, assignment taskassignment.Assignment
if err != nil { if err != nil {
return false, err return false, err
} }
w.event("executor_created", assignment, executor)
} else {
w.event("executor_found", assignment, executor)
} }
if executor.IdentityTarget == "" { if executor.IdentityTarget == "" {
w.event("identity_target_pending", assignment, executor)
return false, nil return false, nil
} }
if err := w.Backend.BindIdentity(ctx, executor, assignment.Identity); err != nil { if err := w.Backend.BindIdentity(ctx, executor, assignment.Identity); err != nil {
return false, err return false, err
} }
w.event("identity_bound", assignment, executor)
return true, nil return true, nil
} }
@@ -178,6 +206,7 @@ func (w Worker) launchSpec(assignment taskassignment.Assignment) (LaunchSpec, er
func BackendMetadata(assignment taskassignment.Assignment) Metadata { func BackendMetadata(assignment taskassignment.Assignment) Metadata {
return Metadata{ return Metadata{
Labels: map[string]string{ Labels: map[string]string{
"ci.ddupan.top/runner": "true",
"ci.ddupan.top/assignment-id": assignment.ID, "ci.ddupan.top/assignment-id": assignment.ID,
"ci.ddupan.top/task-id": strconv.FormatInt(assignment.Task.GetId(), 10), "ci.ddupan.top/task-id": strconv.FormatInt(assignment.Task.GetId(), 10),
"ci.ddupan.top/backend": string(assignment.Backend), "ci.ddupan.top/backend": string(assignment.Backend),
+14 -1
View File
@@ -79,7 +79,8 @@ func TestHandleRecoversExistingExecutorWithoutCreatingAnother(t *testing.T) {
func TestAcceptAcknowledgesAfterBackendAndIdentityAreDurable(t *testing.T) { func TestAcceptAcknowledgesAfterBackendAndIdentityAreDurable(t *testing.T) {
backend := &fakeBackend{} backend := &fakeBackend{}
worker := Worker{Backend: backend, Bootstrap: fakeBootstrap{}} var events []Event
worker := Worker{Backend: backend, Bootstrap: fakeBootstrap{}, OnEvent: func(event Event) { events = append(events, event) }}
accepted, err := worker.Accept(context.Background(), assignment()) accepted, err := worker.Accept(context.Background(), assignment())
if err != nil || !accepted { if err != nil || !accepted {
@@ -88,6 +89,18 @@ func TestAcceptAcknowledgesAfterBackendAndIdentityAreDurable(t *testing.T) {
if backend.created != 1 || backend.bound != 1 || backend.deleted != 0 { if backend.created != 1 || backend.bound != 1 || backend.deleted != 0 {
t.Fatalf("created=%d bound=%d deleted=%d", backend.created, backend.bound, backend.deleted) t.Fatalf("created=%d bound=%d deleted=%d", backend.created, backend.bound, backend.deleted)
} }
want := []string{"executor_absent", "executor_created", "identity_bound"}
if len(events) != len(want) {
t.Fatalf("events = %#v", events)
}
for index := range want {
if events[index].Name != want[index] || events[index].AssignmentID != "gitea-task-42" {
t.Fatalf("event[%d] = %#v", index, events[index])
}
}
if events[1].Executor != "executor" || events[1].Phase != PhaseRunning {
t.Fatalf("created event = %#v", events[1])
}
} }
func TestAcceptRetriesWhileBackendIdentityTargetIsUnavailable(t *testing.T) { func TestAcceptRetriesWhileBackendIdentityTargetIsUnavailable(t *testing.T) {