实现持久化 assignment handoff
This commit is contained in:
@@ -0,0 +1,96 @@
|
||||
// Package assignmentqueue implements the durable assignment handoff with JetStream.
|
||||
package assignmentqueue
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"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
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,130 @@
|
||||
package assignmentqueue
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
runnerv1 "gitea.dev/actionslib/runner/v1"
|
||||
"github.com/nats-io/nats.go"
|
||||
"github.com/nats-io/nats.go/jetstream"
|
||||
"google.golang.org/protobuf/types/known/structpb"
|
||||
|
||||
"git.ddupan.top/panxiao81/gitea-dynamic-runner/internal/taskassignment"
|
||||
)
|
||||
|
||||
func testAssignment(t *testing.T) taskassignment.Assignment {
|
||||
t.Helper()
|
||||
fields, err := structpb.NewStruct(map[string]any{"repository": "owner/repo"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assignment, err := taskassignment.New(&runnerv1.Task{
|
||||
Id: 42,
|
||||
Context: fields,
|
||||
WorkflowPayload: []byte("jobs:\n publish:\n runs-on: [self-hosted, pod]\n steps: []\n"),
|
||||
}, "ddupan.top")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return assignment
|
||||
}
|
||||
|
||||
type fakePublisher struct{ message *nats.Msg }
|
||||
|
||||
func (p *fakePublisher) PublishMsg(_ context.Context, message *nats.Msg, _ ...jetstream.PublishOpt) (*jetstream.PubAck, error) {
|
||||
p.message = message
|
||||
return &jetstream.PubAck{}, nil
|
||||
}
|
||||
|
||||
func TestPublisherUsesBackendSubjectAndAssignmentDeduplication(t *testing.T) {
|
||||
api := &fakePublisher{}
|
||||
publisher := Publisher{JetStream: api, SubjectBase: "ci.assignment"}
|
||||
if err := publisher.Dispatch(context.Background(), testAssignment(t)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if api.message.Subject != "ci.assignment.pod" {
|
||||
t.Fatalf("subject = %q", api.message.Subject)
|
||||
}
|
||||
if api.message.Header.Get(jetstream.MsgIDHeader) != "gitea-task-42" {
|
||||
t.Fatalf("message ID = %q", api.message.Header.Get(jetstream.MsgIDHeader))
|
||||
}
|
||||
}
|
||||
|
||||
type fakeHandler struct {
|
||||
done bool
|
||||
err error
|
||||
}
|
||||
|
||||
func (h fakeHandler) Handle(context.Context, taskassignment.Assignment) (bool, error) {
|
||||
return h.done, h.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 {
|
||||
t.Helper()
|
||||
data, err := taskassignment.Marshal(testAssignment(t))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return data
|
||||
}
|
||||
|
||||
func TestProcessorAcknowledgesOnlyCompletedAssignment(t *testing.T) {
|
||||
message := &fakeMessage{data: encodedAssignment(t)}
|
||||
processor := Processor{TrustDomain: "ddupan.top", Handler: fakeHandler{done: true}}
|
||||
if err := processor.Process(context.Background(), message); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if message.acked != 1 || message.inProgress != 0 || message.nacked != 0 {
|
||||
t.Fatalf("message = %#v", message)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProcessorKeepsRunningAssignmentPending(t *testing.T) {
|
||||
message := &fakeMessage{data: encodedAssignment(t)}
|
||||
processor := Processor{TrustDomain: "ddupan.top", Handler: fakeHandler{}}
|
||||
if err := processor.Process(context.Background(), message); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if message.acked != 0 || message.inProgress != 1 {
|
||||
t.Fatalf("message = %#v", message)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T) {
|
||||
retry := &fakeMessage{data: encodedAssignment(t)}
|
||||
processor := Processor{
|
||||
TrustDomain: "ddupan.top",
|
||||
Handler: fakeHandler{err: errors.New("backend unavailable")},
|
||||
RetryDelay: time.Minute,
|
||||
}
|
||||
if err := processor.Process(context.Background(), retry); err == nil {
|
||||
t.Fatal("expected backend error")
|
||||
}
|
||||
if retry.nacked != time.Minute {
|
||||
t.Fatalf("retry delay = %s", retry.nacked)
|
||||
}
|
||||
|
||||
poison := &fakeMessage{data: []byte("not-json")}
|
||||
if err := processor.Process(context.Background(), poison); err == nil {
|
||||
t.Fatal("expected decode error")
|
||||
}
|
||||
if poison.terminated != 1 || poison.nacked != 0 {
|
||||
t.Fatalf("poison message = %#v", poison)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user