Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
36 kB · 1334 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335package observability
import ( "context" "fmt" "github.com/felixge/httpsnoop" "github.com/go-chi/chi/v5" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/collectors" "github.com/prometheus/client_golang/prometheus/promhttp" "go.opentelemetry.io/otel/trace" "log/slog" "net" "net/http" "sync" "tangled.org/core/spindle/db" "time")
var workflowDurationBuckets = []float64{ 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30, 60, 120, 300, 600, 900, 1800, 3600, 7200,}
var queueDurationBuckets = []float64{ 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30, 60, 120, 300, 600, 900, 1800, 3600, 7200, 21600, 43200, 86400,}
var quotaWaitDurationBuckets = append( prometheus.ExponentialBuckets(0.1, 3, 8), 300, 600, 900, 1800, 3600, 7200,)
type contextKey struct{}
type workflowTerminalObserverKey struct{}
type WorkflowTerminalObserver func(engine, result, failureClass, reason string)
func WithWorkflowTerminalObserver(ctx context.Context, observer WorkflowTerminalObserver) context.Context { return context.WithValue(ctx, workflowTerminalObserverKey{}, observer)}
func NotifyWorkflowTerminal(ctx context.Context, engine, result, failureClass, reason string) bool { if ctx == nil { return false } observer, ok := ctx.Value(workflowTerminalObserverKey{}).(WorkflowTerminalObserver) if !ok || observer == nil { return false } observer(engine, result, failureClass, reason) return true}
func WithMetrics(ctx context.Context, m *Metrics) context.Context { return context.WithValue(ctx, contextKey{}, m)}
func GetMetrics(ctx context.Context) *Metrics { if ctx == nil { return nil } m, _ := ctx.Value(contextKey{}).(*Metrics) return m}
type QuotaUsage struct { Scope string Resource string Used float64}
type QuotaSubjectCount struct { Scope string Resource string Status string Count int64}
type QuotaSnapshot struct { Usage []QuotaUsage Subjects []QuotaSubjectCount}
type QuotaProvider func() QuotaSnapshot
type quotaCollector struct { m *Metrics usageDesc *prometheus.Desc subjectsDesc *prometheus.Desc}
func (c *quotaCollector) Describe(ch chan<- *prometheus.Desc) { ch <- c.usageDesc ch <- c.subjectsDesc}
func (c *quotaCollector) Collect(ch chan<- prometheus.Metric) { if c.m == nil || c.m.quotaProvider == nil { return } snap := c.m.quotaProvider()
usage := make(map[[2]string]float64, len(snap.Usage)) for _, u := range snap.Usage { scope, ok := quotaScope(u.Scope) if !ok { continue } resource, ok := quotaResource(u.Resource) if !ok { continue } usage[[2]string{scope, resource}] += u.Used } for key, used := range usage { ch <- prometheus.MustNewConstMetric(c.usageDesc, prometheus.GaugeValue, used, key[0], key[1]) }
subjects := make(map[[3]string]int64, len(snap.Subjects)) for _, s := range snap.Subjects { scope, ok := quotaScope(s.Scope) if !ok { continue } resource, ok := quotaResource(s.Resource) if !ok { continue } status, ok := quotaStatus(s.Status) if !ok { continue } subjects[[3]string{scope, resource, status}] += s.Count } for key, count := range subjects { ch <- prometheus.MustNewConstMetric(c.subjectsDesc, prometheus.GaugeValue, float64(count), key[0], key[1], key[2]) }}
type Metrics struct { reg *prometheus.Registry
httpRequests *prometheus.CounterVec httpRequestDuration *prometheus.HistogramVec httpInFlight prometheus.Gauge
eventIngestion *prometheus.CounterVec
jobQueueActivity *prometheus.CounterVec jobDequeueLatency prometheus.Histogram
workflowsActive *prometheus.GaugeVec workflowsTotal *prometheus.CounterVec workflowTerminations *prometheus.CounterVec workflowDuration *prometheus.HistogramVec workflowStartupDelay *prometheus.HistogramVec
stepsActive *prometheus.GaugeVec stepsTotal *prometheus.CounterVec stepDuration *prometheus.HistogramVec
poolMemory *prometheus.GaugeVec poolVCPUs *prometheus.GaugeVec poolDisk *prometheus.GaugeVec poolQueueDepth prometheus.Gauge poolAdmission *prometheus.CounterVec
millPlacementAdmission *prometheus.CounterVec millPlacementResults *prometheus.CounterVec millReconnects prometheus.Counter millPlacementWait *prometheus.HistogramVec
jumpActive prometheus.Gauge jumpMax prometheus.Gauge jumpRejections *prometheus.CounterVec
cacheUploads *prometheus.CounterVec cacheUploadBytes *prometheus.CounterVec
quotaDecisions *prometheus.CounterVec quotaDefaultLimit *prometheus.GaugeVec quotaWaitDepth *prometheus.GaugeVec quotaWaitTime *prometheus.HistogramVec quotaProvider QuotaProvider clock clock collectionFailures prometheus.Counter}
func NewMetrics() *Metrics { return newMetricsWithClock(realClock{})}
func newMetricsWithClock(clock clock) *Metrics { reg := prometheus.NewRegistry() reg.MustRegister( collectors.NewGoCollector(), collectors.NewProcessCollector(collectors.ProcessCollectorOpts{}), collectors.NewBuildInfoCollector(), )
m := &Metrics{ reg: reg, clock: clock,
httpRequests: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_http_requests_total", Help: "Total count of HTTP requests.", }, []string{"method", "route", "status_code", "status_class"}),
httpRequestDuration: prometheus.NewHistogramVec(prometheus.HistogramOpts{ Name: "spindle_http_request_duration_seconds", Help: "Duration of HTTP requests in seconds.", Buckets: prometheus.DefBuckets, }, []string{"method", "route", "status_code", "status_class"}),
httpInFlight: prometheus.NewGauge(prometheus.GaugeOpts{ Name: "spindle_http_requests_in_flight", Help: "Current number of in-flight HTTP requests.", }),
eventIngestion: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_event_ingestion_outcomes_total", Help: "Total number of ingested events by consumer and outcome status.", }, []string{"consumer", "status"}),
jobQueueActivity: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_job_queue_activity_total", Help: "Total count of job queue operations.", }, []string{"action"}),
jobDequeueLatency: prometheus.NewHistogram(prometheus.HistogramOpts{ Name: "spindle_job_dequeue_latency_seconds", Help: "Time from durable pipeline admission until a worker dequeues the job.", Buckets: queueDurationBuckets, }),
workflowsActive: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Name: "spindle_workflows_active", Help: "Current number of active workflows by engine.", }, []string{"engine"}),
workflowsTotal: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_workflows_total", Help: "Total number of completed workflows by engine and result.", }, []string{"engine", "result"}),
workflowTerminations: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_workflow_terminations_total", Help: "Total number of completed workflows by engine, result, failure class and bounded reason.", }, []string{"engine", "result", "failure_class", "reason"}),
workflowDuration: prometheus.NewHistogramVec(prometheus.HistogramOpts{ Name: "spindle_workflow_duration_seconds", Help: "Duration of completed workflows by engine and result.", Buckets: workflowDurationBuckets, }, []string{"engine", "result"}),
workflowStartupDelay: prometheus.NewHistogramVec(prometheus.HistogramOpts{ Name: "spindle_workflow_startup_delay_seconds", Help: "Time from pending admission until a workflow starts running.", Buckets: queueDurationBuckets, }, []string{"engine", "placement"}),
stepsActive: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Name: "spindle_steps_active", Help: "Current number of active steps by engine.", }, []string{"engine"}),
stepsTotal: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_steps_total", Help: "Total number of completed steps by engine and status.", }, []string{"engine", "status"}),
stepDuration: prometheus.NewHistogramVec(prometheus.HistogramOpts{ Name: "spindle_step_duration_seconds", Help: "Duration of completed steps by engine and status.", Buckets: workflowDurationBuckets, }, []string{"engine", "status"}),
poolMemory: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Name: "spindle_engine_pool_memory_mib", Help: "Local engine memory pool in MiB by state.", }, []string{"state"}),
poolVCPUs: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Name: "spindle_engine_pool_vcpus", Help: "Local engine vCPU pool by state.", }, []string{"state"}),
poolDisk: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Name: "spindle_engine_pool_disk_mib", Help: "Local engine disk pool in MiB by state.", }, []string{"state"}),
poolQueueDepth: prometheus.NewGauge(prometheus.GaugeOpts{ Name: "spindle_engine_pool_queue_depth", Help: "Current local engine resource scheduler queue depth.", }),
poolAdmission: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_engine_pool_admission_total", Help: "Total count of local engine resource scheduler admissions, host capacity only.", }, []string{"decision", "reason"}),
millPlacementAdmission: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_mill_placement_admission_total", Help: "Total count of mill placement admissions.", }, []string{"status", "reason"}),
millPlacementResults: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_mill_placement_results_total", Help: "Total count of mill placement results.", }, []string{"result"}),
millReconnects: prometheus.NewCounter(prometheus.CounterOpts{ Name: "spindle_mill_reconnects_total", Help: "Total number of executor reconnects to the mill.", }),
millPlacementWait: prometheus.NewHistogramVec(prometheus.HistogramOpts{ Name: "spindle_mill_placement_wait_seconds", Help: "Time spent waiting for mill placement by engine and result.", Buckets: queueDurationBuckets, }, []string{"engine", "result"}),
jumpActive: prometheus.NewGauge(prometheus.GaugeOpts{ Name: "spindle_jump_active_connections", Help: "Current number of active debug jump connections.", }),
jumpMax: prometheus.NewGauge(prometheus.GaugeOpts{ Name: "spindle_jump_max_connections", Help: "Maximum number of debug jump connections allowed.", }),
jumpRejections: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_jump_rejections_total", Help: "Total count of rejected debug jump connections.", }, []string{"reason"}),
cacheUploads: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_cache_uploads_total", Help: "Total count of cache uploads.", }, []string{"backend", "result"}),
cacheUploadBytes: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_cache_upload_bytes_total", Help: "Total bytes of cache uploads.", }, []string{"backend", "result"}),
quotaDecisions: prometheus.NewCounterVec(prometheus.CounterOpts{ Name: "spindle_quota_decisions_total", Help: "Total count of quota admission decisions by kind, resource, decision and reason.", }, []string{"kind", "resource", "decision", "reason"}),
quotaDefaultLimit: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Name: "spindle_quota_default_limit", Help: "Configured default quota limit by scope and resource, in each resource's own unit. Zero means unlimited.", }, []string{"scope", "resource"}),
quotaWaitDepth: prometheus.NewGaugeVec(prometheus.GaugeOpts{ Name: "spindle_quota_wait_depth", Help: "Current number of requests waiting for quota by kind.", }, []string{"kind"}),
quotaWaitTime: prometheus.NewHistogramVec(prometheus.HistogramOpts{ Name: "spindle_quota_wait_duration_seconds", Help: "Time spent waiting for quota before a grant, by kind and the resource that was short.", Buckets: quotaWaitDurationBuckets, }, []string{"kind", "resource"}), collectionFailures: prometheus.NewCounter(prometheus.CounterOpts{ Name: "spindle_collection_failures_total", Help: "Total number of metric collection failures.", }), }
reg.MustRegister( m.httpRequests, m.httpRequestDuration, m.httpInFlight, m.eventIngestion, m.jobQueueActivity, m.jobDequeueLatency, m.workflowsActive, m.workflowsTotal, m.workflowTerminations, m.workflowDuration, m.workflowStartupDelay, m.stepsActive, m.stepsTotal, m.stepDuration, m.poolMemory, m.poolVCPUs, m.poolDisk, m.poolQueueDepth, m.poolAdmission, m.millPlacementAdmission, m.millPlacementResults, m.millReconnects, m.millPlacementWait, m.jumpActive, m.jumpMax, m.jumpRejections, m.cacheUploads, m.cacheUploadBytes, m.quotaDecisions, m.quotaDefaultLimit, m.quotaWaitDepth, m.quotaWaitTime, m.collectionFailures, )
reg.MustRegister("aCollector{ m: m, usageDesc: prometheus.NewDesc( "spindle_quota_usage", "Current committed quota usage by scope and resource, in each resource's own unit.", []string{"scope", "resource"}, nil, ), subjectsDesc: prometheus.NewDesc( "spindle_quota_subjects", "Current count of quota subjects by scope, resource and status.", []string{"scope", "resource", "status"}, nil, ), })
return m}
func boundBackend(backend string) string { switch backend { case "http", "nix_store": return backend default: return "unknown" }}
func boundResult(result string) string { switch result { case "published", "quota_skipped", "failed": return result default: return "unknown" }}
// keep unbounded input out of metric labelsfunc quotaScope(scope string) (string, bool) { switch scope { case "user", "repo": return scope, true default: return "unknown", false }}
func quotaResource(resource string) (string, bool) { switch resource { case "cache_storage_bytes", "workflows", "vcpus", "memory_mib", "disk_mib": return resource, true default: return "unknown", false }}
func quotaStatus(status string) (string, bool) { switch status { case "unlimited", "under_limit", "near_limit", "at_limit", "over_limit": return status, true default: return "unknown", false }}
func quotaKind(kind string) (string, bool) { switch kind { case "workflow", "nix_cache", "generic_cache": return kind, true default: return "unknown", false }}
func quotaReason(reason string) (string, bool) { switch reason { case "allowed_immediately", "allowed_from_queue", "unlimited", "within_limit", "user_limit", "repo_limit", "queued", "no_wait_slot_unavailable", "context_done_in_queue", "context_done_after_ready", "store_error": return reason, true default: return "unknown", false }}
func enginePoolReason(reason string) (string, bool) { switch reason { case "allowed_immediately", "allowed_from_queue", "request_exceeds_budget", "request_exceeds_max", "no_wait_slot_unavailable", "context_done_in_queue", "context_done_after_ready": return reason, true default: return "unknown", false }}
func quotaDecision(allowed, temporary bool) string { switch { case allowed: return "allowed" case temporary: return "deferred" default: return "rejected" }}
func boundWorkflowResult(result string) string { switch result { case "success", "failure", "timeout", "cancelled": return result default: return "unknown" }}
func boundFailureClass(class string) string { switch class { case "none", "user", "infrastructure", "policy": return class default: return "unknown" }}
func boundFailureReason(reason string) string { switch reason { case "success", "command_failed", "resource_exhausted", "configuration_failed", "workflow_invalid", "setup_failed", "runtime_failed", "capacity_unavailable", "quota_denied", "executor_lost", "timeout", "cancelled": return reason default: return "unknown" }}
func boundEngine(engine string) string { switch engine { case "dummy", "microvm", "nixery": return engine default: return "unknown" }}
func boundPlacement(placement string) string { switch placement { case "local", "mill": return placement default: return "unknown" }}
func boundMillPlacementResult(result string) string { switch result { case "success", "timeout", "cancelled", "error", "rejected": return result default: return "unknown" }}
func boundMillPlacementReason(reason string) string { switch reason { case "allowed", "max_pending_reached": return reason default: return "unknown" }}
func (m *Metrics) RecordCacheUpload(backend, result string) { if m == nil { return } backend = boundBackend(backend) result = boundResult(result) m.cacheUploads.WithLabelValues(backend, result).Inc()}
func (m *Metrics) RecordCacheUploadBytes(backend, result string, bytes int64) { if m == nil || bytes < 0 { return } backend = boundBackend(backend) result = boundResult(result) m.cacheUploadBytes.WithLabelValues(backend, result).Add(float64(bytes))}
func (m *Metrics) RecordQuotaDecision(kind, resource string, allowed, temporary bool, reason string) { if m == nil { return } kind, _ = quotaKind(kind) resource, _ = quotaResource(resource) reason, _ = quotaReason(reason) m.quotaDecisions.WithLabelValues(kind, resource, quotaDecision(allowed, temporary), reason).Inc()}
func (m *Metrics) SetQuotaWaitDepth(kind string, depth int64) { if m == nil { return } kind, _ = quotaKind(kind) m.quotaWaitDepth.WithLabelValues(kind).Set(float64(depth))}
func (m *Metrics) RecordQuotaWait(kind, resource string, d time.Duration) { if m == nil || d < 0 { return } kind, _ = quotaKind(kind) resource, _ = quotaResource(resource) m.quotaWaitTime.WithLabelValues(kind, resource).Observe(d.Seconds())}
type clock interface { Now() time.Time}
type realClock struct{}
func (realClock) Now() time.Time { return time.Now()}
type QuotaLoader func() (QuotaSnapshot, error)
type cachedQuotaProvider struct { mu sync.Mutex loader QuotaLoader clock clock ttl time.Duration lastSnap QuotaSnapshot expiry time.Time hasLast bool onError func()}
func newCachedQuotaProvider(loader QuotaLoader, clock clock, ttl time.Duration) *cachedQuotaProvider { return &cachedQuotaProvider{ loader: loader, clock: clock, ttl: ttl, }}
func cloneSnapshot(snap QuotaSnapshot) QuotaSnapshot { var out QuotaSnapshot if snap.Usage != nil { out.Usage = make([]QuotaUsage, len(snap.Usage)) copy(out.Usage, snap.Usage) } if snap.Subjects != nil { out.Subjects = make([]QuotaSubjectCount, len(snap.Subjects)) copy(out.Subjects, snap.Subjects) } return out}
func (c *cachedQuotaProvider) Get() QuotaSnapshot { c.mu.Lock() defer c.mu.Unlock()
now := c.clock.Now() if now.Before(c.expiry) { if c.hasLast { return cloneSnapshot(c.lastSnap) } return QuotaSnapshot{} }
snap, err := c.loader() c.expiry = now.Add(c.ttl) if err != nil { if c.onError != nil { c.onError() } slog.Error("failed to collect quota metrics snapshot", "err", err) if c.hasLast { return cloneSnapshot(c.lastSnap) } return QuotaSnapshot{} } c.lastSnap = cloneSnapshot(snap) c.hasLast = true return cloneSnapshot(c.lastSnap)}
func (m *Metrics) SetQuotaLoader(loader QuotaLoader) { if m == nil { return } if loader == nil { m.quotaProvider = nil return } cache := newCachedQuotaProvider(loader, m.clock, 1*time.Minute) m.quotaProvider = cache.Get cache.onError = m.collectionFailures.Inc}
func (m *Metrics) SetQuotaProvider(provider QuotaProvider) { if m == nil { return } if provider == nil { m.quotaProvider = nil return } m.SetQuotaLoader(func() (QuotaSnapshot, error) { return provider(), nil })}
type QuotaObserver struct { m *Metrics}
func (m *Metrics) QuotaObserver() *QuotaObserver { return &QuotaObserver{m: m}}
func (o *QuotaObserver) RecordDecision(kind, resource string, allowed, temporary bool, reason string) { o.m.RecordQuotaDecision(kind, resource, allowed, temporary, reason)}
func (o *QuotaObserver) SetWaitDepth(kind string, depth int64) { o.m.SetQuotaWaitDepth(kind, depth)}
func (o *QuotaObserver) RecordWait(kind, resource string, d time.Duration) { o.m.RecordQuotaWait(kind, resource, d)}
func (m *Metrics) Registry() *prometheus.Registry { if m == nil { return nil } return m.reg}
func (m *Metrics) RecordEventIngestion(consumer, status string) { if m == nil { return } m.eventIngestion.WithLabelValues(consumer, status).Inc()}
func (m *Metrics) RecordJobQueueActivity(action string) { if m == nil { return } m.jobQueueActivity.WithLabelValues(action).Inc()}
func (m *Metrics) RecordJobDequeueLatency(ctx context.Context, d time.Duration) { if m == nil || d < 0 { return } observeWithExemplar(ctx, m.jobDequeueLatency, d.Seconds())}
func (m *Metrics) SetQuotaDefaultLimit(scope, resource string, limit int64) { if m == nil { return } boundedScope, scopeOK := quotaScope(scope) boundedResource, resourceOK := quotaResource(resource) if !scopeOK || !resourceOK { return } m.quotaDefaultLimit.WithLabelValues(boundedScope, boundedResource).Set(float64(limit))}
func (m *Metrics) RecordWorkflowStart(engine string) { if m == nil { return } m.workflowsActive.WithLabelValues(boundEngine(engine)).Inc()}
func observeWithExemplar(ctx context.Context, observer prometheus.Observer, val float64) { spanContext := trace.SpanContextFromContext(ctx) if spanContext.IsValid() && spanContext.IsSampled() { if exemplarObserver, ok := observer.(prometheus.ExemplarObserver); ok { exemplarObserver.ObserveWithExemplar( val, prometheus.Labels{"traceID": spanContext.TraceID().String()}, ) return } } observer.Observe(val)}
func (m *Metrics) RecordWorkflowEnd( ctx context.Context, engine, result, failureClass, reason string, duration time.Duration,) { if m == nil { return } m.RecordWorkflowExecutionEnd(ctx, engine, result, duration) m.RecordWorkflowTerminal(engine, result, failureClass, reason)}
func (m *Metrics) RecordWorkflowExecutionEnd( ctx context.Context, engine, result string, duration time.Duration,) { if m == nil { return } engine = boundEngine(engine) result = boundWorkflowResult(result) m.workflowsActive.WithLabelValues(engine).Dec() observer := m.workflowDuration.WithLabelValues(engine, result) observeWithExemplar(ctx, observer, duration.Seconds())}
func (m *Metrics) RecordWorkflowTerminal(engine, result, failureClass, reason string) { if m == nil { return } engine = boundEngine(engine) result = boundWorkflowResult(result) failureClass = boundFailureClass(failureClass) reason = boundFailureReason(reason) m.workflowsTotal.WithLabelValues(engine, result).Inc() m.workflowTerminations.WithLabelValues(engine, result, failureClass, reason).Inc()}
func (m *Metrics) RecordWorkflowStartupDelay( ctx context.Context, engine, placement string, d time.Duration,) { if m == nil || d < 0 { return } observer := m.workflowStartupDelay.WithLabelValues(boundEngine(engine), boundPlacement(placement)) observeWithExemplar(ctx, observer, d.Seconds())}
func (m *Metrics) RecordStepStart(engine string) { if m == nil { return } m.stepsActive.WithLabelValues(engine).Inc()}
func (m *Metrics) RecordStepEnd(ctx context.Context, engine, status string, duration time.Duration) { if m == nil { return } m.stepsActive.WithLabelValues(engine).Dec() m.stepsTotal.WithLabelValues(engine, status).Inc() observer := m.stepDuration.WithLabelValues(engine, status) observeWithExemplar(ctx, observer, duration.Seconds())}
func (m *Metrics) RecordEnginePoolSnapshot( usedMem, budgetMem, maxMem int64, usedCPU, budgetCPU, maxCPU int64, usedDisk, budgetDisk, maxDisk int64, queueDepth int,) { if m == nil { return } m.poolMemory.WithLabelValues("used").Set(float64(usedMem)) m.poolMemory.WithLabelValues("limit").Set(float64(budgetMem)) m.poolMemory.WithLabelValues("max_request").Set(float64(maxMem))
m.poolVCPUs.WithLabelValues("used").Set(float64(usedCPU)) m.poolVCPUs.WithLabelValues("limit").Set(float64(budgetCPU)) m.poolVCPUs.WithLabelValues("max_request").Set(float64(maxCPU))
m.poolDisk.WithLabelValues("used").Set(float64(usedDisk)) m.poolDisk.WithLabelValues("limit").Set(float64(budgetDisk)) m.poolDisk.WithLabelValues("max_request").Set(float64(maxDisk))
m.poolQueueDepth.Set(float64(queueDepth))}
func (m *Metrics) RecordEnginePoolAdmission(allowed bool, reason string) { if m == nil { return } decision := "rejected" if allowed { decision = "allowed" } reason, _ = enginePoolReason(reason) m.poolAdmission.WithLabelValues(decision, reason).Inc()}
func (m *Metrics) RegisterMillGauges( pending, maxPending, leases, reservations, activeSessions, disconnectedSessions func() float64,) { if m == nil { return } m.reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Name: "spindle_mill_pending_jobs", Help: "Current number of pending mill jobs.", }, pending)) m.reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Name: "spindle_mill_max_pending_jobs", Help: "Maximum number of pending mill jobs allowed.", }, maxPending)) m.reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Name: "spindle_mill_leases_active", Help: "Current number of active leases on the mill.", }, leases)) m.reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Name: "spindle_mill_reservations_active", Help: "Current number of seat reservations on the mill.", }, reservations)) m.reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Name: "spindle_mill_active_sessions", Help: "Current number of active executor sessions.", }, activeSessions)) m.reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Name: "spindle_mill_disconnected_sessions", Help: "Current number of disconnected executor sessions.", }, disconnectedSessions))}
func (m *Metrics) RecordMillPlacementAdmission(allowed bool, reason string) { if m == nil { return } status := "rejected" if allowed { status = "allowed" } m.millPlacementAdmission.WithLabelValues(status, boundMillPlacementReason(reason)).Inc()}
func (m *Metrics) RecordMillPlacementResult(ctx context.Context, engine, result string, d time.Duration) { if m == nil || d < 0 { return } engine = boundEngine(engine) result = boundMillPlacementResult(result) m.millPlacementResults.WithLabelValues(result).Inc() observeWithExemplar(ctx, m.millPlacementWait.WithLabelValues(engine, result), d.Seconds())}
func (m *Metrics) RecordMillReconnect() { if m == nil { return } m.millReconnects.Inc()}
func (m *Metrics) RegisterExecutorGauges(seats, reservations, jobs, outboxBytes func() float64) { if m == nil { return } m.reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Name: "spindle_executor_seats", Help: "Configured executor seat capacity.", }, seats)) m.reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Name: "spindle_executor_reservations_active", Help: "Current number of seat reservations on the executor.", }, reservations)) m.reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Name: "spindle_executor_jobs_active", Help: "Current number of active jobs on the executor.", }, jobs)) m.reg.MustRegister(prometheus.NewGaugeFunc(prometheus.GaugeOpts{ Name: "spindle_executor_outbox_bytes", Help: "Current outbox size in bytes.", }, outboxBytes))}
func (m *Metrics) RecordJumpRejection(reason string) { if m == nil { return } m.jumpRejections.WithLabelValues(reason).Inc()}
func (m *Metrics) RecordJumpActive(total int64) { if m == nil { return } m.jumpActive.Set(float64(total))}
func (m *Metrics) RecordJumpLimit(limit int64) { if m == nil { return } m.jumpMax.Set(float64(limit))}
type dbCollector struct { db *db.DB queuedJobsDesc *prometheus.Desc workflowStatusDesc *prometheus.Desc collectionFailures prometheus.Counter mu sync.RWMutex queuedJobs int hasQueuedJobs bool workflowStatus map[string]int
maxOpenDesc *prometheus.Desc openDesc *prometheus.Desc inUseDesc *prometheus.Desc idleDesc *prometheus.Desc waitCountDesc *prometheus.Desc waitDurationDesc *prometheus.Desc maxIdleClosedDesc *prometheus.Desc maxLifetimeClosedDesc *prometheus.Desc}
func (c *dbCollector) Describe(ch chan<- *prometheus.Desc) { ch <- c.queuedJobsDesc ch <- c.workflowStatusDesc ch <- c.maxOpenDesc ch <- c.openDesc ch <- c.inUseDesc ch <- c.idleDesc ch <- c.waitCountDesc ch <- c.waitDurationDesc ch <- c.maxIdleClosedDesc ch <- c.maxLifetimeClosedDesc}
func (c *dbCollector) refresh(ctx context.Context) { var queuedJobs int if err := c.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM jobs`).Scan(&queuedJobs); err != nil { c.collectionFailures.Inc() } else { c.mu.Lock() c.queuedJobs = queuedJobs c.hasQueuedJobs = true c.mu.Unlock() }
rows, err := c.db.QueryContext(ctx, ` WITH latest AS ( SELECT status, row_number() OVER ( PARTITION BY pipeline_id, workflow ORDER BY id DESC ) AS rank FROM workflow_statuses ) SELECT CASE status WHEN 'pending' THEN status WHEN 'running' THEN status WHEN 'failed' THEN status WHEN 'timeout' THEN status WHEN 'cancelled' THEN status WHEN 'success' THEN status ELSE 'unknown' END AS bounded_status, COUNT(*) FROM latest WHERE rank = 1 GROUP BY bounded_status `) if err != nil { c.collectionFailures.Inc() return } defer rows.Close()
statuses := make(map[string]int) for rows.Next() { var status string var count int if err := rows.Scan(&status, &count); err != nil { c.collectionFailures.Inc() return } statuses[boundWorkflowStatus(status)] = count } if err := rows.Err(); err != nil { c.collectionFailures.Inc() return } c.mu.Lock() c.workflowStatus = statuses c.mu.Unlock()}
func (c *dbCollector) run(ctx context.Context) { c.refresh(ctx) ticker := time.NewTicker(time.Minute) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: c.refresh(ctx) } }}
func (c *dbCollector) Collect(ch chan<- prometheus.Metric) { if c.db == nil { return } stats := c.db.Stats()
ch <- prometheus.MustNewConstMetric(c.maxOpenDesc, prometheus.GaugeValue, float64(stats.MaxOpenConnections)) ch <- prometheus.MustNewConstMetric(c.openDesc, prometheus.GaugeValue, float64(stats.OpenConnections)) ch <- prometheus.MustNewConstMetric(c.inUseDesc, prometheus.GaugeValue, float64(stats.InUse)) ch <- prometheus.MustNewConstMetric(c.idleDesc, prometheus.GaugeValue, float64(stats.Idle)) ch <- prometheus.MustNewConstMetric(c.waitCountDesc, prometheus.CounterValue, float64(stats.WaitCount)) ch <- prometheus.MustNewConstMetric(c.waitDurationDesc, prometheus.CounterValue, stats.WaitDuration.Seconds()) ch <- prometheus.MustNewConstMetric(c.maxIdleClosedDesc, prometheus.CounterValue, float64(stats.MaxIdleClosed)) ch <- prometheus.MustNewConstMetric(c.maxLifetimeClosedDesc, prometheus.CounterValue, float64(stats.MaxLifetimeClosed))
c.mu.RLock() queuedJobs := c.queuedJobs hasQueuedJobs := c.hasQueuedJobs statuses := make(map[string]int, len(c.workflowStatus)) for status, count := range c.workflowStatus { statuses[status] = count } c.mu.RUnlock()
if hasQueuedJobs { ch <- prometheus.MustNewConstMetric(c.queuedJobsDesc, prometheus.GaugeValue, float64(queuedJobs)) } for status, count := range statuses { ch <- prometheus.MustNewConstMetric(c.workflowStatusDesc, prometheus.GaugeValue, float64(count), status) }}
func (m *Metrics) AttachDB(ctx context.Context, database *db.DB) { if m == nil || database == nil { return } coll := &dbCollector{ db: database, collectionFailures: m.collectionFailures, queuedJobsDesc: prometheus.NewDesc( "spindle_db_queued_jobs", "Current number of queued jobs in database.", nil, nil, ), workflowStatusDesc: prometheus.NewDesc( "spindle_db_workflow_status_count", "Current count of workflows in the database by status.", []string{"status"}, nil, ), maxOpenDesc: prometheus.NewDesc( "spindle_db_max_open_connections", "Maximum number of open connections to the database.", nil, nil, ), openDesc: prometheus.NewDesc( "spindle_db_open_connections", "The number of established connections both in use and idle.", nil, nil, ), inUseDesc: prometheus.NewDesc( "spindle_db_in_use_connections", "The number of connections currently in use.", nil, nil, ), idleDesc: prometheus.NewDesc( "spindle_db_idle_connections", "The number of idle connections.", nil, nil, ), waitCountDesc: prometheus.NewDesc( "spindle_db_wait_count", "The total number of connections waited for.", nil, nil, ), waitDurationDesc: prometheus.NewDesc( "spindle_db_wait_duration_seconds", "The total time blocked waiting for a connection.", nil, nil, ), maxIdleClosedDesc: prometheus.NewDesc( "spindle_db_max_idle_closed", "The total number of connections closed due to SetMaxIdleConns.", nil, nil, ), maxLifetimeClosedDesc: prometheus.NewDesc( "spindle_db_max_lifetime_closed", "The total number of connections closed due to SetConnMaxLifetime.", nil, nil, ), workflowStatus: make(map[string]int), } m.reg.MustRegister(coll) go coll.run(ctx)}
func StartMetricsServer(ctx context.Context, addr string, logger *slog.Logger, reg *prometheus.Registry) (*http.Server, error) { if addr == "" { return nil, nil }
ln, err := net.Listen("tcp", addr) if err != nil { return nil, fmt.Errorf("failed to bind metrics port %s: %w", addr, err) }
mux := http.NewServeMux() mux.Handle("/metrics", promhttp.HandlerFor(reg, promhttp.HandlerOpts{EnableOpenMetrics: true}))
srv := &http.Server{ Addr: ln.Addr().String(), Handler: mux, ReadHeaderTimeout: 5 * time.Second, }
logger.Info("starting dedicated metrics listener", "address", srv.Addr)
go func() { if err := srv.Serve(ln); err != nil && err != http.ErrServerClosed { logger.Error("metrics server error", "err", err) } }()
go func() { <-ctx.Done() shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() _ = srv.Shutdown(shutdownCtx) }()
return srv, nil}
func boundMethod(m string) string { switch m { case "GET", "POST", "PUT", "DELETE", "PATCH", "HEAD", "OPTIONS": return m default: return "UNKNOWN" }}
func boundStatusCode(code int) string { if code >= 100 && code < 600 { return fmt.Sprintf("%d", code) } return "unknown"}
func getStatusClass(code int) string { if code >= 100 && code < 200 { return "1xx" } else if code >= 200 && code < 300 { return "2xx" } else if code >= 300 && code < 400 { return "3xx" } else if code >= 400 && code < 500 { return "4xx" } else if code >= 500 && code < 600 { return "5xx" } return "unknown"}
func boundWorkflowStatus(status string) string { switch status { case "pending", "running", "failed", "timeout", "cancelled", "success": return status default: return "unknown" }}
type recordedKey struct{}
func HTTPMiddleware(m *Metrics) func(http.Handler) http.Handler { return func(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if m == nil { next.ServeHTTP(w, r) return } if r.Context().Value(recordedKey{}) != nil { next.ServeHTTP(w, r) return } r = r.WithContext(context.WithValue(r.Context(), recordedKey{}, true))
m.httpInFlight.Inc() defer m.httpInFlight.Dec()
captured := httpsnoop.CaptureMetrics(next, w, r)
var route string rctx := chi.RouteContext(r.Context()) if rctx != nil { route = rctx.RoutePattern() } if route == "" { route = "unknown" }
method := boundMethod(r.Method) statusCode := boundStatusCode(captured.Code) statusClass := getStatusClass(captured.Code)
m.httpRequests.WithLabelValues(method, route, statusCode, statusClass).Inc() observer := m.httpRequestDuration.WithLabelValues(method, route, statusCode, statusClass) observeWithExemplar(r.Context(), observer, captured.Duration.Seconds()) }) }}