212 lines
6.2 KiB
Go
212 lines
6.2 KiB
Go
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 || manager.config.MaxDeliver != 1000 {
|
|
t.Fatalf("config = %#v", manager.config)
|
|
}
|
|
}
|