check / check (push) Successful in 4m2s
The HTTP metrics recorder, the delivery collectors and the Go and process collectors now register on one prometheus.Registry that fx provides, instead of Prometheus's global default registry, and /metrics serves that registry. A second metrics-enabled router in one process, or the server tests run with -count=2, no longer panics on a duplicate registration. The middleware builds its recorder once, in New, so installing Metrics() on more than one router over the same graph is also safe. The scrape keeps the same series and labels, including go_*, process_* and promhttp_metric_handler_*. Model: opus-5-5
406 lines
13 KiB
Go
406 lines
13 KiB
Go
// Package metrics defines the Prometheus collectors describing
|
|
// webhooker's delivery pipeline: how many events arrive, how many
|
|
// deliveries are attempted, how they end, how long they take, how
|
|
// deep the queues are, and how many circuit breakers are open.
|
|
//
|
|
// It also builds the registry the authenticated /metrics route
|
|
// serves. These collectors, the inbound HTTP metrics recorded in
|
|
// internal/middleware, and the Go runtime and process collectors all
|
|
// register on that one registry, never on Prometheus's global default.
|
|
package metrics
|
|
|
|
import (
|
|
"time"
|
|
|
|
"github.com/prometheus/client_golang/prometheus"
|
|
"github.com/prometheus/client_golang/prometheus/collectors"
|
|
"github.com/prometheus/client_golang/prometheus/promauto"
|
|
"sneak.berlin/go/webhooker/internal/database"
|
|
)
|
|
|
|
// namespace prefixes every collector defined here.
|
|
const namespace = "webhooker"
|
|
|
|
// targetTypeLabel is the only label any delivery metric carries, and
|
|
// cardinality is the whole reason for that.
|
|
//
|
|
// A target type is one of four compile-time constants, so the label
|
|
// domain is bounded by construction. Target ids, event ids and
|
|
// entrypoint ids are not: they are UUIDs minted per operator action
|
|
// or per inbound request, a series is never reclaimed once it exists,
|
|
// and labelling by any of them makes /metrics a memory leak that
|
|
// grows with traffic. normalizeTargetType enforces the bound at every
|
|
// call site — a type the registry does not know collapses into
|
|
// unknownTargetType rather than minting a series of its own.
|
|
const targetTypeLabel = "target_type"
|
|
|
|
// unknownTargetType is the bucket for a target type outside the known
|
|
// set, so an unrecognised value cannot mint a new series.
|
|
const unknownTargetType = "unknown"
|
|
|
|
// Delivery duration buckets, exponential from 5ms so the last bucket
|
|
// (about 98s) sits above the 30s outbound HTTP client timeout.
|
|
const (
|
|
durationBucketStart = 0.005
|
|
durationBucketFactor = 3
|
|
durationBucketCount = 10
|
|
)
|
|
|
|
// knownTargetTypes is the fixed label domain: the target types the
|
|
// delivery engine implements.
|
|
//
|
|
//nolint:gochecknoglobals // the label domain, built once per process
|
|
var knownTargetTypes = []database.TargetType{
|
|
database.TargetTypeHTTP,
|
|
database.TargetTypeDatabase,
|
|
database.TargetTypeLog,
|
|
database.TargetTypeSlack,
|
|
}
|
|
|
|
// NewRegistry returns the registry /metrics serves, carrying the Go
|
|
// runtime and process collectors that Prometheus's global default
|
|
// registry carries, so the go_* and process_* series stay in the
|
|
// scrape.
|
|
//
|
|
// A registry of its own, rather than the global default, is what lets
|
|
// two dependency graphs in one process — two tests, say — each
|
|
// register their collectors without the second registration
|
|
// panicking.
|
|
func NewRegistry() *prometheus.Registry {
|
|
reg := prometheus.NewRegistry()
|
|
reg.MustRegister(
|
|
collectors.NewGoCollector(),
|
|
collectors.NewProcessCollector(
|
|
collectors.ProcessCollectorOpts{},
|
|
),
|
|
)
|
|
|
|
return reg
|
|
}
|
|
|
|
// Set is one registered group of webhooker's delivery collectors.
|
|
// Production builds one on the registry /metrics serves; tests build
|
|
// their own against a private registry so assertions are not
|
|
// disturbed by deliveries other tests are making concurrently.
|
|
type Set struct {
|
|
eventsReceived prometheus.Counter
|
|
deliveryAttempts *prometheus.CounterVec
|
|
deliveriesSucceeded *prometheus.CounterVec
|
|
deliveriesFailed *prometheus.CounterVec
|
|
deliveryRetries *prometheus.CounterVec
|
|
deliveryReplays *prometheus.CounterVec
|
|
eventsResubmitted prometheus.Counter
|
|
deliveryDuration *prometheus.HistogramVec
|
|
deliveriesPending *prometheus.GaugeVec
|
|
deliveriesRetrying *prometheus.GaugeVec
|
|
circuitBreakersOpen *prometheus.GaugeVec
|
|
}
|
|
|
|
// New registers a full set of delivery collectors on reg and returns
|
|
// it. It panics if reg already holds them, which is the intended
|
|
// behaviour for a duplicate registration.
|
|
func New(reg *prometheus.Registry) *Set {
|
|
factory := promauto.With(reg)
|
|
|
|
s := &Set{
|
|
eventsReceived: factory.NewCounter(
|
|
prometheus.CounterOpts{
|
|
Namespace: namespace,
|
|
Name: "events_received_total",
|
|
Help: "Webhook events received and " +
|
|
"stored, so the receive and deliver " +
|
|
"sides can be compared.",
|
|
},
|
|
),
|
|
deliveryDuration: factory.NewHistogramVec(
|
|
prometheus.HistogramOpts{
|
|
Namespace: namespace,
|
|
Name: "delivery_duration_seconds",
|
|
Help: "Wall time of a single delivery " +
|
|
"attempt, by target type.",
|
|
Buckets: prometheus.ExponentialBuckets(
|
|
durationBucketStart,
|
|
durationBucketFactor,
|
|
durationBucketCount,
|
|
),
|
|
},
|
|
[]string{targetTypeLabel},
|
|
),
|
|
}
|
|
|
|
s.registerCounters(factory)
|
|
s.registerGauges(factory)
|
|
s.initSeries()
|
|
|
|
return s
|
|
}
|
|
|
|
// EventReceived counts one inbound webhook event stored.
|
|
func (s *Set) EventReceived() {
|
|
s.eventsReceived.Inc()
|
|
}
|
|
|
|
// DeliveryAttempted counts one delivery attempt dispatched to a
|
|
// target.
|
|
func (s *Set) DeliveryAttempted(t database.TargetType) {
|
|
s.deliveryAttempts.
|
|
WithLabelValues(normalizeTargetType(t)).
|
|
Inc()
|
|
}
|
|
|
|
// ObserveDeliveryDuration records how long one delivery attempt took.
|
|
func (s *Set) ObserveDeliveryDuration(
|
|
t database.TargetType, d time.Duration,
|
|
) {
|
|
s.deliveryDuration.
|
|
WithLabelValues(normalizeTargetType(t)).
|
|
Observe(d.Seconds())
|
|
}
|
|
|
|
// DeliveryReplayed counts one delivery an operator replayed from the
|
|
// event log.
|
|
//
|
|
// A replay runs the ordinary engine path, so it already moves the
|
|
// attempt, outcome and duration series exactly as a first delivery
|
|
// does — deliberately, since a replay is a real delivery and hiding it
|
|
// from those would misreport the pipeline. This counter is the one
|
|
// place the two are distinguishable, and it carries the existing
|
|
// target-type label rather than adding a replay dimension to every
|
|
// other series.
|
|
func (s *Set) DeliveryReplayed(t database.TargetType) {
|
|
s.deliveryReplays.
|
|
WithLabelValues(normalizeTargetType(t)).
|
|
Inc()
|
|
}
|
|
|
|
// EventResubmitted counts one stored event an operator re-injected
|
|
// from the event log.
|
|
//
|
|
// It counts the operator action once, not the deliveries it fans out
|
|
// to: those already move the attempt, outcome and duration series, and
|
|
// the new event moves events_received_total, since it is a stored
|
|
// event that the delivery side will be compared against. This counter
|
|
// is what separates a resubmitted event from a received one.
|
|
//
|
|
// It carries no labels. The only label available at the call site
|
|
// would be the route pattern, which has exactly one value and so would
|
|
// distinguish nothing; the target types the event fans out to belong
|
|
// to the delivery series, not to this one.
|
|
func (s *Set) EventResubmitted() {
|
|
s.eventsResubmitted.Inc()
|
|
}
|
|
|
|
// DeliveryStatusChanged counts a delivery's transition into a new
|
|
// status. The mapping from status to counter lives here, next to the
|
|
// collectors, so the engine has a single call for every transition it
|
|
// persists. A move back to pending is not an outcome and counts
|
|
// nothing.
|
|
func (s *Set) DeliveryStatusChanged(
|
|
t database.TargetType, status database.DeliveryStatus,
|
|
) {
|
|
label := normalizeTargetType(t)
|
|
|
|
switch status {
|
|
case database.DeliveryStatusDelivered:
|
|
s.deliveriesSucceeded.WithLabelValues(label).Inc()
|
|
case database.DeliveryStatusFailed:
|
|
s.deliveriesFailed.WithLabelValues(label).Inc()
|
|
case database.DeliveryStatusRetrying:
|
|
s.deliveryRetries.WithLabelValues(label).Inc()
|
|
case database.DeliveryStatusPending:
|
|
}
|
|
}
|
|
|
|
// SetQueueDepths publishes the pending and retrying queue depths from
|
|
// one sample. Every label in the queue domain is written on every
|
|
// call, so a type whose queue has drained reads zero instead of
|
|
// holding its last value forever.
|
|
func (s *Set) SetQueueDepths(
|
|
pending, retrying map[database.TargetType]int,
|
|
) {
|
|
pendingByLabel := foldToLabels(pending)
|
|
retryingByLabel := foldToLabels(retrying)
|
|
|
|
for _, label := range queueDepthLabels() {
|
|
s.deliveriesPending.WithLabelValues(label).
|
|
Set(float64(pendingByLabel[label]))
|
|
s.deliveriesRetrying.WithLabelValues(label).
|
|
Set(float64(retryingByLabel[label]))
|
|
}
|
|
}
|
|
|
|
// queueDepthLabels is the label domain of the two queue-depth gauges:
|
|
// the known target types plus unknown.
|
|
//
|
|
// Unknown is a real bucket here, not a safety net. A delivery queued
|
|
// against a target that has since been deleted carries a target id no
|
|
// longer in the targets table, so the sample resolves it to the empty
|
|
// type; folding it into unknown is what keeps that backlog visible.
|
|
// Dropping it would hide the one queue nobody is watching.
|
|
func queueDepthLabels() []string {
|
|
labels := make([]string, 0, len(knownTargetTypes)+1)
|
|
|
|
for _, t := range knownTargetTypes {
|
|
labels = append(labels, string(t))
|
|
}
|
|
|
|
return append(labels, unknownTargetType)
|
|
}
|
|
|
|
// foldToLabels collapses a per-target-type count onto the bounded
|
|
// label domain, summing everything outside the known set into
|
|
// unknown.
|
|
func foldToLabels(
|
|
counts map[database.TargetType]int,
|
|
) map[string]int {
|
|
byLabel := make(map[string]int, len(counts))
|
|
|
|
for t, n := range counts {
|
|
byLabel[normalizeTargetType(t)] += n
|
|
}
|
|
|
|
return byLabel
|
|
}
|
|
|
|
// SetCircuitBreakersOpen publishes how many of a target type's
|
|
// circuit breakers are currently open.
|
|
func (s *Set) SetCircuitBreakersOpen(
|
|
t database.TargetType, open int,
|
|
) {
|
|
s.circuitBreakersOpen.
|
|
WithLabelValues(normalizeTargetType(t)).
|
|
Set(float64(open))
|
|
}
|
|
|
|
func (s *Set) registerCounters(factory promauto.Factory) {
|
|
s.deliveryAttempts = factory.NewCounterVec(
|
|
prometheus.CounterOpts{
|
|
Namespace: namespace,
|
|
Name: "delivery_attempts_total",
|
|
Help: "Delivery attempts dispatched to a " +
|
|
"target, by target type.",
|
|
},
|
|
[]string{targetTypeLabel},
|
|
)
|
|
|
|
s.deliveriesSucceeded = factory.NewCounterVec(
|
|
prometheus.CounterOpts{
|
|
Namespace: namespace,
|
|
Name: "deliveries_succeeded_total",
|
|
Help: "Deliveries that reached the delivered " +
|
|
"state, by target type.",
|
|
},
|
|
[]string{targetTypeLabel},
|
|
)
|
|
|
|
s.deliveriesFailed = factory.NewCounterVec(
|
|
prometheus.CounterOpts{
|
|
Namespace: namespace,
|
|
Name: "deliveries_failed_total",
|
|
Help: "Deliveries that failed terminally and " +
|
|
"will not be retried, by target type.",
|
|
},
|
|
[]string{targetTypeLabel},
|
|
)
|
|
|
|
s.deliveryRetries = factory.NewCounterVec(
|
|
prometheus.CounterOpts{
|
|
Namespace: namespace,
|
|
Name: "delivery_retries_total",
|
|
Help: "Deliveries put back into the retrying " +
|
|
"state, by target type.",
|
|
},
|
|
[]string{targetTypeLabel},
|
|
)
|
|
|
|
s.deliveryReplays = factory.NewCounterVec(
|
|
prometheus.CounterOpts{
|
|
Namespace: namespace,
|
|
Name: "delivery_replays_total",
|
|
Help: "Deliveries an operator replayed from the " +
|
|
"event log, by target type.",
|
|
},
|
|
[]string{targetTypeLabel},
|
|
)
|
|
|
|
s.eventsResubmitted = factory.NewCounter(
|
|
prometheus.CounterOpts{
|
|
Namespace: namespace,
|
|
Name: "events_resubmitted_total",
|
|
Help: "Stored events an operator re-injected from " +
|
|
"the event log as new events.",
|
|
},
|
|
)
|
|
}
|
|
|
|
func (s *Set) registerGauges(factory promauto.Factory) {
|
|
s.deliveriesPending = factory.NewGaugeVec(
|
|
prometheus.GaugeOpts{
|
|
Namespace: namespace,
|
|
Name: "deliveries_pending",
|
|
Help: "Deliveries currently in the pending " +
|
|
"state, by target type.",
|
|
},
|
|
[]string{targetTypeLabel},
|
|
)
|
|
|
|
s.deliveriesRetrying = factory.NewGaugeVec(
|
|
prometheus.GaugeOpts{
|
|
Namespace: namespace,
|
|
Name: "deliveries_retrying",
|
|
Help: "Deliveries currently in the retrying " +
|
|
"state, by target type.",
|
|
},
|
|
[]string{targetTypeLabel},
|
|
)
|
|
|
|
s.circuitBreakersOpen = factory.NewGaugeVec(
|
|
prometheus.GaugeOpts{
|
|
Namespace: namespace,
|
|
Name: "circuit_breakers_open",
|
|
Help: "Delivery circuit breakers currently " +
|
|
"open, by target type.",
|
|
},
|
|
[]string{targetTypeLabel},
|
|
)
|
|
}
|
|
|
|
// initSeries materialises every known-target-type series at zero, so
|
|
// a dashboard and an alert rule see a target type that has not
|
|
// delivered yet rather than a missing series.
|
|
//
|
|
// The queue-depth gauges additionally get their unknown series, which
|
|
// holds deliveries queued against a deleted target. That backlog can
|
|
// predate the process — it is read out of the databases, not counted
|
|
// from transitions — so its series has to exist from the first scrape
|
|
// rather than appearing only once a backlog has already built up.
|
|
func (s *Set) initSeries() {
|
|
for _, t := range knownTargetTypes {
|
|
label := string(t)
|
|
|
|
s.deliveryAttempts.WithLabelValues(label)
|
|
s.deliveriesSucceeded.WithLabelValues(label)
|
|
s.deliveriesFailed.WithLabelValues(label)
|
|
s.deliveryRetries.WithLabelValues(label)
|
|
s.deliveryReplays.WithLabelValues(label)
|
|
s.deliveriesPending.WithLabelValues(label)
|
|
s.deliveriesRetrying.WithLabelValues(label)
|
|
s.circuitBreakersOpen.WithLabelValues(label)
|
|
}
|
|
|
|
s.deliveriesPending.WithLabelValues(unknownTargetType)
|
|
s.deliveriesRetrying.WithLabelValues(unknownTargetType)
|
|
}
|
|
|
|
// normalizeTargetType maps a target type onto the bounded label
|
|
// domain, collapsing anything outside it to unknownTargetType.
|
|
func normalizeTargetType(t database.TargetType) string {
|
|
for _, known := range knownTargetTypes {
|
|
if t == known {
|
|
return string(known)
|
|
}
|
|
}
|
|
|
|
return unknownTargetType
|
|
}
|