package 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") }