diff --git a/.gitea/workflows/test.yml b/.gitea/workflows/test.yml index cb40e4d..75d692b 100644 --- a/.gitea/workflows/test.yml +++ b/.gitea/workflows/test.yml @@ -17,7 +17,7 @@ jobs: - run: go vet ./... python: - runs-on: self-hosted + runs-on: [self-hosted, pod] steps: - uses: actions/checkout@v4 - uses: actions/setup-python@v5 @@ -28,7 +28,7 @@ jobs: - run: python -m compileall -q src shell: - runs-on: self-hosted + runs-on: [self-hosted, pod] steps: - uses: actions/checkout@v4 - run: | diff --git a/cmd/gitea-dynamic-runner/controller.go b/cmd/gitea-dynamic-runner/controller.go index c3392d1..0d753f4 100644 --- a/cmd/gitea-dynamic-runner/controller.go +++ b/cmd/gitea-dynamic-runner/controller.go @@ -4,7 +4,7 @@ import ( "context" "errors" "fmt" - "log" + "log/slog" "net/http" "os" "slices" @@ -139,7 +139,7 @@ func runController(ctx context.Context) error { JetStream: producerJS, SubjectBase: config.SubjectBase, }}, 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() 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) } } - 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 { 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 { group, groupContext := errgroup.WithContext(ctx) 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) } } - 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 { 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 { group, groupContext := errgroup.WithContext(ctx) 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{ Consumer: consumer, Capacity: capacity, - Processor: assignmentqueue.Processor{TrustDomain: config.TrustDomain, Accepter: accepter, Claims: claims, Admission: admission}, - OnError: func(err error) { log.Printf("%s worker: %v", backend, err) }, + Processor: assignmentqueue.Processor{TrustDomain: config.TrustDomain, Accepter: accepter, Claims: claims, Admission: admission, OnEvent: func(event assignmentqueue.Event) { + slog.Info("assignment transition", "component", "worker", "event", event.Name, "backend", event.Backend, "assignment", event.AssignmentID, "retry_delay", event.RetryDelay) + }}, + OnError: func(err error) { + slog.Error("assignment processing error", "component", "worker", "backend", backend, "error", err) + }, }, 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) { selection, err := controller.ParseSelection(os.Getenv("COMPONENTS")) if err != nil { diff --git a/cmd/gitea-dynamic-runner/main.go b/cmd/gitea-dynamic-runner/main.go index 5870bee..0029d63 100644 --- a/cmd/gitea-dynamic-runner/main.go +++ b/cmd/gitea-dynamic-runner/main.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "log/slog" "os" "os/signal" "syscall" @@ -12,8 +13,9 @@ import ( ) func main() { + slog.SetDefault(slog.New(slog.NewJSONHandler(os.Stderr, nil))) if err := run(); err != nil { - fmt.Fprintln(os.Stderr, err) + slog.Error("runner stopped", "error", err) os.Exit(1) } } diff --git a/internal/assignmentqueue/jetstream.go b/internal/assignmentqueue/jetstream.go index c3d8f14..74e8ce6 100644 --- a/internal/assignmentqueue/jetstream.go +++ b/internal/assignmentqueue/jetstream.go @@ -78,6 +78,22 @@ type Processor struct { Admission Admission RetryDelay 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 { @@ -88,6 +104,7 @@ func (p Processor) Process(ctx context.Context, message Message) error { if err != nil { return errors.Join(err, message.TermWithReason("invalid assignment")) } + p.event("received", assignment, 0) if _, err := p.Claims.Offer(assignment); err != nil { 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 { delay = 2 * time.Second } + p.event("capacity_wait", assignment, delay) return message.NakWithDelay(delay) } + p.event("capacity_acquired", assignment, 0) accepted, err := p.Accepter.Accept(ctx, assignment) if err != nil { p.Admission.Release(assignment.ID) @@ -105,9 +124,11 @@ func (p Processor) Process(ctx context.Context, message Message) error { if delay <= 0 { delay = 15 * time.Second } + p.event("backend_retry", assignment, delay) return errors.Join(err, message.NakWithDelay(delay)) } if accepted { + p.event("backend_ready", assignment, 0) timeout := p.ClaimTimeout if timeout <= 0 { timeout = 4 * time.Minute @@ -120,17 +141,21 @@ func (p Processor) Process(ctx context.Context, message Message) error { if delay <= 0 { delay = 2 * time.Second } + p.event("claim_timeout", assignment, delay) return errors.Join(err, message.NakWithDelay(delay)) } + p.event("runner_claimed", assignment, 0) if err := message.DoubleAck(ctx); err != nil { return fmt.Errorf("ack assignment %s: %w", assignment.ID, err) } + p.event("acked", assignment, 0) return nil } delay := p.RetryDelay if delay <= 0 { delay = 2 * time.Second } + p.event("backend_pending", assignment, delay) return message.NakWithDelay(delay) } diff --git a/internal/assignmentqueue/jetstream_test.go b/internal/assignmentqueue/jetstream_test.go index 4bc39bf..b053208 100644 --- a/internal/assignmentqueue/jetstream_test.go +++ b/internal/assignmentqueue/jetstream_test.go @@ -126,13 +126,23 @@ func encodedAssignment(t *testing.T) []byte { func TestProcessorAcknowledgesPersistedHandoff(t *testing.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 { t.Fatal(err) } if message.acked != 1 || message.nacked != 0 { 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) { diff --git a/internal/taskworker/worker.go b/internal/taskworker/worker.go index 6c3b703..ea590ad 100644 --- a/internal/taskworker/worker.go +++ b/internal/taskworker/worker.go @@ -71,6 +71,28 @@ type Worker struct { Backend Backend Tasks TaskState 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 @@ -88,6 +110,7 @@ func (w Worker) Accept(ctx context.Context, assignment taskassignment.Assignment return false, err } if executor == nil { + w.event("executor_absent", assignment, nil) launch, launchErr := w.launchSpec(assignment) if launchErr != nil { return false, launchErr @@ -96,13 +119,18 @@ func (w Worker) Accept(ctx context.Context, assignment taskassignment.Assignment if err != nil { return false, err } + w.event("executor_created", assignment, executor) + } else { + w.event("executor_found", assignment, executor) } if executor.IdentityTarget == "" { + w.event("identity_target_pending", assignment, executor) return false, nil } if err := w.Backend.BindIdentity(ctx, executor, assignment.Identity); err != nil { return false, err } + w.event("identity_bound", assignment, executor) return true, nil } diff --git a/internal/taskworker/worker_test.go b/internal/taskworker/worker_test.go index 3729bb6..efdf111 100644 --- a/internal/taskworker/worker_test.go +++ b/internal/taskworker/worker_test.go @@ -79,7 +79,8 @@ func TestHandleRecoversExistingExecutorWithoutCreatingAnother(t *testing.T) { func TestAcceptAcknowledgesAfterBackendAndIdentityAreDurable(t *testing.T) { 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()) if err != nil || !accepted { @@ -88,6 +89,18 @@ func TestAcceptAcknowledgesAfterBackendAndIdentityAreDurable(t *testing.T) { if backend.created != 1 || backend.bound != 1 || backend.deleted != 0 { 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) {