Compare commits

..
Author SHA1 Message Date
panxiao81 87cf359ef6 fix: 标记 OpenSandbox Runner workload
test / shell (pull_request) Successful in 26s
test / go (pull_request) Successful in 2m58s
test / python (pull_request) Failing after 14m15s
2026-09-21 08:22:56 +00:00
panxiao81 ad3934e7e1 Merge OpenSandbox metadata 编码修复
test / python (push) Successful in 21s
test / shell (push) Successful in 24s
test / go (push) Successful in 6m52s
publish images / publish-images (push) Failing after 10m40s
2026-09-21 07:56:50 +00:00
panxiao81 d82687382a fix: 编码 OpenSandbox assignment metadata
test / go (pull_request) Successful in 4m45s
test / shell (pull_request) Failing after 13m27s
test / python (pull_request) Failing after 13m27s
2026-09-21 07:48:10 +00:00
panxiao81 7f52cf393f Merge Pod 与 VM 独立执行容量池
publish images / publish-images (push) Failing after 11m8s
test / python (push) Successful in 10s
test / shell (push) Successful in 26s
test / go (push) Successful in 5m51s
2026-09-21 07:26:10 +00:00
panxiao81 fadc93a0bf feat: 为执行后端增加独立容量池
test / python (pull_request) Successful in 13s
test / shell (pull_request) Successful in 20s
test / go (pull_request) Successful in 5m54s
2026-09-21 07:21:38 +00:00
panxiao81 8a186dcd86 Merge Runner 无中断恢复与容量隔离
test / python (push) Successful in 25s
test / shell (push) Successful in 29s
test / go (push) Successful in 2m28s
publish images / publish-images (push) Failing after 11m13s
2026-09-21 06:48:11 +00:00
panxiao81 d06845bc3c fix: 隔离 Pod 与 VM 并发容量
test / python (pull_request) Successful in 20s
test / shell (pull_request) Successful in 23s
test / go (pull_request) Successful in 3m46s
2026-09-21 06:44:34 +00:00
panxiao81 2a23f0c62e feat: 通过 Lease 单例调度任务
test / python (pull_request) Successful in 31s
test / shell (pull_request) Successful in 34s
test / go (pull_request) Successful in 3m38s
2026-09-21 06:35:06 +00:00
panxiao81 5e94182308 fix: 从后端恢复 Runner claim
test / python (pull_request) Successful in 9s
test / shell (pull_request) Successful in 15s
test / go (pull_request) Successful in 2m18s
2026-09-21 06:18:23 +00:00
panxiao81 327c73e744 Merge Go Runner 与 OpenSandbox VM 集成
test / python (push) Successful in 15s
test / shell (push) Successful in 21s
publish images / publish-images (push) Failing after 13m58s
test / go (push) Failing after 15m46s
2026-09-21 06:11:40 +00:00
22 changed files with 565 additions and 106 deletions
+5 -3
View File
@@ -18,7 +18,8 @@ runs-on: [self-hosted, vm]
目标 Go controller 组件: 目标 Go controller 组件:
- `scheduler`:以常驻 Gitea RunnerService 身份直接领取 task,并把完整 assignment - `scheduler`:以常驻 Gitea RunnerService 身份直接领取 task,并把完整 assignment
持久化到 JetStream;同时提供仅允许 SPIFFE mTLS 的 RunnerService facade。 持久化到 JetStream;一个 registration 下按配置启动多个并发 `FetchTask` goroutine,
同时提供仅允许 SPIFFE mTLS 的 RunnerService facade。
- `pod-worker`:直接在 homelab Kubernetes 创建一次性 Pod。 - `pod-worker`:直接在 homelab Kubernetes 创建一次性 Pod。
- `vm-worker`:通过 OpenSandbox Lifecycle API 从 `ci-vm` Pool 创建 Kata microVM。 - `vm-worker`:通过 OpenSandbox Lifecycle API 从 `ci-vm` Pool 创建 Kata microVM。
- 三个组件默认在同一个 Go 进程启用。首轮集成期间不允许只启动 worker,因为 facade - 三个组件默认在同一个 Go 进程启用。首轮集成期间不允许只启动 worker,因为 facade
@@ -37,8 +38,9 @@ runs-on: [self-hosted, vm]
动态 Pod 或 VM 直接取得自己的 SPIFFE 身份。 动态 Pod 或 VM 直接取得自己的 SPIFFE 身份。
Pod 路径由 homelab 集群中的 `pod-worker` 直接创建 Kubernetes Pod。OpenSandbox 只用于 Pod 路径由 homelab 集群中的 `pod-worker` 直接创建 Kubernetes Pod。OpenSandbox 只用于
VM/Kata workload;两个 backend 使用独立 durable consumer,任一执行层故障不会阻塞另一条 VM/Kata workload;两个 backend 使用独立 durable consumer 和独立容量池。assignment 根据
部署。长期 RunnerService 协议路线见 `runs-on` 进入对应池,池满时留在 JetStream pending,不会创建超出容量的 workload;任一
执行层故障不会阻塞另一条部署。长期 RunnerService 协议路线见
[`docs/runner-protocol-roadmap.md`](docs/runner-protocol-roadmap.md)。 [`docs/runner-protocol-roadmap.md`](docs/runner-protocol-roadmap.md)。
## 开发 ## 开发
+95 -9
View File
@@ -15,8 +15,14 @@ import (
"github.com/nats-io/nats.go" "github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream" "github.com/nats-io/nats.go/jetstream"
"golang.org/x/sync/errgroup" "golang.org/x/sync/errgroup"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/leaderelection"
"k8s.io/client-go/tools/leaderelection/resourcelock"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/assignmentqueue" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/assignmentqueue"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/backendpool"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/controller" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/controller"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/giteaactions" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/giteaactions"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/opensandboxbackend" "git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/opensandboxbackend"
@@ -86,7 +92,8 @@ func runController(ctx context.Context) error {
return err return err
} }
registry := runnerfacade.NewRegistry() registry := runnerfacade.NewRegistry()
gate := taskscheduler.NewSingleFlightGate() podPool := backendpool.New(config.PodCapacity)
vmPool := backendpool.New(config.VMCapacity)
var podExecutorBackend *podbackend.Backend var podExecutorBackend *podbackend.Backend
var vmExecutorBackend *opensandboxbackend.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)
@@ -101,6 +108,7 @@ func runController(ctx context.Context) error {
if err := podExecutorBackend.MarkTerminal(ctx, assignment.ID); err != nil { if err := podExecutorBackend.MarkTerminal(ctx, assignment.ID); err != nil {
return err return err
} }
podPool.Release(assignment.ID)
case taskassignment.BackendVM: case taskassignment.BackendVM:
if vmExecutorBackend == nil { if vmExecutorBackend == nil {
return errors.New("VM lifecycle backend is not configured") return errors.New("VM lifecycle backend is not configured")
@@ -108,8 +116,8 @@ func runController(ctx context.Context) error {
if err := vmExecutorBackend.MarkTerminal(ctx, assignment.ID); err != nil { if err := vmExecutorBackend.MarkTerminal(ctx, assignment.ID); err != nil {
return err return err
} }
vmPool.Release(assignment.ID)
} }
gate.Release()
return nil return nil
}, },
} }
@@ -126,13 +134,25 @@ func runController(ctx context.Context) error {
labels = append(labels, string(taskassignment.BackendVM)) labels = append(labels, string(taskassignment.BackendVM))
} }
poller := taskscheduler.Poller{ poller := taskscheduler.Poller{
Client: giteaClient, Gate: gate, Client: giteaClient,
Scheduler: &taskscheduler.Scheduler{TrustDomain: config.TrustDomain, Dispatcher: assignmentqueue.Publisher{ Scheduler: &taskscheduler.Scheduler{TrustDomain: config.TrustDomain, Dispatcher: assignmentqueue.Publisher{
JetStream: producerJS, SubjectBase: config.SubjectBase, JetStream: producerJS, SubjectBase: config.SubjectBase,
}}, }},
Config: taskscheduler.PollerConfig{Version: "gitea-dynamic-runner/0.4", Labels: labels}, Config: taskscheduler.PollerConfig{Version: "gitea-dynamic-runner/0.4", Labels: labels, Capacity: config.PodCapacity + config.VMCapacity},
OnError: func(err error) { log.Printf("scheduler: %v", err) }, OnError: func(err error) { log.Printf("scheduler: %v", err) },
} }
kubernetesConfig, err := rest.InClusterConfig()
if err != nil {
return fmt.Errorf("load leader election Kubernetes config: %w", err)
}
kubernetesClient, err := kubernetes.NewForConfig(kubernetesConfig)
if err != nil {
return fmt.Errorf("create leader election Kubernetes client: %w", err)
}
leaderIdentity := strings.TrimSpace(os.Getenv("HOSTNAME"))
if leaderIdentity == "" {
return errors.New("HOSTNAME is required for scheduler leader election")
}
facadeServer := runnerfacade.Server{ facadeServer := runnerfacade.Server{
Facade: facade, ListenAddress: config.FacadeListen, TrustDomain: config.TrustDomain, Facade: facade, ListenAddress: config.FacadeListen, TrustDomain: config.TrustDomain,
WorkloadAPIAddr: config.WorkloadAPIAddr, UpstreamURL: config.GiteaURL, WorkloadAPIAddr: config.WorkloadAPIAddr, UpstreamURL: config.GiteaURL,
@@ -141,7 +161,9 @@ func runController(ctx context.Context) error {
controller.Scheduler: runComponent(func(ctx context.Context) error { controller.Scheduler: runComponent(func(ctx context.Context) error {
group, groupContext := errgroup.WithContext(ctx) group, groupContext := errgroup.WithContext(ctx)
group.Go(func() error { return facadeServer.Run(groupContext) }) group.Go(func() error { return facadeServer.Run(groupContext) })
group.Go(func() error { return poller.Run(groupContext) }) group.Go(func() error {
return runSchedulerLeader(groupContext, kubernetesClient, config.PodNamespace, leaderIdentity, poller.Run)
})
return group.Wait() return group.Wait()
}), }),
} }
@@ -158,7 +180,17 @@ func runController(ctx context.Context) error {
SPIREAgentID: config.SPIREAgentID, ExecutorUID: config.PodExecutorUID, SPIREAgentID: config.SPIREAgentID, ExecutorUID: config.PodExecutorUID,
}} }}
podExecutorBackend = &backend podExecutorBackend = &backend
component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendPod, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry) assignments, err := backend.RecoverAssignments(ctx)
if err != nil {
return err
}
for _, assignment := range assignments {
podPool.Restore(assignment.ID)
if err := registry.RecoverClaimed(assignment); err != nil {
return fmt.Errorf("recover Pod facade claim %s: %w", assignment.ID, err)
}
}
component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendPod, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry, podPool)
if err != nil { if err != nil {
return err return err
} }
@@ -178,7 +210,17 @@ func runController(ctx context.Context) error {
Env: map[string]string{"SPIFFE_ENDPOINT_SOCKET": config.WorkloadAPIAddr}, Env: map[string]string{"SPIFFE_ENDPOINT_SOCKET": config.WorkloadAPIAddr},
}} }}
vmExecutorBackend = &backend vmExecutorBackend = &backend
component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendVM, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry) assignments, err := backend.RecoverAssignments(ctx, config.TrustDomain)
if err != nil {
return err
}
for _, assignment := range assignments {
vmPool.Restore(assignment.ID)
if err := registry.RecoverClaimed(assignment); err != nil {
return fmt.Errorf("recover VM facade claim %s: %w", assignment.ID, err)
}
}
component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendVM, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry, vmPool)
if err != nil { if err != nil {
return err return err
} }
@@ -193,6 +235,50 @@ func runController(ctx context.Context) error {
return controller.Run(ctx, config.Components, components) return controller.Run(ctx, config.Components, components)
} }
func runSchedulerLeader(ctx context.Context, client kubernetes.Interface, namespace, identity string, run func(context.Context) error) error {
if client == nil || namespace == "" || identity == "" || run == nil {
return errors.New("leader election client, namespace, identity, and scheduler are required")
}
electionContext, cancel := context.WithCancel(ctx)
defer cancel()
result := make(chan error, 1)
lock := &resourcelock.LeaseLock{
LeaseMeta: metav1.ObjectMeta{Name: "dynamic-runner-scheduler", Namespace: namespace},
Client: client.CoordinationV1(),
LockConfig: resourcelock.ResourceLockConfig{
Identity: identity,
},
}
elector, err := leaderelection.NewLeaderElector(leaderelection.LeaderElectionConfig{
Lock: lock, LeaseDuration: 15 * time.Second, RenewDeadline: 10 * time.Second, RetryPeriod: 2 * time.Second,
ReleaseOnCancel: true,
Callbacks: leaderelection.LeaderCallbacks{
OnStartedLeading: func(leaderContext context.Context) {
result <- run(leaderContext)
cancel()
},
OnStoppedLeading: func() {
if ctx.Err() == nil {
select {
case result <- errors.New("scheduler leadership lost"):
default:
}
}
},
},
})
if err != nil {
return fmt.Errorf("configure scheduler leader election: %w", err)
}
go elector.Run(electionContext)
select {
case err := <-result:
return err
case <-ctx.Done():
return nil
}
}
func connectNATS(server, user, password, caFile, clientName string) (*nats.Conn, error) { func connectNATS(server, user, password, caFile, clientName string) (*nats.Conn, error) {
options := []nats.Option{nats.Name(clientName), nats.UserInfo(user, password)} options := []nats.Option{nats.Name(clientName), nats.UserInfo(user, password)}
if caFile != "" { if caFile != "" {
@@ -201,14 +287,14 @@ 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) (controller.Component, error) { 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) {
consumer, err := assignmentqueue.OpenConsumer(ctx, js, config.Stream, config.SubjectBase, backend, capacity) consumer, err := assignmentqueue.OpenConsumer(ctx, js, config.Stream, config.SubjectBase, backend, 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}, Processor: assignmentqueue.Processor{TrustDomain: config.TrustDomain, Accepter: accepter, Claims: claims, Admission: admission},
OnError: func(err error) { log.Printf("%s worker: %v", backend, err) }, OnError: func(err error) { log.Printf("%s worker: %v", backend, err) },
}, nil }, nil
} }
+13 -7
View File
@@ -3,8 +3,9 @@
## 目标 ## 目标
长期形态不依赖 `workflow_job` webhook 发现工作。controller 本身作为 Gitea Runner 长期形态不依赖 `workflow_job` webhook 发现工作。controller 本身作为 Gitea Runner
协议客户端注册,并声明 `self-hosted`、`pod` 和 `vm` labels;它只在后端存在可用容量 协议客户端注册,并声明 `self-hosted`、`pod` 和 `vm` labels;单个 registration 内按总
时领取 task,然后将该 task 交给一个一次性 Pod 或 microVM 执行。 配置容量启动多个 `FetchTask` goroutine,再将 task 按 `runs-on` 交给 Pod 或 VM 的独立
容量池,由一次性 Pod 或 microVM 执行。
```text ```text
Gitea RunnerService Gitea RunnerService
@@ -47,7 +48,10 @@ facade,并严格校验 facade 的 SPIFFE ID。这样无需修改 runner 或把
## 设计约束 ## 设计约束
- 对 workflow 的接口保持 `[self-hosted, pod]` 和 `[self-hosted, vm]` 不变。 - 对 workflow 的接口保持 `[self-hosted, pod]` 和 `[self-hosted, vm]` 不变。
- scheduler 在没有对应 backend 容量时不领取 task,避免本地形成不可控积压。 - scheduler 使用单一 Gitea runner UUID/token 和一个 `Declare`,不为并发槽位重复注册;
`POD_CAPACITY + VM_CAPACITY` 决定并发 `FetchTask` goroutine 数量。
- task 领取并持久化后按 backend 进入独立 durable consumer;对应容量池已满时延迟 NAK,
assignment 保持 JetStream pending,且不得创建超出配置容量的 workload。
- scheduler Declare 后使用 RunnerService 长轮询;一旦 FetchTask 返回已分配 task,在 - scheduler Declare 后使用 RunnerService 长轮询;一旦 FetchTask 返回已分配 task,在
JetStream publish 成功前只重试该 assignment,不领取下一项。 JetStream publish 成功前只重试该 assignment,不领取下一项。
- 每个 executor 只执行一个 task,完成后销毁。 - 每个 executor 只执行一个 task,完成后销毁。
@@ -72,17 +76,19 @@ 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 和并发上限。consumer 只负责将 assignment - pod 与 vm 使用独立 durable consumer、进程内 admission pool 和并发上限。consumer 只负责将 assignment
幂等落到后端;executor 与身份恢复 metadata 持久化后立即 `DoubleAck`。尚未取得 幂等落到后端;executor 与身份恢复 metadata 持久化后立即 `DoubleAck`。尚未取得
Pod UID 等短暂未就绪状态以及临时后端错误使用延迟 NAK。 Pod UID 等短暂未就绪状态以及临时后端错误使用延迟 NAK。
- admission pool 只保存可重建的并发状态:启动时从 Pod labels/annotations 或 OpenSandbox
metadata 恢复非终态 assignment,terminal update 持久化成功后释放槽位,不引入新存储。
- assignment ACK 后的运行、结果回报和清理由 backend reconciler 根据 Kubernetes、 - assignment ACK 后的运行、结果回报和清理由 backend reconciler 根据 Kubernetes、
OpenSandbox 与 Gitea 的事实状态驱动,不继续占用 JetStream delivery。 OpenSandbox 与 Gitea 的事实状态驱动,不继续占用 JetStream delivery。
- 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 协议,只有环境创建和销毁实现不同。
- 首次生产 canary 使用 controller 进程内 single-flight gate:只有官方 runner 的终态 - scheduler 在 assignment 持久化到 JetStream 后即可继续领取;Pod 与 VM 分别由 durable
`UpdateTask` 已被 Gitea 接受后才允许 Fetch 下一条任务。它把未知故障收敛为停止领取, consumer 的 capacity 限制并发,不共享全局执行槽位。未知后端故障由对应 consumer 的
而不是在 backlog 下连续创建 executor;后续容量调度必须以 backend 实际运行资源为准。 NAK/redelivery 收敛,不能阻塞另一种 backend。
- 两种 backend 都注入同一份 runner bootstrap 环境;Pod 仍由 homelab Kubernetes 原生 - 两种 backend 都注入同一份 runner bootstrap 环境;Pod 仍由 homelab Kubernetes 原生
创建,只有 VM 经 OpenSandbox 创建,bootstrap 机制不改变 backend 边界。 创建,只有 VM 经 OpenSandbox 创建,bootstrap 机制不改变 backend 边界。
+1
View File
@@ -45,6 +45,7 @@ require (
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/nats-io/nkeys v0.4.16 // indirect github.com/nats-io/nkeys v0.4.16 // indirect
github.com/nats-io/nuid v1.0.1 // indirect github.com/nats-io/nuid v1.0.1 // indirect
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect
github.com/sirupsen/logrus v1.10.2 // indirect github.com/sirupsen/logrus v1.10.2 // indirect
github.com/x448/float16 v0.8.4 // indirect github.com/x448/float16 v0.8.4 // indirect
go.yaml.in/yaml/v2 v2.4.4 // indirect go.yaml.in/yaml/v2 v2.4.4 // indirect
+2
View File
@@ -112,6 +112,8 @@ go.opentelemetry.io/otel/sdk/metric v1.39.0 h1:cXMVVFVgsIf2YL6QkRF4Urbr/aMInf+2W
go.opentelemetry.io/otel/sdk/metric v1.39.0/go.mod h1:xq9HEVH7qeX69/JnwEfp6fVq5wosJsY1mt4lLfYdVew= go.opentelemetry.io/otel/sdk/metric v1.39.0/go.mod h1:xq9HEVH7qeX69/JnwEfp6fVq5wosJsY1mt4lLfYdVew=
go.opentelemetry.io/otel/trace v1.39.0 h1:2d2vfpEDmCJ5zVYz7ijaJdOF59xLomrvj7bjt6/qCJI= go.opentelemetry.io/otel/trace v1.39.0 h1:2d2vfpEDmCJ5zVYz7ijaJdOF59xLomrvj7bjt6/qCJI=
go.opentelemetry.io/otel/trace v1.39.0/go.mod h1:88w4/PnZSazkGzz/w84VHpQafiU4EtqqlVdxWy+rNOA= go.opentelemetry.io/otel/trace v1.39.0/go.mod h1:88w4/PnZSazkGzz/w84VHpQafiU4EtqqlVdxWy+rNOA=
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ= go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ=
go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ= go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ=
go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw= go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw=
+16 -2
View File
@@ -57,6 +57,11 @@ type Claims interface {
WaitClaimed(context.Context, string) error WaitClaimed(context.Context, string) error
} }
type Admission interface {
Acquire(string) bool
Release(string)
}
// Message is the subset of jetstream.Msg needed by one reconciliation. // Message is the subset of jetstream.Msg needed by one reconciliation.
type Message interface { type Message interface {
Data() []byte Data() []byte
@@ -70,13 +75,14 @@ type Processor struct {
TrustDomain string TrustDomain string
Accepter Accepter Accepter Accepter
Claims Claims Claims Claims
Admission Admission
RetryDelay time.Duration RetryDelay time.Duration
ClaimTimeout time.Duration ClaimTimeout time.Duration
} }
func (p Processor) Process(ctx context.Context, message Message) error { func (p Processor) Process(ctx context.Context, message Message) error {
if p.Accepter == nil || p.Claims == nil { if p.Accepter == nil || p.Claims == nil || p.Admission == nil {
return errors.New("assignment accepter and claim registry are required") return errors.New("assignment accepter, claim registry, and backend admission pool are required")
} }
assignment, err := taskassignment.Unmarshal(message.Data(), p.TrustDomain) assignment, err := taskassignment.Unmarshal(message.Data(), p.TrustDomain)
if err != nil { if err != nil {
@@ -85,8 +91,16 @@ func (p Processor) Process(ctx context.Context, message Message) error {
if _, err := p.Claims.Offer(assignment); err != nil { if _, err := p.Claims.Offer(assignment); err != nil {
return errors.Join(err, message.TermWithReason("conflicting assignment")) return errors.Join(err, message.TermWithReason("conflicting assignment"))
} }
if !p.Admission.Acquire(assignment.ID) {
delay := p.RetryDelay
if delay <= 0 {
delay = 2 * time.Second
}
return message.NakWithDelay(delay)
}
accepted, err := p.Accepter.Accept(ctx, assignment) accepted, err := p.Accepter.Accept(ctx, assignment)
if err != nil { if err != nil {
p.Admission.Release(assignment.ID)
delay := p.RetryDelay delay := p.RetryDelay
if delay <= 0 { if delay <= 0 {
delay = 15 * time.Second delay = 15 * time.Second
+47 -2
View File
@@ -61,6 +61,28 @@ type fakeClaims struct {
claimed bool claimed bool
} }
type fakeAdmission struct {
allowed bool
active map[string]bool
released int
}
func (a *fakeAdmission) Acquire(assignmentID string) bool {
if !a.allowed {
return false
}
if a.active == nil {
a.active = make(map[string]bool)
}
a.active[assignmentID] = true
return true
}
func (a *fakeAdmission) Release(assignmentID string) {
delete(a.active, assignmentID)
a.released++
}
func (c *fakeClaims) Offer(taskassignment.Assignment) (<-chan struct{}, error) { func (c *fakeClaims) Offer(taskassignment.Assignment) (<-chan struct{}, error) {
ready := make(chan struct{}) ready := make(chan struct{})
if c.claimed { if c.claimed {
@@ -104,7 +126,7 @@ func encodedAssignment(t *testing.T) []byte {
func TestProcessorAcknowledgesPersistedHandoff(t *testing.T) { func TestProcessorAcknowledgesPersistedHandoff(t *testing.T) {
message := &fakeMessage{data: encodedAssignment(t)} message := &fakeMessage{data: encodedAssignment(t)}
processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{accepted: true}, Claims: &fakeClaims{claimed: true}} processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{accepted: true}, Claims: &fakeClaims{claimed: true}, Admission: &fakeAdmission{allowed: true}}
if err := processor.Process(context.Background(), message); err != nil { if err := processor.Process(context.Background(), message); err != nil {
t.Fatal(err) t.Fatal(err)
} }
@@ -115,7 +137,7 @@ func TestProcessorAcknowledgesPersistedHandoff(t *testing.T) {
func TestProcessorRetriesUntilBackendHandoffIsDurable(t *testing.T) { func TestProcessorRetriesUntilBackendHandoffIsDurable(t *testing.T) {
message := &fakeMessage{data: encodedAssignment(t)} message := &fakeMessage{data: encodedAssignment(t)}
processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{}, Claims: &fakeClaims{}, RetryDelay: 2 * time.Second} processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{}, Claims: &fakeClaims{}, Admission: &fakeAdmission{allowed: true}, RetryDelay: 2 * time.Second}
if err := processor.Process(context.Background(), message); err != nil { if err := processor.Process(context.Background(), message); err != nil {
t.Fatal(err) t.Fatal(err)
} }
@@ -126,10 +148,12 @@ func TestProcessorRetriesUntilBackendHandoffIsDurable(t *testing.T) {
func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T) { func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T) {
retry := &fakeMessage{data: encodedAssignment(t)} retry := &fakeMessage{data: encodedAssignment(t)}
admission := &fakeAdmission{allowed: true}
processor := Processor{ processor := Processor{
TrustDomain: "ddupan.top", TrustDomain: "ddupan.top",
Accepter: &fakeAccepter{err: errors.New("backend unavailable")}, Accepter: &fakeAccepter{err: errors.New("backend unavailable")},
Claims: &fakeClaims{}, Claims: &fakeClaims{},
Admission: admission,
RetryDelay: time.Minute, RetryDelay: time.Minute,
} }
if err := processor.Process(context.Background(), retry); err == nil { if err := processor.Process(context.Background(), retry); err == nil {
@@ -138,6 +162,9 @@ func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T)
if retry.nacked != time.Minute { if retry.nacked != time.Minute {
t.Fatalf("retry delay = %s", retry.nacked) t.Fatalf("retry delay = %s", retry.nacked)
} }
if admission.released != 1 {
t.Fatalf("released slots = %d", admission.released)
}
poison := &fakeMessage{data: []byte("not-json")} poison := &fakeMessage{data: []byte("not-json")}
if err := processor.Process(context.Background(), poison); err == nil { if err := processor.Process(context.Background(), poison); err == nil {
@@ -148,6 +175,24 @@ func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T)
} }
} }
func TestProcessorLeavesAssignmentPendingWhenBackendPoolIsFull(t *testing.T) {
message := &fakeMessage{data: encodedAssignment(t)}
accepter := &fakeAccepter{accepted: true}
processor := Processor{
TrustDomain: "ddupan.top",
Accepter: accepter,
Claims: &fakeClaims{},
Admission: &fakeAdmission{},
RetryDelay: 3 * time.Second,
}
if err := processor.Process(context.Background(), message); err != nil {
t.Fatal(err)
}
if message.acked != 0 || message.nacked != 3*time.Second {
t.Fatalf("message = %#v", message)
}
}
type fakeConsumerManager struct{ config jetstream.ConsumerConfig } type fakeConsumerManager struct{ config jetstream.ConsumerConfig }
func (m *fakeConsumerManager) CreateOrUpdateConsumer(_ context.Context, _ string, config jetstream.ConsumerConfig) (jetstream.Consumer, error) { func (m *fakeConsumerManager) CreateOrUpdateConsumer(_ context.Context, _ string, config jetstream.ConsumerConfig) (jetstream.Consumer, error) {
+50
View File
@@ -0,0 +1,50 @@
// Package backendpool manages runtime capacity independently for each executor backend.
package backendpool
import "sync"
type Pool struct {
mu sync.Mutex
capacity int
active map[string]struct{}
}
func New(capacity int) *Pool {
return &Pool{capacity: capacity, active: make(map[string]struct{})}
}
// Acquire reserves a backend slot without blocking. Redelivery of the same
// assignment is idempotent and succeeds even while the pool is full.
func (p *Pool) Acquire(assignmentID string) bool {
p.mu.Lock()
defer p.mu.Unlock()
if _, exists := p.active[assignmentID]; exists {
return true
}
if assignmentID == "" || len(p.active) >= p.capacity {
return false
}
p.active[assignmentID] = struct{}{}
return true
}
func (p *Pool) Restore(assignmentID string) {
if assignmentID == "" {
return
}
p.mu.Lock()
p.active[assignmentID] = struct{}{}
p.mu.Unlock()
}
func (p *Pool) Release(assignmentID string) {
p.mu.Lock()
delete(p.active, assignmentID)
p.mu.Unlock()
}
func (p *Pool) Active() int {
p.mu.Lock()
defer p.mu.Unlock()
return len(p.active)
}
+26
View File
@@ -0,0 +1,26 @@
package backendpool
import "testing"
func TestPoolSeparatesRuntimeCapacityFromDeliveries(t *testing.T) {
pool := New(2)
if !pool.Acquire("one") || !pool.Acquire("two") || pool.Acquire("three") {
t.Fatal("capacity was not enforced")
}
if !pool.Acquire("one") {
t.Fatal("redelivery must be idempotent")
}
pool.Release("one")
if !pool.Acquire("three") || pool.Active() != 2 {
t.Fatalf("active=%d", pool.Active())
}
}
func TestRestoreMayTemporarilyExceedReducedCapacity(t *testing.T) {
pool := New(1)
pool.Restore("one")
pool.Restore("two")
if pool.Active() != 2 || pool.Acquire("three") {
t.Fatalf("active=%d", pool.Active())
}
}
+92 -3
View File
@@ -3,9 +3,13 @@ package opensandboxbackend
import ( import (
"context" "context"
"encoding/base64"
"encoding/json"
"errors" "errors"
"fmt" "fmt"
"net/http" "net/http"
"sort"
"strings"
"time" "time"
opensandbox "github.com/alibaba/OpenSandbox/sdks/sandbox/go" opensandbox "github.com/alibaba/OpenSandbox/sdks/sandbox/go"
@@ -17,6 +21,8 @@ import (
const assignmentMetadata = "ci.ddupan.top/assignment-id" const assignmentMetadata = "ci.ddupan.top/assignment-id"
const terminalMetadata = "ci.ddupan.top/terminal" const terminalMetadata = "ci.ddupan.top/terminal"
const annotationsMetadataPrefix = "ci.ddupan.top/annotations-"
const metadataValueLimit = 63
type Lifecycle interface { type Lifecycle interface {
ListSandboxes(context.Context, opensandbox.ListOptions) (*opensandbox.ListSandboxesResponse, error) ListSandboxes(context.Context, opensandbox.ListOptions) (*opensandbox.ListSandboxesResponse, error)
@@ -108,6 +114,34 @@ type Backend struct {
Config Config Config Config
} }
func (b Backend) RecoverAssignments(ctx context.Context, trustDomain string) ([]taskassignment.Assignment, error) {
if err := b.validate(); err != nil {
return nil, err
}
result, err := b.Lifecycle.ListSandboxes(ctx, opensandbox.ListOptions{
Metadata: map[string]string{"ci.ddupan.top/backend": "vm"}, PageSize: 100,
})
if err != nil {
return nil, fmt.Errorf("list recoverable sandboxes: %w", err)
}
assignments := make([]taskassignment.Assignment, 0, len(result.Items))
for _, sandbox := range result.Items {
if sandbox.Metadata[terminalMetadata] == "true" {
continue
}
annotations, err := decodeAnnotations(sandbox.Metadata)
if err != nil {
return nil, fmt.Errorf("decode sandbox %s annotations: %w", sandbox.ID, err)
}
assignment, err := taskassignment.FromMetadata(sandbox.Metadata, annotations, trustDomain)
if err != nil {
return nil, fmt.Errorf("recover sandbox %s: %w", sandbox.ID, err)
}
assignments = append(assignments, assignment)
}
return assignments, nil
}
func NewLifecycleClient(baseURL, apiKey string, client *http.Client) *opensandbox.LifecycleClient { func NewLifecycleClient(baseURL, apiKey string, client *http.Client) *opensandbox.LifecycleClient {
if client != nil { if client != nil {
return opensandbox.NewLifecycleClient(baseURL, apiKey, opensandbox.WithHTTPClient(client)) return opensandbox.NewLifecycleClient(baseURL, apiKey, opensandbox.WithHTTPClient(client))
@@ -145,8 +179,12 @@ func (b Backend) Create(ctx context.Context, assignment taskassignment.Assignmen
for key, value := range launch.Environment { for key, value := range launch.Environment {
environment[key] = value environment[key] = value
} }
sandboxMetadata := clone(launch.Metadata.Annotations) sandboxMetadata := clone(launch.Metadata.Labels)
for key, value := range launch.Metadata.Labels { annotations, err := encodeAnnotations(launch.Metadata.Annotations)
if err != nil {
return nil, fmt.Errorf("encode sandbox annotations: %w", err)
}
for key, value := range annotations {
sandboxMetadata[key] = value sandboxMetadata[key] = value
} }
request := opensandbox.CreateSandboxRequest{ request := opensandbox.CreateSandboxRequest{
@@ -170,7 +208,11 @@ func (b Backend) BindIdentity(ctx context.Context, executor *taskworker.Executor
if err != nil { if err != nil {
return fmt.Errorf("verify sandbox identity metadata: %w", err) return fmt.Errorf("verify sandbox identity metadata: %w", err)
} }
if sandbox.Metadata["ci.ddupan.top/spiffe-id"] != identity.SPIFFEID { annotations, err := decodeAnnotations(sandbox.Metadata)
if err != nil {
return fmt.Errorf("decode sandbox identity metadata: %w", err)
}
if annotations["ci.ddupan.top/spiffe-id"] != identity.SPIFFEID {
return fmt.Errorf("sandbox %s has inconsistent SPIFFE identity metadata", executor.Name) return fmt.Errorf("sandbox %s has inconsistent SPIFFE identity metadata", executor.Name)
} }
return nil return nil
@@ -219,3 +261,50 @@ func clone(source map[string]string) map[string]string {
} }
return result return result
} }
// OpenSandbox metadata follows Kubernetes label-value constraints, unlike Pod
// annotations. Store the annotation map as deterministic URL-safe base64
// chunks so repository paths and SPIFFE IDs remain lossless and recoverable.
func encodeAnnotations(annotations map[string]string) (map[string]string, error) {
data, err := json.Marshal(annotations)
if err != nil {
return nil, err
}
encoded := base64.RawURLEncoding.EncodeToString(data)
result := make(map[string]string, (len(encoded)+metadataValueLimit-1)/metadataValueLimit)
for index := 0; len(encoded) > 0; index++ {
length := min(metadataValueLimit, len(encoded))
result[fmt.Sprintf("%s%03d", annotationsMetadataPrefix, index)] = encoded[:length]
encoded = encoded[length:]
}
return result, nil
}
func decodeAnnotations(metadata map[string]string) (map[string]string, error) {
keys := make([]string, 0)
for key := range metadata {
if strings.HasPrefix(key, annotationsMetadataPrefix) {
keys = append(keys, key)
}
}
if len(keys) == 0 {
return nil, errors.New("sandbox annotation metadata is missing")
}
sort.Strings(keys)
var encoded strings.Builder
for index, key := range keys {
if key != fmt.Sprintf("%s%03d", annotationsMetadataPrefix, index) {
return nil, errors.New("sandbox annotation metadata chunks are incomplete")
}
encoded.WriteString(metadata[key])
}
data, err := base64.RawURLEncoding.DecodeString(encoded.String())
if err != nil {
return nil, err
}
var annotations map[string]string
if err := json.Unmarshal(data, &annotations); err != nil {
return nil, err
}
return annotations, nil
}
+29 -2
View File
@@ -31,7 +31,8 @@ func (f *fakeLifecycle) GetSandbox(_ context.Context, id string) (*opensandbox.S
return &item, nil return &item, nil
} }
} }
return &opensandbox.SandboxInfo{ID: id, Metadata: map[string]string{"ci.ddupan.top/spiffe-id": assignment().Identity.SPIFFEID}}, nil metadata, _ := encodeAnnotations(map[string]string{"ci.ddupan.top/spiffe-id": assignment().Identity.SPIFFEID})
return &opensandbox.SandboxInfo{ID: id, Metadata: metadata}, nil
} }
func (f *fakeLifecycle) PatchSandboxMetadata(_ context.Context, id string, patch opensandbox.MetadataPatch) (*opensandbox.SandboxInfo, error) { func (f *fakeLifecycle) PatchSandboxMetadata(_ context.Context, id string, patch opensandbox.MetadataPatch) (*opensandbox.SandboxInfo, error) {
for index := range f.items { for index := range f.items {
@@ -89,12 +90,38 @@ func TestCreateUsesPoolAndPersistsRecoveryMetadata(t *testing.T) {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if lifecycle.created.Extensions["poolRef"] != "ci-vm" || lifecycle.created.Metadata[assignmentMetadata] != assignment().ID || lifecycle.created.Env["CI_SPIFFE_ID"] != assignment().Identity.SPIFFEID || executor.Name != "sandbox-42" { if lifecycle.created.Extensions["poolRef"] != "ci-vm" || lifecycle.created.Metadata[assignmentMetadata] != assignment().ID || lifecycle.created.Metadata["ci.ddupan.top/runner"] != "true" || lifecycle.created.Env["CI_SPIFFE_ID"] != assignment().Identity.SPIFFEID || executor.Name != "sandbox-42" {
t.Fatalf("request=%#v executor=%#v", lifecycle.created, executor) t.Fatalf("request=%#v executor=%#v", lifecycle.created, executor)
} }
if lifecycle.created.Env["CI_RUNNER_CAPABILITY"] != "capability" { if lifecycle.created.Env["CI_RUNNER_CAPABILITY"] != "capability" {
t.Fatalf("environment = %#v", lifecycle.created.Env) t.Fatalf("environment = %#v", lifecycle.created.Env)
} }
annotations, err := decodeAnnotations(lifecycle.created.Metadata)
if err != nil || annotations["ci.ddupan.top/repository"] != "owner/repo" || annotations["ci.ddupan.top/spiffe-id"] != assignment().Identity.SPIFFEID {
t.Fatalf("annotations=%#v err=%v", annotations, err)
}
for _, value := range lifecycle.created.Metadata {
if len(value) > metadataValueLimit {
t.Fatalf("metadata value exceeds %d characters: %q", metadataValueLimit, value)
}
}
}
func TestAnnotationMetadataRoundTripPreservesSlashValues(t *testing.T) {
want := taskworker.BackendMetadata(assignment()).Annotations
encoded, err := encodeAnnotations(want)
if err != nil {
t.Fatal(err)
}
got, err := decodeAnnotations(encoded)
if err != nil {
t.Fatal(err)
}
for key, value := range want {
if got[key] != value {
t.Fatalf("%s=%q, want %q", key, got[key], value)
}
}
} }
func TestBindIdentityVerifiesPersistedMetadata(t *testing.T) { func TestBindIdentityVerifiesPersistedMetadata(t *testing.T) {
+22
View File
@@ -116,6 +116,28 @@ type Backend struct {
Config Config Config Config
} }
func (b Backend) RecoverAssignments(ctx context.Context) ([]taskassignment.Assignment, error) {
if err := b.validate(); err != nil {
return nil, err
}
pods, err := b.API.ListPods(ctx, b.Config.Namespace, "ci.ddupan.top/backend=pod")
if err != nil {
return nil, fmt.Errorf("list recoverable assignment Pods: %w", err)
}
assignments := make([]taskassignment.Assignment, 0, len(pods))
for _, pod := range pods {
if pod.Labels[terminalLabel] == "true" {
continue
}
assignment, err := taskassignment.FromMetadata(pod.Labels, pod.Annotations, b.Config.TrustDomain)
if err != nil {
return nil, fmt.Errorf("recover Pod %s: %w", pod.Name, err)
}
assignments = append(assignments, assignment)
}
return assignments, nil
}
func (b Backend) Find(ctx context.Context, assignmentID string) (*taskworker.Executor, error) { func (b Backend) Find(ctx context.Context, assignmentID string) (*taskworker.Executor, error) {
if err := b.validate(); err != nil { if err := b.validate(); err != nil {
return nil, err return nil, err
+12
View File
@@ -225,3 +225,15 @@ func TestCapabilitiesAreDeterministicAndAssignmentScoped(t *testing.T) {
t.Fatal("capability scope is invalid") t.Fatal("capability scope is invalid")
} }
} }
func TestRegistryRecoversClaimedAssignment(t *testing.T) {
registry := NewRegistry()
assignment := facadeAssignment(t)
if err := registry.RecoverClaimed(assignment); err != nil {
t.Fatal(err)
}
resolved, err := registry.Resolve(assignment.ID, assignment.Identity.SPIFFEID)
if err != nil || resolved.Task.GetId() != assignment.Task.GetId() {
t.Fatalf("resolved=%#v err=%v", resolved, err)
}
}
+18
View File
@@ -90,3 +90,21 @@ func (r *Registry) Remove(assignmentID string) {
defer r.mu.Unlock() defer r.mu.Unlock()
delete(r.claims, assignmentID) delete(r.claims, assignmentID)
} }
// RecoverClaimed restores authorization for an executor that already claimed
// its task before the controller restarted.
func (r *Registry) RecoverClaimed(assignment taskassignment.Assignment) error {
ready, err := r.Offer(assignment)
if err != nil {
return err
}
r.mu.Lock()
entry := r.claims[assignment.ID]
if !entry.claimed {
entry.claimed = true
close(entry.ready)
}
r.mu.Unlock()
<-ready
return nil
}
+27
View File
@@ -7,6 +7,7 @@ import (
"errors" "errors"
"fmt" "fmt"
"slices" "slices"
"strconv"
"gitea.dev/actionslib/pkg/model" "gitea.dev/actionslib/pkg/model"
runnerv1 "gitea.dev/actionslib/runner/v1" runnerv1 "gitea.dev/actionslib/runner/v1"
@@ -32,6 +33,32 @@ type Assignment struct {
Identity taskidentity.Identity Identity taskidentity.Identity
} }
// FromMetadata reconstructs the minimal assignment needed to authorize an
// already-running executor after a controller restart. Backend metadata was
// originally derived from the trusted Gitea task and is validated again here.
func FromMetadata(labels, annotations map[string]string, trustDomain string) (Assignment, error) {
taskID, err := strconv.ParseInt(labels["ci.ddupan.top/task-id"], 10, 64)
if err != nil || taskID < 1 {
return Assignment{}, errors.New("backend metadata has invalid task ID")
}
backend := Backend(labels["ci.ddupan.top/backend"])
if backend != BackendPod && backend != BackendVM {
return Assignment{}, errors.New("backend metadata has invalid backend")
}
id := labels["ci.ddupan.top/assignment-id"]
if id != fmt.Sprintf("gitea-task-%d", taskID) {
return Assignment{}, errors.New("backend metadata assignment ID does not match task ID")
}
identity, err := taskidentity.FromMetadata(
annotations["ci.ddupan.top/repository"], annotations["ci.ddupan.top/job-key"],
annotations["ci.ddupan.top/spiffe-id"], trustDomain,
)
if err != nil {
return Assignment{}, err
}
return Assignment{ID: id, Backend: backend, 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"`
@@ -73,3 +73,18 @@ func TestAssignmentWireRoundTripAndValidation(t *testing.T) {
t.Fatal("expected tampered backend to fail") t.Fatal("expected tampered backend to fail")
} }
} }
func TestFromMetadataRecoversMinimalAssignment(t *testing.T) {
assignment, err := FromMetadata(map[string]string{
"ci.ddupan.top/assignment-id": "gitea-task-42",
"ci.ddupan.top/task-id": "42",
"ci.ddupan.top/backend": "vm",
}, map[string]string{
"ci.ddupan.top/repository": "owner/repo",
"ci.ddupan.top/job-key": "publish",
"ci.ddupan.top/spiffe-id": "spiffe://ddupan.top/ci/owner/repo/publish",
}, "ddupan.top")
if err != nil || assignment.Task.GetId() != 42 || assignment.Backend != BackendVM {
t.Fatalf("assignment=%#v err=%v", assignment, err)
}
}
+15
View File
@@ -23,6 +23,21 @@ type Identity struct {
SPIFFEID string SPIFFEID string
} }
// FromMetadata validates identity fields recovered from backend-owned state.
func FromMetadata(repository, task, spiffeID, trustDomain string) (Identity, error) {
parts := strings.Split(repository, "/")
if len(parts) != 2 || parts[0] == "" || parts[1] == "" || !safeTaskKey.MatchString(task) {
return Identity{}, errors.New("invalid recovered repository or task identity")
}
expected := "spiffe://" + trustDomain + "/ci/" + strings.Join([]string{
sanitize(parts[0]), sanitize(parts[1]), task,
}, "/")
if spiffeID != expected {
return Identity{}, fmt.Errorf("recovered SPIFFE ID %q does not match %q", spiffeID, expected)
}
return Identity{Repository: repository, Task: task, SPIFFEID: spiffeID}, nil
}
// FromTask derives the repository/task SPIFFE ID from Gitea's trusted task // FromTask derives the repository/task SPIFFE ID from Gitea's trusted task
// context. Workflow input never supplies or overrides the resulting ID. // context. Workflow input never supplies or overrides the resulting ID.
func FromTask(task *runnerv1.Task, trustDomain string) (Identity, error) { func FromTask(task *runnerv1.Task, trustDomain string) (Identity, error) {
-31
View File
@@ -1,31 +0,0 @@
package taskscheduler
import "context"
// SingleFlightGate keeps at most one fetched task in flight. Release is
// idempotent so repeated terminal updates cannot increase capacity.
type SingleFlightGate struct {
token chan struct{}
}
func NewSingleFlightGate() *SingleFlightGate {
gate := &SingleFlightGate{token: make(chan struct{}, 1)}
gate.token <- struct{}{}
return gate
}
func (g *SingleFlightGate) Acquire(ctx context.Context) error {
select {
case <-g.token:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func (g *SingleFlightGate) Release() {
select {
case g.token <- struct{}{}:
default:
}
}
-29
View File
@@ -1,29 +0,0 @@
package taskscheduler
import (
"context"
"testing"
"time"
)
func TestSingleFlightGateBlocksUntilTerminalRelease(t *testing.T) {
gate := NewSingleFlightGate()
if err := gate.Acquire(context.Background()); err != nil {
t.Fatal(err)
}
blocked, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond)
defer cancel()
if err := gate.Acquire(blocked); err == nil {
t.Fatal("second task acquired capacity before release")
}
gate.Release()
gate.Release()
if err := gate.Acquire(context.Background()); err != nil {
t.Fatal(err)
}
blockedAgain, cancelAgain := context.WithTimeout(context.Background(), 20*time.Millisecond)
defer cancelAgain()
if err := gate.Acquire(blockedAgain); err == nil {
t.Fatal("duplicate terminal release increased capacity")
}
}
+30 -18
View File
@@ -4,8 +4,11 @@ import (
"context" "context"
"errors" "errors"
"fmt" "fmt"
"sync/atomic"
"time" "time"
"golang.org/x/sync/errgroup"
runnerv1 "gitea.dev/actionslib/runner/v1" runnerv1 "gitea.dev/actionslib/runner/v1"
) )
@@ -19,14 +22,15 @@ type PollerConfig struct {
Labels []string Labels []string
EmptyBackoff time.Duration EmptyBackoff time.Duration
ErrorBackoff time.Duration ErrorBackoff time.Duration
Capacity int
} }
// Poller is the scheduler component. Once Gitea assigns a task, it never // Poller is the scheduler component. Each fetcher keeps its assigned task
// fetches another one until the current assignment is durably dispatched. // until that assignment is durably dispatched; all fetchers share one runner
// declaration and a monotonic tasks version.
type Poller struct { type Poller struct {
Client PollClient Client PollClient
Scheduler *Scheduler Scheduler *Scheduler
Gate *SingleFlightGate
Config PollerConfig Config PollerConfig
OnError func(error) OnError func(error)
} }
@@ -50,16 +54,21 @@ func (p Poller) Run(ctx context.Context) error {
if errorBackoff <= 0 { if errorBackoff <= 0 {
errorBackoff = 5 * time.Second errorBackoff = 5 * time.Second
} }
var tasksVersion int64 capacity := p.Config.Capacity
haveLease := false if capacity < 1 {
capacity = 1
}
var tasksVersion atomic.Int64
group, groupContext := errgroup.WithContext(ctx)
for range capacity {
group.Go(func() error { return p.runFetcher(groupContext, &tasksVersion, emptyBackoff, errorBackoff) })
}
return group.Wait()
}
func (p Poller) runFetcher(ctx context.Context, tasksVersion *atomic.Int64, emptyBackoff, errorBackoff time.Duration) error {
for { for {
if p.Gate != nil && !haveLease { response, err := p.Client.FetchTask(ctx, tasksVersion.Load())
if err := p.Gate.Acquire(ctx); err != nil {
return nil
}
haveLease = true
}
response, err := p.Client.FetchTask(ctx, tasksVersion)
if err != nil { if err != nil {
if ctx.Err() != nil { if ctx.Err() != nil {
return nil return nil
@@ -77,13 +86,9 @@ func (p Poller) Run(ctx context.Context) error {
} }
continue continue
} }
tasksVersion = response.GetTasksVersion() storeMaximum(tasksVersion, response.GetTasksVersion())
task := response.GetTask() task := response.GetTask()
if task == nil { if task == nil {
if p.Gate != nil {
p.Gate.Release()
haveLease = false
}
if !wait(ctx, emptyBackoff) { if !wait(ctx, emptyBackoff) {
return nil return nil
} }
@@ -91,7 +96,6 @@ func (p Poller) Run(ctx context.Context) error {
} }
for { for {
if err := p.Scheduler.Run(ctx, task); err == nil { if err := p.Scheduler.Run(ctx, task); err == nil {
haveLease = false
break break
} else { } else {
if ctx.Err() != nil { if ctx.Err() != nil {
@@ -106,6 +110,14 @@ func (p Poller) Run(ctx context.Context) error {
} }
} }
func storeMaximum(value *atomic.Int64, candidate int64) {
for current := value.Load(); candidate > current; current = value.Load() {
if value.CompareAndSwap(current, candidate) {
return
}
}
}
func (p Poller) report(err error) { func (p Poller) report(err error) {
if p.OnError != nil { if p.OnError != nil {
p.OnError(err) p.OnError(err)
+49
View File
@@ -115,3 +115,52 @@ func TestPollerRetriesAssignedTaskBeforeFetchingAnother(t *testing.T) {
t.Fatalf("dispatches=%d fetches-before-dispatch=%d declares=%d", dispatcher.calls, dispatcher.fetchesAtSuccess, client.declared) t.Fatalf("dispatches=%d fetches-before-dispatch=%d declares=%d", dispatcher.calls, dispatcher.fetchesAtSuccess, client.declared)
} }
} }
type blockingPollClient struct {
mu sync.Mutex
declared int
started chan struct{}
}
func (c *blockingPollClient) Declare(context.Context, string, []string) error {
c.mu.Lock()
c.declared++
c.mu.Unlock()
return nil
}
func (c *blockingPollClient) FetchTask(ctx context.Context, _ int64) (*runnerv1.FetchTaskResponse, error) {
c.started <- struct{}{}
<-ctx.Done()
return nil, ctx.Err()
}
func TestPollerStartsConfiguredNumberOfFetchersAfterOneDeclare(t *testing.T) {
client := &blockingPollClient{started: make(chan struct{}, 3)}
poller := Poller{
Client: client,
Scheduler: &Scheduler{TrustDomain: "ddupan.top", Dispatcher: &retryDispatcher{done: make(chan struct{})}},
Config: PollerConfig{
Version: "dev", Labels: []string{"self-hosted:host", "pod:host", "vm:host"}, Capacity: 3,
},
}
ctx, cancel := context.WithCancel(context.Background())
finished := make(chan error, 1)
go func() { finished <- poller.Run(ctx) }()
for range 3 {
select {
case <-client.started:
case <-time.After(time.Second):
t.Fatal("configured fetchers did not start")
}
}
cancel()
if err := <-finished; err != nil {
t.Fatal(err)
}
client.mu.Lock()
defer client.mu.Unlock()
if client.declared != 1 {
t.Fatalf("declares = %d", client.declared)
}
}
+1
View File
@@ -178,6 +178,7 @@ func (w Worker) launchSpec(assignment taskassignment.Assignment) (LaunchSpec, er
func BackendMetadata(assignment taskassignment.Assignment) Metadata { func BackendMetadata(assignment taskassignment.Assignment) Metadata {
return Metadata{ return Metadata{
Labels: map[string]string{ Labels: map[string]string{
"ci.ddupan.top/runner": "true",
"ci.ddupan.top/assignment-id": assignment.ID, "ci.ddupan.top/assignment-id": assignment.ID,
"ci.ddupan.top/task-id": strconv.FormatInt(assignment.Task.GetId(), 10), "ci.ddupan.top/task-id": strconv.FormatInt(assignment.Task.GetId(), 10),
"ci.ddupan.top/backend": string(assignment.Backend), "ci.ddupan.top/backend": string(assignment.Backend),