diff --git a/docs/runner-protocol-roadmap.md b/docs/runner-protocol-roadmap.md index 97a45ab..3e7ebba 100644 --- a/docs/runner-protocol-roadmap.md +++ b/docs/runner-protocol-roadmap.md @@ -51,9 +51,11 @@ controller 使用单一 Go 二进制;默认在同一进程启用 `scheduler` SPIFFE ID写入 sandbox metadata;通过 `extensions.poolRef=ci-vm` 使用既有 Kata Pool。 - 结果处理顺序固定为回报 Gitea、清理后端、ACK assignment;重投时先查询 Gitea 终态, 从而关闭后端已删除但消息尚未 ACK 的崩溃窗口。 -- pod 与 vm 使用独立 durable consumer 和并发上限。worker 持有消息期间持续 reconcile - 后端并发送 `InProgress`;只有完整完成才 `DoubleAck`,进程退出则保留未确认消息供 - 其他实例恢复,临时后端错误使用延迟 NAK。 +- pod 与 vm 使用独立 durable consumer 和并发上限。consumer 只负责将 assignment + 幂等落到后端;executor 与身份恢复 metadata 持久化后立即 `DoubleAck`。尚未取得 + Pod UID 等短暂未就绪状态以及临时后端错误使用延迟 NAK。 +- assignment ACK 后的运行、结果回报和清理由 backend reconciler 根据 Kubernetes、 + OpenSandbox 与 Gitea 的事实状态驱动,不继续占用 JetStream delivery。 - Pod 与 VM 共享 task/executor 协议,只有环境创建和销毁实现不同。 ## 实现顺序 diff --git a/internal/assignmentqueue/jetstream.go b/internal/assignmentqueue/jetstream.go index 59cc49f..0d22884 100644 --- a/internal/assignmentqueue/jetstream.go +++ b/internal/assignmentqueue/jetstream.go @@ -48,8 +48,8 @@ func (p Publisher) Dispatch(ctx context.Context, assignment taskassignment.Assig return nil } -type Handler interface { - Handle(context.Context, taskassignment.Assignment) (bool, error) +type Accepter interface { + Accept(context.Context, taskassignment.Assignment) (bool, error) } // Message is the subset of jetstream.Msg needed by one reconciliation. @@ -57,27 +57,25 @@ type Message interface { Data() []byte DoubleAck(context.Context) error NakWithDelay(time.Duration) error - InProgress() error TermWithReason(string) error } // Processor maps one delivery to one idempotent worker reconciliation. type Processor struct { - TrustDomain string - Handler Handler - RetryDelay time.Duration - PollInterval time.Duration + TrustDomain string + Accepter Accepter + RetryDelay time.Duration } func (p Processor) Process(ctx context.Context, message Message) error { - if p.Handler == nil { - return errors.New("assignment handler is required") + if p.Accepter == nil { + return errors.New("assignment accepter is required") } assignment, err := taskassignment.Unmarshal(message.Data(), p.TrustDomain) if err != nil { return errors.Join(err, message.TermWithReason("invalid assignment")) } - done, err := p.Handler.Handle(ctx, assignment) + accepted, err := p.Accepter.Accept(ctx, assignment) if err != nil { delay := p.RetryDelay if delay <= 0 { @@ -85,60 +83,17 @@ func (p Processor) Process(ctx context.Context, message Message) error { } return errors.Join(err, message.NakWithDelay(delay)) } - if done { + if accepted { if err := message.DoubleAck(ctx); err != nil { return fmt.Errorf("ack assignment %s: %w", assignment.ID, err) } return nil } - if err := message.InProgress(); err != nil { - return fmt.Errorf("extend assignment %s acknowledgement: %w", assignment.ID, err) - } - return nil -} - -// ProcessUntilDone holds one durable delivery while repeatedly reconciling -// backend state. Cancellation leaves it unacknowledged for another process. -func (p Processor) ProcessUntilDone(ctx context.Context, message Message) error { - if p.Handler == nil { - return errors.New("assignment handler is required") - } - assignment, err := taskassignment.Unmarshal(message.Data(), p.TrustDomain) - if err != nil { - return errors.Join(err, message.TermWithReason("invalid assignment")) - } - interval := p.PollInterval - if interval <= 0 { - interval = 2 * time.Second - } - for { - done, handleErr := p.Handler.Handle(ctx, assignment) - if handleErr != nil { - delay := p.RetryDelay - if delay <= 0 { - delay = 15 * time.Second - } - return errors.Join(handleErr, message.NakWithDelay(delay)) - } - if done { - if err := message.DoubleAck(ctx); err != nil { - return fmt.Errorf("ack assignment %s: %w", assignment.ID, err) - } - return nil - } - if err := message.InProgress(); err != nil { - return fmt.Errorf("extend assignment %s acknowledgement: %w", assignment.ID, err) - } - timer := time.NewTimer(interval) - select { - case <-ctx.Done(): - if !timer.Stop() { - <-timer.C - } - return ctx.Err() - case <-timer.C: - } + delay := p.RetryDelay + if delay <= 0 { + delay = 2 * time.Second } + return message.NakWithDelay(delay) } type consumeAPI interface { @@ -198,7 +153,7 @@ func (c ConsumerComponent) Run(ctx context.Context) error { go func() { defer workers.Done() defer func() { <-semaphore }() - if err := c.Processor.ProcessUntilDone(ctx, message); err != nil && !errors.Is(err, context.Canceled) && c.OnError != nil { + if err := c.Processor.Process(ctx, message); err != nil && !errors.Is(err, context.Canceled) && c.OnError != nil { c.OnError(err) } }() diff --git a/internal/assignmentqueue/jetstream_test.go b/internal/assignmentqueue/jetstream_test.go index cdb4f22..df58f46 100644 --- a/internal/assignmentqueue/jetstream_test.go +++ b/internal/assignmentqueue/jetstream_test.go @@ -52,32 +52,25 @@ func TestPublisherUsesBackendSubjectAndAssignmentDeduplication(t *testing.T) { } } -type fakeHandler struct { - done bool - err error - remaining int +type fakeAccepter struct { + accepted bool + err error } -func (h *fakeHandler) Handle(context.Context, taskassignment.Assignment) (bool, error) { - if h.remaining > 0 { - h.remaining-- - return false, nil - } - return h.done, h.err +func (a *fakeAccepter) Accept(context.Context, taskassignment.Assignment) (bool, error) { + return a.accepted, a.err } type fakeMessage struct { data []byte acked int nacked time.Duration - inProgress int terminated int } func (m *fakeMessage) Data() []byte { return m.data } func (m *fakeMessage) DoubleAck(context.Context) error { m.acked++; return nil } func (m *fakeMessage) NakWithDelay(delay time.Duration) error { m.nacked = delay; return nil } -func (m *fakeMessage) InProgress() error { m.inProgress++; return nil } func (m *fakeMessage) TermWithReason(string) error { m.terminated++; return nil } func encodedAssignment(t *testing.T) []byte { @@ -89,24 +82,24 @@ func encodedAssignment(t *testing.T) []byte { return data } -func TestProcessorAcknowledgesOnlyCompletedAssignment(t *testing.T) { +func TestProcessorAcknowledgesPersistedHandoff(t *testing.T) { message := &fakeMessage{data: encodedAssignment(t)} - processor := Processor{TrustDomain: "ddupan.top", Handler: &fakeHandler{done: true}} + processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{accepted: true}} if err := processor.Process(context.Background(), message); err != nil { t.Fatal(err) } - if message.acked != 1 || message.inProgress != 0 || message.nacked != 0 { + if message.acked != 1 || message.nacked != 0 { t.Fatalf("message = %#v", message) } } -func TestProcessorKeepsRunningAssignmentPending(t *testing.T) { +func TestProcessorRetriesUntilBackendHandoffIsDurable(t *testing.T) { message := &fakeMessage{data: encodedAssignment(t)} - processor := Processor{TrustDomain: "ddupan.top", Handler: &fakeHandler{}} + processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{}, RetryDelay: 2 * time.Second} if err := processor.Process(context.Background(), message); err != nil { t.Fatal(err) } - if message.acked != 0 || message.inProgress != 1 { + if message.acked != 0 || message.nacked != 2*time.Second { t.Fatalf("message = %#v", message) } } @@ -115,7 +108,7 @@ func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T) retry := &fakeMessage{data: encodedAssignment(t)} processor := Processor{ TrustDomain: "ddupan.top", - Handler: &fakeHandler{err: errors.New("backend unavailable")}, + Accepter: &fakeAccepter{err: errors.New("backend unavailable")}, RetryDelay: time.Minute, } if err := processor.Process(context.Background(), retry); err == nil { @@ -134,38 +127,6 @@ func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T) } } -func TestProcessUntilDoneReconcilesWithoutRedelivery(t *testing.T) { - message := &fakeMessage{data: encodedAssignment(t)} - processor := Processor{ - TrustDomain: "ddupan.top", - Handler: &fakeHandler{done: true, remaining: 2}, - PollInterval: time.Millisecond, - } - if err := processor.ProcessUntilDone(context.Background(), message); err != nil { - t.Fatal(err) - } - if message.inProgress != 2 || message.acked != 1 || message.nacked != 0 { - t.Fatalf("message = %#v", message) - } -} - -func TestProcessUntilDoneLeavesMessagePendingOnShutdown(t *testing.T) { - message := &fakeMessage{data: encodedAssignment(t)} - processor := Processor{ - TrustDomain: "ddupan.top", - Handler: &fakeHandler{}, - PollInterval: time.Hour, - } - ctx, cancel := context.WithCancel(context.Background()) - cancel() - if err := processor.ProcessUntilDone(ctx, message); !errors.Is(err, context.Canceled) { - t.Fatalf("error = %v", err) - } - if message.acked != 0 || message.nacked != 0 || message.terminated != 0 { - t.Fatalf("message = %#v", message) - } -} - type fakeConsumerManager struct{ config jetstream.ConsumerConfig } func (m *fakeConsumerManager) CreateOrUpdateConsumer(_ context.Context, _ string, config jetstream.ConsumerConfig) (jetstream.Consumer, error) { diff --git a/internal/taskworker/worker.go b/internal/taskworker/worker.go index 1e71ab2..863564d 100644 --- a/internal/taskworker/worker.go +++ b/internal/taskworker/worker.go @@ -57,6 +57,47 @@ type Worker struct { Tasks TaskState } +// Accept completes the durable handoff from JetStream to the backend. Once it +// returns true, all recovery information exists in Kubernetes/OpenSandbox and +// the assignment message can be acknowledged immediately. +func (w Worker) Accept(ctx context.Context, assignment taskassignment.Assignment) (bool, error) { + if w.Backend == nil || w.Tasks == nil { + return false, errors.New("backend and Gitea task state are required") + } + if assignment.ID == "" || assignment.Task == nil { + return false, errors.New("valid assignment is required") + } + terminal, err := w.Tasks.Terminal(ctx, assignment.Task.GetId()) + if err != nil { + return false, err + } + executor, err := w.Backend.Find(ctx, assignment.ID) + if err != nil { + return false, err + } + if terminal { + if executor != nil { + if err := w.Backend.Delete(ctx, executor); err != nil { + return false, err + } + } + return true, nil + } + if executor == nil { + executor, err = w.Backend.Create(ctx, assignment, BackendMetadata(assignment)) + if err != nil { + return false, err + } + } + if executor.IdentityTarget == "" { + return false, nil + } + if err := w.Backend.BindIdentity(ctx, executor, assignment.Identity); err != nil { + return false, err + } + return true, nil +} + // Handle performs one reconciliation. Done means the queue message may be // acknowledged. A false result should remain pending and be reconciled again. func (w Worker) Handle(ctx context.Context, assignment taskassignment.Assignment) (done bool, err error) { diff --git a/internal/taskworker/worker_test.go b/internal/taskworker/worker_test.go index 17b5494..4d1a5fd 100644 --- a/internal/taskworker/worker_test.go +++ b/internal/taskworker/worker_test.go @@ -71,6 +71,32 @@ func TestHandleRecoversExistingExecutorWithoutCreatingAnother(t *testing.T) { } } +func TestAcceptAcknowledgesAfterBackendAndIdentityAreDurable(t *testing.T) { + backend := &fakeBackend{} + worker := Worker{Backend: backend, Tasks: &fakeTasks{}} + + accepted, err := worker.Accept(context.Background(), assignment()) + if err != nil || !accepted { + t.Fatalf("Accept() = (%v, %v), want accepted", accepted, err) + } + if backend.created != 1 || backend.bound != 1 || backend.deleted != 0 { + t.Fatalf("created=%d bound=%d deleted=%d", backend.created, backend.bound, backend.deleted) + } +} + +func TestAcceptRetriesWhileBackendIdentityTargetIsUnavailable(t *testing.T) { + backend := &fakeBackend{executor: &Executor{Name: "pending", Phase: PhasePending}} + worker := Worker{Backend: backend, Tasks: &fakeTasks{}} + + accepted, err := worker.Accept(context.Background(), assignment()) + if err != nil || accepted { + t.Fatalf("Accept() = (%v, %v), want retry", accepted, err) + } + if backend.bound != 0 { + t.Fatalf("identity bindings = %d", backend.bound) + } +} + func TestHandleReportsBeforeCleanupAndBecomesRecoverable(t *testing.T) { backend := &fakeBackend{executor: &Executor{Name: "finished", IdentityTarget: "uid", Phase: PhaseSucceeded}} tasks := &fakeTasks{}