Files
gitea-dynamic-runner/internal/assignmentqueue/jetstream.go
T

221 lines
6.4 KiB
Go

// Package assignmentqueue implements the durable assignment handoff with JetStream.
package assignmentqueue
import (
"context"
"errors"
"fmt"
"strings"
"sync"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskassignment"
)
type publishAPI interface {
PublishMsg(context.Context, *nats.Msg, ...jetstream.PublishOpt) (*jetstream.PubAck, error)
}
// Publisher implements the scheduler dispatcher with one subject per backend.
type Publisher struct {
JetStream publishAPI
SubjectBase string
}
func (p Publisher) Dispatch(ctx context.Context, assignment taskassignment.Assignment) error {
if p.JetStream == nil {
return errors.New("JetStream publisher is required")
}
body, err := taskassignment.Marshal(assignment)
if err != nil {
return err
}
base := strings.TrimSuffix(p.SubjectBase, ".")
if base == "" {
return errors.New("assignment subject base is required")
}
message := &nats.Msg{
Subject: base + "." + string(assignment.Backend),
Header: nats.Header{jetstream.MsgIDHeader: []string{assignment.ID}},
Data: body,
}
if _, err := p.JetStream.PublishMsg(ctx, message); err != nil {
return fmt.Errorf("publish assignment %s: %w", assignment.ID, err)
}
return nil
}
type Handler interface {
Handle(context.Context, taskassignment.Assignment) (bool, error)
}
// Message is the subset of jetstream.Msg needed by one reconciliation.
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
}
func (p Processor) Process(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"))
}
done, err := p.Handler.Handle(ctx, assignment)
if err != nil {
delay := p.RetryDelay
if delay <= 0 {
delay = 15 * time.Second
}
return errors.Join(err, 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)
}
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")
}
}