package observability import ( "context" "encoding/json" "errors" "net/http" "net/http/httptest" "os" "path/filepath" "regexp" "strings" "sync" "sync/atomic" "testing" "time" "github.com/go-chi/chi/v5" "github.com/prometheus/client_golang/prometheus" dto "github.com/prometheus/client_model/go" "go.opentelemetry.io/otel" sdktrace "go.opentelemetry.io/otel/sdk/trace" "go.opentelemetry.io/otel/sdk/trace/tracetest" "tangled.org/core/spindle/db" ) var forbiddenLabelValues = []string{ "did:plc:123", "did:web:alice", "did:web:alice/repo1", "res-abc", "sha256:deadbeef", "dial tcp 10.0.0.1:5432: connection refused", "arbitrary", "invalid_status", "gpu_seconds", } func gather(t *testing.T, m *Metrics) map[string]*dto.MetricFamily { t.Helper() families, err := m.Registry().Gather() if err != nil { t.Fatalf("gather: %v", err) } byName := make(map[string]*dto.MetricFamily, len(families)) for _, f := range families { byName[f.GetName()] = f } return byName } func labelKeys(metric *dto.Metric) []string { keys := make([]string, 0, len(metric.GetLabel())) for _, lp := range metric.GetLabel() { keys = append(keys, lp.GetName()) } return keys } func labelValue(metric *dto.Metric, name string) string { for _, lp := range metric.GetLabel() { if lp.GetName() == name { return lp.GetValue() } } return "" } func TestQuotaMetricsAreBounded(t *testing.T) { m := NewMetrics() if m == nil { t.Fatal("failed to create metrics") } m.SetQuotaProvider(func() QuotaSnapshot { return QuotaSnapshot{ Usage: []QuotaUsage{ {Scope: "user", Resource: "cache_storage_bytes", Used: 5000}, {Scope: "repo", Resource: "cache_storage_bytes", Used: 2000}, {Scope: "user", Resource: "workflows", Used: 3}, {Scope: "user", Resource: "vcpus", Used: 8}, {Scope: "repo", Resource: "memory_mib", Used: 4096}, {Scope: "repo", Resource: "disk_mib", Used: 20480}, {Scope: "user", Resource: "workflows", Used: 2}, {Scope: "did:plc:123", Resource: "cache_storage_bytes", Used: 9999}, {Scope: "did:web:alice/repo1", Resource: "workflows", Used: 7}, {Scope: "user", Resource: "gpu_seconds", Used: 1234}, }, Subjects: []QuotaSubjectCount{ {Scope: "user", Resource: "cache_storage_bytes", Status: "under_limit", Count: 2}, {Scope: "user", Resource: "cache_storage_bytes", Status: "at_limit", Count: 1}, {Scope: "repo", Resource: "cache_storage_bytes", Status: "unlimited", Count: 5}, {Scope: "user", Resource: "workflows", Status: "near_limit", Count: 1}, {Scope: "repo", Resource: "memory_mib", Status: "over_limit", Count: 1}, {Scope: "user", Resource: "cache_storage_bytes", Status: "invalid_status", Count: 42}, {Scope: "did:plc:123", Resource: "cache_storage_bytes", Status: "under_limit", Count: 10}, {Scope: "user", Resource: "gpu_seconds", Status: "under_limit", Count: 3}, }, } }) m.SetQuotaDefaultLimit("user", "workflows", 4) m.SetQuotaDefaultLimit("did:web:alice", "gpu_seconds", 99) m.RecordQuotaDecision("nix_cache", "cache_storage_bytes", true, false, "within_limit") m.RecordQuotaDecision("nix_cache", "cache_storage_bytes", false, false, "repo_limit") m.RecordQuotaDecision("workflow", "workflows", false, true, "queued") m.RecordQuotaDecision("workflow", "vcpus", true, false, "allowed_from_queue") m.RecordQuotaDecision("did:web:alice", "gpu_seconds", false, false, "dial tcp 10.0.0.1:5432: connection refused") m.RecordQuotaDecision("workflow", "workflows", false, false, "request_exceeds_budget") m.SetQuotaWaitDepth("workflow", 4) m.SetQuotaWaitDepth("res-abc", 1) m.RecordQuotaWait("workflow", "memory_mib", 3*time.Second) m.RecordQuotaWait("sha256:deadbeef", "gpu_seconds", time.Second) m.RecordQuotaWait("workflow", "memory_mib", -time.Second) families := gather(t, m) wantLabels := map[string][]string{ "spindle_quota_usage": {"resource", "scope"}, "spindle_quota_default_limit": {"resource", "scope"}, "spindle_quota_subjects": {"resource", "scope", "status"}, "spindle_quota_decisions_total": {"decision", "kind", "reason", "resource"}, "spindle_quota_wait_depth": {"kind"}, "spindle_quota_wait_duration_seconds": {"kind", "resource"}, } allowed := map[string]map[string]bool{ "scope": {"user": true, "repo": true}, "resource": {"cache_storage_bytes": true, "workflows": true, "vcpus": true, "memory_mib": true, "disk_mib": true, "unknown": true}, "status": {"unlimited": true, "under_limit": true, "near_limit": true, "at_limit": true, "over_limit": true}, "kind": {"workflow": true, "nix_cache": true, "generic_cache": true, "unknown": true}, "decision": {"allowed": true, "rejected": true, "deferred": true}, "reason": { "allowed_immediately": true, "allowed_from_queue": true, "unlimited": true, "within_limit": true, "user_limit": true, "repo_limit": true, "queued": true, "no_wait_slot_unavailable": true, "context_done_in_queue": true, "context_done_after_ready": true, "store_error": true, "unknown": true, }, } for name, want := range wantLabels { f, ok := families[name] if !ok { t.Fatalf("missing quota metric family %s", name) } if len(f.GetMetric()) == 0 { t.Fatalf("family %s emitted no series", name) } for _, metric := range f.GetMetric() { keys := labelKeys(metric) if strings.Join(keys, ",") != strings.Join(want, ",") { t.Fatalf("family %s has labels %v, want %v", name, keys, want) } for _, lp := range metric.GetLabel() { if !allowed[lp.GetName()][lp.GetValue()] { t.Fatalf("family %s label %s has unbounded value %q", name, lp.GetName(), lp.GetValue()) } } } } for _, name := range []string{"spindle_quota_usage", "spindle_quota_subjects", "spindle_quota_default_limit"} { for _, metric := range families[name].GetMetric() { for _, lp := range metric.GetLabel() { if lp.GetValue() == "unknown" { t.Fatalf("family %s kept an unknown %s bucket", name, lp.GetName()) } } } } var sawUnknownDecision bool for _, metric := range families["spindle_quota_decisions_total"].GetMetric() { if labelValue(metric, "kind") == "unknown" { sawUnknownDecision = true if labelValue(metric, "resource") != "unknown" || labelValue(metric, "reason") != "unknown" { t.Fatalf("unbounded decision kept a partially unbounded tuple: %v", metric.GetLabel()) } } } if !sawUnknownDecision { t.Fatal("unbounded decision was dropped instead of counted as unknown") } var workflowSeries int for _, metric := range families["spindle_quota_usage"].GetMetric() { if labelValue(metric, "scope") != "user" || labelValue(metric, "resource") != "workflows" { continue } workflowSeries++ if got := metric.GetGauge().GetValue(); got != 5 { t.Fatalf("user workflow usage = %v, want 5", got) } } if workflowSeries != 1 { t.Fatalf("user workflow usage emitted %d series, want 1", workflowSeries) } if got := len(families["spindle_quota_usage"].GetMetric()); got != 6 { t.Fatalf("quota usage emitted %d series, want 6", got) } if got := len(families["spindle_quota_subjects"].GetMetric()); got != 5 { t.Fatalf("quota subjects emitted %d series, want 5", got) } var waits uint64 for _, metric := range families["spindle_quota_wait_duration_seconds"].GetMetric() { waits += metric.GetHistogram().GetSampleCount() } if waits != 2 { t.Fatalf("observed %d quota waits, want 2", waits) } for name, f := range families { if !strings.HasPrefix(name, "spindle_quota_") { continue } if _, ok := wantLabels[name]; !ok { t.Fatalf("unexpected quota metric family %s", name) } if f.GetHelp() == "" { t.Fatalf("family %s has no help text", name) } } } func TestProvisionedDashboardMatchesEmittedQuotaMetrics(t *testing.T) { const dashboardPath = "../../localinfra/observability/grafana/provisioning/dashboards/spindle.json" raw, err := os.ReadFile(dashboardPath) if err != nil { t.Fatalf("read dashboard: %v", err) } var dashboard struct { UID string `json:"uid"` } if err := json.Unmarshal(raw, &dashboard); err != nil { t.Fatalf("parse dashboard: %v", err) } if dashboard.UID != "efv22ywqutn28f" { t.Fatalf("dashboard uid = %q, provisioned uid is efv22ywqutn28f", dashboard.UID) } m := NewMetrics() m.SetQuotaProvider(func() QuotaSnapshot { return QuotaSnapshot{ Usage: []QuotaUsage{{Scope: "user", Resource: "cache_storage_bytes", Used: 1}}, Subjects: []QuotaSubjectCount{{Scope: "user", Resource: "cache_storage_bytes", Status: "under_limit", Count: 1}}, } }) m.SetQuotaDefaultLimit("user", "cache_storage_bytes", 1) obs := m.QuotaObserver() obs.RecordDecision("workflow", "workflows", true, false, "allowed_immediately") obs.SetWaitDepth("workflow", 0) obs.RecordWait("workflow", "memory_mib", time.Second) m.RecordEnginePoolSnapshot(0, 1, 1, 0, 1, 1, 0, 1, 1, 0) m.RecordEnginePoolAdmission(true, "allowed_immediately") m.RecordCacheUpload("http", "published") m.RecordCacheUploadBytes("http", "published", 1) emitted := gather(t, m) referenced := regexp.MustCompile(`spindle_(?:quota|engine_pool|cache)_[a-z0-9_]+`).FindAllString(string(raw), -1) if len(referenced) == 0 { t.Fatal("dashboard references no quota or cache metrics") } for _, name := range referenced { base := strings.TrimSuffix(name, "_bucket") if _, ok := emitted[base]; !ok { t.Fatalf("dashboard queries %s, which this package does not emit", name) } } } func TestCacheUploadMetricsAreBounded(t *testing.T) { m := NewMetrics() m.RecordCacheUpload("bad_backend", "bad_result") m.RecordCacheUpload("http", "published") m.RecordCacheUploadBytes("http", "published", 1024) m.RecordCacheUploadBytes("http", "published", -500) m.RecordCacheUploadBytes("nix_store", "quota_skipped", 2048) m.RecordCacheUploadBytes("did:web:alice", "published", 2048) families := gather(t, m) for _, name := range []string{"spindle_cache_uploads_total", "spindle_cache_upload_bytes_total"} { f, ok := families[name] if !ok { t.Fatalf("missing family %s", name) } for _, metric := range f.GetMetric() { keys := labelKeys(metric) if strings.Join(keys, ",") != "backend,result" { t.Fatalf("family %s has labels %v, want [backend result]", name, keys) } backend := labelValue(metric, "backend") if backend != "http" && backend != "nix_store" && backend != "unknown" { t.Fatalf("family %s has unbounded backend %q", name, backend) } result := labelValue(metric, "result") if result != "published" && result != "quota_skipped" && result != "failed" && result != "unknown" { t.Fatalf("family %s has unbounded result %q", name, result) } } } // negative byte counts are ignored for _, metric := range families["spindle_cache_upload_bytes_total"].GetMetric() { if labelValue(metric, "backend") == "http" && metric.GetCounter().GetValue() != 1024 { t.Fatalf("http upload bytes = %v, want 1024", metric.GetCounter().GetValue()) } } } // the cache-only quota families are gone, nothing may resurrect them func TestObsoleteCacheQuotaFamiliesAreGone(t *testing.T) { m := NewMetrics() m.SetQuotaProvider(func() QuotaSnapshot { return QuotaSnapshot{ Usage: []QuotaUsage{{Scope: "user", Resource: "cache_storage_bytes", Used: 1}}, Subjects: []QuotaSubjectCount{{Scope: "user", Resource: "cache_storage_bytes", Status: "under_limit", Count: 1}}, } }) m.RecordQuotaDecision("nix_cache", "cache_storage_bytes", true, false, "within_limit") families := gather(t, m) for _, name := range []string{ "spindle_cache_quota_usage_bytes", "spindle_cache_quota_subjects", "spindle_cache_quota_decisions_total", "spindle_quota_memory_mib", "spindle_quota_vcpus", "spindle_quota_disk_mib", "spindle_quota_queue_depth", } { if _, ok := families[name]; ok { t.Fatalf("obsolete metric family %s is still registered", name) } } } // the local engine pool is physical capacity, both its levels and its // admissions stay outside the quota namespace func TestEnginePoolMetricsAreSeparateFromQuota(t *testing.T) { m := NewMetrics() m.RecordEnginePoolSnapshot(2048, 8192, 4096, 2, 8, 4, 10240, 40960, 20480, 3) m.RecordEnginePoolAdmission(true, "allowed_immediately") m.RecordEnginePoolAdmission(false, "request_exceeds_budget") // a tenant quota reason is not a host capacity reason m.RecordEnginePoolAdmission(false, "user_limit") m.RecordEnginePoolAdmission(false, "did:web:alice") families := gather(t, m) for _, name := range []string{ "spindle_engine_pool_memory_mib", "spindle_engine_pool_vcpus", "spindle_engine_pool_disk_mib", } { f, ok := families[name] if !ok { t.Fatalf("missing family %s", name) } for _, metric := range f.GetMetric() { keys := labelKeys(metric) if strings.Join(keys, ",") != "state" { t.Fatalf("family %s has labels %v, want [state]", name, keys) } state := labelValue(metric, "state") if state != "used" && state != "limit" && state != "max_request" { t.Fatalf("family %s has unbounded state %q", name, state) } } } depth, ok := families["spindle_engine_pool_queue_depth"] if !ok { t.Fatal("missing family spindle_engine_pool_queue_depth") } if got := depth.GetMetric()[0].GetGauge().GetValue(); got != 3 { t.Fatalf("engine pool queue depth = %v, want 3", got) } admission, ok := families["spindle_engine_pool_admission_total"] if !ok { t.Fatal("missing family spindle_engine_pool_admission_total") } poolReasons := map[string]bool{ "allowed_immediately": true, "allowed_from_queue": true, "request_exceeds_budget": true, "request_exceeds_max": true, "no_wait_slot_unavailable": true, "context_done_in_queue": true, "context_done_after_ready": true, "unknown": true, } var unknownPoolReasons float64 for _, metric := range admission.GetMetric() { keys := labelKeys(metric) if strings.Join(keys, ",") != "decision,reason" { t.Fatalf("engine pool admission has labels %v, want [decision reason]", keys) } if decision := labelValue(metric, "decision"); decision != "allowed" && decision != "rejected" { t.Fatalf("engine pool admission has unbounded decision %q", decision) } reason := labelValue(metric, "reason") if !poolReasons[reason] { t.Fatalf("engine pool admission has unbounded reason %q", reason) } if reason == "unknown" { unknownPoolReasons += metric.GetCounter().GetValue() } } if unknownPoolReasons != 2 { t.Fatalf("counted %v unknown pool reasons, want 2", unknownPoolReasons) } // host capacity admissions must never reach the tenant quota family if _, ok := families["spindle_quota_decisions_total"]; ok { t.Fatal("engine pool admissions leaked into spindle_quota_decisions_total") } } // the two reason vocabularies do not bleed into each other func TestQuotaAndEnginePoolReasonsAreDisjointWhereTheyMustBe(t *testing.T) { for _, reason := range []string{"request_exceeds_budget", "request_exceeds_max"} { if _, ok := quotaReason(reason); ok { t.Fatalf("host capacity reason %q is accepted as a quota reason", reason) } } for _, reason := range []string{"user_limit", "repo_limit", "within_limit", "unlimited", "queued", "store_error"} { if _, ok := enginePoolReason(reason); ok { t.Fatalf("tenant quota reason %q is accepted as a host capacity reason", reason) } } } // no forbidden identifier survives any bound function func TestQuotaBoundFunctionsRejectIdentifiers(t *testing.T) { for _, value := range forbiddenLabelValues { if got, ok := quotaScope(value); ok || got != "unknown" { t.Fatalf("quotaScope(%q) = %q, %v", value, got, ok) } if got, ok := quotaStatus(value); ok || got != "unknown" { t.Fatalf("quotaStatus(%q) = %q, %v", value, got, ok) } if got, ok := enginePoolReason(value); ok || got != "unknown" { t.Fatalf("enginePoolReason(%q) = %q, %v", value, got, ok) } if got, ok := quotaKind(value); ok || got != "unknown" { t.Fatalf("quotaKind(%q) = %q, %v", value, got, ok) } if got, ok := quotaResource(value); ok || got != "unknown" { t.Fatalf("quotaResource(%q) = %q, %v", value, got, ok) } if got, ok := quotaReason(value); ok || got != "unknown" { t.Fatalf("quotaReason(%q) = %q, %v", value, got, ok) } } if got := quotaDecision(true, true); got != "allowed" { t.Fatalf("allowed decision = %q", got) } if got := quotaDecision(false, true); got != "deferred" { t.Fatalf("temporary denial = %q", got) } if got := quotaDecision(false, false); got != "rejected" { t.Fatalf("permanent denial = %q", got) } } type mockClock struct { mu sync.Mutex t time.Time } func (m *mockClock) Now() time.Time { m.mu.Lock() defer m.mu.Unlock() return m.t } func (m *mockClock) Advance(d time.Duration) { m.mu.Lock() defer m.mu.Unlock() m.t = m.t.Add(d) } func TestCachedQuotaProvider(t *testing.T) { clk := &mockClock{t: time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC)} t.Run("first load", func(t *testing.T) { var calls int32 loader := func() (QuotaSnapshot, error) { atomic.AddInt32(&calls, 1) return QuotaSnapshot{ Usage: []QuotaUsage{{Scope: "user", Resource: "vcpus", Used: 10}}, }, nil } cache := newCachedQuotaProvider(loader, clk, 1*time.Minute) snap := cache.Get() if atomic.LoadInt32(&calls) != 1 { t.Fatalf("expected 1 loader call, got %d", calls) } if len(snap.Usage) != 1 || snap.Usage[0].Used != 10 { t.Fatalf("unexpected snapshot content: %v", snap) } }) t.Run("cache hit", func(t *testing.T) { var calls int32 loader := func() (QuotaSnapshot, error) { atomic.AddInt32(&calls, 1) return QuotaSnapshot{ Usage: []QuotaUsage{{Scope: "user", Resource: "vcpus", Used: 10}}, }, nil } cache := newCachedQuotaProvider(loader, clk, 1*time.Minute) cache.Get() snap := cache.Get() if atomic.LoadInt32(&calls) != 1 { t.Fatalf("expected exactly 1 loader call (first load), got %d", calls) } if len(snap.Usage) != 1 || snap.Usage[0].Used != 10 { t.Fatalf("unexpected snapshot content: %v", snap) } }) t.Run("TTL expiry", func(t *testing.T) { var calls int32 loader := func() (QuotaSnapshot, error) { atomic.AddInt32(&calls, 1) return QuotaSnapshot{ Usage: []QuotaUsage{{Scope: "user", Resource: "vcpus", Used: float64(atomic.LoadInt32(&calls) * 10)}}, }, nil } cache := newCachedQuotaProvider(loader, clk, 1*time.Minute) snap1 := cache.Get() if snap1.Usage[0].Used != 10 { t.Fatalf("expected 10, got %g", snap1.Usage[0].Used) } clk.Advance(59 * time.Second) snap2 := cache.Get() if snap2.Usage[0].Used != 10 { t.Fatalf("expected cached 10, got %g", snap2.Usage[0].Used) } if atomic.LoadInt32(&calls) != 1 { t.Fatalf("expected exactly 1 call, got %d", calls) } clk.Advance(2 * time.Second) snap3 := cache.Get() if snap3.Usage[0].Used != 20 { t.Fatalf("expected refreshed 20, got %g", snap3.Usage[0].Used) } if atomic.LoadInt32(&calls) != 2 { t.Fatalf("expected exactly 2 calls, got %d", calls) } }) t.Run("concurrent miss coalescing", func(t *testing.T) { var calls int32 startChan := make(chan struct{}) blockChan := make(chan struct{}) loader := func() (QuotaSnapshot, error) { atomic.AddInt32(&calls, 1) close(startChan) <-blockChan return QuotaSnapshot{ Usage: []QuotaUsage{{Scope: "user", Resource: "vcpus", Used: 42}}, }, nil } cache := newCachedQuotaProvider(loader, clk, 1*time.Minute) var wg sync.WaitGroup const numCallers = 10 results := make([]QuotaSnapshot, numCallers) for i := 0; i < numCallers; i++ { wg.Add(1) go func(idx int) { defer wg.Done() results[idx] = cache.Get() }(i) } <-startChan time.Sleep(50 * time.Millisecond) close(blockChan) wg.Wait() if atomic.LoadInt32(&calls) != 1 { t.Fatalf("expected concurrent scrapes to coalesce into exactly 1 loader call, got %d", calls) } for idx, snap := range results { if len(snap.Usage) != 1 || snap.Usage[0].Used != 42 { t.Fatalf("caller %d got unexpected snapshot: %v", idx, snap) } } }) t.Run("last-good preservation after error", func(t *testing.T) { var calls int32 var shouldFail int32 loader := func() (QuotaSnapshot, error) { atomic.AddInt32(&calls, 1) if atomic.LoadInt32(&shouldFail) == 1 { return QuotaSnapshot{}, errors.New("database connection lost") } return QuotaSnapshot{ Usage: []QuotaUsage{{Scope: "user", Resource: "vcpus", Used: 100}}, }, nil } cache := newCachedQuotaProvider(loader, clk, 1*time.Minute) snap1 := cache.Get() if len(snap1.Usage) != 1 || snap1.Usage[0].Used != 100 { t.Fatalf("expected initial load success, got %v", snap1) } atomic.StoreInt32(&shouldFail, 1) clk.Advance(2 * time.Minute) snap2 := cache.Get() if len(snap2.Usage) != 1 || snap2.Usage[0].Used != 100 { t.Fatalf("expected last-good preservation, got %v", snap2) } clk.Advance(2 * time.Minute) snap3 := cache.Get() if len(snap3.Usage) != 1 || snap3.Usage[0].Used != 100 { t.Fatalf("expected last-good preservation on second failure, got %v", snap3) } if atomic.LoadInt32(&calls) != 3 { t.Fatalf("expected 3 load attempts, got %d", calls) } }) t.Run("cold error is cached until expiry", func(t *testing.T) { var calls int32 cache := newCachedQuotaProvider(func() (QuotaSnapshot, error) { atomic.AddInt32(&calls, 1) return QuotaSnapshot{}, errors.New("database unavailable") }, clk, time.Minute) cache.Get() cache.Get() if got := atomic.LoadInt32(&calls); got != 1 { t.Fatalf("cold failure loader calls = %d, want 1", got) } clk.Advance(time.Minute) cache.Get() if got := atomic.LoadInt32(&calls); got != 2 { t.Fatalf("loader calls at expiry = %d, want 2", got) } }) } func BenchmarkRecordWorkflow(b *testing.B) { m := NewMetrics() b.ReportAllocs() b.ResetTimer() for i := 0; i < b.N; i++ { m.RecordWorkflowStart("nixery") m.RecordWorkflowEnd(context.Background(), "nixery", "success", "none", "success", 123*time.Millisecond) } } func BenchmarkRecordStep(b *testing.B) { m := NewMetrics() b.ReportAllocs() b.ResetTimer() for i := 0; i < b.N; i++ { m.RecordStepStart("nixery") m.RecordStepEnd(context.Background(), "nixery", "success", 45*time.Millisecond) } } func BenchmarkRecordQuotaDecision(b *testing.B) { m := NewMetrics() b.ReportAllocs() b.ResetTimer() for i := 0; i < b.N; i++ { m.RecordQuotaDecision("workflow", "vcpus", true, false, "within_limit") } } func BenchmarkHTTPMiddleware(b *testing.B) { m := NewMetrics() handler := HTTPMiddleware(m)(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })) req, err := http.NewRequest("GET", "/test", nil) if err != nil { b.Fatal(err) } rctx := chi.NewRouteContext() rctx.RoutePatterns = []string{"/test"} req = req.WithContext(context.WithValue(req.Context(), chi.RouteCtxKey, rctx)) w := httptest.NewRecorder() b.ReportAllocs() b.ResetTimer() for i := 0; i < b.N; i++ { handler.ServeHTTP(w, req) } } func TestRecordWorkflowEnd_Exemplar(t *testing.T) { oldProvider := otel.GetTracerProvider() recorder := tracetest.NewSpanRecorder() provider := sdktrace.NewTracerProvider( sdktrace.WithSpanProcessor(recorder), sdktrace.WithSampler(sdktrace.AlwaysSample()), ) otel.SetTracerProvider(provider) t.Cleanup(func() { _ = provider.Shutdown(context.Background()) otel.SetTracerProvider(oldProvider) }) m := NewMetrics() ctx, span := Tracer().Start(context.Background(), "test-workflow") defer span.End() m.RecordWorkflowEnd(ctx, "nixery", "success", "none", "success", 123*time.Millisecond) m.RecordStepEnd(ctx, "nixery", "success", 45*time.Millisecond) families, err := m.Registry().Gather() if err != nil { t.Fatal(err) } findExemplar := func(name string) *dto.Exemplar { for _, f := range families { if f.GetName() == name { for _, metric := range f.GetMetric() { h := metric.GetHistogram() if h == nil { continue } for _, b := range h.GetBucket() { if e := b.GetExemplar(); e != nil { return e } } } } } return nil } wfExemplar := findExemplar("spindle_workflow_duration_seconds") if wfExemplar == nil { t.Error("spindle_workflow_duration_seconds missing exemplar") } else { hasTraceID := false for _, label := range wfExemplar.GetLabel() { if label.GetName() == "traceID" && label.GetValue() == span.SpanContext().TraceID().String() { hasTraceID = true } } if !hasTraceID { t.Errorf("workflow exemplar missing correct traceID, got label: %+v", wfExemplar.GetLabel()) } } stepExemplar := findExemplar("spindle_step_duration_seconds") if stepExemplar == nil { t.Error("spindle_step_duration_seconds missing exemplar") } else { hasTraceID := false for _, label := range stepExemplar.GetLabel() { if label.GetName() == "traceID" && label.GetValue() == span.SpanContext().TraceID().String() { hasTraceID = true } } if !hasTraceID { t.Errorf("step exemplar missing correct traceID, got label: %+v", stepExemplar.GetLabel()) } } } func TestWorkflowQueueAndExecutorMetrics(t *testing.T) { m := NewMetrics() m.RecordWorkflowStart("microvm") m.RecordWorkflowEnd( context.Background(), "microvm", "failure", "user", "command_failed", 15*time.Minute, ) m.RecordWorkflowTerminal("arbitrary", "invalid_status", "did:plc:123", "dial tcp 10.0.0.1:5432: connection refused") m.RecordJobDequeueLatency(context.Background(), 20*time.Minute) m.RecordWorkflowStartupDelay(context.Background(), "microvm", "mill", 12*time.Minute) m.RecordQuotaWait("workflow", "workflows", time.Second) m.RecordMillPlacementResult(context.Background(), "microvm", "success", 11*time.Minute) m.RegisterExecutorGauges( func() float64 { return 8 }, func() float64 { return 2 }, func() float64 { return 5 }, func() float64 { return 1024 }, ) families := gather(t, m) workflows := families["spindle_workflows_total"] if workflows == nil { t.Fatal("spindle_workflows_total missing") } for _, metric := range workflows.GetMetric() { if got := len(metric.GetLabel()); got != 2 { t.Fatalf("spindle_workflows_total label count = %d, want existing engine/result schema", got) } } terminations := families["spindle_workflow_terminations_total"] if terminations == nil { t.Fatal("spindle_workflow_terminations_total missing") } if len(terminations.GetMetric()) != 2 { t.Fatalf("spindle_workflow_terminations_total series = %d, want 2", len(terminations.GetMetric())) } var sawClassified, sawBounded bool for _, metric := range terminations.GetMetric() { switch labelValue(metric, "engine") { case "microvm": sawClassified = labelValue(metric, "result") == "failure" && labelValue(metric, "failure_class") == "user" && labelValue(metric, "reason") == "command_failed" case "unknown": sawBounded = labelValue(metric, "result") == "unknown" && labelValue(metric, "failure_class") == "unknown" && labelValue(metric, "reason") == "unknown" } } if !sawClassified || !sawBounded { t.Fatalf("workflow classification: classified=%t bounded=%t", sawClassified, sawBounded) } for _, name := range []string{ "spindle_workflow_duration_seconds", "spindle_job_dequeue_latency_seconds", "spindle_workflow_startup_delay_seconds", "spindle_mill_placement_wait_seconds", } { family := families[name] if family == nil || len(family.GetMetric()) == 0 { t.Fatalf("%s missing", name) } histogram := family.GetMetric()[0].GetHistogram() if histogram == nil { t.Fatalf("%s is not a histogram", name) } maxBucket := float64(0) for _, bucket := range histogram.GetBucket() { if bucket.GetUpperBound() > maxBucket { maxBucket = bucket.GetUpperBound() } } wantMin := float64(7200) if name != "spindle_workflow_duration_seconds" { wantMin = 86400 } if maxBucket < wantMin { t.Fatalf("%s max finite bucket = %v, want at least %v", name, maxBucket, wantMin) } } workflowHistogram := families["spindle_workflow_duration_seconds"].GetMetric()[0].GetHistogram() for _, required := range prometheus.DefBuckets { found := false for _, bucket := range workflowHistogram.GetBucket() { if bucket.GetUpperBound() == required { found = true break } } if !found { t.Fatalf("workflow histogram dropped deployed bucket %v", required) } } quotaHistogram := families["spindle_quota_wait_duration_seconds"].GetMetric()[0].GetHistogram() for _, required := range prometheus.ExponentialBuckets(0.1, 3, 8) { found := false for _, bucket := range quotaHistogram.GetBucket() { if bucket.GetUpperBound() == required { found = true break } } if !found { t.Fatalf("quota wait histogram dropped deployed bucket %v", required) } } seats := families["spindle_executor_seats"] if seats == nil || len(seats.GetMetric()) != 1 || seats.GetMetric()[0].GetGauge().GetValue() != 8 { t.Fatalf("spindle_executor_seats = %v, want 8", seats) } } func TestDBCollectorUsesLatestNormalizedWorkflowStatus(t *testing.T) { ctx := context.Background() database, err := db.Make(ctx, filepath.Join(t.TempDir(), "metrics.db")) if err != nil { t.Fatal(err) } defer database.Close() for _, row := range []struct { pipeline, workflow, status string }{ {"p1", "build", "pending"}, {"p1", "build", "running"}, {"p2", "test", "failed"}, {"p1", "build", "success"}, } { if _, err := database.Exec(`insert into workflow_statuses (pipeline_id, workflow, status, created_at) values (?, ?, ?, ?)`, row.pipeline, row.workflow, row.status, time.Now().Format(time.RFC3339)); err != nil { t.Fatal(err) } } collector := &dbCollector{ db: database, collectionFailures: prometheus.NewCounter(prometheus.CounterOpts{Name: "test_collection_failures_total"}), workflowStatus: make(map[string]int), } collector.refresh(ctx) if collector.workflowStatus["success"] != 1 || collector.workflowStatus["failed"] != 1 || len(collector.workflowStatus) != 2 { t.Fatalf("workflow statuses = %#v", collector.workflowStatus) } }