Expose delivery metrics on /metrics (closes #209)
All checks were successful
check / check (push) Successful in 4m13s
All checks were successful
check / check (push) Successful in 4m13s
/metrics carried only the inbound HTTP surface, so a destination failing for an hour, a growing retry backlog and a stuck-open circuit breaker were all invisible: the receive side stays healthy in each case because it is. New internal/metrics registers, on the existing default registry that the go-http-metrics recorder and the promhttp handler already share: - webhooker_events_received_total - webhooker_delivery_attempts_total - webhooker_deliveries_succeeded_total - webhooker_deliveries_failed_total - webhooker_delivery_retries_total - webhooker_delivery_duration_seconds - webhooker_deliveries_pending / _retrying - webhooker_circuit_breakers_open The route mounting is untouched. Every delivery metric carries one label, target_type, whose domain is the four target-type constants; anything outside it collapses to "unknown" so no series can be minted from a UUID. Target ids, event ids and entrypoint ids are deliberately not labels. Instrumentation sits at the points every target type already passes through: processDelivery for the attempt counter and the duration histogram, updateDeliveryStatus for the outcome counters. The queue-depth gauges are counted out of the per-webhook databases by a 30s sampler rather than tracked as deltas, which would need seeding at startup and would drift on any transition that failed to persist. The open-breaker gauge is recounted from the target's breaker registry on every state change.
This commit is contained in:
@@ -15,6 +15,7 @@ import (
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
"sneak.berlin/go/webhooker/internal/lifecycle"
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
"sneak.berlin/go/webhooker/internal/metrics"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -139,6 +140,12 @@ type Engine struct {
|
||||
retryCh chan Task
|
||||
workers int
|
||||
|
||||
// mx is the delivery metric set. Production wires the
|
||||
// process-wide one; a test can substitute a set registered on
|
||||
// a private registry so its assertions are not disturbed by
|
||||
// deliveries other tests are making at the same time.
|
||||
mx *metrics.Set
|
||||
|
||||
// targets maps each target type to its implementation.
|
||||
targets map[database.TargetType]Target
|
||||
|
||||
@@ -164,6 +171,7 @@ func New(
|
||||
deliveryCh: make(chan Task, deliveryChannelSize),
|
||||
retryCh: make(chan Task, retryChannelSize),
|
||||
workers: defaultWorkers,
|
||||
mx: metrics.Default(),
|
||||
}
|
||||
|
||||
e.initTargets(&http.Client{
|
||||
@@ -283,6 +291,10 @@ func (e *Engine) start() {
|
||||
|
||||
go e.retrySweep(ctx)
|
||||
|
||||
e.wg.Add(1)
|
||||
|
||||
go e.queueDepthSampler(ctx)
|
||||
|
||||
e.log.Info(
|
||||
"delivery engine started",
|
||||
"workers", e.workers,
|
||||
@@ -826,6 +838,11 @@ func (e *Engine) failUnretryableRetry(
|
||||
target.Type,
|
||||
)
|
||||
|
||||
// The delivery was loaded without its target relation, so
|
||||
// attach it: updateDeliveryStatus reads the type to label the
|
||||
// terminal failure it is about to count.
|
||||
d.Target = *target
|
||||
|
||||
e.recordResult(
|
||||
webhookDB,
|
||||
d,
|
||||
@@ -844,12 +861,21 @@ func (e *Engine) failUnretryableRetry(
|
||||
|
||||
// processDelivery dispatches a delivery to the target that
|
||||
// owns its type. Unknown target types fail the delivery.
|
||||
//
|
||||
// It is also where the attempt counter and the duration histogram
|
||||
// are recorded, because it is the one point every target type
|
||||
// passes through on every attempt: a target added later is
|
||||
// instrumented without touching it, and the duration measured is
|
||||
// the whole cost of the attempt rather than whatever each target
|
||||
// happens to time for itself.
|
||||
func (e *Engine) processDelivery(
|
||||
ctx context.Context,
|
||||
webhookDB *gorm.DB,
|
||||
d *database.Delivery,
|
||||
task *Task,
|
||||
) {
|
||||
e.mx.DeliveryAttempted(d.Target.Type)
|
||||
|
||||
target, ok := e.targets[d.Target.Type]
|
||||
if !ok {
|
||||
e.log.Error(
|
||||
@@ -865,7 +891,13 @@ func (e *Engine) processDelivery(
|
||||
return
|
||||
}
|
||||
|
||||
start := time.Now()
|
||||
|
||||
target.Deliver(ctx, webhookDB, d, task, e)
|
||||
|
||||
e.mx.ObserveDeliveryDuration(
|
||||
d.Target.Type, time.Since(start),
|
||||
)
|
||||
}
|
||||
|
||||
// recordResult persists a DeliveryResult row describing a
|
||||
@@ -901,12 +933,16 @@ func (e *Engine) recordResult(
|
||||
}
|
||||
|
||||
// updateDeliveryStatus persists a new status for a delivery.
|
||||
// It is a cross-target helper the targets call.
|
||||
// It is a cross-target helper the targets call, and therefore the
|
||||
// single point where a delivery's outcome — delivered, terminally
|
||||
// failed, or put back into retry — is counted.
|
||||
func (e *Engine) updateDeliveryStatus(
|
||||
webhookDB *gorm.DB,
|
||||
d *database.Delivery,
|
||||
status database.DeliveryStatus,
|
||||
) {
|
||||
e.mx.DeliveryStatusChanged(d.Target.Type, status)
|
||||
|
||||
err := webhookDB.Model(d).
|
||||
Update("status", status).Error
|
||||
if err != nil {
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"go.uber.org/fx"
|
||||
"gorm.io/gorm"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
"sneak.berlin/go/webhooker/internal/metrics"
|
||||
)
|
||||
|
||||
// ErrExportArchiveWriterEvicted exposes the sentinel returned by
|
||||
@@ -253,6 +254,7 @@ func NewTestEngine(
|
||||
deliveryCh: make(chan Task, deliveryChannelSize),
|
||||
retryCh: make(chan Task, retryChannelSize),
|
||||
workers: workers,
|
||||
mx: metrics.Default(),
|
||||
}
|
||||
e.initTargets(client)
|
||||
|
||||
@@ -267,6 +269,7 @@ func NewTestEngineSmallRetry(
|
||||
e := &Engine{
|
||||
log: log,
|
||||
retryCh: make(chan Task, 1),
|
||||
mx: metrics.Default(),
|
||||
}
|
||||
e.initTargets(nil)
|
||||
|
||||
@@ -289,12 +292,25 @@ func NewTestEngineWithDB(
|
||||
deliveryCh: make(chan Task, deliveryChannelSize),
|
||||
retryCh: make(chan Task, retryChannelSize),
|
||||
workers: workers,
|
||||
mx: metrics.Default(),
|
||||
}
|
||||
e.initTargets(client)
|
||||
|
||||
return e
|
||||
}
|
||||
|
||||
// ExportSetMetrics substitutes the engine's metric set, so a test can
|
||||
// assert on collectors registered on a private registry instead of
|
||||
// the process-wide ones every other test is also moving.
|
||||
func (e *Engine) ExportSetMetrics(mx *metrics.Set) {
|
||||
e.mx = mx
|
||||
}
|
||||
|
||||
// ExportSampleQueueDepths runs one queue depth sample synchronously.
|
||||
func (e *Engine) ExportSampleQueueDepths(ctx context.Context) {
|
||||
e.sampleQueueDepths(ctx)
|
||||
}
|
||||
|
||||
// NewTestCircuitBreaker creates a CircuitBreaker with
|
||||
// custom settings for testing.
|
||||
func NewTestCircuitBreaker(
|
||||
|
||||
378
internal/delivery/metrics_test.go
Normal file
378
internal/delivery/metrics_test.go
Normal file
@@ -0,0 +1,378 @@
|
||||
package delivery_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
dto "github.com/prometheus/client_model/go"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
"sneak.berlin/go/webhooker/internal/delivery"
|
||||
"sneak.berlin/go/webhooker/internal/metrics"
|
||||
)
|
||||
|
||||
// Metric names as exposed on /metrics.
|
||||
const (
|
||||
mAttempts = "webhooker_delivery_attempts_total"
|
||||
mSucceeded = "webhooker_deliveries_succeeded_total"
|
||||
mFailed = "webhooker_deliveries_failed_total"
|
||||
mRetries = "webhooker_delivery_retries_total"
|
||||
mDuration = "webhooker_delivery_duration_seconds"
|
||||
mPending = "webhooker_deliveries_pending"
|
||||
mRetrying = "webhooker_deliveries_retrying"
|
||||
mBreakers = "webhooker_circuit_breakers_open"
|
||||
)
|
||||
|
||||
const (
|
||||
mTypeHTTP = "http"
|
||||
mTypeLog = "log"
|
||||
)
|
||||
|
||||
// mIsolate gives the setup's engine a metric set registered on a
|
||||
// private registry. The process-wide collectors are moved by every
|
||||
// other delivery test running in parallel, so exact assertions are
|
||||
// only possible against a registry this test owns.
|
||||
func mIsolate(
|
||||
t *testing.T, s iSetup,
|
||||
) *prometheus.Registry {
|
||||
t.Helper()
|
||||
|
||||
reg := prometheus.NewRegistry()
|
||||
s.Engine.ExportSetMetrics(metrics.New(reg))
|
||||
|
||||
return reg
|
||||
}
|
||||
|
||||
// mFind returns the series of the named metric carrying the given
|
||||
// target_type label.
|
||||
func mFind(
|
||||
t *testing.T,
|
||||
reg *prometheus.Registry,
|
||||
name, targetType string,
|
||||
) *dto.Metric {
|
||||
t.Helper()
|
||||
|
||||
families, err := reg.Gather()
|
||||
require.NoError(t, err)
|
||||
|
||||
for _, fam := range families {
|
||||
if fam.GetName() != name {
|
||||
continue
|
||||
}
|
||||
|
||||
for _, m := range fam.GetMetric() {
|
||||
for _, label := range m.GetLabel() {
|
||||
if label.GetName() == "target_type" &&
|
||||
label.GetValue() == targetType {
|
||||
return m
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
t.Fatalf(
|
||||
"metric %s{target_type=%q} not found",
|
||||
name, targetType,
|
||||
)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func mCounter(
|
||||
t *testing.T,
|
||||
reg *prometheus.Registry,
|
||||
name, targetType string,
|
||||
) float64 {
|
||||
t.Helper()
|
||||
|
||||
return mFind(t, reg, name, targetType).
|
||||
GetCounter().GetValue()
|
||||
}
|
||||
|
||||
func mGauge(
|
||||
t *testing.T,
|
||||
reg *prometheus.Registry,
|
||||
name, targetType string,
|
||||
) float64 {
|
||||
t.Helper()
|
||||
|
||||
return mFind(t, reg, name, targetType).
|
||||
GetGauge().GetValue()
|
||||
}
|
||||
|
||||
func mObservations(
|
||||
t *testing.T,
|
||||
reg *prometheus.Registry,
|
||||
name, targetType string,
|
||||
) uint64 {
|
||||
t.Helper()
|
||||
|
||||
return mFind(t, reg, name, targetType).
|
||||
GetHistogram().GetSampleCount()
|
||||
}
|
||||
|
||||
// TestDeliveryMetrics_SuccessAndRetryExhaustion drives one delivery
|
||||
// that succeeds and one that fails every attempt until its retries
|
||||
// are exhausted, and asserts every delivery counter across both.
|
||||
func TestDeliveryMetrics_SuccessAndRetryExhaustion(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
s := newISetup(t)
|
||||
reg := mIsolate(t, s)
|
||||
|
||||
mDeliverOK(t, s)
|
||||
|
||||
assert.InDelta(t, 1.0,
|
||||
mCounter(t, reg, mAttempts, mTypeHTTP), 0)
|
||||
assert.InDelta(t, 1.0,
|
||||
mCounter(t, reg, mSucceeded, mTypeHTTP), 0)
|
||||
assert.InDelta(t, 0.0,
|
||||
mCounter(t, reg, mFailed, mTypeHTTP), 0)
|
||||
assert.InDelta(t, 0.0,
|
||||
mCounter(t, reg, mRetries, mTypeHTTP), 0)
|
||||
assert.Equal(t, uint64(1),
|
||||
mObservations(t, reg, mDuration, mTypeHTTP))
|
||||
|
||||
mExhaustRetries(t, s)
|
||||
|
||||
// Two further attempts: the first is retried, the second is
|
||||
// the last one allowed and fails the delivery terminally.
|
||||
assert.InDelta(t, 3.0,
|
||||
mCounter(t, reg, mAttempts, mTypeHTTP), 0)
|
||||
assert.InDelta(t, 1.0,
|
||||
mCounter(t, reg, mSucceeded, mTypeHTTP), 0)
|
||||
assert.InDelta(t, 1.0,
|
||||
mCounter(t, reg, mRetries, mTypeHTTP), 0)
|
||||
assert.InDelta(t, 1.0,
|
||||
mCounter(t, reg, mFailed, mTypeHTTP), 0)
|
||||
assert.Equal(t, uint64(3),
|
||||
mObservations(t, reg, mDuration, mTypeHTTP))
|
||||
|
||||
// Two consecutive failures are below the trip threshold.
|
||||
assert.InDelta(t, 0.0,
|
||||
mGauge(t, reg, mBreakers, mTypeHTTP), 0)
|
||||
|
||||
// The label is the target type and nothing finer: two http
|
||||
// targets shared one series, and no other type's moved.
|
||||
assert.InDelta(t, 0.0,
|
||||
mCounter(t, reg, mAttempts, mTypeLog), 0)
|
||||
assert.InDelta(t, 0.0,
|
||||
mCounter(t, reg, mFailed, mTypeLog), 0)
|
||||
}
|
||||
|
||||
// mDeliverOK delivers one event to a target that answers 200.
|
||||
func mDeliverOK(t *testing.T, s iSetup) {
|
||||
t.Helper()
|
||||
|
||||
ts := httptest.NewServer(http.HandlerFunc(
|
||||
func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
},
|
||||
))
|
||||
defer ts.Close()
|
||||
|
||||
event := iSeedEvent(
|
||||
t, s.WebhookDB, s.WebhookID, `{"ok":true}`,
|
||||
)
|
||||
targetID := uuid.New().String()
|
||||
|
||||
d := iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, targetID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
body := event.Body
|
||||
task := iTask(
|
||||
d, event, s.WebhookID, targetID,
|
||||
"metrics-ok", iHTTPConfig(ts.URL), 3, 1, &body,
|
||||
)
|
||||
|
||||
s.Engine.ExportProcessNewTask(context.TODO(), &task)
|
||||
|
||||
iAssertStatus(t, s.WebhookDB, d.ID,
|
||||
database.DeliveryStatusDelivered,
|
||||
)
|
||||
}
|
||||
|
||||
// mExhaustRetries delivers to a target that answers 500 with a
|
||||
// two-attempt budget, driving both attempts so the delivery ends
|
||||
// terminally failed.
|
||||
func mExhaustRetries(t *testing.T, s iSetup) {
|
||||
t.Helper()
|
||||
|
||||
ts := httptest.NewServer(http.HandlerFunc(
|
||||
func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
},
|
||||
))
|
||||
defer ts.Close()
|
||||
|
||||
event := iSeedEvent(
|
||||
t, s.WebhookDB, s.WebhookID, `{"ok":false}`,
|
||||
)
|
||||
targetID := uuid.New().String()
|
||||
|
||||
d := iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, targetID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
body := event.Body
|
||||
cfg := iHTTPConfig(ts.URL)
|
||||
|
||||
first := iTask(
|
||||
d, event, s.WebhookID, targetID,
|
||||
"metrics-fail", cfg, 2, 1, &body,
|
||||
)
|
||||
|
||||
s.Engine.ExportProcessNewTask(context.TODO(), &first)
|
||||
|
||||
iAssertStatus(t, s.WebhookDB, d.ID,
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
|
||||
// The engine's own scheduler would re-enqueue this after the
|
||||
// backoff; driving the second attempt directly keeps the test
|
||||
// deterministic and off the wall clock.
|
||||
second := iTask(
|
||||
d, event, s.WebhookID, targetID,
|
||||
"metrics-fail", cfg, 2, 2, &body,
|
||||
)
|
||||
|
||||
s.Engine.ExportProcessRetryTask(
|
||||
context.TODO(), &second,
|
||||
)
|
||||
|
||||
iAssertStatus(t, s.WebhookDB, d.ID,
|
||||
database.DeliveryStatusFailed,
|
||||
)
|
||||
}
|
||||
|
||||
// TestDeliveryMetrics_CircuitBreakerGauge proves the open-breaker
|
||||
// gauge follows a breaker that trips.
|
||||
func TestDeliveryMetrics_CircuitBreakerGauge(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
s := newISetup(t)
|
||||
reg := mIsolate(t, s)
|
||||
|
||||
ts := httptest.NewServer(http.HandlerFunc(
|
||||
func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
},
|
||||
))
|
||||
defer ts.Close()
|
||||
|
||||
event := iSeedEvent(
|
||||
t, s.WebhookDB, s.WebhookID, `{"trip":true}`,
|
||||
)
|
||||
targetID := uuid.New().String()
|
||||
|
||||
d := iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, targetID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
body := event.Body
|
||||
cfg := iHTTPConfig(ts.URL)
|
||||
|
||||
// A retry budget above the failure threshold, so the breaker
|
||||
// rather than the budget is what stops the delivery.
|
||||
maxRetries := delivery.ExportDefaultFailureThreshold + 5
|
||||
|
||||
first := iTask(
|
||||
d, event, s.WebhookID, targetID,
|
||||
"metrics-trip", cfg, maxRetries, 1, &body,
|
||||
)
|
||||
|
||||
s.Engine.ExportProcessNewTask(context.TODO(), &first)
|
||||
|
||||
assert.InDelta(t, 0.0,
|
||||
mGauge(t, reg, mBreakers, mTypeHTTP), 0)
|
||||
|
||||
for attempt := 2; attempt <= delivery.
|
||||
ExportDefaultFailureThreshold; attempt++ {
|
||||
task := iTask(
|
||||
d, event, s.WebhookID, targetID,
|
||||
"metrics-trip", cfg, maxRetries, attempt, &body,
|
||||
)
|
||||
|
||||
s.Engine.ExportProcessRetryTask(
|
||||
context.TODO(), &task,
|
||||
)
|
||||
}
|
||||
|
||||
assert.InDelta(t, 1.0,
|
||||
mGauge(t, reg, mBreakers, mTypeHTTP), 0)
|
||||
}
|
||||
|
||||
// TestDeliveryMetrics_QueueDepthGauges proves the sampler publishes
|
||||
// the queued deliveries it finds in the per-webhook databases, and
|
||||
// that a drained queue reads zero rather than keeping its last
|
||||
// value.
|
||||
func TestDeliveryMetrics_QueueDepthGauges(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
s := newISetup(t)
|
||||
reg := mIsolate(t, s)
|
||||
|
||||
iCreateWebhook(
|
||||
t, s.MainDB, s.WebhookID, "queue-depth",
|
||||
)
|
||||
|
||||
targetID := uuid.New().String()
|
||||
|
||||
iCreateTarget(t, s.MainDB, targetID, s.WebhookID,
|
||||
"queue-depth-target", database.TargetTypeHTTP,
|
||||
iHTTPConfig("https://example.com/hook"), 3,
|
||||
)
|
||||
|
||||
event := iSeedEvent(
|
||||
t, s.WebhookDB, s.WebhookID, `{"queued":true}`,
|
||||
)
|
||||
|
||||
pending := iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, targetID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, targetID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
retrying := iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, targetID,
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
|
||||
s.Engine.ExportSampleQueueDepths(context.Background())
|
||||
|
||||
assert.InDelta(t, 2.0,
|
||||
mGauge(t, reg, mPending, mTypeHTTP), 0)
|
||||
assert.InDelta(t, 1.0,
|
||||
mGauge(t, reg, mRetrying, mTypeHTTP), 0)
|
||||
assert.InDelta(t, 0.0,
|
||||
mGauge(t, reg, mPending, mTypeLog), 0)
|
||||
|
||||
require.NoError(t, s.WebhookDB.
|
||||
Model(&database.Delivery{}).
|
||||
Where("id IN ?", []string{pending.ID, retrying.ID}).
|
||||
Update(
|
||||
"status", database.DeliveryStatusDelivered,
|
||||
).Error)
|
||||
|
||||
s.Engine.ExportSampleQueueDepths(context.Background())
|
||||
|
||||
assert.InDelta(t, 1.0,
|
||||
mGauge(t, reg, mPending, mTypeHTTP), 0)
|
||||
assert.InDelta(t, 0.0,
|
||||
mGauge(t, reg, mRetrying, mTypeHTTP), 0)
|
||||
}
|
||||
183
internal/delivery/queue_depth.go
Normal file
183
internal/delivery/queue_depth.go
Normal file
@@ -0,0 +1,183 @@
|
||||
package delivery
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
)
|
||||
|
||||
// queueDepthSampleInterval is how often the pending and retrying
|
||||
// queue depths are counted and published as gauges.
|
||||
const queueDepthSampleInterval = 30 * time.Second
|
||||
|
||||
// queueDepthSampler publishes the pending and retrying queue depths
|
||||
// on a timer for as long as the engine runs.
|
||||
//
|
||||
// The depths are counted out of the databases rather than tracked as
|
||||
// deltas alongside the status transitions. A delta counter would have
|
||||
// to be seeded correctly at startup from rows written by a previous
|
||||
// process, and would drift permanently on any transition that failed
|
||||
// to persist. Counting is the measurement that cannot go wrong, and
|
||||
// it is the same whole-database walk the retry sweep already makes.
|
||||
func (e *Engine) queueDepthSampler(ctx context.Context) {
|
||||
defer e.wg.Done()
|
||||
|
||||
ticker := time.NewTicker(queueDepthSampleInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
e.sampleQueueDepths(ctx)
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
e.sampleQueueDepths(ctx)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// sampleQueueDepths counts every queued delivery across all
|
||||
// per-webhook databases and publishes the result.
|
||||
func (e *Engine) sampleQueueDepths(ctx context.Context) {
|
||||
if e.database == nil || e.dbManager == nil {
|
||||
return
|
||||
}
|
||||
|
||||
types, err := e.targetTypesByID()
|
||||
if err != nil {
|
||||
e.log.Error(
|
||||
"queue depth sample: failed to load target types",
|
||||
"error", err,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
var webhookIDs []string
|
||||
|
||||
err = e.database.DB().
|
||||
Model(&database.Webhook{}).
|
||||
Pluck("id", &webhookIDs).Error
|
||||
if err != nil {
|
||||
e.log.Error(
|
||||
"queue depth sample: failed to query webhook IDs",
|
||||
"error", err,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
pending := make(map[database.TargetType]int)
|
||||
retrying := make(map[database.TargetType]int)
|
||||
|
||||
for _, webhookID := range webhookIDs {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
default:
|
||||
}
|
||||
|
||||
if !e.dbManager.DBExists(webhookID) {
|
||||
continue
|
||||
}
|
||||
|
||||
e.sampleWebhookQueueDepths(
|
||||
webhookID, types, pending, retrying,
|
||||
)
|
||||
}
|
||||
|
||||
e.mx.SetQueueDepths(pending, retrying)
|
||||
}
|
||||
|
||||
// targetTypesByID maps every configured target id to its type. The
|
||||
// deliveries live in the per-webhook databases but carry only a
|
||||
// target id, so the type label has to come from the main database.
|
||||
func (e *Engine) targetTypesByID() (
|
||||
map[string]database.TargetType, error,
|
||||
) {
|
||||
var rows []struct {
|
||||
ID string
|
||||
Type database.TargetType
|
||||
}
|
||||
|
||||
err := e.database.DB().
|
||||
Model(&database.Target{}).
|
||||
Select("id", "type").
|
||||
Scan(&rows).Error
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("loading targets: %w", err)
|
||||
}
|
||||
|
||||
types := make(map[string]database.TargetType, len(rows))
|
||||
|
||||
for _, row := range rows {
|
||||
types[row.ID] = row.Type
|
||||
}
|
||||
|
||||
return types, nil
|
||||
}
|
||||
|
||||
// sampleWebhookQueueDepths adds one webhook's queued deliveries into
|
||||
// the running totals. A delivery whose target has since been deleted
|
||||
// resolves to the empty type and lands in the unknown bucket rather
|
||||
// than being dropped.
|
||||
func (e *Engine) sampleWebhookQueueDepths(
|
||||
webhookID string,
|
||||
types map[string]database.TargetType,
|
||||
pending, retrying map[database.TargetType]int,
|
||||
) {
|
||||
webhookDB, err := e.dbManager.GetDB(webhookID)
|
||||
if err != nil {
|
||||
e.log.Error(
|
||||
"queue depth sample: failed to get webhook database",
|
||||
"webhook_id", webhookID,
|
||||
"error", err,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
var rows []struct {
|
||||
TargetID string
|
||||
Status database.DeliveryStatus
|
||||
Depth int
|
||||
}
|
||||
|
||||
err = webhookDB.
|
||||
Model(&database.Delivery{}).
|
||||
Select("target_id", "status", "count(*) as depth").
|
||||
Where("status IN ?", []database.DeliveryStatus{
|
||||
database.DeliveryStatusPending,
|
||||
database.DeliveryStatusRetrying,
|
||||
}).
|
||||
Group("target_id, status").
|
||||
Scan(&rows).Error
|
||||
if err != nil {
|
||||
e.log.Error(
|
||||
"queue depth sample: "+
|
||||
"failed to count queued deliveries",
|
||||
"webhook_id", webhookID,
|
||||
"error", err,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
for _, row := range rows {
|
||||
targetType := types[row.TargetID]
|
||||
|
||||
switch row.Status {
|
||||
case database.DeliveryStatusPending:
|
||||
pending[targetType] += row.Depth
|
||||
case database.DeliveryStatusRetrying:
|
||||
retrying[targetType] += row.Depth
|
||||
case database.DeliveryStatusDelivered,
|
||||
database.DeliveryStatusFailed:
|
||||
// Excluded by the query above: a delivery that has
|
||||
// reached a terminal state is not queued.
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -107,6 +107,11 @@ func (c *httpCore) withRetry(
|
||||
return
|
||||
}
|
||||
|
||||
// Allow may have moved the breaker to half-open, and the
|
||||
// attempt below may open or close it, so the gauge is
|
||||
// republished on every exit from here.
|
||||
defer c.publishCircuitState(d.Target.Type)
|
||||
|
||||
attemptNum := task.AttemptNum
|
||||
|
||||
res := attempt()
|
||||
@@ -146,6 +151,8 @@ func (c *httpCore) circuitBreakerBlock(
|
||||
return false
|
||||
}
|
||||
|
||||
defer c.publishCircuitState(d.Target.Type)
|
||||
|
||||
remaining := cb.CooldownRemaining()
|
||||
|
||||
c.eng.log.Info(
|
||||
@@ -215,6 +222,28 @@ func (c *httpCore) getCircuitBreaker(
|
||||
return cb
|
||||
}
|
||||
|
||||
// publishCircuitState recounts this core's open breakers and
|
||||
// publishes the gauge. Each core holds the breakers of exactly one
|
||||
// target type, so the recount is over that type's targets alone.
|
||||
// Counting rather than adjusting a delta keeps the gauge honest
|
||||
// however a breaker changed state.
|
||||
func (c *httpCore) publishCircuitState(
|
||||
targetType database.TargetType,
|
||||
) {
|
||||
open := 0
|
||||
|
||||
c.circuitBreakers.Range(func(_, val any) bool {
|
||||
cb, ok := val.(*CircuitBreaker)
|
||||
if ok && cb.State() == CircuitOpen {
|
||||
open++
|
||||
}
|
||||
|
||||
return true
|
||||
})
|
||||
|
||||
c.eng.mx.SetCircuitBreakersOpen(targetType, open)
|
||||
}
|
||||
|
||||
// remainingBackoff returns how long remains of the backoff
|
||||
// window for the last attempt of a recovered retrying
|
||||
// delivery. It implements rescheduler.
|
||||
|
||||
Reference in New Issue
Block a user