diff --git a/docs/runner-protocol-roadmap.md b/docs/runner-protocol-roadmap.md index 5b66a88..97a45ab 100644 --- a/docs/runner-protocol-roadmap.md +++ b/docs/runner-protocol-roadmap.md @@ -51,6 +51,9 @@ 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 共享 task/executor 协议,只有环境创建和销毁实现不同。 ## 实现顺序 diff --git a/internal/assignmentqueue/jetstream.go b/internal/assignmentqueue/jetstream.go index 25e0c57..59cc49f 100644 --- a/internal/assignmentqueue/jetstream.go +++ b/internal/assignmentqueue/jetstream.go @@ -6,6 +6,7 @@ import ( "errors" "fmt" "strings" + "sync" "time" "github.com/nats-io/nats.go" @@ -62,9 +63,10 @@ type Message interface { // Processor maps one delivery to one idempotent worker reconciliation. type Processor struct { - TrustDomain string - Handler Handler - RetryDelay time.Duration + TrustDomain string + Handler Handler + RetryDelay time.Duration + PollInterval time.Duration } func (p Processor) Process(ctx context.Context, message Message) error { @@ -94,3 +96,125 @@ func (p Processor) Process(ctx context.Context, message Message) error { } 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: + } + } +} + +type consumeAPI interface { + Consume(jetstream.MessageHandler, ...jetstream.PullConsumeOpt) (jetstream.ConsumeContext, error) +} + +// ConsumerComponent runs bounded reconciliation goroutines for one durable +// backend consumer. The goroutine set is operational state, not task storage. +type ConsumerComponent struct { + Consumer consumeAPI + Processor Processor + Capacity int + OnError func(error) +} + +type consumerManager interface { + CreateOrUpdateConsumer(context.Context, string, jetstream.ConsumerConfig) (jetstream.Consumer, error) +} + +// OpenConsumer creates the durable backend cursor. Capacity is enforced both +// server-side and by ConsumerComponent's local semaphore. +func OpenConsumer(ctx context.Context, manager consumerManager, stream, subjectBase string, backend taskassignment.Backend, capacity int) (jetstream.Consumer, error) { + if manager == nil || stream == "" || strings.TrimSuffix(subjectBase, ".") == "" || capacity < 1 { + return nil, errors.New("JetStream manager, stream, subject base, and positive capacity are required") + } + if backend != taskassignment.BackendPod && backend != taskassignment.BackendVM { + return nil, fmt.Errorf("unsupported assignment backend %q", backend) + } + consumer, err := manager.CreateOrUpdateConsumer(ctx, stream, jetstream.ConsumerConfig{ + Name: string(backend), + Durable: string(backend), + FilterSubject: strings.TrimSuffix(subjectBase, ".") + "." + string(backend), + AckPolicy: jetstream.AckExplicitPolicy, + AckWait: 5 * time.Minute, + MaxAckPending: capacity, + MaxDeliver: 20, + }) + if err != nil { + return nil, fmt.Errorf("open %s assignment consumer: %w", backend, err) + } + return consumer, nil +} + +func (c ConsumerComponent) Run(ctx context.Context) error { + if c.Consumer == nil || c.Capacity < 1 { + return errors.New("JetStream consumer and positive capacity are required") + } + semaphore := make(chan struct{}, c.Capacity) + var workers sync.WaitGroup + consumeContext, err := c.Consumer.Consume(func(message jetstream.Msg) { + select { + case semaphore <- struct{}{}: + case <-ctx.Done(): + return + } + workers.Add(1) + 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 { + c.OnError(err) + } + }() + }, jetstream.PullMaxMessages(c.Capacity)) + if err != nil { + return fmt.Errorf("start JetStream consumer: %w", err) + } + + select { + case <-ctx.Done(): + consumeContext.Stop() + <-consumeContext.Closed() + workers.Wait() + return nil + case <-consumeContext.Closed(): + workers.Wait() + return errors.New("JetStream consumer stopped unexpectedly") + } +} diff --git a/internal/assignmentqueue/jetstream_test.go b/internal/assignmentqueue/jetstream_test.go index 13619de..cdb4f22 100644 --- a/internal/assignmentqueue/jetstream_test.go +++ b/internal/assignmentqueue/jetstream_test.go @@ -53,11 +53,16 @@ func TestPublisherUsesBackendSubjectAndAssignmentDeduplication(t *testing.T) { } type fakeHandler struct { - done bool - err error + done bool + err error + remaining int } -func (h fakeHandler) Handle(context.Context, taskassignment.Assignment) (bool, 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 } @@ -86,7 +91,7 @@ func encodedAssignment(t *testing.T) []byte { func TestProcessorAcknowledgesOnlyCompletedAssignment(t *testing.T) { message := &fakeMessage{data: encodedAssignment(t)} - processor := Processor{TrustDomain: "ddupan.top", Handler: fakeHandler{done: true}} + processor := Processor{TrustDomain: "ddupan.top", Handler: &fakeHandler{done: true}} if err := processor.Process(context.Background(), message); err != nil { t.Fatal(err) } @@ -97,7 +102,7 @@ func TestProcessorAcknowledgesOnlyCompletedAssignment(t *testing.T) { func TestProcessorKeepsRunningAssignmentPending(t *testing.T) { message := &fakeMessage{data: encodedAssignment(t)} - processor := Processor{TrustDomain: "ddupan.top", Handler: fakeHandler{}} + processor := Processor{TrustDomain: "ddupan.top", Handler: &fakeHandler{}} if err := processor.Process(context.Background(), message); err != nil { t.Fatal(err) } @@ -110,7 +115,7 @@ func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T) retry := &fakeMessage{data: encodedAssignment(t)} processor := Processor{ TrustDomain: "ddupan.top", - Handler: fakeHandler{err: errors.New("backend unavailable")}, + Handler: &fakeHandler{err: errors.New("backend unavailable")}, RetryDelay: time.Minute, } if err := processor.Process(context.Background(), retry); err == nil { @@ -128,3 +133,52 @@ func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T) t.Fatalf("poison message = %#v", poison) } } + +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) { + m.config = config + return nil, nil +} + +func TestOpenConsumerUsesIndependentDurablePerBackend(t *testing.T) { + manager := &fakeConsumerManager{} + if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.BackendPod, 4); err != nil { + t.Fatal(err) + } + if manager.config.Durable != "pod" || manager.config.FilterSubject != "ci.assignment.pod" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 { + t.Fatalf("config = %#v", manager.config) + } +}