Compare commits

...
Author SHA1 Message Date
panxiao81 48b6b8038e feat: 隔离 VM 集成测试标签
test / python (pull_request) Failing after 11m44s
test / shell (pull_request) Successful in 21s
test / go (pull_request) Successful in 5m12s
2026-09-21 08:50:18 +00:00
panxiao81 70c5ff422f Merge assignment 重投预算修复
publish images / publish-images (push) Failing after 12m48s
test / python (push) Successful in 31s
test / shell (push) Successful in 34s
test / go (push) Successful in 5m40s
2026-09-21 08:44:42 +00:00
panxiao81 906e6a2e18 fix: 保持后端 assignment 可重试
test / shell (pull_request) Failing after 13m28s
test / python (pull_request) Failing after 13m29s
test / go (pull_request) Successful in 6m15s
2026-09-21 08:44:00 +00:00
panxiao81 ace84373f6 Merge OpenSandbox Runner workload 标记
test / python (push) Successful in 25s
test / shell (push) Successful in 30s
test / go (push) Successful in 4m30s
publish images / publish-images (push) Failing after 13m58s
2026-09-21 08:28:32 +00:00
panxiao81 87cf359ef6 fix: 标记 OpenSandbox Runner workload
test / shell (pull_request) Successful in 26s
test / go (pull_request) Successful in 2m58s
test / python (pull_request) Failing after 14m15s
2026-09-21 08:22:56 +00:00
panxiao81 ad3934e7e1 Merge OpenSandbox metadata 编码修复
test / python (push) Successful in 21s
test / shell (push) Successful in 24s
test / go (push) Successful in 6m52s
publish images / publish-images (push) Failing after 10m40s
2026-09-21 07:56:50 +00:00
panxiao81 d82687382a fix: 编码 OpenSandbox assignment metadata
test / go (pull_request) Successful in 4m45s
test / shell (pull_request) Failing after 13m27s
test / python (pull_request) Failing after 13m27s
2026-09-21 07:48:10 +00:00
panxiao81 7f52cf393f Merge Pod 与 VM 独立执行容量池
publish images / publish-images (push) Failing after 11m8s
test / python (push) Successful in 10s
test / shell (push) Successful in 26s
test / go (push) Successful in 5m51s
2026-09-21 07:26:10 +00:00
panxiao81 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 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
16 changed files with 360 additions and 34 deletions
+1 -1
View File
@@ -5,7 +5,7 @@ on:
jobs: jobs:
kind: kind:
runs-on: [self-hosted, vm] runs-on: [self-hosted, vm-dev]
steps: steps:
- name: Verify Docker - name: Verify Docker
run: docker info run: docker info
+9 -3
View File
@@ -15,10 +15,15 @@ runs-on: [self-hosted, vm]
只执行一个 job,并在 job 结束后连同本地状态一起销毁。完整的设计约束见 只执行一个 job,并在 job 结束后连同本地状态一起销毁。完整的设计约束见
[`docs/design-principles.md`](docs/design-principles.md)。 [`docs/design-principles.md`](docs/design-principles.md)。
集成期间可将 `VM_RUNNER_LABEL=vm-dev`,只接取显式使用
`runs-on: [self-hosted, vm-dev]` 的测试任务;生产 `vm` job 将保持在 Gitea pending,
不会在 backend 修复过程中继续涌入。
目标 Go controller 组件: 目标 Go controller 组件:
- `scheduler`:以常驻 Gitea RunnerService 身份直接领取 task,并把完整 assignment - `scheduler`:以常驻 Gitea RunnerService 身份直接领取 task,并把完整 assignment
持久化到 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 +42,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)。
## 开发 ## 开发
+18 -8
View File
@@ -22,6 +22,7 @@ import (
"k8s.io/client-go/tools/leaderelection/resourcelock" "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"
@@ -48,7 +49,7 @@ type controllerConfig struct {
PodNamespace, PodImage, PodServiceAccount, SPIRECluster, SPIREClass string PodNamespace, PodImage, PodServiceAccount, SPIRECluster, SPIREClass string
SPIREAgentID string SPIREAgentID string
PodExecutorUID, PodCapacity int PodExecutorUID, PodCapacity int
OpenSandboxURL, OpenSandboxAPIKey, OpenSandboxPool string OpenSandboxURL, OpenSandboxAPIKey, OpenSandboxPool, VMRunnerLabel string
VMTimeout, VMCapacity int VMTimeout, VMCapacity int
} }
@@ -91,6 +92,8 @@ func runController(ctx context.Context) error {
return err return err
} }
registry := runnerfacade.NewRegistry() registry := runnerfacade.NewRegistry()
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)
@@ -105,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")
@@ -112,6 +116,7 @@ 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)
} }
return nil return nil
}, },
@@ -126,14 +131,14 @@ func runController(ctx context.Context) error {
labels = append(labels, string(taskassignment.BackendPod)) labels = append(labels, string(taskassignment.BackendPod))
} }
if slices.Contains(config.Components, controller.VMWorker) { if slices.Contains(config.Components, controller.VMWorker) {
labels = append(labels, string(taskassignment.BackendVM)) labels = append(labels, config.VMRunnerLabel)
} }
poller := taskscheduler.Poller{ poller := taskscheduler.Poller{
Client: giteaClient, Client: giteaClient,
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() kubernetesConfig, err := rest.InClusterConfig()
@@ -180,11 +185,12 @@ func runController(ctx context.Context) error {
return err return err
} }
for _, assignment := range assignments { for _, assignment := range assignments {
podPool.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}, registry) 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
} }
@@ -209,11 +215,12 @@ func runController(ctx context.Context) error {
return err return err
} }
for _, assignment := range assignments { for _, assignment := range assignments {
vmPool.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}, registry) 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
} }
@@ -280,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
} }
@@ -338,7 +345,7 @@ func loadControllerConfig() (controllerConfig, error) {
FacadeListen: env("RUNNER_FACADE_LISTEN", ":8443"), FacadeURL: os.Getenv("RUNNER_FACADE_URL"), FacadeSPIFFEID: os.Getenv("RUNNER_FACADE_SPIFFE_ID"), CapabilityKey: []byte(capabilityKey), FacadeListen: env("RUNNER_FACADE_LISTEN", ":8443"), FacadeURL: os.Getenv("RUNNER_FACADE_URL"), FacadeSPIFFEID: os.Getenv("RUNNER_FACADE_SPIFFE_ID"), CapabilityKey: []byte(capabilityKey),
PodNamespace: env("POD_NAMESPACE", "gitea-actions"), PodImage: os.Getenv("POD_EXECUTOR_IMAGE"), PodServiceAccount: env("POD_SERVICE_ACCOUNT", "gitea-task-executor"), PodNamespace: env("POD_NAMESPACE", "gitea-actions"), PodImage: os.Getenv("POD_EXECUTOR_IMAGE"), PodServiceAccount: env("POD_SERVICE_ACCOUNT", "gitea-task-executor"),
SPIRECluster: env("SPIRE_CLUSTER", "homelab"), SPIREClass: env("SPIRE_CLASS", "spire-mgmt-spire"), SPIREAgentID: os.Getenv("SPIRE_AGENT_ID"), PodExecutorUID: envInt("POD_EXECUTOR_UID", 2000), PodCapacity: envInt("POD_CAPACITY", 4), SPIRECluster: env("SPIRE_CLUSTER", "homelab"), SPIREClass: env("SPIRE_CLASS", "spire-mgmt-spire"), SPIREAgentID: os.Getenv("SPIRE_AGENT_ID"), PodExecutorUID: envInt("POD_EXECUTOR_UID", 2000), PodCapacity: envInt("POD_CAPACITY", 4),
OpenSandboxURL: os.Getenv("OPENSANDBOX_API"), OpenSandboxPool: env("OPENSANDBOX_POOL", "ci-vm"), VMTimeout: envInt("VM_TIMEOUT_SECONDS", 14400), VMCapacity: envInt("VM_CAPACITY", 1), OpenSandboxURL: os.Getenv("OPENSANDBOX_API"), OpenSandboxPool: env("OPENSANDBOX_POOL", "ci-vm"), VMRunnerLabel: env("VM_RUNNER_LABEL", "vm"), VMTimeout: envInt("VM_TIMEOUT_SECONDS", 14400), VMCapacity: envInt("VM_CAPACITY", 1),
} }
if config.WorkloadAPIAddr == "" || config.FacadeURL == "" || config.FacadeSPIFFEID == "" { if config.WorkloadAPIAddr == "" || config.FacadeURL == "" || config.FacadeSPIFFEID == "" {
return controllerConfig{}, errors.New("SPIFFE_ENDPOINT_SOCKET, RUNNER_FACADE_URL, and RUNNER_FACADE_SPIFFE_ID are required") return controllerConfig{}, errors.New("SPIFFE_ENDPOINT_SOCKET, RUNNER_FACADE_URL, and RUNNER_FACADE_SPIFFE_ID are required")
@@ -350,6 +357,9 @@ func loadControllerConfig() (controllerConfig, error) {
return controllerConfig{}, errors.New("SPIRE_AGENT_ID is required for pod-worker") return controllerConfig{}, errors.New("SPIRE_AGENT_ID is required for pod-worker")
} }
if slices.Contains(selection, controller.VMWorker) { if slices.Contains(selection, controller.VMWorker) {
if config.VMRunnerLabel != "vm" && config.VMRunnerLabel != "vm-dev" {
return controllerConfig{}, errors.New("VM_RUNNER_LABEL must be vm or vm-dev")
}
if config.OpenSandboxURL == "" { if config.OpenSandboxURL == "" {
return controllerConfig{}, errors.New("OPENSANDBOX_API is required for vm-worker") return controllerConfig{}, errors.New("OPENSANDBOX_API is required for vm-worker")
} }
@@ -56,6 +56,7 @@ func TestLoadControllerConfigRequiresOpenSandboxSecretOnlyForVM(t *testing.T) {
t.Setenv("RUNNER_FACADE_URL", "https://facade:8443") t.Setenv("RUNNER_FACADE_URL", "https://facade:8443")
t.Setenv("RUNNER_FACADE_SPIFFE_ID", "spiffe://ddupan.top/controller") t.Setenv("RUNNER_FACADE_SPIFFE_ID", "spiffe://ddupan.top/controller")
t.Setenv("OPENSANDBOX_API", "http://opensandbox.internal") t.Setenv("OPENSANDBOX_API", "http://opensandbox.internal")
t.Setenv("VM_RUNNER_LABEL", "vm-dev")
if _, err := loadControllerConfig(); err == nil { if _, err := loadControllerConfig(); err == nil {
t.Fatal("expected missing OpenSandbox API key file error") t.Fatal("expected missing OpenSandbox API key file error")
} }
+10 -4
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,9 +76,11 @@ 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;无需把完整
+17 -3
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
@@ -153,7 +167,7 @@ func OpenConsumer(ctx context.Context, manager consumerManager, stream, subjectB
AckPolicy: jetstream.AckExplicitPolicy, AckPolicy: jetstream.AckExplicitPolicy,
AckWait: 5 * time.Minute, AckWait: 5 * time.Minute,
MaxAckPending: capacity, MaxAckPending: capacity,
MaxDeliver: 20, MaxDeliver: 1000,
}) })
if err != nil { if err != nil {
return nil, fmt.Errorf("open %s assignment consumer: %w", backend, err) return nil, fmt.Errorf("open %s assignment consumer: %w", backend, err)
+48 -3
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) {
@@ -160,7 +205,7 @@ func TestOpenConsumerUsesIndependentDurablePerBackend(t *testing.T) {
if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.BackendPod, 4); err != nil { if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.BackendPod, 4); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if manager.config.Durable != "pod" || manager.config.FilterSubject != "ci.assignment.pod" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 { if manager.config.Durable != "pod" || manager.config.FilterSubject != "ci.assignment.pod" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 || manager.config.MaxDeliver != 1000 {
t.Fatalf("config = %#v", manager.config) t.Fatalf("config = %#v", manager.config)
} }
} }
+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())
}
}
+69 -4
View File
@@ -3,9 +3,13 @@ package opensandboxbackend
import ( import (
"context" "context"
"encoding/base64"
"encoding/json"
"errors" "errors"
"fmt" "fmt"
"net/http" "net/http"
"sort"
"strings"
"time" "time"
opensandbox "github.com/alibaba/OpenSandbox/sdks/sandbox/go" opensandbox "github.com/alibaba/OpenSandbox/sdks/sandbox/go"
@@ -17,6 +21,8 @@ import (
const assignmentMetadata = "ci.ddupan.top/assignment-id" const assignmentMetadata = "ci.ddupan.top/assignment-id"
const terminalMetadata = "ci.ddupan.top/terminal" const terminalMetadata = "ci.ddupan.top/terminal"
const annotationsMetadataPrefix = "ci.ddupan.top/annotations-"
const metadataValueLimit = 63
type Lifecycle interface { type Lifecycle interface {
ListSandboxes(context.Context, opensandbox.ListOptions) (*opensandbox.ListSandboxesResponse, error) ListSandboxes(context.Context, opensandbox.ListOptions) (*opensandbox.ListSandboxesResponse, error)
@@ -123,7 +129,11 @@ func (b Backend) RecoverAssignments(ctx context.Context, trustDomain string) ([]
if sandbox.Metadata[terminalMetadata] == "true" { if sandbox.Metadata[terminalMetadata] == "true" {
continue continue
} }
assignment, err := taskassignment.FromMetadata(sandbox.Metadata, sandbox.Metadata, trustDomain) annotations, err := decodeAnnotations(sandbox.Metadata)
if err != nil {
return nil, fmt.Errorf("decode sandbox %s annotations: %w", sandbox.ID, err)
}
assignment, err := taskassignment.FromMetadata(sandbox.Metadata, annotations, trustDomain)
if err != nil { if err != nil {
return nil, fmt.Errorf("recover sandbox %s: %w", sandbox.ID, err) return nil, fmt.Errorf("recover sandbox %s: %w", sandbox.ID, err)
} }
@@ -169,8 +179,12 @@ func (b Backend) Create(ctx context.Context, assignment taskassignment.Assignmen
for key, value := range launch.Environment { for key, value := range launch.Environment {
environment[key] = value environment[key] = value
} }
sandboxMetadata := clone(launch.Metadata.Annotations) sandboxMetadata := clone(launch.Metadata.Labels)
for key, value := range launch.Metadata.Labels { annotations, err := encodeAnnotations(launch.Metadata.Annotations)
if err != nil {
return nil, fmt.Errorf("encode sandbox annotations: %w", err)
}
for key, value := range annotations {
sandboxMetadata[key] = value sandboxMetadata[key] = value
} }
request := opensandbox.CreateSandboxRequest{ request := opensandbox.CreateSandboxRequest{
@@ -194,7 +208,11 @@ func (b Backend) BindIdentity(ctx context.Context, executor *taskworker.Executor
if err != nil { if err != nil {
return fmt.Errorf("verify sandbox identity metadata: %w", err) return fmt.Errorf("verify sandbox identity metadata: %w", err)
} }
if sandbox.Metadata["ci.ddupan.top/spiffe-id"] != identity.SPIFFEID { annotations, err := decodeAnnotations(sandbox.Metadata)
if err != nil {
return fmt.Errorf("decode sandbox identity metadata: %w", err)
}
if annotations["ci.ddupan.top/spiffe-id"] != identity.SPIFFEID {
return fmt.Errorf("sandbox %s has inconsistent SPIFFE identity metadata", executor.Name) return fmt.Errorf("sandbox %s has inconsistent SPIFFE identity metadata", executor.Name)
} }
return nil return nil
@@ -243,3 +261,50 @@ func clone(source map[string]string) map[string]string {
} }
return result return result
} }
// OpenSandbox metadata follows Kubernetes label-value constraints, unlike Pod
// annotations. Store the annotation map as deterministic URL-safe base64
// chunks so repository paths and SPIFFE IDs remain lossless and recoverable.
func encodeAnnotations(annotations map[string]string) (map[string]string, error) {
data, err := json.Marshal(annotations)
if err != nil {
return nil, err
}
encoded := base64.RawURLEncoding.EncodeToString(data)
result := make(map[string]string, (len(encoded)+metadataValueLimit-1)/metadataValueLimit)
for index := 0; len(encoded) > 0; index++ {
length := min(metadataValueLimit, len(encoded))
result[fmt.Sprintf("%s%03d", annotationsMetadataPrefix, index)] = encoded[:length]
encoded = encoded[length:]
}
return result, nil
}
func decodeAnnotations(metadata map[string]string) (map[string]string, error) {
keys := make([]string, 0)
for key := range metadata {
if strings.HasPrefix(key, annotationsMetadataPrefix) {
keys = append(keys, key)
}
}
if len(keys) == 0 {
return nil, errors.New("sandbox annotation metadata is missing")
}
sort.Strings(keys)
var encoded strings.Builder
for index, key := range keys {
if key != fmt.Sprintf("%s%03d", annotationsMetadataPrefix, index) {
return nil, errors.New("sandbox annotation metadata chunks are incomplete")
}
encoded.WriteString(metadata[key])
}
data, err := base64.RawURLEncoding.DecodeString(encoded.String())
if err != nil {
return nil, err
}
var annotations map[string]string
if err := json.Unmarshal(data, &annotations); err != nil {
return nil, err
}
return annotations, nil
}
+29 -2
View File
@@ -31,7 +31,8 @@ func (f *fakeLifecycle) GetSandbox(_ context.Context, id string) (*opensandbox.S
return &item, nil return &item, nil
} }
} }
return &opensandbox.SandboxInfo{ID: id, Metadata: map[string]string{"ci.ddupan.top/spiffe-id": assignment().Identity.SPIFFEID}}, nil metadata, _ := encodeAnnotations(map[string]string{"ci.ddupan.top/spiffe-id": assignment().Identity.SPIFFEID})
return &opensandbox.SandboxInfo{ID: id, Metadata: metadata}, nil
} }
func (f *fakeLifecycle) PatchSandboxMetadata(_ context.Context, id string, patch opensandbox.MetadataPatch) (*opensandbox.SandboxInfo, error) { func (f *fakeLifecycle) PatchSandboxMetadata(_ context.Context, id string, patch opensandbox.MetadataPatch) (*opensandbox.SandboxInfo, error) {
for index := range f.items { for index := range f.items {
@@ -89,12 +90,38 @@ func TestCreateUsesPoolAndPersistsRecoveryMetadata(t *testing.T) {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if lifecycle.created.Extensions["poolRef"] != "ci-vm" || lifecycle.created.Metadata[assignmentMetadata] != assignment().ID || lifecycle.created.Env["CI_SPIFFE_ID"] != assignment().Identity.SPIFFEID || executor.Name != "sandbox-42" { if lifecycle.created.Extensions["poolRef"] != "ci-vm" || lifecycle.created.Metadata[assignmentMetadata] != assignment().ID || lifecycle.created.Metadata["ci.ddupan.top/runner"] != "true" || lifecycle.created.Env["CI_SPIFFE_ID"] != assignment().Identity.SPIFFEID || executor.Name != "sandbox-42" {
t.Fatalf("request=%#v executor=%#v", lifecycle.created, executor) t.Fatalf("request=%#v executor=%#v", lifecycle.created, executor)
} }
if lifecycle.created.Env["CI_RUNNER_CAPABILITY"] != "capability" { if lifecycle.created.Env["CI_RUNNER_CAPABILITY"] != "capability" {
t.Fatalf("environment = %#v", lifecycle.created.Env) t.Fatalf("environment = %#v", lifecycle.created.Env)
} }
annotations, err := decodeAnnotations(lifecycle.created.Metadata)
if err != nil || annotations["ci.ddupan.top/repository"] != "owner/repo" || annotations["ci.ddupan.top/spiffe-id"] != assignment().Identity.SPIFFEID {
t.Fatalf("annotations=%#v err=%v", annotations, err)
}
for _, value := range lifecycle.created.Metadata {
if len(value) > metadataValueLimit {
t.Fatalf("metadata value exceeds %d characters: %q", metadataValueLimit, value)
}
}
}
func TestAnnotationMetadataRoundTripPreservesSlashValues(t *testing.T) {
want := taskworker.BackendMetadata(assignment()).Annotations
encoded, err := encodeAnnotations(want)
if err != nil {
t.Fatal(err)
}
got, err := decodeAnnotations(encoded)
if err != nil {
t.Fatal(err)
}
for key, value := range want {
if got[key] != value {
t.Fatalf("%s=%q, want %q", key, got[key], value)
}
}
} }
func TestBindIdentityVerifiesPersistedMetadata(t *testing.T) { func TestBindIdentityVerifiesPersistedMetadata(t *testing.T) {
+1 -1
View File
@@ -102,7 +102,7 @@ func backendFromTask(task *runnerv1.Task) (Backend, error) {
return "", fmt.Errorf("task runs-on labels must include self-hosted: %v", labels) return "", fmt.Errorf("task runs-on labels must include self-hosted: %v", labels)
} }
hasPod := slices.Contains(labels, string(BackendPod)) hasPod := slices.Contains(labels, string(BackendPod))
hasVM := slices.Contains(labels, string(BackendVM)) hasVM := slices.Contains(labels, string(BackendVM)) || slices.Contains(labels, "vm-dev")
if hasPod == hasVM { if hasPod == hasVM {
return "", fmt.Errorf("task runs-on labels must select exactly one of pod or vm: %v", labels) return "", fmt.Errorf("task runs-on labels must select exactly one of pod or vm: %v", labels)
} }
@@ -28,6 +28,7 @@ func TestNewSelectsBackendFromRunsOn(t *testing.T) {
}{ }{
{"[self-hosted, pod]", BackendPod}, {"[self-hosted, pod]", BackendPod},
{"[self-hosted, vm]", BackendVM}, {"[self-hosted, vm]", BackendVM},
{"[self-hosted, vm-dev]", BackendVM},
} { } {
assignment, err := New(task(t, test.labels), "ddupan.top") assignment, err := New(task(t, test.labels), "ddupan.top")
if err != nil { if err != nil {
+30 -5
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,10 +22,12 @@ 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
@@ -49,9 +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
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 {
response, err := p.Client.FetchTask(ctx, tasksVersion) response, err := p.Client.FetchTask(ctx, tasksVersion.Load())
if err != nil { if err != nil {
if ctx.Err() != nil { if ctx.Err() != nil {
return nil return nil
@@ -69,7 +86,7 @@ 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 !wait(ctx, emptyBackoff) { if !wait(ctx, emptyBackoff) {
@@ -93,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),