Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259package engine
import ( "context" "errors" "io" "log/slog" "sync" "testing" "time"
"tangled.org/core/spindle/config" "tangled.org/core/spindle/models" "tangled.org/core/spindle/observability" "tangled.org/core/spindle/quota")
type spyStore struct { mu sync.Mutex reserveCalls []quota.ReserveRequest releaseCalls []string reserveErr error}
func (s *spyStore) Reserve(ctx context.Context, req quota.ReserveRequest) (quota.Reservation, error) { s.mu.Lock() defer s.mu.Unlock() s.reserveCalls = append(s.reserveCalls, req) if s.reserveErr != nil { return quota.Reservation{}, s.reserveErr } return quota.Reservation{ ID: "res-123", Allowed: true, Temporary: false, }, nil}
func (s *spyStore) BeginCommit(ctx context.Context, id string) error { return nil }func (s *spyStore) Commit(ctx context.Context, id string) error { return nil }func (s *spyStore) Release(ctx context.Context, id string) error { s.mu.Lock() defer s.mu.Unlock() s.releaseCalls = append(s.releaseCalls, id) return nil}
func (s *spyStore) ListLimits(ctx context.Context) ([]quota.Limit, error) { return nil, nil }func (s *spyStore) ListUsage(ctx context.Context) ([]quota.Usage, error) { return nil, nil }func (s *spyStore) MetricsSnapshot(ctx context.Context) (quota.MetricsSnapshot, error) { return quota.MetricsSnapshot{}, nil}func (s *spyStore) Recover(ctx context.Context, liveIDs []string) error { return nil }func (s *spyStore) SetLimit(ctx context.Context, scope quota.Scope, did string, resource string, limit int64) error { return nil}func (s *spyStore) UnsetLimit(ctx context.Context, scope quota.Scope, did string, resource string) error { return nil}
type quotaReporterEngine struct { *mockEngine workflows, memory, vcpus, disk int64}
func (e *quotaReporterEngine) QuotaResources(wf *models.Workflow) quota.Resources { return quota.Resources{ quota.ResourceWorkflows: e.workflows, quota.ResourceMemoryMiB: e.memory, quota.ResourceVCPUs: e.vcpus, quota.ResourceDiskMiB: e.disk, }}
type noReporterEngine struct { *mockEngine}
func TestStartWorkflows_NoDoubleAcquisitionByNoReporterEngine(t *testing.T) { testDB := newTestDB(t) logger := slog.New(slog.NewTextHandler(io.Discard, nil)) cfg := &config.Config{Server: config.Server{LogDir: t.TempDir()}}
store := &spyStore{} qm := quota.NewManager(store, 50*time.Millisecond, nil) defer qm.Close()
eng := &noReporterEngine{mockEngine: &mockEngine{}}
wf := models.Workflow{ Name: "job1", Steps: []models.Step{mockStep{name: "step1"}}, OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo", }
pipeline := &models.Pipeline{ Workflows: map[models.Engine][]models.Workflow{ eng: {wf}, }, } pipelineId := models.PipelineId{Knot: "knot", Rkey: "rkey"}
StartWorkflows(logger, nil, cfg, qm, nil, testDB, nil, nil, nil, context.Background(), pipeline, pipelineId)
store.mu.Lock() resCount := len(store.reserveCalls) store.mu.Unlock()
if resCount != 0 { t.Fatalf("expected 0 quota reservations for engine without reporter, got %d", resCount) }}
func TestStartWorkflows_QuotaAcquisitionAndIdentity(t *testing.T) { testDB := newTestDB(t) logger := slog.New(slog.NewTextHandler(io.Discard, nil)) cfg := &config.Config{Server: config.Server{LogDir: t.TempDir()}}
store := &spyStore{} qm := quota.NewManager(store, 50*time.Millisecond, nil) defer qm.Close()
eng := "aReporterEngine{ mockEngine: &mockEngine{}, workflows: 1, memory: 256, vcpus: 2, disk: 512, }
wf := models.Workflow{ Name: "job1", Steps: []models.Step{mockStep{name: "step1"}}, OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo", }
pipeline := &models.Pipeline{ Workflows: map[models.Engine][]models.Workflow{ eng: {wf}, }, } pipelineId := models.PipelineId{Knot: "knot", Rkey: "rkey"}
StartWorkflows(logger, nil, cfg, qm, nil, testDB, nil, nil, nil, context.Background(), pipeline, pipelineId)
store.mu.Lock() resCount := len(store.reserveCalls) store.mu.Unlock()
if resCount != 1 { t.Fatalf("expected 1 quota reservation call, got %d", resCount) }
req := store.reserveCalls[0] if req.Identity.OwnerDID != "did:web:alice" || req.Identity.RepoDID != "did:web:alice/repo" { t.Errorf("unexpected identity in reservation: %+v", req.Identity) }
expectedID := quota.WorkflowReservationID("", "did:web:alice", "did:web:alice/repo", "knot", "rkey", "job1") if req.ID != expectedID { t.Errorf("expected reservation ID %q, got %q", expectedID, req.ID) } if req.Key != expectedID { t.Errorf("expected reservation key %q, got %q", expectedID, req.Key) }
if req.Resources[quota.ResourceWorkflows] != 1 || req.Resources[quota.ResourceMemoryMiB] != 256 || req.Resources[quota.ResourceVCPUs] != 2 || req.Resources[quota.ResourceDiskMiB] != 512 { t.Errorf("unexpected resource vectors: %+v", req.Resources) }
engNoDisk := "aReporterEngine{ mockEngine: &mockEngine{}, workflows: 1, memory: 256, vcpus: 0, disk: 0, } store.mu.Lock() store.reserveCalls = nil store.mu.Unlock()
wf2 := wf wf2.Name = "job2"
pipelineNoDisk := &models.Pipeline{ Workflows: map[models.Engine][]models.Workflow{ engNoDisk: {wf2}, }, } StartWorkflows(logger, nil, cfg, qm, nil, testDB, nil, nil, nil, context.Background(), pipelineNoDisk, pipelineId)
store.mu.Lock() resCount2 := len(store.reserveCalls) store.mu.Unlock()
if resCount2 != 1 { t.Fatalf("expected 1 quota reservation call, got %d", resCount2) } req2 := store.reserveCalls[0] if _, exists := req2.Resources[quota.ResourceVCPUs]; exists { t.Error("vcpus = 0 should be omitted from resources map") } if _, exists := req2.Resources[quota.ResourceDiskMiB]; exists { t.Error("disk_mib = 0 should be omitted from resources map") }}
func TestStartWorkflows_QuotaFailureRecordsWorkflowFailure(t *testing.T) { testDB := newTestDB(t) logger := slog.New(slog.NewTextHandler(io.Discard, nil)) cfg := &config.Config{Server: config.Server{LogDir: t.TempDir()}}
store := &spyStore{reserveErr: errors.New("quota store unavailable")} qm := quota.NewManager(store, 50*time.Millisecond, nil) defer qm.Close()
eng := "aReporterEngine{ mockEngine: &mockEngine{}, workflows: 1, } pipeline := &models.Pipeline{ RepoDid: "did:web:alice/repo", Workflows: map[models.Engine][]models.Workflow{ eng: {{ Name: "job1", Steps: []models.Step{mockStep{name: "step1"}}, OwnerDID: "did:web:alice", RepoDID: "did:web:alice/repo", }}, }, } metrics := observability.NewMetrics() ctx := observability.WithMetrics(context.Background(), metrics) StartWorkflows(logger, nil, cfg, qm, nil, testDB, nil, nil, nil, ctx, pipeline, models.PipelineId{Knot: "knot", Rkey: "rkey"})
families, err := metrics.Registry().Gather() if err != nil { t.Fatal(err) } for _, family := range families { if family.GetName() != "spindle_workflows_total" { continue } for _, metric := range family.GetMetric() { for _, label := range metric.GetLabel() { if label.GetName() == "result" && label.GetValue() == "failure" && metric.GetCounter().GetValue() == 1 { return } } } } t.Fatal("quota acquisition failure was not recorded as a failed workflow")}