Merge pull request '增加 Runner 生命周期结构化日志' (#37) from fix/structured-runner-lifecycle-logs into main
This commit was merged in pull request #37.
This commit is contained in:
@@ -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: |
|
||||||
|
|||||||
@@ -4,7 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
"os"
|
"os"
|
||||||
"slices"
|
"slices"
|
||||||
@@ -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 {
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
@@ -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
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
Reference in New Issue
Block a user