Compare commits

...
Author SHA1 Message Date
panxiao81 53a080b310 ci: 为测试显式选择 Pod runner
test / shell (pull_request) Successful in 1m34s
test / python (pull_request) Successful in 1m13s
test / go (pull_request) Successful in 3m14s
2026-09-21 09:34:30 +00:00
panxiao81 38e8d59541 feat: 增加 Runner 生命周期结构化日志
test / go (pull_request) Successful in 3m10s
test / shell (pull_request) Failing after 10m4s
test / python (pull_request) Failing after 10m4s
2026-09-21 09:31:46 +00:00
panxiao81 661b5e9218 Merge vm-dev 集成测试隔离标签
publish images / publish-images (push) Failing after 11m12s
test / python (push) Successful in 24s
test / shell (push) Successful in 28s
test / go (push) Successful in 6m44s
2026-09-21 08:51:17 +00:00
panxiao81 48b6b8038e feat: 隔离 VM 集成测试标签
test / python (pull_request) Failing after 11m44s
test / shell (pull_request) Successful in 21s
test / go (pull_request) Successful in 5m12s
2026-09-21 08:50:18 +00:00
panxiao81 70c5ff422f Merge assignment 重投预算修复
publish images / publish-images (push) Failing after 12m48s
test / python (push) Successful in 31s
test / shell (push) Successful in 34s
test / go (push) Successful in 5m40s
2026-09-21 08:44:42 +00:00
panxiao81 906e6a2e18 fix: 保持后端 assignment 可重试
test / shell (pull_request) Failing after 13m28s
test / python (pull_request) Failing after 13m29s
test / go (pull_request) Successful in 6m15s
2026-09-21 08:44:00 +00:00
panxiao81 ace84373f6 Merge OpenSandbox Runner workload 标记
test / python (push) Successful in 25s
test / shell (push) Successful in 30s
test / go (push) Successful in 4m30s
publish images / publish-images (push) Failing after 13m58s
2026-09-21 08:28:32 +00:00
12 changed files with 121 additions and 20 deletions
+2 -2
View File
@@ -17,7 +17,7 @@ jobs:
- run: go vet ./... - run: go vet ./...
python: python:
runs-on: self-hosted runs-on: [self-hosted, pod]
steps: steps:
- uses: actions/checkout@v4 - uses: actions/checkout@v4
- uses: actions/setup-python@v5 - uses: actions/setup-python@v5
@@ -28,7 +28,7 @@ jobs:
- run: python -m compileall -q src - run: python -m compileall -q src
shell: shell:
runs-on: self-hosted runs-on: [self-hosted, pod]
steps: steps:
- uses: actions/checkout@v4 - uses: actions/checkout@v4
- run: | - run: |
+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
+4
View File
@@ -15,6 +15,10 @@ runs-on: [self-hosted, vm]
只执行一个 job,并在 job 结束后连同本地状态一起销毁。完整的设计约束见 只执行一个 job,并在 job 结束后连同本地状态一起销毁。完整的设计约束见
[`docs/design-principles.md`](docs/design-principles.md)。 [`docs/design-principles.md`](docs/design-principles.md)。
集成期间可将 `VM_RUNNER_LABEL=vm-dev`,只接取显式使用
`runs-on: [self-hosted, vm-dev]` 的测试任务;生产 `vm` job 将保持在 Gitea pending,
不会在 backend 修复过程中继续涌入。
目标 Go controller 组件: 目标 Go controller 组件:
- `scheduler`:以常驻 Gitea RunnerService 身份直接领取 task,并把完整 assignment - `scheduler`:以常驻 Gitea RunnerService 身份直接领取 task,并把完整 assignment
+28 -11
View File
@@ -4,7 +4,7 @@ import (
"context" "context"
"errors" "errors"
"fmt" "fmt"
"log" "log/slog"
"net/http" "net/http"
"os" "os"
"slices" "slices"
@@ -49,7 +49,7 @@ type controllerConfig struct {
PodNamespace, PodImage, PodServiceAccount, SPIRECluster, SPIREClass string PodNamespace, PodImage, PodServiceAccount, SPIRECluster, SPIREClass string
SPIREAgentID string SPIREAgentID string
PodExecutorUID, PodCapacity int PodExecutorUID, PodCapacity int
OpenSandboxURL, OpenSandboxAPIKey, OpenSandboxPool string OpenSandboxURL, OpenSandboxAPIKey, OpenSandboxPool, VMRunnerLabel string
VMTimeout, VMCapacity int VMTimeout, VMCapacity int
} }
@@ -131,7 +131,7 @@ func runController(ctx context.Context) error {
labels = append(labels, string(taskassignment.BackendPod)) labels = append(labels, string(taskassignment.BackendPod))
} }
if slices.Contains(config.Components, controller.VMWorker) { if slices.Contains(config.Components, controller.VMWorker) {
labels = append(labels, string(taskassignment.BackendVM)) labels = append(labels, config.VMRunnerLabel)
} }
poller := taskscheduler.Poller{ poller := taskscheduler.Poller{
Client: giteaClient, Client: giteaClient,
@@ -139,7 +139,7 @@ func runController(ctx context.Context) error {
JetStream: producerJS, SubjectBase: config.SubjectBase, JetStream: producerJS, SubjectBase: config.SubjectBase,
}}, }},
Config: taskscheduler.PollerConfig{Version: "gitea-dynamic-runner/0.4", Labels: labels, Capacity: config.PodCapacity + config.VMCapacity}, Config: taskscheduler.PollerConfig{Version: "gitea-dynamic-runner/0.4", Labels: labels, Capacity: config.PodCapacity + config.VMCapacity},
OnError: func(err error) { log.Printf("scheduler: %v", err) }, OnError: func(err error) { slog.Error("scheduler error", "component", "scheduler", "error", err) },
} }
kubernetesConfig, err := rest.InClusterConfig() kubernetesConfig, err := rest.InClusterConfig()
if err != nil { if err != nil {
@@ -190,11 +190,13 @@ func runController(ctx context.Context) error {
return fmt.Errorf("recover Pod facade claim %s: %w", assignment.ID, err) return fmt.Errorf("recover Pod facade claim %s: %w", assignment.ID, err)
} }
} }
component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendPod, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry, podPool) component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendPod, config.PodCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap, OnEvent: workerEventLogger(taskassignment.BackendPod)}, registry, podPool)
if err != nil { if err != nil {
return err return err
} }
lifecycle := podbackend.Lifecycle{Backend: backend, OnError: func(err error) { log.Printf("pod lifecycle: %v", err) }} lifecycle := podbackend.Lifecycle{Backend: backend, OnError: func(err error) {
slog.Error("backend lifecycle error", "component", "lifecycle", "backend", taskassignment.BackendPod, "error", err)
}}
components[controller.PodWorker] = runComponent(func(ctx context.Context) error { components[controller.PodWorker] = runComponent(func(ctx context.Context) error {
group, groupContext := errgroup.WithContext(ctx) group, groupContext := errgroup.WithContext(ctx)
group.Go(func() error { return component.Run(groupContext) }) group.Go(func() error { return component.Run(groupContext) })
@@ -220,11 +222,13 @@ func runController(ctx context.Context) error {
return fmt.Errorf("recover VM facade claim %s: %w", assignment.ID, err) return fmt.Errorf("recover VM facade claim %s: %w", assignment.ID, err)
} }
} }
component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendVM, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap}, registry, vmPool) component, err := workerComponent(ctx, workerJS, config, taskassignment.BackendVM, config.VMCapacity, taskworker.Worker{Backend: backend, Bootstrap: bootstrap, OnEvent: workerEventLogger(taskassignment.BackendVM)}, registry, vmPool)
if err != nil { if err != nil {
return err return err
} }
lifecycleReconciler := opensandboxbackend.LifecycleReconciler{Backend: backend, OnError: func(err error) { log.Printf("VM lifecycle: %v", err) }} lifecycleReconciler := opensandboxbackend.LifecycleReconciler{Backend: backend, OnError: func(err error) {
slog.Error("backend lifecycle error", "component", "lifecycle", "backend", taskassignment.BackendVM, "error", err)
}}
components[controller.VMWorker] = runComponent(func(ctx context.Context) error { components[controller.VMWorker] = runComponent(func(ctx context.Context) error {
group, groupContext := errgroup.WithContext(ctx) group, groupContext := errgroup.WithContext(ctx)
group.Go(func() error { return component.Run(groupContext) }) group.Go(func() error { return component.Run(groupContext) })
@@ -294,11 +298,21 @@ func workerComponent(ctx context.Context, js jetstream.JetStream, config control
} }
return assignmentqueue.ConsumerComponent{ return assignmentqueue.ConsumerComponent{
Consumer: consumer, Capacity: capacity, Consumer: consumer, Capacity: capacity,
Processor: assignmentqueue.Processor{TrustDomain: config.TrustDomain, Accepter: accepter, Claims: claims, Admission: admission}, Processor: assignmentqueue.Processor{TrustDomain: config.TrustDomain, Accepter: accepter, Claims: claims, Admission: admission, OnEvent: func(event assignmentqueue.Event) {
OnError: func(err error) { log.Printf("%s worker: %v", backend, err) }, slog.Info("assignment transition", "component", "worker", "event", event.Name, "backend", event.Backend, "assignment", event.AssignmentID, "retry_delay", event.RetryDelay)
}},
OnError: func(err error) {
slog.Error("assignment processing error", "component", "worker", "backend", backend, "error", err)
},
}, nil }, nil
} }
func workerEventLogger(backend taskassignment.Backend) func(taskworker.Event) {
return func(event taskworker.Event) {
slog.Info("executor transition", "component", "worker", "event", event.Name, "backend", backend, "assignment", event.AssignmentID, "executor", event.Executor, "phase", event.Phase)
}
}
func loadControllerConfig() (controllerConfig, error) { func loadControllerConfig() (controllerConfig, error) {
selection, err := controller.ParseSelection(os.Getenv("COMPONENTS")) selection, err := controller.ParseSelection(os.Getenv("COMPONENTS"))
if err != nil { if err != nil {
@@ -345,7 +359,7 @@ func loadControllerConfig() (controllerConfig, error) {
FacadeListen: env("RUNNER_FACADE_LISTEN", ":8443"), FacadeURL: os.Getenv("RUNNER_FACADE_URL"), FacadeSPIFFEID: os.Getenv("RUNNER_FACADE_SPIFFE_ID"), CapabilityKey: []byte(capabilityKey), FacadeListen: env("RUNNER_FACADE_LISTEN", ":8443"), FacadeURL: os.Getenv("RUNNER_FACADE_URL"), FacadeSPIFFEID: os.Getenv("RUNNER_FACADE_SPIFFE_ID"), CapabilityKey: []byte(capabilityKey),
PodNamespace: env("POD_NAMESPACE", "gitea-actions"), PodImage: os.Getenv("POD_EXECUTOR_IMAGE"), PodServiceAccount: env("POD_SERVICE_ACCOUNT", "gitea-task-executor"), PodNamespace: env("POD_NAMESPACE", "gitea-actions"), PodImage: os.Getenv("POD_EXECUTOR_IMAGE"), PodServiceAccount: env("POD_SERVICE_ACCOUNT", "gitea-task-executor"),
SPIRECluster: env("SPIRE_CLUSTER", "homelab"), SPIREClass: env("SPIRE_CLASS", "spire-mgmt-spire"), SPIREAgentID: os.Getenv("SPIRE_AGENT_ID"), PodExecutorUID: envInt("POD_EXECUTOR_UID", 2000), PodCapacity: envInt("POD_CAPACITY", 4), SPIRECluster: env("SPIRE_CLUSTER", "homelab"), SPIREClass: env("SPIRE_CLASS", "spire-mgmt-spire"), SPIREAgentID: os.Getenv("SPIRE_AGENT_ID"), PodExecutorUID: envInt("POD_EXECUTOR_UID", 2000), PodCapacity: envInt("POD_CAPACITY", 4),
OpenSandboxURL: os.Getenv("OPENSANDBOX_API"), OpenSandboxPool: env("OPENSANDBOX_POOL", "ci-vm"), VMTimeout: envInt("VM_TIMEOUT_SECONDS", 14400), VMCapacity: envInt("VM_CAPACITY", 1), OpenSandboxURL: os.Getenv("OPENSANDBOX_API"), OpenSandboxPool: env("OPENSANDBOX_POOL", "ci-vm"), VMRunnerLabel: env("VM_RUNNER_LABEL", "vm"), VMTimeout: envInt("VM_TIMEOUT_SECONDS", 14400), VMCapacity: envInt("VM_CAPACITY", 1),
} }
if config.WorkloadAPIAddr == "" || config.FacadeURL == "" || config.FacadeSPIFFEID == "" { if config.WorkloadAPIAddr == "" || config.FacadeURL == "" || config.FacadeSPIFFEID == "" {
return controllerConfig{}, errors.New("SPIFFE_ENDPOINT_SOCKET, RUNNER_FACADE_URL, and RUNNER_FACADE_SPIFFE_ID are required") return controllerConfig{}, errors.New("SPIFFE_ENDPOINT_SOCKET, RUNNER_FACADE_URL, and RUNNER_FACADE_SPIFFE_ID are required")
@@ -357,6 +371,9 @@ func loadControllerConfig() (controllerConfig, error) {
return controllerConfig{}, errors.New("SPIRE_AGENT_ID is required for pod-worker") return controllerConfig{}, errors.New("SPIRE_AGENT_ID is required for pod-worker")
} }
if slices.Contains(selection, controller.VMWorker) { if slices.Contains(selection, controller.VMWorker) {
if config.VMRunnerLabel != "vm" && config.VMRunnerLabel != "vm-dev" {
return controllerConfig{}, errors.New("VM_RUNNER_LABEL must be vm or vm-dev")
}
if config.OpenSandboxURL == "" { if config.OpenSandboxURL == "" {
return controllerConfig{}, errors.New("OPENSANDBOX_API is required for vm-worker") return controllerConfig{}, errors.New("OPENSANDBOX_API is required for vm-worker")
} }
@@ -56,6 +56,7 @@ func TestLoadControllerConfigRequiresOpenSandboxSecretOnlyForVM(t *testing.T) {
t.Setenv("RUNNER_FACADE_URL", "https://facade:8443") t.Setenv("RUNNER_FACADE_URL", "https://facade:8443")
t.Setenv("RUNNER_FACADE_SPIFFE_ID", "spiffe://ddupan.top/controller") t.Setenv("RUNNER_FACADE_SPIFFE_ID", "spiffe://ddupan.top/controller")
t.Setenv("OPENSANDBOX_API", "http://opensandbox.internal") t.Setenv("OPENSANDBOX_API", "http://opensandbox.internal")
t.Setenv("VM_RUNNER_LABEL", "vm-dev")
if _, err := loadControllerConfig(); err == nil { if _, err := loadControllerConfig(); err == nil {
t.Fatal("expected missing OpenSandbox API key file error") t.Fatal("expected missing OpenSandbox API key file error")
} }
+3 -1
View File
@@ -4,6 +4,7 @@ import (
"context" "context"
"errors" "errors"
"fmt" "fmt"
"log/slog"
"os" "os"
"os/signal" "os/signal"
"syscall" "syscall"
@@ -12,8 +13,9 @@ import (
) )
func main() { func main() {
slog.SetDefault(slog.New(slog.NewJSONHandler(os.Stderr, nil)))
if err := run(); err != nil { if err := run(); err != nil {
fmt.Fprintln(os.Stderr, err) slog.Error("runner stopped", "error", err)
os.Exit(1) os.Exit(1)
} }
} }
+26 -1
View File
@@ -78,6 +78,22 @@ type Processor struct {
Admission Admission Admission Admission
RetryDelay time.Duration RetryDelay time.Duration
ClaimTimeout time.Duration ClaimTimeout time.Duration
OnEvent func(Event)
}
// Event describes a non-sensitive assignment handoff transition. It never
// contains task payloads, credentials, capabilities, or workload identities.
type Event struct {
Name string
AssignmentID string
Backend taskassignment.Backend
RetryDelay time.Duration
}
func (p Processor) event(name string, assignment taskassignment.Assignment, retryDelay time.Duration) {
if p.OnEvent != nil {
p.OnEvent(Event{Name: name, AssignmentID: assignment.ID, Backend: assignment.Backend, RetryDelay: retryDelay})
}
} }
func (p Processor) Process(ctx context.Context, message Message) error { func (p Processor) Process(ctx context.Context, message Message) error {
@@ -88,6 +104,7 @@ func (p Processor) Process(ctx context.Context, message Message) error {
if err != nil { if err != nil {
return errors.Join(err, message.TermWithReason("invalid assignment")) return errors.Join(err, message.TermWithReason("invalid assignment"))
} }
p.event("received", assignment, 0)
if _, err := p.Claims.Offer(assignment); err != nil { if _, err := p.Claims.Offer(assignment); err != nil {
return errors.Join(err, message.TermWithReason("conflicting assignment")) return errors.Join(err, message.TermWithReason("conflicting assignment"))
} }
@@ -96,8 +113,10 @@ func (p Processor) Process(ctx context.Context, message Message) error {
if delay <= 0 { if delay <= 0 {
delay = 2 * time.Second delay = 2 * time.Second
} }
p.event("capacity_wait", assignment, delay)
return message.NakWithDelay(delay) return message.NakWithDelay(delay)
} }
p.event("capacity_acquired", assignment, 0)
accepted, err := p.Accepter.Accept(ctx, assignment) accepted, err := p.Accepter.Accept(ctx, assignment)
if err != nil { if err != nil {
p.Admission.Release(assignment.ID) p.Admission.Release(assignment.ID)
@@ -105,9 +124,11 @@ func (p Processor) Process(ctx context.Context, message Message) error {
if delay <= 0 { if delay <= 0 {
delay = 15 * time.Second delay = 15 * time.Second
} }
p.event("backend_retry", assignment, delay)
return errors.Join(err, message.NakWithDelay(delay)) return errors.Join(err, message.NakWithDelay(delay))
} }
if accepted { if accepted {
p.event("backend_ready", assignment, 0)
timeout := p.ClaimTimeout timeout := p.ClaimTimeout
if timeout <= 0 { if timeout <= 0 {
timeout = 4 * time.Minute timeout = 4 * time.Minute
@@ -120,17 +141,21 @@ func (p Processor) Process(ctx context.Context, message Message) error {
if delay <= 0 { if delay <= 0 {
delay = 2 * time.Second delay = 2 * time.Second
} }
p.event("claim_timeout", assignment, delay)
return errors.Join(err, message.NakWithDelay(delay)) return errors.Join(err, message.NakWithDelay(delay))
} }
p.event("runner_claimed", assignment, 0)
if err := message.DoubleAck(ctx); err != nil { if err := message.DoubleAck(ctx); err != nil {
return fmt.Errorf("ack assignment %s: %w", assignment.ID, err) return fmt.Errorf("ack assignment %s: %w", assignment.ID, err)
} }
p.event("acked", assignment, 0)
return nil return nil
} }
delay := p.RetryDelay delay := p.RetryDelay
if delay <= 0 { if delay <= 0 {
delay = 2 * time.Second delay = 2 * time.Second
} }
p.event("backend_pending", assignment, delay)
return message.NakWithDelay(delay) return message.NakWithDelay(delay)
} }
@@ -167,7 +192,7 @@ func OpenConsumer(ctx context.Context, manager consumerManager, stream, subjectB
AckPolicy: jetstream.AckExplicitPolicy, AckPolicy: jetstream.AckExplicitPolicy,
AckWait: 5 * time.Minute, AckWait: 5 * time.Minute,
MaxAckPending: capacity, MaxAckPending: capacity,
MaxDeliver: 20, MaxDeliver: 1000,
}) })
if err != nil { if err != nil {
return nil, fmt.Errorf("open %s assignment consumer: %w", backend, err) return nil, fmt.Errorf("open %s assignment consumer: %w", backend, err)
+12 -2
View File
@@ -126,13 +126,23 @@ func encodedAssignment(t *testing.T) []byte {
func TestProcessorAcknowledgesPersistedHandoff(t *testing.T) { func TestProcessorAcknowledgesPersistedHandoff(t *testing.T) {
message := &fakeMessage{data: encodedAssignment(t)} message := &fakeMessage{data: encodedAssignment(t)}
processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{accepted: true}, Claims: &fakeClaims{claimed: true}, Admission: &fakeAdmission{allowed: true}} var events []Event
processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{accepted: true}, Claims: &fakeClaims{claimed: true}, Admission: &fakeAdmission{allowed: true}, OnEvent: func(event Event) { events = append(events, event) }}
if err := processor.Process(context.Background(), message); err != nil { if err := processor.Process(context.Background(), message); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if message.acked != 1 || message.nacked != 0 { if message.acked != 1 || message.nacked != 0 {
t.Fatalf("message = %#v", message) t.Fatalf("message = %#v", message)
} }
want := []string{"received", "capacity_acquired", "backend_ready", "runner_claimed", "acked"}
if len(events) != len(want) {
t.Fatalf("events = %#v", events)
}
for index := range want {
if events[index].Name != want[index] || events[index].AssignmentID != "gitea-task-42" || events[index].Backend != taskassignment.BackendPod {
t.Fatalf("event[%d] = %#v", index, events[index])
}
}
} }
func TestProcessorRetriesUntilBackendHandoffIsDurable(t *testing.T) { func TestProcessorRetriesUntilBackendHandoffIsDurable(t *testing.T) {
@@ -205,7 +215,7 @@ func TestOpenConsumerUsesIndependentDurablePerBackend(t *testing.T) {
if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.BackendPod, 4); err != nil { if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.BackendPod, 4); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if manager.config.Durable != "pod" || manager.config.FilterSubject != "ci.assignment.pod" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 { if manager.config.Durable != "pod" || manager.config.FilterSubject != "ci.assignment.pod" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 || manager.config.MaxDeliver != 1000 {
t.Fatalf("config = %#v", manager.config) t.Fatalf("config = %#v", manager.config)
} }
} }
+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 {
+28
View File
@@ -71,6 +71,28 @@ type Worker struct {
Backend Backend Backend Backend
Tasks TaskState Tasks TaskState
Bootstrap Bootstrap Bootstrap Bootstrap
OnEvent func(Event)
}
// Event describes a backend lifecycle transition without exposing launch
// environment values or other credentials.
type Event struct {
Name string
AssignmentID string
Executor string
Phase Phase
}
func (w Worker) event(name string, assignment taskassignment.Assignment, executor *Executor) {
if w.OnEvent == nil {
return
}
event := Event{Name: name, AssignmentID: assignment.ID}
if executor != nil {
event.Executor = executor.Name
event.Phase = executor.Phase
}
w.OnEvent(event)
} }
// Accept completes the durable handoff from JetStream to the backend. Once it // Accept completes the durable handoff from JetStream to the backend. Once it
@@ -88,6 +110,7 @@ func (w Worker) Accept(ctx context.Context, assignment taskassignment.Assignment
return false, err return false, err
} }
if executor == nil { if executor == nil {
w.event("executor_absent", assignment, nil)
launch, launchErr := w.launchSpec(assignment) launch, launchErr := w.launchSpec(assignment)
if launchErr != nil { if launchErr != nil {
return false, launchErr return false, launchErr
@@ -96,13 +119,18 @@ func (w Worker) Accept(ctx context.Context, assignment taskassignment.Assignment
if err != nil { if err != nil {
return false, err return false, err
} }
w.event("executor_created", assignment, executor)
} else {
w.event("executor_found", assignment, executor)
} }
if executor.IdentityTarget == "" { if executor.IdentityTarget == "" {
w.event("identity_target_pending", assignment, executor)
return false, nil return false, nil
} }
if err := w.Backend.BindIdentity(ctx, executor, assignment.Identity); err != nil { if err := w.Backend.BindIdentity(ctx, executor, assignment.Identity); err != nil {
return false, err return false, err
} }
w.event("identity_bound", assignment, executor)
return true, nil return true, nil
} }
+14 -1
View File
@@ -79,7 +79,8 @@ func TestHandleRecoversExistingExecutorWithoutCreatingAnother(t *testing.T) {
func TestAcceptAcknowledgesAfterBackendAndIdentityAreDurable(t *testing.T) { func TestAcceptAcknowledgesAfterBackendAndIdentityAreDurable(t *testing.T) {
backend := &fakeBackend{} backend := &fakeBackend{}
worker := Worker{Backend: backend, Bootstrap: fakeBootstrap{}} var events []Event
worker := Worker{Backend: backend, Bootstrap: fakeBootstrap{}, OnEvent: func(event Event) { events = append(events, event) }}
accepted, err := worker.Accept(context.Background(), assignment()) accepted, err := worker.Accept(context.Background(), assignment())
if err != nil || !accepted { if err != nil || !accepted {
@@ -88,6 +89,18 @@ func TestAcceptAcknowledgesAfterBackendAndIdentityAreDurable(t *testing.T) {
if backend.created != 1 || backend.bound != 1 || backend.deleted != 0 { if backend.created != 1 || backend.bound != 1 || backend.deleted != 0 {
t.Fatalf("created=%d bound=%d deleted=%d", backend.created, backend.bound, backend.deleted) t.Fatalf("created=%d bound=%d deleted=%d", backend.created, backend.bound, backend.deleted)
} }
want := []string{"executor_absent", "executor_created", "identity_bound"}
if len(events) != len(want) {
t.Fatalf("events = %#v", events)
}
for index := range want {
if events[index].Name != want[index] || events[index].AssignmentID != "gitea-task-42" {
t.Fatalf("event[%d] = %#v", index, events[index])
}
}
if events[1].Executor != "executor" || events[1].Phase != PhaseRunning {
t.Fatalf("created event = %#v", events[1])
}
} }
func TestAcceptRetriesWhileBackendIdentityTargetIsUnavailable(t *testing.T) { func TestAcceptRetriesWhileBackendIdentityTargetIsUnavailable(t *testing.T) {