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 fakeAccepter struct { accepted bool err error } type fakeClaims struct { claimed bool } type fakeAdmission struct { allowed bool active map[string]bool released int } func (a *fakeAdmission) Acquire(assignmentID string) bool { if !a.allowed { return false } if a.active == nil { a.active = make(map[string]bool) } a.active[assignmentID] = true return true } func (a *fakeAdmission) Release(assignmentID string) { delete(a.active, assignmentID) a.released++ } func (c *fakeClaims) Offer(taskassignment.Assignment) (<-chan struct{}, error) { ready := make(chan struct{}) if c.claimed { close(ready) } return ready, nil } func (c *fakeClaims) WaitClaimed(ctx context.Context, _ string) error { if c.claimed { return nil } <-ctx.Done() return ctx.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 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) 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 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}} if err := processor.Process(context.Background(), message); err != nil { t.Fatal(err) } if message.acked != 1 || message.nacked != 0 { t.Fatalf("message = %#v", message) } } func TestProcessorRetriesUntilBackendHandoffIsDurable(t *testing.T) { message := &fakeMessage{data: encodedAssignment(t)} processor := Processor{TrustDomain: "ddupan.top", Accepter: &fakeAccepter{}, Claims: &fakeClaims{}, Admission: &fakeAdmission{allowed: true}, RetryDelay: 2 * time.Second} if err := processor.Process(context.Background(), message); err != nil { t.Fatal(err) } if message.acked != 0 || message.nacked != 2*time.Second { t.Fatalf("message = %#v", message) } } func TestProcessorRetriesBackendFailureAndTerminatesPoisonMessage(t *testing.T) { retry := &fakeMessage{data: encodedAssignment(t)} admission := &fakeAdmission{allowed: true} processor := Processor{ TrustDomain: "ddupan.top", Accepter: &fakeAccepter{err: errors.New("backend unavailable")}, Claims: &fakeClaims{}, Admission: admission, 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) } if admission.released != 1 { t.Fatalf("released slots = %d", admission.released) } 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) } } func TestProcessorLeavesAssignmentPendingWhenBackendPoolIsFull(t *testing.T) { message := &fakeMessage{data: encodedAssignment(t)} accepter := &fakeAccepter{accepted: true} processor := Processor{ TrustDomain: "ddupan.top", Accepter: accepter, Claims: &fakeClaims{}, Admission: &fakeAdmission{}, RetryDelay: 3 * time.Second, } if err := processor.Process(context.Background(), message); err != nil { t.Fatal(err) } if message.acked != 0 || message.nacked != 3*time.Second { 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) } }