Compare commits

...
Author SHA1 Message Date
panxiao81 452be27ea9 Merge 解耦 runner 与 workload placement
test / shell (push) Successful in 3m21s
publish controller image / publish-controller (push) Successful in 7m15s
test / python (push) Successful in 6m42s
test / go (push) Successful in 7m41s
2026-09-25 19:12:06 +00:00
panxiao81 53732699df fix: 解耦 runner 与 placement 实现
test / shell (pull_request) Successful in 3m31s
test / python (pull_request) Failing after 10m54s
test / go (pull_request) Failing after 12m12s
2026-09-25 19:03:13 +00:00
panxiao81 011c1e4c83 Merge 修复 placement consumer 名称导致的启动失败
test / shell (push) Successful in 43s
test / python (push) Successful in 1m31s
publish controller image / publish-controller (push) Successful in 5m8s
test / go (push) Successful in 5m36s
2026-09-25 18:46:30 +00:00
panxiao81 1f9577c816 fix: 使用合法的 JetStream consumer 名称
test / shell (pull_request) Successful in 52s
test / python (pull_request) Successful in 1m24s
test / go (pull_request) Successful in 3m57s
2026-09-25 18:41:47 +00:00
panxiao81 6ea78d4704 Merge pull request '重构 workload class 与 placement driver' (#54) from refactor/workload-placement into main
test / shell (push) Successful in 1m7s
test / python (push) Successful in 3m3s
publish controller image / publish-controller (push) Successful in 5m44s
test / go (push) Successful in 9m7s
Reviewed-on: #54
2026-09-25 17:11:01 +00:00
panxiao81 897d327a26 refactor: separate workload class from placement driver
test / shell (pull_request) Successful in 28s
test / python (pull_request) Successful in 59s
test / go (pull_request) Successful in 3m35s
2026-09-25 16:48:32 +00:00
panxiao81 162f742880 Merge 切换镜像发布任务到 Pod 后端
test / python (push) Successful in 2m57s
test / shell (push) Successful in 3m44s
publish controller image / publish-controller (push) Successful in 5m8s
test / go (push) Successful in 5m31s
2026-09-24 07:45:00 +00:00
panxiao81 0571e64b1a ci: 使用 Pod 后端发布镜像
test / shell (pull_request) Successful in 1m7s
test / python (pull_request) Successful in 2m5s
test / go (pull_request) Successful in 5m32s
2026-09-24 07:44:46 +00:00
panxiao81 7921840604 Merge 在 runner 中暴露内置 SPIRE CLI
test / shell (push) Successful in 55s
test / python (push) Successful in 1m46s
test / go (push) Successful in 4m37s
2026-09-24 07:08:41 +00:00
panxiao81 45be4a9fd9 fix(runner): 暴露内置 SPIRE CLI
test / shell (pull_request) Successful in 40s
test / python (pull_request) Successful in 1m26s
test / go (pull_request) Successful in 4m46s
2026-09-24 07:08:08 +00:00
panxiao81 02ccd5b8d7 Merge 修复 runner 工作目录权限
test / go (push) Failing after 2s
test / python (push) Failing after 5s
test / shell (push) Failing after 5s
2026-09-24 06:25:47 +00:00
panxiao81 2edb7b2b82 fix(runner): 创建可写的工作目录
test / python (pull_request) Failing after 4s
test / go (pull_request) Failing after 4s
test / shell (pull_request) Failing after 4s
2026-09-24 06:25:23 +00:00
panxiao81 d472a906ac Merge job Docker daemon detach 修复
test / shell (push) Successful in 26s
test / python (push) Successful in 1m5s
test / go (push) Successful in 5m24s
2026-09-21 18:08:45 +00:00
panxiao81 0151b6faf6 fix(runner): fully detach job Docker daemon
test / shell (pull_request) Successful in 29s
test / python (pull_request) Successful in 59s
test / go (pull_request) Successful in 5m1s
2026-09-21 18:08:05 +00:00
panxiao81 713d9a922a Merge runner job hooks 配置加载修复
test / shell (push) Successful in 33s
test / python (push) Successful in 1m14s
test / go (push) Successful in 3m15s
publish controller image / publish-controller (push) Failing after 11m1s
2026-09-21 17:54:27 +00:00
panxiao81 537c620051 fix(runner): load job hooks configuration
test / python (pull_request) Successful in 51s
test / shell (pull_request) Successful in 40s
test / go (pull_request) Successful in 3m1s
2026-09-21 17:50:47 +00:00
panxiao81 a7b62868b6 Merge 镜像发布 SPIFFE 身份兼容修复
test / python (push) Successful in 1m11s
test / go (push) Successful in 3m20s
test / shell (push) Successful in 25s
publish controller image / publish-controller (push) Failing after 10m57s
2026-09-21 17:38:01 +00:00
panxiao81 84aa607c86 fix(ci): preserve published image job identity
test / python (pull_request) Successful in 1m27s
test / shell (pull_request) Successful in 34s
test / go (pull_request) Successful in 2m59s
2026-09-21 17:36:53 +00:00
panxiao81 57524be30f Merge controller 与 runner 镜像发布拆分
test / go (push) Successful in 2m55s
test / shell (push) Successful in 30s
test / python (push) Successful in 57s
publish controller image / publish-controller (push) Failing after 14m2s
2026-09-21 17:19:44 +00:00
panxiao81 2ffc45c766 ci: publish runner image only on manual dispatch
test / go (pull_request) Successful in 3m3s
test / shell (pull_request) Successful in 36s
test / python (pull_request) Successful in 1m5s
2026-09-21 17:15:33 +00:00
panxiao81 68c3771dc8 Merge Docker cgroup 初始化修复
test / shell (push) Successful in 28s
test / python (push) Successful in 57s
test / go (push) Successful in 3m4s
publish images / publish-images (push) Failing after 35m46s
2026-09-21 16:46:04 +00:00
panxiao81 16054d78e3 fix(runner): correct cgroup controller expression
test / shell (pull_request) Successful in 32s
test / python (pull_request) Successful in 1m6s
test / go (pull_request) Successful in 3m10s
2026-09-21 16:42:47 +00:00
panxiao81 5d3d2a94bd Merge runner 默认 Docker 与镜像发布工具链
test / shell (push) Successful in 29s
test / python (push) Successful in 1m0s
test / go (push) Successful in 3m6s
publish images / publish-images (push) Failing after 15m15s
2026-09-21 16:24:18 +00:00
panxiao81 0bf39b4751 fix(runner): redirect dockerd logs as root
test / shell (pull_request) Successful in 37s
test / python (pull_request) Successful in 1m13s
test / go (pull_request) Successful in 3m9s
2026-09-21 16:20:02 +00:00
panxiao81 d776fa71e9 feat(runner): provide Docker before workflow steps
test / python (pull_request) Successful in 12m25s
test / shell (pull_request) Failing after 13m32s
test / go (pull_request) Successful in 22m28s
2026-09-21 15:53:03 +00:00
panxiao81 adb5af1486 fix: 为镜像发布准备构建工具链
test / shell (pull_request) Successful in 30s
test / python (pull_request) Successful in 1m2s
test / go (pull_request) Successful in 3m23s
2026-09-21 15:43:26 +00:00
panxiao81 94fc84a47c Merge Pod Docker 数据卷修复
test / shell (push) Successful in 27s
test / python (push) Successful in 1m3s
publish images / publish-images (push) Failing after 1m36s
test / go (push) Successful in 3m0s
2026-09-21 15:39:13 +00:00
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
35 changed files with 673 additions and 343 deletions
+23
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,24 @@ 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: Build and run image
shell: bash
run: |
set -euo pipefail
findmnt /var/lib/docker
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
+9 -101
View File
@@ -1,128 +1,36 @@
--- ---
name: publish images name: publish controller image
on: on:
push: push:
branches: [main] branches: [main]
paths: paths:
- '.gitea/workflows/publish-images.yml' - '.gitea/workflows/publish-images.yml'
- 'config/**' - 'container/controller.Dockerfile'
- 'container/**'
- 'cmd/**' - 'cmd/**'
- 'internal/**' - 'internal/**'
- 'scripts/**' - 'scripts/publish-image'
- 'src/**'
- 'go.mod' - 'go.mod'
- 'go.sum' - 'go.sum'
- 'pyproject.toml'
- 'README.md'
workflow_dispatch: workflow_dispatch:
jobs: jobs:
publish-images: publish-images:
name: publish-images name: publish-controller
runs-on: [self-hosted, vm] runs-on: [self-hosted, pod]
timeout-minutes: 45 timeout-minutes: 45
permissions: permissions:
contents: read contents: read
env: env:
PUSH_REGISTRY: zot-push.ad.ddupan.top PUSH_REGISTRY: zot-push.ad.ddupan.top
PULL_REGISTRY: zot.ad.ddupan.top PULL_REGISTRY: zot.ad.ddupan.top
CONTROLLER_REPOSITORY: panxiao81/gitea-dynamic-runner-controller IMAGE_NAME: controller
RUNNER_REPOSITORY: panxiao81/gitea-dynamic-runner-runner IMAGE_REPOSITORY: panxiao81/gitea-dynamic-runner-controller
IMAGE_DOCKERFILE: container/controller.Dockerfile
SPIRE_AGENT_SOCKET: /run/spire/agent-sockets/spire-agent.sock SPIRE_AGENT_SOCKET: /run/spire/agent-sockets/spire-agent.sock
steps: steps:
- uses: actions/checkout@v4 - uses: actions/checkout@v4
- name: Test source
shell: bash
run: |
set -euo pipefail
go test ./...
go vet ./...
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 - name: Build and publish
shell: bash shell: bash
run: | run: scripts/publish-image
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
+26
View File
@@ -0,0 +1,26 @@
---
name: publish runner image
on:
workflow_dispatch:
jobs:
publish-images:
name: publish-runner
runs-on: [self-hosted, pod]
timeout-minutes: 45
permissions:
contents: read
env:
PUSH_REGISTRY: zot-push.ad.ddupan.top
PULL_REGISTRY: zot.ad.ddupan.top
IMAGE_NAME: runner
IMAGE_REPOSITORY: panxiao81/gitea-dynamic-runner-runner
IMAGE_DOCKERFILE: container/runner.Dockerfile
SPIRE_AGENT_SOCKET: /run/spire/agent-sockets/spire-agent.sock
steps:
- uses: actions/checkout@v4
- name: Build and publish
shell: bash
run: scripts/publish-image
+5 -14
View File
@@ -5,29 +5,20 @@ on:
jobs: jobs:
jwt-svid: jwt-svid:
runs-on: self-hosted runs-on: [self-hosted, pod]
container:
volumes:
- /run/spire/agent-sockets:/run/spire/agent-sockets:ro
steps: steps:
- name: Fetch pinned SPIRE CLI - name: Verify bundled SPIRE CLI
shell: bash shell: bash
run: | run: |
set -euo pipefail set -euo pipefail
archive=/tmp/spire.tar.gz command -v spire-agent
curl --fail --location --silent --show-error \ spire-agent -version
--output "$archive" \
https://github.com/spiffe/spire/releases/download/v1.15.3/spire-1.15.3-linux-amd64-musl.tar.gz
printf '%s %s\n' \
ca1a4d1155317bdd2afc7f36663828a10410c7c840e54725b90b4064b0a301c7 \
"$archive" | sha256sum --check --status
tar -xzf "$archive" -C /tmp spire-1.15.3/bin/spire-agent
- name: Fetch short-lived zot JWT-SVID - name: Fetch short-lived zot JWT-SVID
shell: bash shell: bash
run: | run: |
set -euo pipefail set -euo pipefail
/tmp/spire-1.15.3/bin/spire-agent api fetch jwt \ spire-agent api fetch jwt \
-audience zot \ -audience zot \
-socketPath /run/spire/agent-sockets/spire-agent.sock \ -socketPath /run/spire/agent-sockets/spire-agent.sock \
>/dev/null >/dev/null
+107 -7
View File
@@ -1,31 +1,131 @@
---
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-dev] runs-on: [self-hosted, vm]
steps: steps:
- name: Prepare nested kubelet device
shell: bash
run: |
set -euo pipefail
if [[ ! -e /dev/kmsg ]]; then
sudo mknod /dev/kmsg c 1 11
fi
- name: Verify Docker - name: Verify Docker
run: docker info shell: bash
run: |
set -euo pipefail
findmnt /var/lib/docker
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
+1 -1
View File
@@ -7,7 +7,7 @@ on:
jobs: jobs:
runtime: runtime:
runs-on: [self-hosted, vm-dev] runs-on: [self-hosted, vm]
steps: steps:
- name: Verify workload identity socket - name: Verify workload identity socket
shell: bash shell: bash
+17 -13
View File
@@ -1,17 +1,19 @@
# Gitea dynamic runner # Gitea dynamic runner
为 Gitea Actions 按需创建一次性执行环境。对 workflow 提供两种稳定的 runner 为 Gitea Actions 按需创建一次性执行环境。workflow 分别声明 workload class 与
接口: placement driver:
```yaml ```yaml
runs-on: [self-hosted, pod] runs-on: [self-hosted, container, kubernetes]
``` ```
```yaml ```yaml
runs-on: [self-hosted, vm] runs-on: [self-hosted, vm, opensandbox]
``` ```
`pod` 使用动态 Kubernetes Pod,`vm` 使用动态 Cloud Hypervisor microVM。每个环境 兼容标签 `[self-hosted, pod]` 严格映射为 `container+kubernetes`,
`[self-hosted, vm]` 和 `vm-dev` 严格映射为 `vm+opensandbox`;显式 driver 容量耗尽时
不会回退到其他 driver。每个环境
只执行一个 job,并在 job 结束后连同本地状态一起销毁。完整的设计约束见 只执行一个 job,并在 job 结束后连同本地状态一起销毁。完整的设计约束见
[`docs/design-principles.md`](docs/design-principles.md)。 [`docs/design-principles.md`](docs/design-principles.md)。
@@ -24,8 +26,8 @@ runs-on: [self-hosted, vm]
- `scheduler`:以常驻 Gitea RunnerService 身份直接领取 task,并把完整 assignment - `scheduler`:以常驻 Gitea RunnerService 身份直接领取 task,并把完整 assignment
持久化到 JetStream;一个 registration 下按配置启动多个并发 `FetchTask` goroutine, 持久化到 JetStream;一个 registration 下按配置启动多个并发 `FetchTask` goroutine,
同时提供仅允许 SPIFFE mTLS 的 RunnerService facade。 同时提供仅允许 SPIFFE mTLS 的 RunnerService facade。
- `pod-worker`:直接在 homelab Kubernetes 创建一次性 Pod。 - `kubernetes-worker`:直接在 homelab Kubernetes 创建一次性 container workload。
- `vm-worker`:通过 OpenSandbox Lifecycle API 从 `ci-vm` Pool 创建 Kata microVM。 - `opensandbox-worker`:通过 OpenSandbox Lifecycle API 从 `ci-vm` Pool 创建 VM workload。
- 三个组件默认在同一个 Go 进程启用。首轮集成期间不允许只启动 worker,因为 facade - 三个组件默认在同一个 Go 进程启用。首轮集成期间不允许只启动 worker,因为 facade
的 assignment claim registry 仍是进程内状态;支持安全拆分前进程会明确拒绝该配置。 的 assignment claim registry 仍是进程内状态;支持安全拆分前进程会明确拒绝该配置。
- `microvm-runner-launch`:为每个任务以 direct I/O 转换出 flat qcow2 root disk、创建 NoCloud seed 和 TAP,运行 - `microvm-runner-launch`:为每个任务以 direct I/O 转换出 flat qcow2 root disk、创建 NoCloud seed 和 TAP,运行
@@ -36,13 +38,14 @@ runs-on: [self-hosted, vm]
entry;不持有 OpenSandbox API key、Gitea token 或 Bao 凭据。身份与 Pool 契约见 entry;不持有 OpenSandbox API key、Gitea token 或 Bao 凭据。身份与 Pool 契约见
[`docs/opensandbox-runner.md`](docs/opensandbox-runner.md)。 [`docs/opensandbox-runner.md`](docs/opensandbox-runner.md)。
- Pod executor:在 Kubernetes 中创建一次性 privileged Pod;Pod 内的 workflow 使用 - Pod executor:在 Kubernetes 中创建一次性 privileged Pod;Pod 内的 workflow 使用
host executor,Docker、BuildKit 和 kind 等工具由 pipeline 按需 setup。Runner 固定在 host executor。Runner 固定在支持原生 job hooks 的 3.x 版本,在 workflow 第一步前
支持原生 job hooks 的 3.x 版本,在 workflow 第一步前等待实际任务对应的 SVID。 等待实际任务对应的 SVID,并启动 job-local Docker daemon;workflow 可直接使用与
GitHub-hosted runner 相同的 Docker/BuildKit action。
- `jwt-broker`:早期共享 Kubernetes runner 的过渡实验;目标架构不部署它,每个 - `jwt-broker`:早期共享 Kubernetes runner 的过渡实验;目标架构不部署它,每个
动态 Pod 或 VM 直接取得自己的 SPIFFE 身份。 动态 Pod 或 VM 直接取得自己的 SPIFFE 身份。
Pod 路径由 homelab 集群中的 `pod-worker` 直接创建 Kubernetes Pod。OpenSandbox 只用于 Container/Kubernetes 路径由 homelab 集群中的 `kubernetes-worker` 直接创建 Pod。
VM/Kata workload;两个 backend 使用独立 durable consumer 和独立容量池。assignment 根据 OpenSandbox 只用于 VM/Kata workload;每个 placement 使用独立 durable consumer 和容量池。assignment 根据
`runs-on` 进入对应池,池满时留在 JetStream pending,不会创建超出容量的 workload;任一 `runs-on` 进入对应池,池满时留在 JetStream pending,不会创建超出容量的 workload;任一
执行层故障不会阻塞另一条部署。长期 RunnerService 协议路线见 执行层故障不会阻塞另一条部署。长期 RunnerService 协议路线见
[`docs/runner-protocol-roadmap.md`](docs/runner-protocol-roadmap.md)。 [`docs/runner-protocol-roadmap.md`](docs/runner-protocol-roadmap.md)。
@@ -69,12 +72,13 @@ credential 都从挂载文件读取,不接受明文环境变量:
- `NATS_PRODUCER_PASSWORD_FILE`、`NATS_WORKER_PASSWORD_FILE`:分别使用现有最小权限的 - `NATS_PRODUCER_PASSWORD_FILE`、`NATS_WORKER_PASSWORD_FILE`:分别使用现有最小权限的
`ci-producer` publish 连接和 `ci-worker` pull/ACK 连接,controller 不合并权限。 `ci-producer` publish 连接和 `ci-worker` pull/ACK 连接,controller 不合并权限。
- `RUNNER_FACADE_CAPABILITY_KEY_FILE`:至少 32 字节的 controller HMAC key。 - `RUNNER_FACADE_CAPABILITY_KEY_FILE`:至少 32 字节的 controller HMAC key。
- `OPENSANDBOX_API_KEY_FILE`:仅启用 `vm-worker` 时读取。 - `OPENSANDBOX_API_KEY_FILE`:仅启用 `opensandbox-worker` 时读取。
必要的非 secret 配置包括 `POD_EXECUTOR_IMAGE`(应使用 digest)、`SPIRE_AGENT_ID`、 必要的非 secret 配置包括 `POD_EXECUTOR_IMAGE`(应使用 digest)、`SPIRE_AGENT_ID`、
`RUNNER_FACADE_URL`、`RUNNER_FACADE_SPIFFE_ID` 和 `SPIFFE_ENDPOINT_SOCKET`。默认 `RUNNER_FACADE_URL`、`RUNNER_FACADE_SPIFFE_ID` 和 `SPIFFE_ENDPOINT_SOCKET`。默认
`COMPONENTS=all`、Pod 并发 4、VM 并发 1;首次 smoke test 应显式设为 `COMPONENTS=all`、Pod 并发 4、VM 并发 1;首次 smoke test 应显式设为
`COMPONENTS=scheduler,pod-worker`,先验证 Pod 链路,避免同时消耗 VM 容量。 `COMPONENTS=scheduler,kubernetes-worker`,先验证 Kubernetes 链路,避免同时消耗 VM
容量。旧组件名 `pod-worker`、`vm-worker` 只作为配置兼容别名保留。
Pod task 的 terminal update 被 Gitea 接受后,controller 会在 Pod 上持久写入 Pod task 的 terminal update 被 Gitea 接受后,controller 会在 Pod 上持久写入
`ci.ddupan.top/terminal=true` label。生命周期 reconciler 只清理同时带该 label 且已经 `ci.ddupan.top/terminal=true` label。生命周期 reconciler 只清理同时带该 label 且已经
+48 -48
View File
@@ -38,6 +38,15 @@ type runComponent func(context.Context) error
func (function runComponent) Run(ctx context.Context) error { return function(ctx) } func (function runComponent) Run(ctx context.Context) error { return function(ctx) }
type terminalBackend interface {
MarkTerminal(context.Context, string) error
}
type placementRuntime struct {
backend terminalBackend
pool *backendpool.Pool
}
type controllerConfig struct { type controllerConfig struct {
Components controller.Selection Components controller.Selection
TrustDomain, WorkloadAPIAddr string TrustDomain, WorkloadAPIAddr string
@@ -61,7 +70,7 @@ func runController(ctx context.Context) error {
if !slices.Contains(config.Components, controller.Scheduler) { if !slices.Contains(config.Components, controller.Scheduler) {
return errors.New("split worker deployment is not yet safe: scheduler/facade must be enabled with workers") return errors.New("split worker deployment is not yet safe: scheduler/facade must be enabled with workers")
} }
if !slices.Contains(config.Components, controller.PodWorker) && !slices.Contains(config.Components, controller.VMWorker) { if !slices.Contains(config.Components, controller.KubernetesWorker) && !slices.Contains(config.Components, controller.OpenSandboxWorker) {
return errors.New("scheduler requires at least one local backend worker") return errors.New("scheduler requires at least one local backend worker")
} }
@@ -92,32 +101,19 @@ func runController(ctx context.Context) error {
return err return err
} }
registry := runnerfacade.NewRegistry() registry := runnerfacade.NewRegistry()
podPool := backendpool.New(config.PodCapacity) runtimes := make(map[taskassignment.Placement]placementRuntime)
vmPool := backendpool.New(config.VMCapacity)
var podExecutorBackend *podbackend.Backend
var vmExecutorBackend *opensandboxbackend.Backend
giteaClient := giteaactions.NewClient(giteaactions.DefaultHTTPClient(), config.GiteaURL, config.GiteaUUID, config.GiteaToken) giteaClient := giteaactions.NewClient(giteaactions.DefaultHTTPClient(), config.GiteaURL, config.GiteaUUID, config.GiteaToken)
facade := &runnerfacade.Facade{ facade := &runnerfacade.Facade{
Registry: registry, Capabilities: capabilities, Upstream: giteaClient, Registry: registry, Capabilities: capabilities, Upstream: giteaClient,
OnTerminal: func(ctx context.Context, assignment taskassignment.Assignment) error { OnTerminal: func(ctx context.Context, assignment taskassignment.Assignment) error {
switch assignment.Backend { runtime, ok := runtimes[assignment.Placement]
case taskassignment.BackendPod: if !ok {
if podExecutorBackend == nil { return fmt.Errorf("placement runtime %s is not configured", assignment.Placement.Key())
return errors.New("Pod lifecycle backend is not configured")
} }
if err := podExecutorBackend.MarkTerminal(ctx, assignment.ID); err != nil { if err := runtime.backend.MarkTerminal(ctx, assignment.ID); err != nil {
return err return err
} }
podPool.Release(assignment.ID) runtime.pool.Release(assignment.ID)
case taskassignment.BackendVM:
if vmExecutorBackend == nil {
return errors.New("VM lifecycle backend is not configured")
}
if err := vmExecutorBackend.MarkTerminal(ctx, assignment.ID); err != nil {
return err
}
vmPool.Release(assignment.ID)
}
return nil return nil
}, },
} }
@@ -127,11 +123,11 @@ func runController(ctx context.Context) error {
} }
labels := []string{"self-hosted"} labels := []string{"self-hosted"}
if slices.Contains(config.Components, controller.PodWorker) { if slices.Contains(config.Components, controller.KubernetesWorker) {
labels = append(labels, string(taskassignment.BackendPod)) labels = append(labels, "pod", string(taskassignment.WorkloadContainer), string(taskassignment.DriverKubernetes))
} }
if slices.Contains(config.Components, controller.VMWorker) { if slices.Contains(config.Components, controller.OpenSandboxWorker) {
labels = append(labels, config.VMRunnerLabel) labels = append(labels, config.VMRunnerLabel, string(taskassignment.DriverOpenSandbox))
} }
poller := taskscheduler.Poller{ poller := taskscheduler.Poller{
Client: giteaClient, Client: giteaClient,
@@ -168,7 +164,9 @@ func runController(ctx context.Context) error {
}), }),
} }
if slices.Contains(config.Components, controller.PodWorker) { if slices.Contains(config.Components, controller.KubernetesWorker) {
placement := taskassignment.KubernetesContainer
pool := backendpool.New(config.PodCapacity)
client, err := podbackend.NewInClusterClient() client, err := podbackend.NewInClusterClient()
if err != nil { if err != nil {
return err return err
@@ -179,57 +177,59 @@ func runController(ctx context.Context) error {
SPIRECluster: config.SPIRECluster, SPIREClass: config.SPIREClass, SPIRECluster: config.SPIRECluster, SPIREClass: config.SPIREClass,
SPIREAgentID: config.SPIREAgentID, ExecutorUID: config.PodExecutorUID, SPIREAgentID: config.SPIREAgentID, ExecutorUID: config.PodExecutorUID,
}} }}
podExecutorBackend = &backend runtimes[placement] = placementRuntime{backend: backend, pool: pool}
assignments, err := backend.RecoverAssignments(ctx) assignments, err := backend.RecoverAssignments(ctx)
if err != nil { if err != nil {
return err return err
} }
for _, assignment := range assignments { for _, assignment := range assignments {
podPool.Restore(assignment.ID) pool.Restore(assignment.ID)
if err := registry.RecoverClaimed(assignment); err != nil { if err := registry.RecoverClaimed(assignment); err != nil {
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, OnEvent: workerEventLogger(taskassignment.BackendPod)}, registry, podPool) component, err := workerComponent(ctx, workerJS, config, placement, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap, OnEvent: workerEventLogger(placement)}, registry, pool)
if err != nil { if err != nil {
return err return err
} }
lifecycle := podbackend.Lifecycle{Backend: backend, OnError: func(err error) { lifecycle := podbackend.Lifecycle{Backend: backend, OnError: func(err error) {
slog.Error("backend lifecycle error", "component", "lifecycle", "backend", taskassignment.BackendPod, "error", err) slog.Error("backend lifecycle error", "component", "lifecycle", "placement", placement.Key(), "error", err)
}} }}
components[controller.PodWorker] = runComponent(func(ctx context.Context) error { components[controller.KubernetesWorker] = 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) })
group.Go(func() error { return lifecycle.Run(groupContext) }) group.Go(func() error { return lifecycle.Run(groupContext) })
return group.Wait() return group.Wait()
}) })
} }
if slices.Contains(config.Components, controller.VMWorker) { if slices.Contains(config.Components, controller.OpenSandboxWorker) {
placement := taskassignment.OpenSandboxVM
pool := backendpool.New(config.VMCapacity)
lifecycle := opensandboxbackend.NewLifecycleClient(config.OpenSandboxURL, config.OpenSandboxAPIKey, &http.Client{Timeout: 60 * time.Second}) lifecycle := opensandboxbackend.NewLifecycleClient(config.OpenSandboxURL, config.OpenSandboxAPIKey, &http.Client{Timeout: 60 * time.Second})
backend := opensandboxbackend.Backend{Lifecycle: lifecycle, Config: opensandboxbackend.Config{ backend := opensandboxbackend.Backend{Lifecycle: lifecycle, Config: opensandboxbackend.Config{
Pool: config.OpenSandboxPool, Timeout: config.VMTimeout, Pool: config.OpenSandboxPool, Timeout: config.VMTimeout,
Entrypoint: []string{"/usr/local/bin/gitea-dynamic-runner", "executor"}, Entrypoint: []string{"/usr/local/bin/gitea-dynamic-runner", "executor"},
Env: map[string]string{"SPIFFE_ENDPOINT_SOCKET": config.WorkloadAPIAddr}, Env: map[string]string{"SPIFFE_ENDPOINT_SOCKET": config.WorkloadAPIAddr},
}} }}
vmExecutorBackend = &backend runtimes[placement] = placementRuntime{backend: backend, pool: pool}
assignments, err := backend.RecoverAssignments(ctx, config.TrustDomain) assignments, err := backend.RecoverAssignments(ctx, config.TrustDomain)
if err != nil { if err != nil {
return err return err
} }
for _, assignment := range assignments { for _, assignment := range assignments {
vmPool.Restore(assignment.ID) pool.Restore(assignment.ID)
if err := registry.RecoverClaimed(assignment); err != nil { if err := registry.RecoverClaimed(assignment); err != nil {
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, OnEvent: workerEventLogger(taskassignment.BackendVM)}, registry, vmPool) component, err := workerComponent(ctx, workerJS, config, placement, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap, OnEvent: workerEventLogger(placement)}, registry, pool)
if err != nil { if err != nil {
return err return err
} }
lifecycleReconciler := opensandboxbackend.LifecycleReconciler{Backend: backend, OnError: func(err error) { lifecycleReconciler := opensandboxbackend.LifecycleReconciler{Backend: backend, OnError: func(err error) {
slog.Error("backend lifecycle error", "component", "lifecycle", "backend", taskassignment.BackendVM, "error", err) slog.Error("backend lifecycle error", "component", "lifecycle", "placement", placement.Key(), "error", err)
}} }}
components[controller.VMWorker] = runComponent(func(ctx context.Context) error { components[controller.OpenSandboxWorker] = 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) })
group.Go(func() error { return lifecycleReconciler.Run(groupContext) }) group.Go(func() error { return lifecycleReconciler.Run(groupContext) })
@@ -291,25 +291,25 @@ func connectNATS(server, user, password, caFile, clientName string) (*nats.Conn,
return nats.Connect(server, options...) return nats.Connect(server, options...)
} }
func workerComponent(ctx context.Context, js jetstream.JetStream, config controllerConfig, backend taskassignment.Backend, capacity int, accepter assignmentqueue.Accepter, claims assignmentqueue.Claims, admission assignmentqueue.Admission) (controller.Component, error) { func workerComponent(ctx context.Context, js jetstream.JetStream, config controllerConfig, placement taskassignment.Placement, capacity int, accepter assignmentqueue.Accepter, claims assignmentqueue.Claims, admission assignmentqueue.Admission) (controller.Component, error) {
consumer, err := assignmentqueue.OpenConsumer(ctx, js, config.Stream, config.SubjectBase, backend, capacity) consumer, err := assignmentqueue.OpenConsumer(ctx, js, config.Stream, config.SubjectBase, placement, capacity)
if err != nil { if err != nil {
return nil, err return nil, err
} }
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, OnEvent: func(event assignmentqueue.Event) { Processor: assignmentqueue.Processor{TrustDomain: config.TrustDomain, Accepter: accepter, Claims: claims, Admission: admission, OnEvent: func(event assignmentqueue.Event) {
slog.Info("assignment transition", "component", "worker", "event", event.Name, "backend", event.Backend, "assignment", event.AssignmentID, "retry_delay", event.RetryDelay) slog.Info("assignment transition", "component", "worker", "event", event.Name, "placement", event.Placement.Key(), "assignment", event.AssignmentID, "retry_delay", event.RetryDelay)
}}, }},
OnError: func(err error) { OnError: func(err error) {
slog.Error("assignment processing error", "component", "worker", "backend", backend, "error", err) slog.Error("assignment processing error", "component", "worker", "placement", placement.Key(), "error", err)
}, },
}, nil }, nil
} }
func workerEventLogger(backend taskassignment.Backend) func(taskworker.Event) { func workerEventLogger(placement taskassignment.Placement) func(taskworker.Event) {
return func(event 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) slog.Info("executor transition", "component", "worker", "event", event.Name, "placement", placement.Key(), "assignment", event.AssignmentID, "executor", event.Executor, "phase", event.Phase)
} }
} }
@@ -364,18 +364,18 @@ func loadControllerConfig() (controllerConfig, error) {
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")
} }
if slices.Contains(selection, controller.PodWorker) && config.PodImage == "" { if slices.Contains(selection, controller.KubernetesWorker) && config.PodImage == "" {
return controllerConfig{}, errors.New("POD_EXECUTOR_IMAGE is required for pod-worker") return controllerConfig{}, errors.New("POD_EXECUTOR_IMAGE is required for kubernetes-worker")
} }
if slices.Contains(selection, controller.PodWorker) && config.SPIREAgentID == "" { if slices.Contains(selection, controller.KubernetesWorker) && config.SPIREAgentID == "" {
return controllerConfig{}, errors.New("SPIRE_AGENT_ID is required for pod-worker") return controllerConfig{}, errors.New("SPIRE_AGENT_ID is required for kubernetes-worker")
} }
if slices.Contains(selection, controller.VMWorker) { if slices.Contains(selection, controller.OpenSandboxWorker) {
if config.VMRunnerLabel != "vm" && config.VMRunnerLabel != "vm-dev" { if config.VMRunnerLabel != "vm" && config.VMRunnerLabel != "vm-dev" {
return controllerConfig{}, errors.New("VM_RUNNER_LABEL must be vm or 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 opensandbox-worker")
} }
config.OpenSandboxAPIKey, err = read("OPENSANDBOX_API_KEY_FILE") config.OpenSandboxAPIKey, err = read("OPENSANDBOX_API_KEY_FILE")
if err != nil { if err != nil {
+1 -1
View File
@@ -34,7 +34,7 @@ func TestLoadControllerConfigUsesFileSecrets(t *testing.T) {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if len(config.Components) != 2 || config.Components[0] != controller.Scheduler || config.Components[1] != controller.PodWorker { if len(config.Components) != 2 || config.Components[0] != controller.Scheduler || config.Components[1] != controller.KubernetesWorker {
t.Fatalf("components = %#v", config.Components) t.Fatalf("components = %#v", config.Components)
} }
if config.GiteaUUID != "scheduler-uuid" || config.GiteaToken != "scheduler-token" || config.NATSProducerPassword != "producer-password" || config.NATSWorkerPassword != "worker-password" { if config.GiteaUUID != "scheduler-uuid" || config.GiteaToken != "scheduler-token" || config.NATSProducerPassword != "producer-password" || config.NATSWorkerPassword != "worker-password" {
+3 -1
View File
@@ -21,14 +21,16 @@ RUN groupadd --gid 2000 runner \
&& useradd --uid 2000 --gid 2000 --groups docker --create-home --shell /bin/bash runner \ && useradd --uid 2000 --gid 2000 --groups docker --create-home --shell /bin/bash runner \
&& printf 'runner ALL=(ALL) NOPASSWD:ALL\n' >/etc/sudoers.d/runner \ && printf 'runner ALL=(ALL) NOPASSWD:ALL\n' >/etc/sudoers.d/runner \
&& chmod 0440 /etc/sudoers.d/runner \ && chmod 0440 /etc/sudoers.d/runner \
&& install -d -o 2000 -g 2000 /data && install -d -o 2000 -g 2000 /data /workspace
COPY --from=runner /usr/local/bin/gitea-runner /usr/local/bin/gitea-runner COPY --from=runner /usr/local/bin/gitea-runner /usr/local/bin/gitea-runner
COPY --from=controller /out/gitea-dynamic-runner /usr/local/bin/gitea-dynamic-runner COPY --from=controller /out/gitea-dynamic-runner /usr/local/bin/gitea-dynamic-runner
COPY --from=spire /opt/spire/bin/spire-agent /opt/spire/bin/spire-agent COPY --from=spire /opt/spire/bin/spire-agent /opt/spire/bin/spire-agent
RUN ln -s /opt/spire/bin/spire-agent /usr/local/bin/spire-agent
COPY config/runner.yaml /etc/gitea-runner/config.yaml COPY config/runner.yaml /etc/gitea-runner/config.yaml
COPY --chmod=0755 scripts/gitea-job-started /usr/local/libexec/gitea-job-started COPY --chmod=0755 scripts/gitea-job-started /usr/local/libexec/gitea-job-started
COPY --chmod=0755 scripts/gitea-opensandbox-runner /usr/local/libexec/gitea-opensandbox-runner COPY --chmod=0755 scripts/gitea-opensandbox-runner /usr/local/libexec/gitea-opensandbox-runner
COPY --chmod=0755 scripts/setup-job-docker /usr/local/libexec/setup-job-docker
VOLUME ["/data"] VOLUME ["/data"]
ENV HOME=/home/runner ENV HOME=/home/runner
+8 -5
View File
@@ -2,19 +2,22 @@
## 对 workflow 的接口 ## 对 workflow 的接口
Runner 只向 workflow 暴露两个执行环境: Runner 将执行环境的能力类型与实现 driver 分开声明:
```yaml ```yaml
runs-on: [self-hosted, pod] runs-on: [self-hosted, container, kubernetes]
``` ```
```yaml ```yaml
runs-on: [self-hosted, vm] runs-on: [self-hosted, vm, opensandbox]
``` ```
- `self-hosted` 是固定前缀。 - `self-hosted` 是固定前缀。
- `pod` 表示一次性 Kubernetes Pod,承担常规 CI、镜像构建和 kind 等任务。 - `container`、`vm` 是 workload class;`kubernetes`、`opensandbox` 是 placement driver。
- `vm` 表示一次性 microVM,承担需要独立内核、KVM、systemd 或更强隔离的任务。 - 当前支持 `container+kubernetes` 和 `vm+opensandbox`。旧 `pod` 与 `vm` 标签分别是
两个组合的严格兼容别名,显式指定 driver 后不得因容量或故障回退到另一 driver。
- container 承担常规 CI、镜像构建和 kind 等任务;VM 承担需要独立内核、KVM、
systemd 或更强隔离的任务。
执行后端是基础设施选择,不是权限角色。workflow 不需要额外声明由 controller 执行后端是基础设施选择,不是权限角色。workflow 不需要额外声明由 controller
维护的 role 或权限 label。 维护的 role 或权限 label。
+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 独立运行;job-started hook 在第一步 workflow 之前
启动 job-local Docker daemon,业务 workflow 不负责 runner 基础设施初始化;
- Kata VM 中 Docker 数据使用 guest 内的 loop-backed ext4,并随 sandbox 一起删除;
- `self-hosted` 必须是所有 runner labels 的前缀; - `self-hosted` 必须是所有 runner labels 的前缀;
- ephemeral/once runner 完成一项任务后退出。 - ephemeral/once runner 完成一项任务后退出。
+16 -13
View File
@@ -3,9 +3,9 @@
## 目标 ## 目标
长期形态不依赖 `workflow_job` webhook 发现工作。controller 本身作为 Gitea Runner 长期形态不依赖 `workflow_job` webhook 发现工作。controller 本身作为 Gitea Runner
协议客户端注册,并声明 `self-hosted`、`pod` 和 `vm` labels;单个 registration 内按总 协议客户端注册,并声明 workload class 与 driver labels;单个 registration 内按总
配置容量启动多个 `FetchTask` goroutine,再将 task 按 `runs-on` 交给 Pod 或 VM 的独立 配置容量启动多个 `FetchTask` goroutine,再将 task 按 `runs-on` 交给 placement 对应的
容量池,由一次性 Pod 或 microVM 执行。 独立容量池,由一次性 Pod 或 microVM 执行。
```text ```text
Gitea RunnerService Gitea RunnerService
@@ -13,14 +13,14 @@ Gitea RunnerService
▼ ▼
dynamic-runner scheduler dynamic-runner scheduler
│ 已领取的 task + lease │ 已领取的 task + lease
├── Pod executor ├── container.kubernetes executor
└── microVM executor └── vm.opensandbox executor
│ logs / state / result │ logs / state / result
└──────────────────────► Gitea └──────────────────────► Gitea
``` ```
controller 使用单一 Go 二进制;默认在同一进程启用 `scheduler`、`pod-worker` 和 controller 使用单一 Go 二进制;默认在同一进程启用 `scheduler`、`kubernetes-worker`
`vm-worker`,也可通过 `--components` 只启用其中一部分。组件是独立应用服务边界, 和 `opensandbox-worker`,也可通过 `--components` 只启用其中一部分。组件是独立应用服务边界,
共享进程不意味着共享后端状态或把 assignment 降级为内存 channel。 共享进程不意味着共享后端状态或把 assignment 降级为内存 channel。
首轮集成的 facade pending/claimed registry 与三个组件同进程。虽然二进制保留组件选择 首轮集成的 facade pending/claimed registry 与三个组件同进程。虽然二进制保留组件选择
@@ -47,10 +47,13 @@ facade,并严格校验 facade 的 SPIFFE ID。这样无需修改 runner 或把
## 设计约束 ## 设计约束
- 对 workflow 的接口保持 `[self-hosted, pod]` 和 `[self-hosted, vm]` 不变。 - 规范接口为 `[self-hosted, container, kubernetes]` 和
`[self-hosted, vm, opensandbox]`。旧 `pod`、`vm`、`vm-dev` 标签保留严格映射,不能与
冲突 class/driver 混用;显式 driver 不允许自动回退。
- scheduler 使用单一 Gitea runner UUID/token 和一个 `Declare`,不为并发槽位重复注册; - scheduler 使用单一 Gitea runner UUID/token 和一个 `Declare`,不为并发槽位重复注册;
`POD_CAPACITY + VM_CAPACITY` 决定并发 `FetchTask` goroutine 数量。 `POD_CAPACITY + VM_CAPACITY` 决定并发 `FetchTask` goroutine 数量。
- task 领取并持久化后按 backend 进入独立 durable consumer;对应容量池已满时延迟 NAK, - task 领取并持久化后按 placement 进入独立 durable consumer;v2 subject 为
`<subject-base>.<workload-class>.<driver>`。对应容量池已满时延迟 NAK,
assignment 保持 JetStream pending,且不得创建超出配置容量的 workload。 assignment 保持 JetStream pending,且不得创建超出配置容量的 workload。
- scheduler Declare 后使用 RunnerService 长轮询;一旦 FetchTask 返回已分配 task,在 - scheduler Declare 后使用 RunnerService 长轮询;一旦 FetchTask 返回已分配 task,在
JetStream publish 成功前只重试该 assignment,不领取下一项。 JetStream publish 成功前只重试该 assignment,不领取下一项。
@@ -66,7 +69,7 @@ facade,并严格校验 facade 的 SPIFFE ID。这样无需修改 runner 或把
- JetStream 只持久化和投递 assignment,不保存 executor 生命周期状态。Pod labels/annotations - JetStream 只持久化和投递 assignment,不保存 executor 生命周期状态。Pod labels/annotations
与 OpenSandbox metadata 是后端运行状态的权威来源,Gitea 是 task 终态的权威来源。 与 OpenSandbox metadata 是后端运行状态的权威来源,Gitea 是 task 终态的权威来源。
- assignment 使用版本化 envelope 保存完整 Gitea protobuf task,并从 workflow `runs-on` - assignment 使用版本化 envelope 保存完整 Gitea protobuf task,并从 workflow `runs-on`
严格选择 pod 或 vm subject;消费者解码后重新派生 backend 与身份,拒绝被篡改的冗余字段。 严格选择 workload class 与 driver;消费者解码后重新派生 placement 与身份,拒绝被篡改的冗余字段。
- JetStream 的 message ID 等于稳定 assignment ID `gitea-task-<task-id>`,仅用于发布去重, - JetStream 的 message ID 等于稳定 assignment ID `gitea-task-<task-id>`,仅用于发布去重,
不承担 executor 生命周期记录。 不承担 executor 生命周期记录。
- worker 按稳定 assignment ID reconcile 后端资源,进程内只保留并发控制等可丢弃状态; - worker 按稳定 assignment ID reconcile 后端资源,进程内只保留并发控制等可丢弃状态;
@@ -76,7 +79,7 @@ facade,并严格校验 facade 的 SPIFFE ID。这样无需修改 runner 或把
- executor 成功 claim 后 ACK assignment。Gitea 接受 terminal update 后,facade 在后端 - executor 成功 claim 后 ACK assignment。Gitea 接受 terminal update 后,facade 在后端
metadata 写入持久 terminal marker;backend reconciler 仅在执行环境也进入终态后清理, metadata 写入持久 terminal marker;backend reconciler 仅在执行环境也进入终态后清理,
从而关闭进程重启窗口且避免删除尚未完成结果上报的环境。 从而关闭进程重启窗口且避免删除尚未完成结果上报的环境。
- pod 与 vm 使用独立 durable consumer、进程内 admission pool 和并发上限。consumer 只负责将 assignment - 每个 placement 使用独立 durable consumer、进程内 admission pool 和并发上限。consumer 只负责将 assignment
幂等落到后端;executor 与身份恢复 metadata 持久化后立即 `DoubleAck`。尚未取得 幂等落到后端;executor 与身份恢复 metadata 持久化后立即 `DoubleAck`。尚未取得
Pod UID 等短暂未就绪状态以及临时后端错误使用延迟 NAK。 Pod UID 等短暂未就绪状态以及临时后端错误使用延迟 NAK。
- admission pool 只保存可重建的并发状态:启动时从 Pod labels/annotations 或 OpenSandbox - admission pool 只保存可重建的并发状态:启动时从 Pod labels/annotations 或 OpenSandbox
@@ -86,10 +89,10 @@ facade,并严格校验 facade 的 SPIFFE ID。这样无需修改 runner 或把
- consumer 在 executor 使用上述 facade 成功 claim task 后确认 assignment;无需把完整 - consumer 在 executor 使用上述 facade 成功 claim task 后确认 assignment;无需把完整
task 写入 Pod annotation、OpenSandbox metadata 或环境变量。 task 写入 Pod annotation、OpenSandbox metadata 或环境变量。
- Pod 与 VM 共享 task/executor 协议,只有环境创建和销毁实现不同。 - Pod 与 VM 共享 task/executor 协议,只有环境创建和销毁实现不同。
- scheduler 在 assignment 持久化到 JetStream 后即可继续领取;Pod 与 VM 分别由 durable - scheduler 在 assignment 持久化到 JetStream 后即可继续领取;各 placement 分别由 durable
consumer 的 capacity 限制并发,不共享全局执行槽位。未知后端故障由对应 consumer 的 consumer 的 capacity 限制并发,不共享全局执行槽位。未知后端故障由对应 consumer 的
NAK/redelivery 收敛,不能阻塞另一种 backend。 NAK/redelivery 收敛,不能阻塞另一种 backend。
- 两种 backend 都注入同一份 runner bootstrap 环境;Pod 仍由 homelab Kubernetes 原生 - 两种现有 driver 都注入同一份 runner bootstrap 环境;container 仍由 homelab Kubernetes 原生
创建,只有 VM 经 OpenSandbox 创建,bootstrap 机制不改变 backend 边界。 创建,只有 VM 经 OpenSandbox 创建,bootstrap 机制不改变 backend 边界。
## 实现顺序 ## 实现顺序
+16 -11
View File
@@ -19,7 +19,7 @@ type publishAPI interface {
PublishMsg(context.Context, *nats.Msg, ...jetstream.PublishOpt) (*jetstream.PubAck, error) PublishMsg(context.Context, *nats.Msg, ...jetstream.PublishOpt) (*jetstream.PubAck, error)
} }
// Publisher implements the scheduler dispatcher with one subject per backend. // Publisher implements the scheduler dispatcher with one subject per placement.
type Publisher struct { type Publisher struct {
JetStream publishAPI JetStream publishAPI
SubjectBase string SubjectBase string
@@ -38,7 +38,7 @@ func (p Publisher) Dispatch(ctx context.Context, assignment taskassignment.Assig
return errors.New("assignment subject base is required") return errors.New("assignment subject base is required")
} }
message := &nats.Msg{ message := &nats.Msg{
Subject: base + "." + string(assignment.Backend), Subject: base + "." + assignment.Placement.Key(),
Header: nats.Header{jetstream.MsgIDHeader: []string{assignment.ID}}, Header: nats.Header{jetstream.MsgIDHeader: []string{assignment.ID}},
Data: body, Data: body,
} }
@@ -86,13 +86,13 @@ type Processor struct {
type Event struct { type Event struct {
Name string Name string
AssignmentID string AssignmentID string
Backend taskassignment.Backend Placement taskassignment.Placement
RetryDelay time.Duration RetryDelay time.Duration
} }
func (p Processor) event(name string, assignment taskassignment.Assignment, retryDelay time.Duration) { func (p Processor) event(name string, assignment taskassignment.Assignment, retryDelay time.Duration) {
if p.OnEvent != nil { if p.OnEvent != nil {
p.OnEvent(Event{Name: name, AssignmentID: assignment.ID, Backend: assignment.Backend, RetryDelay: retryDelay}) p.OnEvent(Event{Name: name, AssignmentID: assignment.ID, Placement: assignment.Placement, RetryDelay: retryDelay})
} }
} }
@@ -178,24 +178,29 @@ type consumerManager interface {
// OpenConsumer creates the durable backend cursor. Capacity is enforced both // OpenConsumer creates the durable backend cursor. Capacity is enforced both
// server-side and by ConsumerComponent's local semaphore. // server-side and by ConsumerComponent's local semaphore.
func OpenConsumer(ctx context.Context, manager consumerManager, stream, subjectBase string, backend taskassignment.Backend, capacity int) (jetstream.Consumer, error) { func OpenConsumer(ctx context.Context, manager consumerManager, stream, subjectBase string, placement taskassignment.Placement, capacity int) (jetstream.Consumer, error) {
if manager == nil || stream == "" || strings.TrimSuffix(subjectBase, ".") == "" || capacity < 1 { if manager == nil || stream == "" || strings.TrimSuffix(subjectBase, ".") == "" || capacity < 1 {
return nil, errors.New("JetStream manager, stream, subject base, and positive capacity are required") return nil, errors.New("JetStream manager, stream, subject base, and positive capacity are required")
} }
if backend != taskassignment.BackendPod && backend != taskassignment.BackendVM { if err := placement.Validate(); err != nil {
return nil, fmt.Errorf("unsupported assignment backend %q", backend) return nil, err
} }
key := placement.Key()
// JetStream consumer names may not contain dots even though subjects do.
// Keep the placement subject hierarchical while using a stable, legal name
// for the durable cursor.
consumerName := strings.ReplaceAll(key, ".", "-")
consumer, err := manager.CreateOrUpdateConsumer(ctx, stream, jetstream.ConsumerConfig{ consumer, err := manager.CreateOrUpdateConsumer(ctx, stream, jetstream.ConsumerConfig{
Name: string(backend), Name: consumerName,
Durable: string(backend), Durable: consumerName,
FilterSubject: strings.TrimSuffix(subjectBase, ".") + "." + string(backend), FilterSubject: strings.TrimSuffix(subjectBase, ".") + "." + key,
AckPolicy: jetstream.AckExplicitPolicy, AckPolicy: jetstream.AckExplicitPolicy,
AckWait: 5 * time.Minute, AckWait: 5 * time.Minute,
MaxAckPending: capacity, MaxAckPending: capacity,
MaxDeliver: 1000, 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", key, err)
} }
return consumer, nil return consumer, nil
} }
+6 -6
View File
@@ -38,13 +38,13 @@ func (p *fakePublisher) PublishMsg(_ context.Context, message *nats.Msg, _ ...je
return &jetstream.PubAck{}, nil return &jetstream.PubAck{}, nil
} }
func TestPublisherUsesBackendSubjectAndAssignmentDeduplication(t *testing.T) { func TestPublisherUsesPlacementSubjectAndAssignmentDeduplication(t *testing.T) {
api := &fakePublisher{} api := &fakePublisher{}
publisher := Publisher{JetStream: api, SubjectBase: "ci.assignment"} publisher := Publisher{JetStream: api, SubjectBase: "ci.assignment"}
if err := publisher.Dispatch(context.Background(), testAssignment(t)); err != nil { if err := publisher.Dispatch(context.Background(), testAssignment(t)); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if api.message.Subject != "ci.assignment.pod" { if api.message.Subject != "ci.assignment.container.kubernetes" {
t.Fatalf("subject = %q", api.message.Subject) t.Fatalf("subject = %q", api.message.Subject)
} }
if api.message.Header.Get(jetstream.MsgIDHeader) != "gitea-task-42" { if api.message.Header.Get(jetstream.MsgIDHeader) != "gitea-task-42" {
@@ -139,7 +139,7 @@ func TestProcessorAcknowledgesPersistedHandoff(t *testing.T) {
t.Fatalf("events = %#v", events) t.Fatalf("events = %#v", events)
} }
for index := range want { for index := range want {
if events[index].Name != want[index] || events[index].AssignmentID != "gitea-task-42" || events[index].Backend != taskassignment.BackendPod { if events[index].Name != want[index] || events[index].AssignmentID != "gitea-task-42" || events[index].Placement != taskassignment.KubernetesContainer {
t.Fatalf("event[%d] = %#v", index, events[index]) t.Fatalf("event[%d] = %#v", index, events[index])
} }
} }
@@ -210,12 +210,12 @@ func (m *fakeConsumerManager) CreateOrUpdateConsumer(_ context.Context, _ string
return nil, nil return nil, nil
} }
func TestOpenConsumerUsesIndependentDurablePerBackend(t *testing.T) { func TestOpenConsumerUsesIndependentDurablePerPlacement(t *testing.T) {
manager := &fakeConsumerManager{} manager := &fakeConsumerManager{}
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.KubernetesContainer, 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 || manager.config.MaxDeliver != 1000 { if manager.config.Name != "container-kubernetes" || manager.config.Durable != "container-kubernetes" || manager.config.FilterSubject != "ci.assignment.container.kubernetes" || 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)
} }
} }
+9 -3
View File
@@ -15,11 +15,11 @@ type ComponentName string
const ( const (
Scheduler ComponentName = "scheduler" Scheduler ComponentName = "scheduler"
PodWorker ComponentName = "pod-worker" KubernetesWorker ComponentName = "kubernetes-worker"
VMWorker ComponentName = "vm-worker" OpenSandboxWorker ComponentName = "opensandbox-worker"
) )
var defaultComponents = []ComponentName{Scheduler, PodWorker, VMWorker} var defaultComponents = []ComponentName{Scheduler, KubernetesWorker, OpenSandboxWorker}
// Selection parses --components. An empty value enables all components. // Selection parses --components. An empty value enables all components.
type Selection []ComponentName type Selection []ComponentName
@@ -31,6 +31,12 @@ func ParseSelection(value string) (Selection, error) {
var selected Selection var selected Selection
for _, raw := range strings.Split(value, ",") { for _, raw := range strings.Split(value, ",") {
name := ComponentName(strings.TrimSpace(raw)) name := ComponentName(strings.TrimSpace(raw))
switch name {
case "pod-worker":
name = KubernetesWorker
case "vm-worker":
name = OpenSandboxWorker
}
if !slices.Contains(defaultComponents, name) { if !slices.Contains(defaultComponents, name) {
return nil, fmt.Errorf("unknown controller component %q", name) return nil, fmt.Errorf("unknown controller component %q", name)
} }
+16 -6
View File
@@ -13,18 +13,18 @@ func TestParseSelectionDefaultsToAll(t *testing.T) {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if len(selection) != 3 || selection[0] != Scheduler || selection[1] != PodWorker || selection[2] != VMWorker { if len(selection) != 3 || selection[0] != Scheduler || selection[1] != KubernetesWorker || selection[2] != OpenSandboxWorker {
t.Fatalf("selection = %v", selection) t.Fatalf("selection = %v", selection)
} }
} }
} }
func TestParseSelectionAllowsOneOrMoreComponents(t *testing.T) { func TestParseSelectionAllowsOneOrMoreComponents(t *testing.T) {
selection, err := ParseSelection("vm-worker,scheduler,vm-worker") selection, err := ParseSelection("opensandbox-worker,scheduler,opensandbox-worker")
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if len(selection) != 2 || selection[0] != VMWorker || selection[1] != Scheduler { if len(selection) != 2 || selection[0] != OpenSandboxWorker || selection[1] != Scheduler {
t.Fatalf("selection = %v", selection) t.Fatalf("selection = %v", selection)
} }
if _, err := ParseSelection("webhook"); err == nil { if _, err := ParseSelection("webhook"); err == nil {
@@ -32,6 +32,16 @@ func TestParseSelectionAllowsOneOrMoreComponents(t *testing.T) {
} }
} }
func TestParseSelectionNormalizesLegacyWorkerNames(t *testing.T) {
selection, err := ParseSelection("pod-worker,vm-worker")
if err != nil {
t.Fatal(err)
}
if len(selection) != 2 || selection[0] != KubernetesWorker || selection[1] != OpenSandboxWorker {
t.Fatalf("selection = %v", selection)
}
}
type componentFunc func(context.Context) error type componentFunc func(context.Context) error
func (f componentFunc) Run(ctx context.Context) error { return f(ctx) } func (f componentFunc) Run(ctx context.Context) error { return f(ctx) }
@@ -45,14 +55,14 @@ func TestRunStartsSelectedComponentsAndCancelsPeers(t *testing.T) {
started <- Scheduler started <- Scheduler
return errors.New("poll failed") return errors.New("poll failed")
}), }),
PodWorker: componentFunc(func(ctx context.Context) error { KubernetesWorker: componentFunc(func(ctx context.Context) error {
started <- PodWorker started <- KubernetesWorker
<-ctx.Done() <-ctx.Done()
once.Do(func() { close(peerStopped) }) once.Do(func() { close(peerStopped) })
return ctx.Err() return ctx.Err()
}), }),
} }
err := Run(context.Background(), Selection{Scheduler, PodWorker}, registry) err := Run(context.Background(), Selection{Scheduler, KubernetesWorker}, registry)
if err == nil || !errors.Is(err, context.Canceled) && err.Error() != "component scheduler: poll failed" { if err == nil || !errors.Is(err, context.Canceled) && err.Error() != "component scheduler: poll failed" {
t.Fatalf("Run() error = %v", err) t.Fatalf("Run() error = %v", err)
} }
+6 -3
View File
@@ -119,7 +119,10 @@ func (b Backend) RecoverAssignments(ctx context.Context, trustDomain string) ([]
return nil, err return nil, err
} }
result, err := b.Lifecycle.ListSandboxes(ctx, opensandbox.ListOptions{ result, err := b.Lifecycle.ListSandboxes(ctx, opensandbox.ListOptions{
Metadata: map[string]string{"ci.ddupan.top/backend": "vm"}, PageSize: 100, Metadata: map[string]string{
"ci.ddupan.top/workload-class": "vm",
"ci.ddupan.top/driver": "opensandbox",
}, PageSize: 100,
}) })
if err != nil { if err != nil {
return nil, fmt.Errorf("list recoverable sandboxes: %w", err) return nil, fmt.Errorf("list recoverable sandboxes: %w", err)
@@ -172,8 +175,8 @@ func (b Backend) Create(ctx context.Context, assignment taskassignment.Assignmen
if err := b.validate(); err != nil { if err := b.validate(); err != nil {
return nil, err return nil, err
} }
if assignment.Backend != taskassignment.BackendVM { if assignment.Placement != taskassignment.OpenSandboxVM {
return nil, fmt.Errorf("OpenSandbox backend cannot create %q assignment", assignment.Backend) return nil, fmt.Errorf("OpenSandbox backend cannot create %q", assignment.Placement.Key())
} }
environment := clone(b.Config.Env) environment := clone(b.Config.Env)
for key, value := range launch.Environment { for key, value := range launch.Environment {
+1 -1
View File
@@ -63,7 +63,7 @@ func backend(lifecycle Lifecycle) Backend {
func assignment() taskassignment.Assignment { func assignment() taskassignment.Assignment {
return taskassignment.Assignment{ return taskassignment.Assignment{
ID: "gitea-task-42", Backend: taskassignment.BackendVM, ID: "gitea-task-42", Placement: taskassignment.OpenSandboxVM,
Task: &runnerv1.Task{Id: 42}, Task: &runnerv1.Task{Id: 42},
Identity: taskidentity.Identity{Repository: "owner/repo", Task: "publish", SPIFFEID: "spiffe://ddupan.top/ci/owner/repo/publish"}, Identity: taskidentity.Identity{Repository: "owner/repo", Task: "publish", SPIFFEID: "spiffe://ddupan.top/ci/owner/repo/publish"},
} }
+3 -3
View File
@@ -120,7 +120,7 @@ func (b Backend) RecoverAssignments(ctx context.Context) ([]taskassignment.Assig
if err := b.validate(); err != nil { if err := b.validate(); err != nil {
return nil, err return nil, err
} }
pods, err := b.API.ListPods(ctx, b.Config.Namespace, "ci.ddupan.top/backend=pod") pods, err := b.API.ListPods(ctx, b.Config.Namespace, "ci.ddupan.top/workload-class=container,ci.ddupan.top/driver=kubernetes")
if err != nil { if err != nil {
return nil, fmt.Errorf("list recoverable assignment Pods: %w", err) return nil, fmt.Errorf("list recoverable assignment Pods: %w", err)
} }
@@ -159,8 +159,8 @@ func (b Backend) Create(ctx context.Context, assignment taskassignment.Assignmen
if err := b.validate(); err != nil { if err := b.validate(); err != nil {
return nil, err return nil, err
} }
if assignment.Backend != taskassignment.BackendPod { if assignment.Placement != taskassignment.KubernetesContainer {
return nil, fmt.Errorf("Pod backend cannot create %q assignment", assignment.Backend) return nil, fmt.Errorf("Kubernetes backend cannot create %q", assignment.Placement.Key())
} }
labels := clone(launch.Metadata.Labels) labels := clone(launch.Metadata.Labels)
labels["app.kubernetes.io/name"] = "gitea-dynamic-runner" labels["app.kubernetes.io/name"] = "gitea-dynamic-runner"
+1 -1
View File
@@ -59,7 +59,7 @@ func backend(api API) Backend {
func assignment() taskassignment.Assignment { func assignment() taskassignment.Assignment {
return taskassignment.Assignment{ return taskassignment.Assignment{
ID: "gitea-task-42", Backend: taskassignment.BackendPod, ID: "gitea-task-42", Placement: taskassignment.KubernetesContainer,
Task: &runnerv1.Task{Id: 42}, Task: &runnerv1.Task{Id: 42},
Identity: taskidentity.Identity{ Identity: taskidentity.Identity{
Repository: "owner/repo", Task: "publish", Repository: "owner/repo", Task: "publish",
+19 -4
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,17 +77,26 @@ 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{
Volumes: []corev1.Volume{{ {
Name: "spire-agent-socket", Name: "spire-agent-socket",
VolumeSource: corev1.VolumeSource{CSI: &corev1.CSIVolumeSource{ VolumeSource: corev1.VolumeSource{CSI: &corev1.CSIVolumeSource{
Driver: "csi.spiffe.io", ReadOnly: boolPointer(true), 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{})
if err != nil { if err != nil {
@@ -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)
} }
+15 -6
View File
@@ -18,7 +18,7 @@ const (
EnvFacadeURL = "CI_RUNNER_FACADE_URL" EnvFacadeURL = "CI_RUNNER_FACADE_URL"
EnvFacadeID = "CI_RUNNER_FACADE_SPIFFE_ID" EnvFacadeID = "CI_RUNNER_FACADE_SPIFFE_ID"
EnvSPIFFEID = "CI_SPIFFE_ID" EnvSPIFFEID = "CI_SPIFFE_ID"
EnvBackend = "CI_RUNNER_BACKEND" EnvRunnerLabels = "CI_RUNNER_LABELS_JSON"
) )
// Bootstrap emits assignment-scoped launch configuration. FacadeURL is the // Bootstrap emits assignment-scoped launch configuration. FacadeURL is the
@@ -42,13 +42,17 @@ func (b Bootstrap) Environment(assignment taskassignment.Assignment) (map[string
if capability == "" { if capability == "" {
return nil, errors.New("runner capability issuer is not configured") return nil, errors.New("runner capability issuer is not configured")
} }
labels, err := json.Marshal([]string{"self-hosted", string(assignment.Placement.Class), string(assignment.Placement.Driver)})
if err != nil {
return nil, fmt.Errorf("encode runner labels: %w", err)
}
return map[string]string{ return map[string]string{
EnvAssignmentID: assignment.ID, EnvAssignmentID: assignment.ID,
EnvCapability: capability, EnvCapability: capability,
EnvFacadeURL: b.FacadeURL, EnvFacadeURL: b.FacadeURL,
EnvFacadeID: b.FacadeSPIFFEID, EnvFacadeID: b.FacadeSPIFFEID,
EnvSPIFFEID: assignment.Identity.SPIFFEID, EnvSPIFFEID: assignment.Identity.SPIFFEID,
EnvBackend: string(assignment.Backend), EnvRunnerLabels: string(labels),
"SPIFFE_ENDPOINT_SOCKET": b.WorkloadAPIAddr, "SPIFFE_ENDPOINT_SOCKET": b.WorkloadAPIAddr,
}, nil }, nil
} }
@@ -67,7 +71,7 @@ type Registration struct {
Ephemeral bool `json:"ephemeral"` Ephemeral bool `json:"ephemeral"`
} }
func RegistrationJSON(assignmentID, capability, localProxyURL string, backend taskassignment.Backend) ([]byte, error) { func RegistrationJSON(assignmentID, capability, localProxyURL string, labels []string) ([]byte, error) {
if assignmentID == "" || capability == "" { if assignmentID == "" || capability == "" {
return nil, errors.New("assignment ID and runner capability are required") return nil, errors.New("assignment ID and runner capability are required")
} }
@@ -75,13 +79,18 @@ func RegistrationJSON(assignmentID, capability, localProxyURL string, backend ta
if err != nil || parsed.Scheme != "http" || parsed.Host == "" { if err != nil || parsed.Scheme != "http" || parsed.Host == "" {
return nil, errors.New("local runner proxy URL must be an absolute http URL") return nil, errors.New("local runner proxy URL must be an absolute http URL")
} }
if backend != taskassignment.BackendPod && backend != taskassignment.BackendVM { if len(labels) == 0 {
return nil, fmt.Errorf("unsupported runner backend %q", backend) return nil, errors.New("runner labels are required")
}
for _, label := range labels {
if label == "" {
return nil, errors.New("runner labels must not be empty")
}
} }
registration := Registration{ registration := Registration{
Warning: "Generated for one preassigned task by gitea-dynamic-runner.", Warning: "Generated for one preassigned task by gitea-dynamic-runner.",
UUID: assignmentID, Name: assignmentID, Token: capability, UUID: assignmentID, Name: assignmentID, Token: capability,
Address: localProxyURL, Labels: []string{"self-hosted", string(backend)}, Ephemeral: true, Address: localProxyURL, Labels: append([]string(nil), labels...), Ephemeral: true,
} }
data, err := json.MarshalIndent(registration, "", " ") data, err := json.MarshalIndent(registration, "", " ")
if err != nil { if err != nil {
+5 -5
View File
@@ -25,7 +25,7 @@ func testBootstrap(t *testing.T) Bootstrap {
func testAssignment() taskassignment.Assignment { func testAssignment() taskassignment.Assignment {
return taskassignment.Assignment{ return taskassignment.Assignment{
ID: "gitea-task-42", Backend: taskassignment.BackendPod, ID: "gitea-task-42", Placement: taskassignment.KubernetesContainer,
Identity: taskidentity.Identity{SPIFFEID: "spiffe://ddupan.top/ci/owner/repo/publish"}, Identity: taskidentity.Identity{SPIFFEID: "spiffe://ddupan.top/ci/owner/repo/publish"},
} }
} }
@@ -46,7 +46,7 @@ func TestEnvironmentIsDeterministicAndAssignmentScoped(t *testing.T) {
if first[EnvAssignmentID] != "gitea-task-42" || first[EnvSPIFFEID] != testAssignment().Identity.SPIFFEID { if first[EnvAssignmentID] != "gitea-task-42" || first[EnvSPIFFEID] != testAssignment().Identity.SPIFFEID {
t.Fatalf("environment = %#v", first) t.Fatalf("environment = %#v", first)
} }
if first[EnvBackend] != "pod" || first[EnvFacadeID] == "" { if first[EnvRunnerLabels] != `["self-hosted","container","kubernetes"]` || first[EnvFacadeID] == "" {
t.Fatalf("environment = %#v", first) t.Fatalf("environment = %#v", first)
} }
if first["SPIFFE_ENDPOINT_SOCKET"] != "unix:///run/spire/agent-sockets/spire-agent.sock" { if first["SPIFFE_ENDPOINT_SOCKET"] != "unix:///run/spire/agent-sockets/spire-agent.sock" {
@@ -55,7 +55,7 @@ func TestEnvironmentIsDeterministicAndAssignmentScoped(t *testing.T) {
} }
func TestRegistrationMatchesOfficialRunnerSchema(t *testing.T) { func TestRegistrationMatchesOfficialRunnerSchema(t *testing.T) {
data, err := RegistrationJSON("gitea-task-42", "capability", "http://127.0.0.1:8080", taskassignment.BackendVM) data, err := RegistrationJSON("gitea-task-42", "capability", "http://127.0.0.1:8080", []string{"self-hosted", "vm", "opensandbox"})
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
@@ -66,13 +66,13 @@ func TestRegistrationMatchesOfficialRunnerSchema(t *testing.T) {
if registration.UUID != "gitea-task-42" || registration.Token != "capability" || registration.Address != "http://127.0.0.1:8080" || !registration.Ephemeral { if registration.UUID != "gitea-task-42" || registration.Token != "capability" || registration.Address != "http://127.0.0.1:8080" || !registration.Ephemeral {
t.Fatalf("registration = %#v", registration) t.Fatalf("registration = %#v", registration)
} }
if len(registration.Labels) != 2 || registration.Labels[0] != "self-hosted" || registration.Labels[1] != "vm" { if len(registration.Labels) != 3 || registration.Labels[0] != "self-hosted" || registration.Labels[1] != "vm" || registration.Labels[2] != "opensandbox" {
t.Fatalf("labels = %#v", registration.Labels) t.Fatalf("labels = %#v", registration.Labels)
} }
} }
func TestRegistrationRejectsNonLocalTLSAddress(t *testing.T) { func TestRegistrationRejectsNonLocalTLSAddress(t *testing.T) {
if _, err := RegistrationJSON("id", "capability", "https://facade.example", taskassignment.BackendPod); err == nil { if _, err := RegistrationJSON("id", "capability", "https://facade.example", []string{"self-hosted"}); err == nil {
t.Fatal("expected local proxy URL validation error") t.Fatal("expected local proxy URL validation error")
} }
} }
+18 -10
View File
@@ -2,6 +2,7 @@ package runnerbootstrap
import ( import (
"context" "context"
"encoding/json"
"errors" "errors"
"fmt" "fmt"
"net" "net"
@@ -10,18 +11,17 @@ import (
"os/exec" "os/exec"
"path/filepath" "path/filepath"
"time" "time"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskassignment"
) )
type ExecutorConfig struct { type ExecutorConfig struct {
AssignmentID string AssignmentID string
Capability string Capability string
Backend taskassignment.Backend RunnerLabels []string
FacadeURL string FacadeURL string
FacadeSPIFFEID string FacadeSPIFFEID string
WorkloadAPIAddr string WorkloadAPIAddr string
RunnerBinary string RunnerBinary string
RunnerConfig string
ListenAddress string ListenAddress string
WorkDir string WorkDir string
Stdout *os.File Stdout *os.File
@@ -35,6 +35,9 @@ func RunExecutor(ctx context.Context, config ExecutorConfig) error {
if config.RunnerBinary == "" { if config.RunnerBinary == "" {
config.RunnerBinary = "gitea-runner" config.RunnerBinary = "gitea-runner"
} }
if config.RunnerConfig == "" {
config.RunnerConfig = "/etc/gitea-runner/config.yaml"
}
if config.ListenAddress == "" { if config.ListenAddress == "" {
config.ListenAddress = "127.0.0.1:0" config.ListenAddress = "127.0.0.1:0"
} }
@@ -67,7 +70,7 @@ func RunExecutor(ctx context.Context, config ExecutorConfig) error {
defer os.RemoveAll(workDir) defer os.RemoveAll(workDir)
} }
registration, err := RegistrationJSON( registration, err := RegistrationJSON(
config.AssignmentID, config.Capability, "http://"+listener.Addr().String(), config.Backend, config.AssignmentID, config.Capability, "http://"+listener.Addr().String(), config.RunnerLabels,
) )
if err != nil { if err != nil {
return err return err
@@ -93,7 +96,7 @@ func RunExecutor(ctx context.Context, config ExecutorConfig) error {
return errors.Join(readyErr, shutdownErr, serverErr) return errors.Join(readyErr, shutdownErr, serverErr)
} }
command := exec.CommandContext(ctx, config.RunnerBinary, "daemon", "--once") command := exec.CommandContext(ctx, config.RunnerBinary, runnerArguments(config.RunnerConfig)...)
command.Dir = workDir command.Dir = workDir
command.Stdout = config.Stdout command.Stdout = config.Stdout
command.Stderr = config.Stderr command.Stderr = config.Stderr
@@ -108,6 +111,10 @@ func RunExecutor(ctx context.Context, config ExecutorConfig) error {
return errors.Join(runnerErr, shutdownErr, serverErr) return errors.Join(runnerErr, shutdownErr, serverErr)
} }
func runnerArguments(configFile string) []string {
return []string{"daemon", "--config", configFile, "--once"}
}
func waitForFacade(ctx context.Context, endpoint string) error { func waitForFacade(ctx context.Context, endpoint string) error {
client := &http.Client{Timeout: 2 * time.Second} client := &http.Client{Timeout: 2 * time.Second}
ticker := time.NewTicker(250 * time.Millisecond) ticker := time.NewTicker(250 * time.Millisecond)
@@ -136,14 +143,15 @@ func waitForFacade(ctx context.Context, endpoint string) error {
// the assignment-scoped values injected by the backend. The Workload API // the assignment-scoped values injected by the backend. The Workload API
// address follows SPIFFE_ENDPOINT_SOCKET through go-spiffe when not set here. // address follows SPIFFE_ENDPOINT_SOCKET through go-spiffe when not set here.
func ExecutorConfigFromEnvironment() (ExecutorConfig, error) { func ExecutorConfigFromEnvironment() (ExecutorConfig, error) {
backend := taskassignment.Backend(os.Getenv(EnvBackend)) var labels []string
if backend != taskassignment.BackendPod && backend != taskassignment.BackendVM { if err := json.Unmarshal([]byte(os.Getenv(EnvRunnerLabels)), &labels); err != nil || len(labels) == 0 {
return ExecutorConfig{}, fmt.Errorf("invalid %s %q", EnvBackend, backend) return ExecutorConfig{}, errors.New("valid runner labels are required")
} }
config := ExecutorConfig{ config := ExecutorConfig{
AssignmentID: os.Getenv(EnvAssignmentID), Capability: os.Getenv(EnvCapability), AssignmentID: os.Getenv(EnvAssignmentID), Capability: os.Getenv(EnvCapability),
Backend: backend, FacadeURL: os.Getenv(EnvFacadeURL), FacadeSPIFFEID: os.Getenv(EnvFacadeID), RunnerLabels: labels, FacadeURL: os.Getenv(EnvFacadeURL), FacadeSPIFFEID: os.Getenv(EnvFacadeID),
RunnerBinary: os.Getenv("GITEA_RUNNER_BINARY"), ListenAddress: "127.0.0.1:0", RunnerBinary: os.Getenv("GITEA_RUNNER_BINARY"), RunnerConfig: os.Getenv("GITEA_RUNNER_CONFIG_FILE"),
ListenAddress: "127.0.0.1:0",
Stdout: os.Stdout, Stderr: os.Stderr, Stdout: os.Stdout, Stderr: os.Stderr,
} }
if config.AssignmentID == "" || config.Capability == "" || config.FacadeURL == "" || config.FacadeSPIFFEID == "" { if config.AssignmentID == "" || config.Capability == "" || config.FacadeURL == "" || config.FacadeSPIFFEID == "" {
@@ -4,11 +4,19 @@ import (
"context" "context"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"reflect"
"sync/atomic" "sync/atomic"
"testing" "testing"
"time" "time"
) )
func TestRunnerArgumentsLoadJobHooksConfig(t *testing.T) {
want := []string{"daemon", "--config", "/etc/gitea-runner/config.yaml", "--once"}
if got := runnerArguments("/etc/gitea-runner/config.yaml"); !reflect.DeepEqual(got, want) {
t.Fatalf("runner arguments = %q, want %q", got, want)
}
}
func TestWaitForFacadeRetriesTransientGatewayFailure(t *testing.T) { func TestWaitForFacadeRetriesTransientGatewayFailure(t *testing.T) {
var requests atomic.Int32 var requests atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, _ *http.Request) { server := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, _ *http.Request) {
+19 -36
View File
@@ -6,7 +6,6 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"slices"
"strconv" "strconv"
"gitea.dev/actionslib/pkg/model" "gitea.dev/actionslib/pkg/model"
@@ -16,34 +15,27 @@ import (
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskidentity" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskidentity"
) )
const wireVersion = 1 const wireVersion = 2
type Backend string
const (
BackendPod Backend = "pod"
BackendVM Backend = "vm"
)
// Assignment is the only document persisted in the handoff queue. // Assignment is the only document persisted in the handoff queue.
type Assignment struct { type Assignment struct {
ID string ID string
Backend Backend Placement Placement
Task *runnerv1.Task Task *runnerv1.Task
Identity taskidentity.Identity Identity taskidentity.Identity
} }
// FromMetadata reconstructs the minimal assignment needed to authorize an // FromMetadata reconstructs the minimal assignment needed to authorize an
// already-running executor after a controller restart. Backend metadata was // already-running executor after a controller restart. Placement metadata was
// originally derived from the trusted Gitea task and is validated again here. // originally derived from the trusted Gitea task and is validated again here.
func FromMetadata(labels, annotations map[string]string, trustDomain string) (Assignment, error) { func FromMetadata(labels, annotations map[string]string, trustDomain string) (Assignment, error) {
taskID, err := strconv.ParseInt(labels["ci.ddupan.top/task-id"], 10, 64) taskID, err := strconv.ParseInt(labels["ci.ddupan.top/task-id"], 10, 64)
if err != nil || taskID < 1 { if err != nil || taskID < 1 {
return Assignment{}, errors.New("backend metadata has invalid task ID") return Assignment{}, errors.New("backend metadata has invalid task ID")
} }
backend := Backend(labels["ci.ddupan.top/backend"]) placement := Placement{Class: WorkloadClass(labels["ci.ddupan.top/workload-class"]), Driver: Driver(labels["ci.ddupan.top/driver"])}
if backend != BackendPod && backend != BackendVM { if err := placement.Validate(); err != nil {
return Assignment{}, errors.New("backend metadata has invalid backend") return Assignment{}, err
} }
id := labels["ci.ddupan.top/assignment-id"] id := labels["ci.ddupan.top/assignment-id"]
if id != fmt.Sprintf("gitea-task-%d", taskID) { if id != fmt.Sprintf("gitea-task-%d", taskID) {
@@ -56,13 +48,13 @@ func FromMetadata(labels, annotations map[string]string, trustDomain string) (As
if err != nil { if err != nil {
return Assignment{}, err return Assignment{}, err
} }
return Assignment{ID: id, Backend: backend, Task: &runnerv1.Task{Id: taskID}, Identity: identity}, nil return Assignment{ID: id, Placement: placement, Task: &runnerv1.Task{Id: taskID}, Identity: identity}, nil
} }
type envelope struct { type envelope struct {
Version int `json:"version"` Version int `json:"version"`
ID string `json:"id"` ID string `json:"id"`
Backend Backend `json:"backend"` Placement Placement `json:"placement"`
Task []byte `json:"task"` Task []byte `json:"task"`
Identity taskidentity.Identity `json:"identity"` Identity taskidentity.Identity `json:"identity"`
} }
@@ -76,40 +68,28 @@ func New(task *runnerv1.Task, trustDomain string) (Assignment, error) {
if err != nil { if err != nil {
return Assignment{}, err return Assignment{}, err
} }
backend, err := backendFromTask(task) placement, err := placementFromTask(task)
if err != nil { if err != nil {
return Assignment{}, err return Assignment{}, err
} }
return Assignment{ return Assignment{
ID: fmt.Sprintf("gitea-task-%d", task.GetId()), ID: fmt.Sprintf("gitea-task-%d", task.GetId()),
Backend: backend, Placement: placement,
Task: task, Task: task,
Identity: identity, Identity: identity,
}, nil }, nil
} }
func backendFromTask(task *runnerv1.Task) (Backend, error) { func placementFromTask(task *runnerv1.Task) (Placement, error) {
workflow, err := model.ReadWorkflow(bytes.NewReader(task.GetWorkflowPayload())) workflow, err := model.ReadWorkflow(bytes.NewReader(task.GetWorkflowPayload()))
if err != nil { if err != nil {
return "", fmt.Errorf("parse task workflow for backend: %w", err) return Placement{}, fmt.Errorf("parse task workflow for placement: %w", err)
} }
jobIDs := workflow.GetJobIDs() jobIDs := workflow.GetJobIDs()
if len(jobIDs) != 1 || workflow.GetJob(jobIDs[0]) == nil { if len(jobIDs) != 1 || workflow.GetJob(jobIDs[0]) == nil {
return "", fmt.Errorf("task workflow must contain exactly one non-empty job") return Placement{}, fmt.Errorf("task workflow must contain exactly one non-empty job")
} }
labels := workflow.GetJob(jobIDs[0]).RunsOnLabels() return PlacementFromLabels(workflow.GetJob(jobIDs[0]).RunsOnLabels())
if !slices.Contains(labels, "self-hosted") {
return "", fmt.Errorf("task runs-on labels must include self-hosted: %v", labels)
}
hasPod := slices.Contains(labels, string(BackendPod))
hasVM := slices.Contains(labels, string(BackendVM)) || slices.Contains(labels, "vm-dev")
if hasPod == hasVM {
return "", fmt.Errorf("task runs-on labels must select exactly one of pod or vm: %v", labels)
}
if hasPod {
return BackendPod, nil
}
return BackendVM, nil
} }
// Marshal encodes a versioned assignment. Protobuf preserves the exact Gitea task. // Marshal encodes a versioned assignment. Protobuf preserves the exact Gitea task.
@@ -117,13 +97,16 @@ func Marshal(assignment Assignment) ([]byte, error) {
if assignment.Task == nil { if assignment.Task == nil {
return nil, errors.New("assignment task is required") return nil, errors.New("assignment task is required")
} }
if err := assignment.Placement.Validate(); err != nil {
return nil, err
}
task, err := proto.Marshal(assignment.Task) task, err := proto.Marshal(assignment.Task)
if err != nil { if err != nil {
return nil, fmt.Errorf("marshal Gitea task: %w", err) return nil, fmt.Errorf("marshal Gitea task: %w", err)
} }
return json.Marshal(envelope{ return json.Marshal(envelope{
Version: wireVersion, Version: wireVersion,
ID: assignment.ID, Backend: assignment.Backend, ID: assignment.ID, Placement: assignment.Placement,
Task: task, Identity: assignment.Identity, Task: task, Identity: assignment.Identity,
}) })
} }
@@ -145,7 +128,7 @@ func Unmarshal(data []byte, trustDomain string) (Assignment, error) {
if err != nil { if err != nil {
return Assignment{}, err return Assignment{}, err
} }
if wire.ID != canonical.ID || wire.Backend != canonical.Backend || wire.Identity != canonical.Identity { if wire.ID != canonical.ID || wire.Placement != canonical.Placement || wire.Identity != canonical.Identity {
return Assignment{}, errors.New("assignment metadata does not match its Gitea task") return Assignment{}, errors.New("assignment metadata does not match its Gitea task")
} }
return canonical, nil return canonical, nil
+21 -12
View File
@@ -21,29 +21,37 @@ func task(t *testing.T, labels string) *runnerv1.Task {
} }
} }
func TestNewSelectsBackendFromRunsOn(t *testing.T) { func TestNewSelectsPlacementFromRunsOn(t *testing.T) {
for _, test := range []struct { for _, test := range []struct {
labels string labels string
backend Backend placement Placement
}{ }{
{"[self-hosted, pod]", BackendPod}, {"[self-hosted, pod]", KubernetesContainer},
{"[self-hosted, vm]", BackendVM}, {"[self-hosted, container]", KubernetesContainer},
{"[self-hosted, vm-dev]", BackendVM}, {"[self-hosted, container, kubernetes]", KubernetesContainer},
{"[self-hosted, vm]", OpenSandboxVM},
{"[self-hosted, vm-dev]", OpenSandboxVM},
{"[self-hosted, vm, opensandbox]", OpenSandboxVM},
} { } {
assignment, err := New(task(t, test.labels), "ddupan.top") assignment, err := New(task(t, test.labels), "ddupan.top")
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if assignment.Backend != test.backend || assignment.ID != "gitea-task-42" { if assignment.Placement != test.placement || assignment.ID != "gitea-task-42" {
t.Fatalf("assignment = %#v", assignment) t.Fatalf("assignment = %#v", assignment)
} }
} }
} }
func TestNewRejectsAmbiguousBackend(t *testing.T) { func TestNewRejectsInvalidPlacement(t *testing.T) {
for _, labels := range []string{ for _, labels := range []string{
"[self-hosted]", "[self-hosted]",
"[self-hosted, pod, vm]", "[self-hosted, pod, vm]",
"[self-hosted, pod, container]",
"[self-hosted, pod, kubernetes]",
"[self-hosted, container, opensandbox]",
"[self-hosted, vm, kubernetes]",
"[self-hosted, container, kubernetes, opensandbox]",
"[pod]", "[pod]",
} { } {
if _, err := New(task(t, labels), "ddupan.top"); err == nil { if _, err := New(task(t, labels), "ddupan.top"); err == nil {
@@ -65,13 +73,13 @@ func TestAssignmentWireRoundTripAndValidation(t *testing.T) {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if got.ID != want.ID || got.Backend != want.Backend || got.Identity != want.Identity || !bytes.Equal(got.Task.WorkflowPayload, want.Task.WorkflowPayload) { if got.ID != want.ID || got.Placement != want.Placement || got.Identity != want.Identity || !bytes.Equal(got.Task.WorkflowPayload, want.Task.WorkflowPayload) {
t.Fatalf("round trip = %#v, want %#v", got, want) t.Fatalf("round trip = %#v, want %#v", got, want)
} }
tampered := bytes.Replace(data, []byte(`"backend":"pod"`), []byte(`"backend":"vm"`), 1) tampered := bytes.Replace(data, []byte(`"driver":"kubernetes"`), []byte(`"driver":"opensandbox"`), 1)
if _, err := Unmarshal(tampered, "ddupan.top"); err == nil { if _, err := Unmarshal(tampered, "ddupan.top"); err == nil {
t.Fatal("expected tampered backend to fail") t.Fatal("expected tampered placement to fail")
} }
} }
@@ -79,13 +87,14 @@ func TestFromMetadataRecoversMinimalAssignment(t *testing.T) {
assignment, err := FromMetadata(map[string]string{ assignment, err := FromMetadata(map[string]string{
"ci.ddupan.top/assignment-id": "gitea-task-42", "ci.ddupan.top/assignment-id": "gitea-task-42",
"ci.ddupan.top/task-id": "42", "ci.ddupan.top/task-id": "42",
"ci.ddupan.top/backend": "vm", "ci.ddupan.top/workload-class": "vm",
"ci.ddupan.top/driver": "opensandbox",
}, map[string]string{ }, map[string]string{
"ci.ddupan.top/repository": "owner/repo", "ci.ddupan.top/repository": "owner/repo",
"ci.ddupan.top/job-key": "publish", "ci.ddupan.top/job-key": "publish",
"ci.ddupan.top/spiffe-id": "spiffe://ddupan.top/ci/owner/repo/publish", "ci.ddupan.top/spiffe-id": "spiffe://ddupan.top/ci/owner/repo/publish",
}, "ddupan.top") }, "ddupan.top")
if err != nil || assignment.Task.GetId() != 42 || assignment.Backend != BackendVM { if err != nil || assignment.Task.GetId() != 42 || assignment.Placement != OpenSandboxVM {
t.Fatalf("assignment=%#v err=%v", assignment, err) t.Fatalf("assignment=%#v err=%v", assignment, err)
} }
} }
+76
View File
@@ -0,0 +1,76 @@
package taskassignment
import (
"errors"
"fmt"
"slices"
)
type WorkloadClass string
const (
WorkloadContainer WorkloadClass = "container"
WorkloadVM WorkloadClass = "vm"
)
type Driver string
const (
DriverKubernetes Driver = "kubernetes"
DriverOpenSandbox Driver = "opensandbox"
)
type Placement struct {
Class WorkloadClass `json:"workload_class"`
Driver Driver `json:"driver"`
}
var (
KubernetesContainer = Placement{Class: WorkloadContainer, Driver: DriverKubernetes}
OpenSandboxVM = Placement{Class: WorkloadVM, Driver: DriverOpenSandbox}
)
func (p Placement) Validate() error {
switch p {
case KubernetesContainer, OpenSandboxVM:
return nil
default:
return fmt.Errorf("unsupported workload placement %s/%s", p.Class, p.Driver)
}
}
func (p Placement) Key() string { return string(p.Class) + "." + string(p.Driver) }
func PlacementFromLabels(labels []string) (Placement, error) {
if !slices.Contains(labels, "self-hosted") {
return Placement{}, fmt.Errorf("task runs-on labels must include self-hosted: %v", labels)
}
legacyPod := slices.Contains(labels, "pod")
container := slices.Contains(labels, string(WorkloadContainer))
vm := slices.Contains(labels, string(WorkloadVM)) || slices.Contains(labels, "vm-dev")
kubernetes := slices.Contains(labels, string(DriverKubernetes))
opensandbox := slices.Contains(labels, string(DriverOpenSandbox))
if legacyPod {
if container || vm || kubernetes || opensandbox {
return Placement{}, errors.New("legacy pod label cannot be combined with workload or VM driver labels")
}
return KubernetesContainer, nil
}
if container == vm {
return Placement{}, fmt.Errorf("task runs-on labels must select exactly one workload class: %v", labels)
}
if kubernetes && opensandbox {
return Placement{}, fmt.Errorf("task runs-on labels select multiple drivers: %v", labels)
}
if container {
if opensandbox {
return Placement{}, fmt.Errorf("opensandbox does not support container workloads")
}
return KubernetesContainer, nil
}
if kubernetes {
return Placement{}, fmt.Errorf("kubernetes does not support VM workloads")
}
return OpenSandboxVM, nil
}
+2 -1
View File
@@ -209,7 +209,8 @@ func BackendMetadata(assignment taskassignment.Assignment) Metadata {
"ci.ddupan.top/runner": "true", "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/workload-class": string(assignment.Placement.Class),
"ci.ddupan.top/driver": string(assignment.Placement.Driver),
}, },
Annotations: map[string]string{ Annotations: map[string]string{
"ci.ddupan.top/repository": assignment.Identity.Repository, "ci.ddupan.top/repository": assignment.Identity.Repository,
+1 -1
View File
@@ -54,7 +54,7 @@ func (t *fakeTasks) Report(_ context.Context, _ int64, phase Phase) error {
func assignment() taskassignment.Assignment { func assignment() taskassignment.Assignment {
return taskassignment.Assignment{ return taskassignment.Assignment{
ID: "gitea-task-42", ID: "gitea-task-42",
Backend: taskassignment.BackendPod, Placement: taskassignment.KubernetesContainer,
Task: &runnerv1.Task{Id: 42}, Task: &runnerv1.Task{Id: 42},
Identity: taskidentity.Identity{ Identity: taskidentity.Identity{
Repository: "owner/repo", Repository: "owner/repo",
+1
View File
@@ -10,6 +10,7 @@ while [ "$(date +%s)" -lt "$deadline" ]; do
-audience ci-job-ready \ -audience ci-job-ready \
-socketPath "$socket" \ -socketPath "$socket" \
>/dev/null 2>&1; then >/dev/null 2>&1; then
/usr/local/libexec/setup-job-docker
exit 0 exit 0
fi fi
sleep 1 sleep 1
+68
View File
@@ -0,0 +1,68 @@
#!/usr/bin/env bash
set -euo pipefail
set +x
: "${GITHUB_SHA:?GITHUB_SHA is required}"
: "${IMAGE_NAME:?IMAGE_NAME is required}"
: "${IMAGE_REPOSITORY:?IMAGE_REPOSITORY is required}"
: "${IMAGE_DOCKERFILE:?IMAGE_DOCKERFILE is required}"
: "${PUSH_REGISTRY:?PUSH_REGISTRY is required}"
: "${PULL_REGISTRY:?PULL_REGISTRY is required}"
: "${SPIRE_AGENT_SOCKET:?SPIRE_AGENT_SOCKET is required}"
# Bootstrap the image that first introduces automatic Docker setup. Once that
# runner is deployed, the job-started hook makes this an idempotent no-op.
sudo scripts/setup-job-docker
image_tag="sha-${GITHUB_SHA}"
metadata="${IMAGE_NAME}-metadata.json"
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
docker buildx build \
--builder ci-builder \
--platform linux/amd64 \
--file "$IMAGE_DOCKERFILE" \
--tag "${PUSH_REGISTRY}/${IMAGE_REPOSITORY}:${image_tag}" \
--tag "${PUSH_REGISTRY}/${IMAGE_REPOSITORY}:main" \
--provenance=mode=max \
--sbom=true \
--metadata-file "$metadata" \
--push \
.
image_digest=$(
# shellcheck disable=SC2016
jq -er '."containerimage.digest"' "$metadata"
)
image_ref="${PULL_REGISTRY}/${IMAGE_REPOSITORY}@${image_digest}"
printf '%s=%s\n' "$IMAGE_NAME" "$image_ref"
if [[ -n "${GITHUB_STEP_SUMMARY:-}" ]]; then
{
printf '## Published image\n\n'
# shellcheck disable=SC2016
printf -- '- %s: `%s`\n' "$IMAGE_NAME" "$image_ref"
# shellcheck disable=SC2016
printf -- '- Source: `%s`\n' "$GITHUB_SHA"
} >>"$GITHUB_STEP_SUMMARY"
fi
+60
View File
@@ -0,0 +1,60 @@
#!/usr/bin/env bash
set -euo pipefail
if docker info >/dev/null 2>&1; then
printf '%s\n' 'job Docker daemon is already ready'
exit 0
fi
storage_size=${DOCKER_DATA_SIZE:-20G}
storage_driver=${DOCKER_STORAGE_DRIVER:-overlay2}
wait_seconds=${DOCKER_START_WAIT_SECONDS:-60}
sudo install -d /var/lib/docker
if ! mountpoint --quiet /var/lib/docker; then
printf '%s\n' "preparing ${storage_size} loop-backed Docker storage"
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
sudo truncate -s "$storage_size" /tmp/docker-data.img
sudo mkfs.ext4 -F /tmp/docker-data.img
loop_device=$(sudo losetup --find --show /tmp/docker-data.img)
sudo mount "$loop_device" /var/lib/docker
else
printf '%s\n' 'using mounted Docker storage at /var/lib/docker'
fi
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 sh -c 'nohup dockerd "$@" </dev/null >/tmp/dockerd.log 2>&1 &' sh \
--host=unix:///var/run/docker.sock \
--storage-driver="$storage_driver"
printf '%s\n' 'waiting for job Docker daemon'
for ((attempt = 0; attempt < wait_seconds; attempt++)); do
if docker info >/dev/null 2>&1; then
findmnt /var/lib/docker
docker info --format \
'{{json .ServerVersion}} {{json .Driver}} {{json .CgroupVersion}}'
exit 0
fi
sleep 1
done
cat /tmp/dockerd.log
exit 1