feat: 增加 Runner 生命周期结构化日志
This commit is contained in:
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user