Compare commits

..
Author SHA1 Message Date
clawbot 4a24b5f889 Serve /metrics from a registry of its own (closes #227)
check / check (push) Waiting to run
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;
NewForTest gives its Middleware a recorder on a fresh registry. The
scrape keeps the same series and labels, including go_*, process_*
and promhttp_metric_handler_*.

Model: opus-5-5
2026-10-01 21:00:29 +00:00
35 changed files with 294 additions and 1679 deletions
+18 -61
View File
@@ -1065,7 +1065,7 @@ unconditionally against whatever files it finds:
- the main database on connect — `Setting`, `User`, `APIKey`, `Webhook`, - the main database on connect — `Setting`, `User`, `APIKey`, `Webhook`,
`Entrypoint`, `Target` `Entrypoint`, `Target`
- each event database when it is lazily opened — `Event`, `Delivery`, - each event database when it is lazily opened — `Event`, `Delivery`,
`DeliveryResult`, `EventTotals`, `TargetTotals` `DeliveryResult`
- each archive database on every open and reopen - each archive database on every open and reopen
There is no schema version table, no migration ledger, and no down There is no schema version table, no migration ledger, and no down
@@ -1381,7 +1381,7 @@ The codebase uses consistent naming throughout (rename completed in
### Data Model ### Data Model
webhooker's data model has eleven entities organized into two tiers: the webhooker's data model has nine entities organized into two tiers: the
**application tier** (user and webhook configuration) and the **event **application tier** (user and webhook configuration) and the **event
tier** (event ingestion, delivery, and logging). tier** (event ingestion, delivery, and logging).
@@ -1410,13 +1410,6 @@ tier** (event ingestion, delivery, and logging).
│ ┌──────────┐ ┌──────────┐ ┌─────────────────┐ │ │ ┌──────────┐ ┌──────────┐ ┌─────────────────┐ │
│ │ Event │──1:N──│ Delivery │──1:N──│ DeliveryResult │ │ │ │ Event │──1:N──│ Delivery │──1:N──│ DeliveryResult │ │
│ └──────────┘ └──────────┘ └─────────────────┘ │ │ └──────────┘ └──────────┘ └─────────────────┘ │
│ │
│ ┌──────────────┐ (one row: running counts of events) │
│ │ EventTotals │ │
│ └──────────────┘ │
│ ┌──────────────┐ (one row per target: running counts │
│ │ TargetTotals │ of its deliveries) │
│ └──────────────┘ │
└─────────────────────────────────────────────────────────────┘ └─────────────────────────────────────────────────────────────┘
``` ```
@@ -1669,7 +1662,6 @@ status across potentially multiple attempts.
| `event_id` | UUID | Foreign key → Event | | `event_id` | UUID | Foreign key → Event |
| `target_id`| UUID | Foreign key → Target | | `target_id`| UUID | Foreign key → Target |
| `status` | DeliveryStatus | One of: `pending`, `delivered`, `failed`, `retrying` | | `status` | DeliveryStatus | One of: `pending`, `delivered`, `failed`, `retrying` |
| `finished_at` | timestamp | When the delivery became `delivered` or `failed` (nullable; empty while `pending` or `retrying`) |
**Relations:** Belongs to Event. Belongs to Target. Has many **Relations:** Belongs to Event. Belongs to Target. Has many
DeliveryResults. DeliveryResults.
@@ -1737,65 +1729,33 @@ retries) is individually logged for full observability.
**Relations:** Belongs to Delivery. **Relations:** Belongs to Delivery.
#### EventTotals and TargetTotals
Running counts in each event database, read by the statistics pane at the
top of the webhook page. `EventTotals` is one row:
| Field | Type | Description |
| ---------------- | ------- | ----------- |
| `events` | integer | Events ever stored, resubmitted copies included |
| `events_removed` | integer | Events retention has deleted |
`TargetTotals` is one row per target, created by the first delivery to it:
| Field | Type | Description |
| -------------------- | ------- | ----------- |
| `target_id` | UUID | The target (primary key) |
| `deliveries` | integer | Deliveries to it ever created, replays included |
| `delivered` | integer | Of those, how many became `delivered` |
| `failed` | integer | Of those, how many became `failed` |
| `deliveries_removed` | integer | Its deliveries retention has deleted |
| `failed_removed` | integer | Its failed deliveries retention has deleted |
Each count changes in the transaction that writes or deletes the rows it
counts. The pane's lifetime events are `events`, and its lifetime
deliveries and failures are `deliveries` and `failed` summed over the
targets; each figure within retention is the same less what retention
removed, so neither needs the rows themselves. Its last-10-minutes and
last-24-hours figures are counted from the `events` and `deliveries`
indexes over just that window, the deliveries in one query grouped by
target. Its failure percentage for a window is the deliveries that became
`failed` in it out of all that became `delivered` or `failed` in it, and
a dash when none did.
#### Event-tier indexes #### Event-tier indexes
These indexes on the per-webhook event databases are declared in the model These indexes on the per-webhook event databases are declared in the model
tags, so `AutoMigrate` creates them on a fresh database: tags, so `AutoMigrate` creates them on a fresh and on an existing database:
| Table | Columns | Serves | | Table | Columns | Serves |
| ------------------ | --------------------------- | ------ | | ------------------ | --------------------------- | ------ |
| `deliveries` | `status`, `deleted_at`, `finished_at`, `target_id` | Startup recovery, the retry and pending sweeps every 60 seconds and the queue-depth sampler every 30 seconds, which select deliveries by status, and the webhook page's statistics, which count each target's deliveries by status and when they finished | | `deliveries` | `status`, `deleted_at` | Startup recovery, the retry and pending sweeps every 60 seconds and the queue-depth sampler every 30 seconds, which select deliveries by status |
| `deliveries` | `event_id`, `deleted_at` | The event log, which loads each event's deliveries, and retention, which counts and deletes the deliveries of expired events | | `deliveries` | `event_id`, `deleted_at` | The event log, which loads each event's deliveries, and retention, which selects and deletes the deliveries of expired events |
| `delivery_results` | `delivery_id`, `deleted_at` | The event log, which loads the attempts of a page's deliveries, and retention, which deletes the attempts of expired events | | `delivery_results` | `delivery_id`, `deleted_at` | The event log, which loads the attempts of a page's deliveries, and retention, which deletes the attempts of expired events |
| `events` | `deleted_at`, `created_at` | The webhook page's statistics, which count recent events and find the newest | | `events` | `deleted_at`, `created_at` | Retention, which selects expired events by age |
| `events` | `created_at` | Retention, which selects expired events by age | | `events` | `created_at` | Retention's delete of the expired events themselves |
GORM's soft delete adds `deleted_at IS NULL` to these queries; retention GORM's soft delete adds `deleted_at IS NULL` to these queries; retention's
leaves it out. SQLite keeps no statistics on these tables, and without them it deletes leave it out, but their lookups of expired rows keep it. SQLite keeps
rates the `deleted_at` index, which every live row matches, above an index on no statistics on these tables, and without them it rates the `deleted_at`
a column matched against several values or compared with `<`. So every index index, which every live row matches, above an index on a column matched
but the last also covers `deleted_at`. It comes second, so that retention can against several values or compared with `<`. So every index but the last also
use the index without it, except in `events`, where `created_at` is compared covers `deleted_at`. It comes second, so that retention's deletes can use the
with `<` and SQLite narrows by a `<` only on the last column it uses. index without it, except in `events`, where `created_at` is compared with `<`
and SQLite narrows by a `<` only on the last column it uses.
#### Common Fields #### Common Fields
Every entity except `Setting`, `EventTotals` and `TargetTotals` includes Every entity except `Setting` includes these fields from `BaseModel`.
these fields from `BaseModel`. `Setting` is a bare key-value row with no `Setting` is a bare key-value row with no `id`, no timestamps and no
`id`, no timestamps and no soft delete, and the two totals tables hold soft delete:
only counts, keyed by a numeric `id` and by `target_id`:
| Field | Type | Description | | Field | Type | Description |
| ------------ | --------- | ----------- | | ------------ | --------- | ----------- |
@@ -1837,8 +1797,6 @@ encryption key is generated and stored, and an `admin` user is created.
- **Events** — captured incoming webhook payloads - **Events** — captured incoming webhook payloads
- **Deliveries** — event-to-target pairings and their status - **Deliveries** — event-to-target pairings and their status
- **DeliveryResults** — individual delivery attempt logs - **DeliveryResults** — individual delivery attempt logs
- **EventTotals** and **TargetTotals** — running counts of the above,
the deliveries per target, kept through retention
Per-webhook databases are created automatically when a webhook is Per-webhook databases are created automatically when a webhook is
created (and lazily on first access for webhooks that predate this created (and lazily on first access for webhooks that predate this
@@ -2819,7 +2777,6 @@ webhooker/
│ │ ├── model_event.go # Event entity (per-webhook DB) │ │ ├── model_event.go # Event entity (per-webhook DB)
│ │ ├── model_delivery.go # Delivery entity (per-webhook DB) │ │ ├── model_delivery.go # Delivery entity (per-webhook DB)
│ │ ├── model_delivery_result.go # DeliveryResult entity (per-webhook DB) │ │ ├── model_delivery_result.go # DeliveryResult entity (per-webhook DB)
│ │ ├── model_totals.go # EventTotals and TargetTotals (per-webhook DB)
│ │ ├── model_apikey.go # APIKey entity │ │ ├── model_apikey.go # APIKey entity
│ │ ├── password.go # Argon2id hashing and verification │ │ ├── password.go # Argon2id hashing and verification
│ │ ├── retention.go # Retention reaper (per-webhook event expiry) │ │ ├── retention.go # Retention reaper (per-webhook event expiry)
+5
View File
@@ -16,6 +16,7 @@ import (
"sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck" "sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware" "sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/resetpw" "sneak.berlin/go/webhooker/internal/resetpw"
"sneak.berlin/go/webhooker/internal/server" "sneak.berlin/go/webhooker/internal/server"
@@ -177,6 +178,10 @@ func newApp() *fx.App {
healthcheck.New, healthcheck.New,
session.New, session.New,
handlers.New, handlers.New,
// The registry /metrics serves, and the delivery
// collectors registered on it.
metrics.NewRegistry,
metrics.New,
middleware.New, middleware.New,
// The one SSRF guard both target-creation validation // The one SSRF guard both target-creation validation
// and the delivery dialer consult, so they cannot // and the delivery dialer consult, so they cannot
+21 -94
View File
@@ -93,11 +93,11 @@ func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
deliveries []database.Delivery deliveries []database.Delivery
results []database.DeliveryResult results []database.DeliveryResult
depths []struct{ Depth int } depths []struct{ Depth int }
removed []database.TargetTotals
) )
byStatus := "idx_deliveries_status (status=? AND deleted_at=?)" byStatus := "idx_deliveries_status (status=? AND deleted_at=?)"
byEvent := "idx_deliveries_event_id (event_id=? AND deleted_at=?)" byEvent := "idx_deliveries_event_id (event_id=? AND deleted_at=?)"
byAge := "idx_events_deleted_at_created_at (deleted_at=? AND created_at<?)"
// The delivery engine: recovery and the retry sweep, the sweep for // The delivery engine: recovery and the retry sweep, the sweep for
// stranded pending deliveries, and the queue depth count. // stranded pending deliveries, and the queue depth count.
@@ -123,89 +123,25 @@ func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
Order("attempt_num ASC").Find(&results), Order("attempt_num ASC").Find(&results),
"idx_delivery_results_delivery_id (delivery_id=? AND deleted_at=?)") "idx_delivery_results_delivery_id (delivery_id=? AND deleted_at=?)")
// Retention (reapExpired, deleteEvents): one batch of expired // Retention's three deletes (reapExpired), whose subqueries are built
// events, then their attempts, deliveries and the events. // afresh for each statement as it builds them.
var expired []string expiredEventIDs := func() *gorm.DB {
return dry.Model(&database.Event{}).Select("id").
Where("created_at < ?", cutoff)
}
assertPlanUses(t, db, dry.Unscoped().Model(&database.Event{}).
Where("created_at < ?", cutoff).
Limit(database.ExportReapBatchSize).Pluck("id", &expired),
"idx_events_created_at (created_at<?)")
assertPlanUses(t, db, dry.Unscoped().Where( assertPlanUses(t, db, dry.Unscoped().Where(
"delivery_id IN (?)", dry.Unscoped().Model(&database.Delivery{}). "delivery_id IN (?)", dry.Model(&database.Delivery{}).
Select("id").Where("event_id IN ?", ids), Select("id").Where("event_id IN (?)", expiredEventIDs()),
).Delete(&database.DeliveryResult{}), ).Delete(&database.DeliveryResult{}),
"idx_delivery_results_delivery_id (delivery_id=?)", "idx_delivery_results_delivery_id (delivery_id=?)", byEvent, byAge)
"idx_deliveries_event_id (event_id=?)") assertPlanUses(t, db, dry.Unscoped().Where(
assertPlanUses(t, db, dry.Unscoped().Model(&database.Delivery{}). "event_id IN (?)", expiredEventIDs(),
Select("target_id, count(*) AS deliveries_removed, "+ ).Delete(&database.Delivery{}),
"count(CASE WHEN status = ? THEN 1 END) AS failed_removed", "idx_deliveries_event_id (event_id=?)", byAge)
database.DeliveryStatusFailed). assertPlanUses(t, db, dry.Unscoped().Where(
Where("event_id IN ?", ids).Group("target_id").Find(&removed), "created_at < ?", cutoff,
"idx_deliveries_event_id (event_id=?)") ).Delete(&database.Event{}), "idx_events_created_at (created_at<?)")
assertPlanUses(t, db, dry.Unscoped().Where("event_id IN ?", ids).
Delete(&database.Delivery{}), "idx_deliveries_event_id (event_id=?)")
assertPlanUses(t, db, dry.Unscoped().Where("id IN ?", ids).
Delete(&database.Event{}), "sqlite_autoindex_events_1 (id=?)")
}
// TestStatisticsQueriesUseTheirIndexes does the same for the webhook
// page's statistics (readEventStats in the handlers): deliveries in
// progress, each target's deliveries finished since a time, which must
// come from the index alone, events received since a time, and the
// newest event, which must come straight off an index rather than from
// sorting every event.
func TestStatisticsQueriesUseTheirIndexes(t *testing.T) {
t.Parallel()
mgr, lc := setupTestWebhookDBManager(t)
ctx := context.Background()
require.NoError(t, lc.Start(ctx))
defer func() { require.NoError(t, lc.Stop(ctx)) }()
db, err := mgr.GetDB(uuid.New().String())
require.NoError(t, err)
dry := db.Session(&gorm.Session{DryRun: true})
since := time.Now()
var (
count int64
newest []time.Time
byTarget []struct{ TargetID string }
)
assertPlanUses(t, db, dry.Model(&database.Delivery{}).
Where("status IN ?", []database.DeliveryStatus{
database.DeliveryStatusPending,
database.DeliveryStatusRetrying,
}).Count(&count),
"idx_deliveries_status (status=? AND deleted_at=?)")
assertPlanUses(t, db, dry.Model(&database.Delivery{}).
Select("target_id, "+
"count(CASE WHEN status = ? THEN 1 END) AS delivered, "+
"count(CASE WHEN status = ? THEN 1 END) AS failed",
database.DeliveryStatusDelivered,
database.DeliveryStatusFailed).
Where("status IN ? AND finished_at >= ?",
[]database.DeliveryStatus{
database.DeliveryStatusDelivered,
database.DeliveryStatusFailed,
}, since).
Group("target_id").Find(&byTarget),
"COVERING INDEX idx_deliveries_status "+
"(status=? AND deleted_at=? AND finished_at>?)")
assertPlanUses(t, db, dry.Model(&database.Event{}).
Where("created_at >= ?", since).Count(&count),
"idx_events_deleted_at_created_at "+
"(deleted_at=? AND created_at>?)")
newestEvent := dry.Model(&database.Event{}).
Order("created_at DESC").Limit(1).Pluck("created_at", &newest)
assertPlanUses(t, db, newestEvent,
"idx_events_deleted_at_created_at (deleted_at=?)")
assert.NotContains(t, queryPlan(t, db, newestEvent), "TEMP B-TREE")
} }
// assertPlanUses asserts that SQLite's plan for a statement GORM built // assertPlanUses asserts that SQLite's plan for a statement GORM built
@@ -216,18 +152,6 @@ func assertPlanUses(
) { ) {
t.Helper() t.Helper()
plan := queryPlan(t, db, built)
for _, index := range indexes {
assert.Contains(t, plan, index, built.Statement.SQL.String())
}
}
// queryPlan returns SQLite's plan for a statement GORM built in a dry
// run, run with the same SQL and arguments GORM would send.
func queryPlan(t *testing.T, db, built *gorm.DB) string {
t.Helper()
var plan []struct{ Detail string } var plan []struct{ Detail string }
require.NoError(t, db.Raw( require.NoError(t, db.Raw(
@@ -235,5 +159,8 @@ func queryPlan(t *testing.T, db, built *gorm.DB) string {
built.Statement.Vars..., built.Statement.Vars...,
).Scan(&plan).Error) ).Scan(&plan).Error)
return fmt.Sprint(plan) for _, index := range indexes {
assert.Contains(t, fmt.Sprint(plan), index,
built.Statement.SQL.String())
}
} }
-4
View File
@@ -28,10 +28,6 @@ func NewTestRetentionReaper(
} }
} }
// ExportReapBatchSize exposes how many expired events one retention
// transaction deletes.
const ExportReapBatchSize = reapBatchSize
// ExportSweep runs a single retention sweep synchronously for tests. // ExportSweep runs a single retention sweep synchronously for tests.
func (r *RetentionReaper) ExportSweep(ctx context.Context) { func (r *RetentionReaper) ExportSweep(ctx context.Context) {
r.sweep(ctx) r.sweep(ctx)
+2 -13
View File
@@ -1,10 +1,6 @@
package database package database
import ( import "gorm.io/gorm"
"time"
"gorm.io/gorm"
)
// DeliveryStatus represents the status of a delivery // DeliveryStatus represents the status of a delivery
type DeliveryStatus string type DeliveryStatus string
@@ -41,7 +37,7 @@ type Delivery struct {
BaseModel BaseModel
EventID string `gorm:"type:uuid;not null;index:idx_deliveries_event_id,priority:1" json:"eventId"` EventID string `gorm:"type:uuid;not null;index:idx_deliveries_event_id,priority:1" json:"eventId"`
TargetID string `gorm:"type:uuid;not null;index:idx_deliveries_status,priority:4" json:"targetId"` TargetID string `gorm:"type:uuid;not null" json:"targetId"`
Status DeliveryStatus `gorm:"not null;default:'pending';index:idx_deliveries_status,priority:1" json:"status"` Status DeliveryStatus `gorm:"not null;default:'pending';index:idx_deliveries_status,priority:1" json:"status"`
// DeletedAt repeats the BaseModel field only to be the second column // DeletedAt repeats the BaseModel field only to be the second column
@@ -49,13 +45,6 @@ type Delivery struct {
// gives. // gives.
DeletedAt gorm.DeletedAt `gorm:"index:idx_deliveries_event_id,priority:2;index:idx_deliveries_status,priority:2" json:"deletedAt,omitzero"` DeletedAt gorm.DeletedAt `gorm:"index:idx_deliveries_event_id,priority:2;index:idx_deliveries_status,priority:2" json:"deletedAt,omitzero"`
// FinishedAt is when the delivery became delivered or failed, and
// nil while it is pending or retrying. It and then TargetID end the
// status index, so the webhook page counts each target's deliveries
// that finished in a recent window by reading just that window from
// the index.
FinishedAt *time.Time `gorm:"index:idx_deliveries_status,priority:3" json:"finishedAt,omitempty"`
// Relations // Relations
Event Event `json:"event,omitzero"` Event Event `json:"event,omitzero"`
Target Target `json:"target,omitzero"` Target Target `json:"target,omitzero"`
-91
View File
@@ -1,91 +0,0 @@
package database
import (
"fmt"
"gorm.io/gorm"
)
// The running totals in a webhook's event database keep the webhook
// page's lifetime figures right after retention has removed the rows
// they count, and let the page show them without counting every row.
// Each total changes in the transaction that writes or deletes the
// rows it counts.
// EventTotals is the single row counting a webhook's events: every
// event ever stored, and how many of them retention has deleted.
type EventTotals struct {
ID int64 `gorm:"primaryKey"`
Events int64 `gorm:"not null"`
EventsRemoved int64 `gorm:"not null"`
}
// TableName names the table AddEventTotals updates.
func (EventTotals) TableName() string {
return "event_totals"
}
// TargetTotals is one row per target counting its deliveries: every
// delivery ever created, how many became delivered and how many
// failed, and how many deliveries and failed deliveries retention has
// deleted. The webhook's delivery figures are these rows summed.
type TargetTotals struct {
TargetID string `gorm:"type:uuid;primaryKey"`
Deliveries int64 `gorm:"not null"`
Delivered int64 `gorm:"not null"`
Failed int64 `gorm:"not null"`
DeliveriesRemoved int64 `gorm:"not null"`
FailedRemoved int64 `gorm:"not null"`
}
// TableName names the table AddTargetTotals updates.
func (TargetTotals) TableName() string {
return "target_totals"
}
// AddEventTotals adds each count in add to the webhook's event totals.
// Call it on the transaction that writes or deletes the events it
// counts.
func AddEventTotals(tx *gorm.DB, add EventTotals) error {
err := tx.Exec(
`UPDATE event_totals SET
events = events + ?,
events_removed = events_removed + ?`,
add.Events, add.EventsRemoved,
).Error
if err != nil {
return fmt.Errorf("adding to event totals: %w", err)
}
return nil
}
// AddTargetTotals adds each count in add to the totals of the target
// add.TargetID names, creating its row the first time. Call it on the
// transaction that writes or deletes the deliveries it counts.
func AddTargetTotals(tx *gorm.DB, add TargetTotals) error {
err := tx.Exec(
`INSERT INTO target_totals (target_id, deliveries, delivered,
failed, deliveries_removed, failed_removed)
VALUES (?, ?, ?, ?, ?, ?)
ON CONFLICT (target_id) DO UPDATE SET
deliveries = deliveries + excluded.deliveries,
delivered = delivered + excluded.delivered,
failed = failed + excluded.failed,
deliveries_removed =
deliveries_removed + excluded.deliveries_removed,
failed_removed = failed_removed + excluded.failed_removed`,
add.TargetID, add.Deliveries, add.Delivered,
add.Failed, add.DeliveriesRemoved, add.FailedRemoved,
).Error
if err != nil {
return fmt.Errorf(
"adding to totals of target %s: %w", add.TargetID, err,
)
}
return nil
}
+1 -2
View File
@@ -2,8 +2,7 @@ package database
// Migrate runs database migrations for the main application database. // Migrate runs database migrations for the main application database.
// Only configuration-tier models are stored in the main database. // Only configuration-tier models are stored in the main database.
// Event-tier models (Event, Delivery, DeliveryResult, EventTotals, // Event-tier models (Event, Delivery, DeliveryResult) live in
// TargetTotals) live in
// per-webhook dedicated databases managed by WebhookDBManager. // per-webhook dedicated databases managed by WebhookDBManager.
func (d *Database) Migrate() error { func (d *Database) Migrate() error {
return d.db.AutoMigrate( return d.db.AutoMigrate(
+41 -88
View File
@@ -18,13 +18,6 @@ import (
// computation. // computation.
const hoursPerDay = 24 const hoursPerDay = 24
// reapBatchSize is how many expired events one retention transaction
// deletes. A transaction holds the event database's write lock, which
// the receiver and the delivery workers wait for, so a large prune is
// split into transactions each short enough to finish well inside the
// busy timeout.
const reapBatchSize = 1000
// RetentionReaperParams holds the fx dependencies for the // RetentionReaperParams holds the fx dependencies for the
// RetentionReaper. // RetentionReaper.
type RetentionReaperParams struct { type RetentionReaperParams struct {
@@ -272,97 +265,57 @@ func retentionCutoff(
), true ), true
} }
// reapExpired hard-deletes the events older than cutoff, with their // reapExpired hard-deletes, in foreign-key-safe order, the delivery
// deliveries and delivery results, reapBatchSize events per // results, deliveries, and events associated with events older than
// transaction until none is left. It returns the number of events // cutoff. Deletes are unscoped so rows are physically removed rather
// than soft-deleted, reclaiming disk. It returns the number of events
// deleted. // deleted.
func reapExpired(db *gorm.DB, cutoff time.Time) (int64, error) { func reapExpired(db *gorm.DB, cutoff time.Time) (int64, error) {
var total int64 // Fresh subqueries are built per statement to avoid reusing a
// mutated builder across executions.
for { expiredEventIDs := func() *gorm.DB {
var eventIDs []string return db.Model(&Event{}).
err := db.Transaction(func(tx *gorm.DB) error {
err := tx.Unscoped().Model(&Event{}).
Where("created_at < ?", cutoff).
Limit(reapBatchSize).
Pluck("id", &eventIDs).Error
if err != nil {
return fmt.Errorf("selecting expired events: %w", err)
}
if len(eventIDs) == 0 {
return nil
}
return deleteEvents(tx, eventIDs)
})
if err != nil {
return total, err
}
total += int64(len(eventIDs))
if len(eventIDs) < reapBatchSize {
return total, nil
}
}
}
// deleteEvents hard-deletes the given events and, in foreign-key-safe
// order before them, their delivery results and deliveries, then adds
// what it deleted to the running totals. It runs on reapExpired's
// transaction, so the totals change exactly when the rows do. Deletes
// are unscoped so rows are physically removed rather than
// soft-deleted, reclaiming disk.
func deleteEvents(tx *gorm.DB, eventIDs []string) error {
// 1. The delivery results of the events' deliveries.
err := tx.Unscoped().
Where("delivery_id IN (?)", tx.Unscoped().Model(&Delivery{}).
Select("id"). Select("id").
Where("event_id IN ?", eventIDs)). Where("created_at < ?", cutoff)
Delete(&DeliveryResult{}).Error }
if err != nil { expiredDeliveryIDs := func() *gorm.DB {
return fmt.Errorf("deleting expired delivery results: %w", err) return db.Model(&Delivery{}).
Select("id").
Where("event_id IN (?)", expiredEventIDs())
} }
// 2. The events' deliveries, after counting them, and the failed // 1. Delivery results whose delivery belongs to an expired event.
// ones among them, per target. The status is tested in the select res := db.Unscoped().
// list rather than the WHERE clause: there, SQLite would read every Where("delivery_id IN (?)", expiredDeliveryIDs()).
// failed delivery the webhook has through the status index, Delete(&DeliveryResult{})
// instead of only these through the event_id index. if res.Error != nil {
var removed []TargetTotals return 0, fmt.Errorf(
"deleting expired delivery results: %w",
err = tx.Unscoped().Model(&Delivery{}). res.Error,
Select("target_id, count(*) AS deliveries_removed, "+ )
"count(CASE WHEN status = ? THEN 1 END) AS failed_removed",
DeliveryStatusFailed).
Where("event_id IN ?", eventIDs).
Group("target_id").
Find(&removed).Error
if err != nil {
return fmt.Errorf("counting expired deliveries: %w", err)
} }
err = tx.Unscoped(). // 2. Deliveries belonging to an expired event.
Where("event_id IN ?", eventIDs). del := db.Unscoped().
Delete(&Delivery{}).Error Where("event_id IN (?)", expiredEventIDs()).
if err != nil { Delete(&Delivery{})
return fmt.Errorf("deleting expired deliveries: %w", err) if del.Error != nil {
return 0, fmt.Errorf(
"deleting expired deliveries: %w",
del.Error,
)
} }
// 3. The events themselves. // 3. The expired events themselves.
ev := tx.Unscoped().Where("id IN ?", eventIDs).Delete(&Event{}) ev := db.Unscoped().
Where("created_at < ?", cutoff).
Delete(&Event{})
if ev.Error != nil { if ev.Error != nil {
return fmt.Errorf("deleting expired events: %w", ev.Error) return 0, fmt.Errorf(
"deleting expired events: %w",
ev.Error,
)
} }
for i := range removed { return ev.RowsAffected, nil
err = AddTargetTotals(tx, removed[i])
if err != nil {
return err
}
}
return AddEventTotals(tx, EventTotals{EventsRemoved: ev.RowsAffected})
} }
-215
View File
@@ -1,215 +0,0 @@
package database_test
import (
"context"
"net/http"
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
)
// readEventTotals reads a webhook database's row of event totals,
// asserting that it has exactly one.
func readEventTotals(t *testing.T, db *gorm.DB) database.EventTotals {
t.Helper()
var rows []database.EventTotals
require.NoError(t, db.Find(&rows).Error)
require.Len(t, rows, 1)
return rows[0]
}
// readTargetTotals reads a webhook database's target totals, keyed by
// target.
func readTargetTotals(
t *testing.T, db *gorm.DB,
) map[string]database.TargetTotals {
t.Helper()
var rows []database.TargetTotals
require.NoError(t, db.Find(&rows).Error)
byTarget := make(map[string]database.TargetTotals, len(rows))
for _, row := range rows {
byTarget[row.TargetID] = row
}
return byTarget
}
// TestWebhookDBManager_TotalsSurviveReopen verifies that a new event
// database starts with one row of zero event totals and no target
// totals, that adding to a target twice adds to the one row, and that
// opening the database again keeps everything added.
func TestWebhookDBManager_TotalsSurviveReopen(t *testing.T) {
t.Parallel()
mgr, lc := setupTestWebhookDBManager(t)
ctx := context.Background()
require.NoError(t, lc.Start(ctx))
defer func() { require.NoError(t, lc.Stop(ctx)) }()
webhookID := uuid.New().String()
db, err := mgr.GetDB(webhookID)
require.NoError(t, err)
fresh := readEventTotals(t, db)
assert.Equal(t, database.EventTotals{ID: fresh.ID}, fresh)
assert.Empty(t, readTargetTotals(t, db))
first, second := uuid.New().String(), uuid.New().String()
require.NoError(t, database.AddEventTotals(db, database.EventTotals{
Events: 2,
}))
require.NoError(t, database.AddTargetTotals(db, database.TargetTotals{
TargetID: first, Deliveries: 2, Delivered: 1,
}))
require.NoError(t, database.AddTargetTotals(db, database.TargetTotals{
TargetID: first, Failed: 1,
}))
require.NoError(t, database.AddTargetTotals(db, database.TargetTotals{
TargetID: second, Deliveries: 1,
}))
// Drop the cached connection so the next open reopens the file,
// as a restart would.
require.NoError(t, mgr.CloseAll())
db, err = mgr.GetDB(webhookID)
require.NoError(t, err)
assert.Equal(t, database.EventTotals{ID: fresh.ID, Events: 2},
readEventTotals(t, db))
assert.Equal(t, map[string]database.TargetTotals{
first: {
TargetID: first, Deliveries: 2, Delivered: 1, Failed: 1,
},
second: {TargetID: second, Deliveries: 1},
}, readTargetTotals(t, db))
}
// TestRetentionReaper_PrunesMoreThanOneBatch verifies that a prune
// larger than one transaction's batch removes every expired event with
// its deliveries and delivery results, keeps the recent event, and
// adds what it removed to the event and target totals, so the totals
// within retention match the rows still stored.
func TestRetentionReaper_PrunesMoreThanOneBatch(t *testing.T) {
t.Parallel()
env := setupRetentionTest(t)
webhookID := createWebhook(t, env.mainDB.DB(), 30)
db, err := env.mgr.GetDB(webhookID)
require.NoError(t, err)
// Every expired event has a delivered delivery to one target and a
// failed one to the other, each with one attempt.
expired := database.ExportReapBatchSize + 1
delivered, failed := uuid.New().String(), uuid.New().String()
old := time.Now().Add(-40 * 24 * time.Hour)
events := make([]database.Event, expired)
deliveries := make([]database.Delivery, 0, 2*expired)
for i := range events {
events[i] = database.Event{
WebhookID: webhookID,
EntrypointID: uuid.New().String(),
Method: http.MethodPost,
}
events[i].ID = uuid.New().String()
events[i].CreatedAt = old
deliveries = append(deliveries,
database.Delivery{
EventID: events[i].ID,
TargetID: delivered,
Status: database.DeliveryStatusDelivered,
},
database.Delivery{
EventID: events[i].ID,
TargetID: failed,
Status: database.DeliveryStatusFailed,
},
)
}
require.NoError(t, db.CreateInBatches(events, 500).Error)
require.NoError(t, db.CreateInBatches(deliveries, 500).Error)
results := make([]database.DeliveryResult, len(deliveries))
for i := range deliveries {
results[i] = database.DeliveryResult{
DeliveryID: deliveries[i].ID, AttemptNum: 1,
}
}
require.NoError(t, db.CreateInBatches(results, 500).Error)
// One recent event, delivered to the first target.
recent := seedEventChain(t, db, webhookID, time.Now())
require.NoError(t, db.Model(&database.Delivery{}).
Where("id = ?", recent.deliveryID).
Update("target_id", delivered).Error)
// The totals storing those rows would have left.
n := int64(expired)
require.NoError(t, database.AddEventTotals(db, database.EventTotals{
Events: n + 1,
}))
require.NoError(t, database.AddTargetTotals(db, database.TargetTotals{
TargetID: delivered, Deliveries: n + 1, Delivered: n + 1,
}))
require.NoError(t, database.AddTargetTotals(db, database.TargetTotals{
TargetID: failed, Deliveries: n, Failed: n,
}))
env.reaper.ExportSweep(context.Background())
// Only the recent event's rows are left.
for _, model := range []any{
&database.Event{}, &database.Delivery{}, &database.DeliveryResult{},
} {
var count int64
require.NoError(t, db.Model(model).Count(&count).Error)
assert.Equal(t, int64(1), count, "%T rows left", model)
}
assertChainPresent(t, db, recent)
eventTotals := readEventTotals(t, db)
assert.Equal(t, database.EventTotals{
ID: eventTotals.ID, Events: n + 1, EventsRemoved: n,
}, eventTotals)
targetTotals := readTargetTotals(t, db)
assert.Equal(t, map[string]database.TargetTotals{
delivered: {
TargetID: delivered, Deliveries: n + 1, Delivered: n + 1,
DeliveriesRemoved: n,
},
failed: {
TargetID: failed, Deliveries: n, Failed: n,
DeliveriesRemoved: n, FailedRemoved: n,
},
}, targetTotals)
// A sweep with nothing left to remove changes nothing.
env.reaper.ExportSweep(context.Background())
assert.Equal(t, eventTotals, readEventTotals(t, db))
assert.Equal(t, targetTotals, readTargetTotals(t, db))
}
+1 -15
View File
@@ -35,8 +35,7 @@ var errInvalidCachedDBType = errors.New(
// WebhookDBManager manages per-webhook SQLite database files // WebhookDBManager manages per-webhook SQLite database files
// for event storage. Each webhook gets its own dedicated // for event storage. Each webhook gets its own dedicated
// database containing Events, Deliveries, DeliveryResults and the // database containing Events, Deliveries, and DeliveryResults.
// running totals of them (EventTotals, TargetTotals).
// Database connections are opened lazily and cached. // Database connections are opened lazily and cached.
type WebhookDBManager struct { type WebhookDBManager struct {
dataDir string dataDir string
@@ -296,7 +295,6 @@ func (m *WebhookDBManager) openDB(
// Run migrations for event-tier models only // Run migrations for event-tier models only
err = db.AutoMigrate( err = db.AutoMigrate(
&Event{}, &Delivery{}, &DeliveryResult{}, &Event{}, &Delivery{}, &DeliveryResult{},
&EventTotals{}, &TargetTotals{},
) )
if err != nil { if err != nil {
_ = sqlDB.Close() _ = sqlDB.Close()
@@ -307,18 +305,6 @@ func (m *WebhookDBManager) openDB(
) )
} }
// A new database gets its row of event totals, all zero. Target
// totals rows are created by the first delivery to each target.
err = db.FirstOrCreate(&EventTotals{}).Error
if err != nil {
_ = sqlDB.Close()
return nil, fmt.Errorf(
"creating event totals for webhook database %s: %w",
webhookID, err,
)
}
m.log.Info( m.log.Info(
"opened per-webhook database", "opened per-webhook database",
"webhook_id", webhookID, "webhook_id", webhookID,
-116
View File
@@ -1,116 +0,0 @@
package delivery_test
import (
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
)
// targetTotals reads one target's totals from a webhook database, all
// zero when it has no row.
func targetTotals(
t *testing.T, db *gorm.DB, targetID string,
) database.TargetTotals {
t.Helper()
var rows []database.TargetTotals
require.NoError(t, db.Where("target_id = ?", targetID).
Find(&rows).Error)
if len(rows) == 0 {
return database.TargetTotals{TargetID: targetID}
}
return rows[0]
}
// TestUpdateDeliveryStatus_FinishTimeAndTargetTotals pins what a status
// write records for the webhook page's statistics: the time a delivery
// finished, set only when it becomes delivered or failed, and one more
// on its target's delivered or failed total.
func TestUpdateDeliveryStatus_FinishTimeAndTargetTotals(t *testing.T) {
t.Parallel()
tests := []struct {
status database.DeliveryStatus
finished bool
delivered int64
failed int64
}{
{database.DeliveryStatusRetrying, false, 0, 0},
{database.DeliveryStatusDelivered, true, 1, 0},
{database.DeliveryStatusFailed, true, 0, 1},
}
for _, tt := range tests {
t.Run(string(tt.status), func(t *testing.T) {
t.Parallel()
db := testWebhookDB(t)
e := testEngine(t, 1)
event := seedEvent(t, db, `{}`)
targetID := uuid.New().String()
d := seedDelivery(
t, db, event.ID, targetID,
database.DeliveryStatusPending,
)
before := time.Now()
require.NoError(t, e.ExportUpdateDeliveryStatus(
db, &d, tt.status,
))
var stored database.Delivery
require.NoError(t, db.First(&stored, "id = ?", d.ID).Error)
assert.Equal(t, tt.status, stored.Status)
if tt.finished {
require.NotNil(t, stored.FinishedAt)
assert.False(t, stored.FinishedAt.Before(before))
} else {
assert.Nil(t, stored.FinishedAt)
}
assert.Equal(t, database.TargetTotals{
TargetID: targetID,
Delivered: tt.delivered,
Failed: tt.failed,
}, targetTotals(t, db, targetID))
})
}
}
// TestUpdateDeliveryStatus_DeletedDeliveryIsNotCounted covers a
// delivery retention deleted while the engine still held it. Failing
// it afterwards writes no row, so it adds no failure either: retention
// has already counted what it removed.
func TestUpdateDeliveryStatus_DeletedDeliveryIsNotCounted(t *testing.T) {
t.Parallel()
db := testWebhookDB(t)
e := testEngine(t, 1)
event := seedEvent(t, db, `{}`)
targetID := uuid.New().String()
d := seedDelivery(
t, db, event.ID, targetID,
database.DeliveryStatusRetrying,
)
require.NoError(t, db.Unscoped().
Delete(&database.Delivery{}, "id = ?", d.ID).Error)
require.NoError(t, e.ExportUpdateDeliveryStatus(
db, &d, database.DeliveryStatusFailed,
))
assert.Equal(t, database.TargetTotals{TargetID: targetID},
targetTotals(t, db, targetID))
}
+8 -38
View File
@@ -148,6 +148,7 @@ type EngineParams struct {
DBManager *database.WebhookDBManager DBManager *database.WebhookDBManager
Logger *logger.Logger Logger *logger.Logger
SSRFGuard *Guard SSRFGuard *Guard
Metrics *metrics.Set
} }
// Engine processes queued deliveries in the background // Engine processes queued deliveries in the background
@@ -167,10 +168,10 @@ type Engine struct {
retryCh chan Task retryCh chan Task
workers int workers int
// mtr is the delivery metric set. Production wires the // mtr is the delivery metric set. Production wires the one
// process-wide one; a test can substitute a set registered on // registered on the registry /metrics serves; a test can
// a private registry so its assertions are not disturbed by // substitute a set registered on a registry it holds, so it can
// deliveries other tests are making at the same time. // gather what its own deliveries recorded.
mtr *metrics.Set mtr *metrics.Set
// targets maps each target type to its implementation. // targets maps each target type to its implementation.
@@ -204,7 +205,7 @@ func New(
deliveryCh: make(chan Task, deliveryChannelSize), deliveryCh: make(chan Task, deliveryChannelSize),
retryCh: make(chan Task, retryChannelSize), retryCh: make(chan Task, retryChannelSize),
workers: defaultWorkers, workers: defaultWorkers,
mtr: metrics.Default(), mtr: params.Metrics,
} }
e.initTargets(&http.Client{ e.initTargets(&http.Client{
@@ -1554,9 +1555,8 @@ func (e *Engine) updateDeliveryStatus(
targetType database.TargetType, targetType database.TargetType,
status database.DeliveryStatus, status database.DeliveryStatus,
) error { ) error {
err := webhookDB.Transaction(func(tx *gorm.DB) error { err := webhookDB.Model(d).
return writeDeliveryStatus(tx, d, status) Update("status", status).Error
})
if err != nil { if err != nil {
return fmt.Errorf( return fmt.Errorf(
"updating delivery %s to status %s: %w", "updating delivery %s to status %s: %w",
@@ -1575,36 +1575,6 @@ func (e *Engine) updateDeliveryStatus(
return nil return nil
} }
// writeDeliveryStatus writes a delivery's new status. A delivery that
// becomes delivered or failed also gets the time it finished, and is
// added to its target's delivered or failed total. It is counted only
// if the row was still there to update: retention may have deleted it
// while the engine was working on it.
func writeDeliveryStatus(
tx *gorm.DB,
d *database.Delivery,
status database.DeliveryStatus,
) error {
if !status.Terminal() {
return tx.Model(d).Update("status", status).Error
}
res := tx.Model(d).Updates(map[string]any{
"status": status,
"finished_at": time.Now(),
})
if res.Error != nil || res.RowsAffected == 0 {
return res.Error
}
add := database.TargetTotals{TargetID: d.TargetID, Delivered: 1}
if status == database.DeliveryStatusFailed {
add = database.TargetTotals{TargetID: d.TargetID, Failed: 1}
}
return database.AddTargetTotals(tx, add)
}
// settleStatus moves a delivery to its outcome status and reports a // settleStatus moves a delivery to its outcome status and reports a
// failed write through bookkeepingFailed, which leaves the row // failed write through bookkeepingFailed, which leaves the row
// recoverable. It exists so the target call sites read as one // recoverable. It exists so the target call sites read as one
-3
View File
@@ -57,10 +57,7 @@ func testWebhookDB(t *testing.T) *gorm.DB {
&database.Event{}, &database.Event{},
&database.Delivery{}, &database.Delivery{},
&database.DeliveryResult{}, &database.DeliveryResult{},
&database.EventTotals{},
&database.TargetTotals{},
)) ))
require.NoError(t, db.Create(&database.EventTotals{}).Error)
return db return db
} }
+5 -15
View File
@@ -9,6 +9,7 @@ import (
"net/url" "net/url"
"time" "time"
"github.com/prometheus/client_golang/prometheus"
"go.uber.org/fx" "go.uber.org/fx"
"gorm.io/gorm" "gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
@@ -150,16 +151,6 @@ func (e *Engine) ExportDeliverSlack(
) )
} }
// ExportUpdateDeliveryStatus exposes updateDeliveryStatus. It passes no
// target type, so no metric moves.
func (e *Engine) ExportUpdateDeliveryStatus(
webhookDB *gorm.DB,
d *database.Delivery,
status database.DeliveryStatus,
) error {
return e.updateDeliveryStatus(webhookDB, d, "", status)
}
// ExportProcessNewTask exposes processNewTask. // ExportProcessNewTask exposes processNewTask.
func (e *Engine) ExportProcessNewTask( func (e *Engine) ExportProcessNewTask(
ctx context.Context, task *Task, ctx context.Context, task *Task,
@@ -399,7 +390,7 @@ func NewTestEngine(
deliveryCh: make(chan Task, deliveryChannelSize), deliveryCh: make(chan Task, deliveryChannelSize),
retryCh: make(chan Task, retryChannelSize), retryCh: make(chan Task, retryChannelSize),
workers: workers, workers: workers,
mtr: metrics.Default(), mtr: metrics.New(prometheus.NewRegistry()),
} }
e.initTargets(client) e.initTargets(client)
@@ -414,7 +405,7 @@ func NewTestEngineSmallRetry(
e := &Engine{ e := &Engine{
log: log, log: log,
retryCh: make(chan Task, 1), retryCh: make(chan Task, 1),
mtr: metrics.Default(), mtr: metrics.New(prometheus.NewRegistry()),
} }
e.initTargets(nil) e.initTargets(nil)
@@ -437,7 +428,7 @@ func NewTestEngineWithDB(
deliveryCh: make(chan Task, deliveryChannelSize), deliveryCh: make(chan Task, deliveryChannelSize),
retryCh: make(chan Task, retryChannelSize), retryCh: make(chan Task, retryChannelSize),
workers: workers, workers: workers,
mtr: metrics.Default(), mtr: metrics.New(prometheus.NewRegistry()),
} }
e.initTargets(client) e.initTargets(client)
@@ -445,8 +436,7 @@ func NewTestEngineWithDB(
} }
// ExportSetMetrics substitutes the engine's metric set, so a test can // ExportSetMetrics substitutes the engine's metric set, so a test can
// assert on collectors registered on a private registry instead of // assert on collectors registered on a registry it holds.
// the process-wide ones every other test is also moving.
func (e *Engine) ExportSetMetrics(mtr *metrics.Set) { func (e *Engine) ExportSetMetrics(mtr *metrics.Set) {
e.mtr = mtr e.mtr = mtr
} }
+2 -3
View File
@@ -35,9 +35,8 @@ const (
) )
// mIsolate gives the setup's engine a metric set registered on a // mIsolate gives the setup's engine a metric set registered on a
// private registry. The process-wide collectors are moved by every // registry this test holds, so its exact assertions can gather from
// other delivery test running in parallel, so exact assertions are // it.
// only possible against a registry this test owns.
func mIsolate( func mIsolate(
t *testing.T, s iSetup, t *testing.T, s iSetup,
) *prometheus.Registry { ) *prometheus.Registry {
+3 -13
View File
@@ -299,9 +299,8 @@ func countInFlightDeliveries(
return count, err return count, err
} }
// createReplayDelivery writes the new pending delivery row, adds it to // createReplayDelivery writes the new pending delivery row and returns
// its target's totals in the same transaction, and returns the task // the task that carries it to the delivery engine.
// that carries it to the delivery engine.
// //
// The row is written with associations omitted, and neither Event nor // The row is written with associations omitted, and neither Event nor
// Target is populated on it: GORM's SaveBeforeAssociations would // Target is populated on it: GORM's SaveBeforeAssociations would
@@ -320,16 +319,7 @@ func createReplayDelivery(
Status: database.DeliveryStatusPending, Status: database.DeliveryStatusPending,
} }
err := webhookDB.Transaction(func(tx *gorm.DB) error { err := webhookDB.Omit(clause.Associations).Create(dlv).Error
err := tx.Omit(clause.Associations).Create(dlv).Error
if err != nil {
return err
}
return database.AddTargetTotals(tx, database.TargetTotals{
TargetID: dlv.TargetID, Deliveries: 1,
})
})
if err != nil { if err != nil {
return delivery.Task{}, err return delivery.Task{}, err
} }
-25
View File
@@ -4,9 +4,7 @@ import (
"html/template" "html/template"
"log/slog" "log/slog"
"net/http" "net/http"
"time"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
) )
@@ -71,29 +69,6 @@ func (s *Handlers) LoadEventLogViewsForTest(
return views return views
} }
// WebhookStatsForTest returns the figures the statistics pane on a
// webhook's page shows, from the webhook's entrypoints and targets
// loaded as that page loads them.
func (s *Handlers) WebhookStatsForTest(webhookID string) *WebhookStats {
var entrypoints []database.Entrypoint
s.db.DB().Where("webhook_id = ?", webhookID).Find(&entrypoints)
var targets []database.Target
s.db.DB().Where("webhook_id = ?", webhookID).Find(&targets)
return s.loadWebhookStats(webhookID, entrypoints, targets)
}
// FinishedByTargetForTest exposes finishedByTarget for use in the
// handlers_test package.
func FinishedByTargetForTest(
webhookDB *gorm.DB, since time.Time,
) ([]TargetFinished, error) {
return finishedByTarget(webhookDB, since)
}
// AddTemplateForTest registers a template under a page name so that // AddTemplateForTest registers a template under a page name so that
// the handlers_test package can drive the render path with a // the handlers_test package can drive the render path with a
// template of its own. // template of its own.
+16 -17
View File
@@ -12,6 +12,7 @@ import (
"net/http" "net/http"
"sync/atomic" "sync/atomic"
"github.com/prometheus/client_golang/prometheus"
"go.uber.org/fx" "go.uber.org/fx"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
@@ -61,6 +62,8 @@ type HandlersParams struct {
Notifier delivery.Notifier Notifier delivery.Notifier
Evictor delivery.WebhookEvictor Evictor delivery.WebhookEvictor
SSRFGuard *delivery.Guard SSRFGuard *delivery.Guard
Metrics *metrics.Set
Registry *prometheus.Registry
} }
// Handlers provides HTTP handler methods for all application // Handlers provides HTTP handler methods for all application
@@ -91,22 +94,18 @@ type Handlers struct {
// parsePageTemplate parses a page-specific template set from the // parsePageTemplate parses a page-specific template set from the
// embedded FS. Each page template is combined with the shared // embedded FS. Each page template is combined with the shared
// base, htmlheader, and navbar templates, and with any further files // base, htmlheader, and navbar templates. The page file must be
// the page includes. The page file must be listed first so that its // listed first so that its root action ({{template "base" .}})
// root action ({{template "base" .}}) becomes the template set's entry // becomes the template set's entry point.
// point. func parsePageTemplate(pageFile string) *template.Template {
func parsePageTemplate(
pageFile string, included ...string,
) *template.Template {
files := append([]string{
pageFile,
"base.html",
"htmlheader.html",
"navbar.html",
}, included...)
return template.Must( return template.Must(
template.ParseFS(templates.Templates, files...), template.ParseFS(
templates.Templates,
pageFile,
"base.html",
"htmlheader.html",
"navbar.html",
),
) )
} }
@@ -126,7 +125,7 @@ func New(
s.mw = params.Middleware s.mw = params.Middleware
s.notifier = params.Notifier s.notifier = params.Notifier
s.evictor = params.Evictor s.evictor = params.Evictor
s.mtr = metrics.Default() s.mtr = params.Metrics
s.ssrf = params.SSRFGuard s.ssrf = params.SSRFGuard
// Parse all page templates once at startup // Parse all page templates once at startup
@@ -135,7 +134,7 @@ func New(
"profile.html": parsePageTemplate("profile.html"), "profile.html": parsePageTemplate("profile.html"),
"sources_list.html": parsePageTemplate("sources_list.html"), "sources_list.html": parsePageTemplate("sources_list.html"),
"sources_new.html": parsePageTemplate("sources_new.html"), "sources_new.html": parsePageTemplate("sources_new.html"),
"source_detail.html": parsePageTemplate("source_detail.html", "webhook_stats.html"), "source_detail.html": parsePageTemplate("source_detail.html"),
"source_edit.html": parsePageTemplate("source_edit.html"), "source_edit.html": parsePageTemplate("source_edit.html"),
"source_logs.html": parsePageTemplate("source_logs.html"), "source_logs.html": parsePageTemplate("source_logs.html"),
"target_edit.html": parsePageTemplate("target_edit.html"), "target_edit.html": parsePageTemplate("target_edit.html"),
+3
View File
@@ -20,6 +20,7 @@ import (
"sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck" "sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware" "sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/session" "sneak.berlin/go/webhooker/internal/session"
) )
@@ -109,6 +110,8 @@ func newTestApp(
func(r *recordingEvictor) delivery.WebhookEvictor { func(r *recordingEvictor) delivery.WebhookEvictor {
return r return r
}, },
metrics.NewRegistry,
metrics.New,
middleware.New, middleware.New,
delivery.NewGuard, delivery.NewGuard,
handlers.New, handlers.New,
+21
View File
@@ -0,0 +1,21 @@
package handlers
import (
"net/http"
"github.com/prometheus/client_golang/prometheus/promhttp"
)
// HandleMetrics returns the Prometheus scrape handler for the
// registry built by metrics.NewRegistry, which the HTTP, delivery, Go
// runtime and process collectors register on. It is what
// promhttp.Handler builds for the global default registry, including
// the promhttp_metric_handler_* series that count scrapes, pointed at
// that registry instead.
func (s *Handlers) HandleMetrics() http.HandlerFunc {
reg := s.params.Registry
return promhttp.InstrumentMetricHandler(
reg, promhttp.HandlerFor(reg, promhttp.HandlerOpts{}),
).ServeHTTP
}
-1
View File
@@ -450,7 +450,6 @@ func (h *Handlers) renderSourceDetail(
"Targets": delivery.NewTargetViews(targets), "Targets": delivery.NewTargetViews(targets),
"Events": events, "Events": events,
"BaseURL": baseURL, "BaseURL": baseURL,
"Stats": h.loadWebhookStats(webhook.ID, entrypoints, targets),
} }
h.renderTemplate(w, r, "source_detail.html", data) h.renderTemplate(w, r, "source_detail.html", data)
+7 -23
View File
@@ -252,12 +252,11 @@ func requestEventSource(
} }
} }
// createAndFanOut writes the event and one pending delivery per target, // createAndFanOut writes the event and one pending delivery per target
// and adds them to the webhook's running totals, in a single // in a single transaction, then hands the tasks to the delivery
// transaction, then hands the tasks to the delivery engine. It is the // engine. It is the only path by which an event and its deliveries are
// only path by which an event and its deliveries are created, so a // created, so a resubmitted event is retried, SSRF-guarded and
// resubmitted event is retried, SSRF-guarded and circuit-broken // circuit-broken exactly as a received one is.
// exactly as a received one is.
// //
// The tasks are returned as well as queued, so a caller can report how // The tasks are returned as well as queued, so a caller can report how
// many targets the event went to. // many targets the event went to.
@@ -297,13 +296,6 @@ func (h *Handlers) createAndFanOut(
return nil, nil, err return nil, nil, err
} }
err = database.AddEventTotals(tx, database.EventTotals{Events: 1})
if err != nil {
tx.Rollback()
return nil, nil, err
}
err = tx.Commit().Error err = tx.Commit().Error
if err != nil { if err != nil {
return nil, nil, fmt.Errorf( return nil, nil, fmt.Errorf(
@@ -362,9 +354,8 @@ func (h *Handlers) finishWebhookResponse(
} }
// buildDeliveryTasks creates one pending delivery per target in the // buildDeliveryTasks creates one pending delivery per target in the
// transaction, adds each to its target's totals, and returns the tasks // transaction and returns the tasks for the delivery engine. The
// for the delivery engine. The caller owns the transaction and rolls // caller owns the transaction and rolls it back on error.
// it back on error.
func buildDeliveryTasks( func buildDeliveryTasks(
tx *gorm.DB, tx *gorm.DB,
event *database.Event, event *database.Event,
@@ -388,13 +379,6 @@ func buildDeliveryTasks(
) )
} }
err = database.AddTargetTotals(tx, database.TargetTotals{
TargetID: targets[i].ID, Deliveries: 1,
})
if err != nil {
return nil, err
}
tasks = append(tasks, delivery.Task{ tasks = append(tasks, delivery.Task{
DeliveryID: dlv.ID, DeliveryID: dlv.ID,
EventID: event.ID, EventID: event.ID,
-271
View File
@@ -1,271 +0,0 @@
package handlers
import (
"fmt"
"time"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/database"
)
// The spans of the two recent windows the statistics pane reports on:
// the last 10 minutes and the last 24 hours.
const (
shortWindow = 10 * time.Minute
longWindow = 24 * time.Hour
)
// percent turns a fraction into a percentage.
const percent = 100
// WebhookStats holds the figures in the statistics pane at the top of
// the webhook page.
type WebhookStats struct {
Entrypoints int
ActiveEntrypoints int
Targets int
ActiveTargets int
// Lifetime counts every event, delivery and failure the webhook
// has had, and WithinRetention those still stored.
Lifetime Counts
WithinRetention Counts
// InProgress counts the deliveries still pending or retrying.
InProgress int64
// LastEventAt is when the newest stored event arrived, or nil when
// none is stored.
LastEventAt *time.Time
Last10Minutes RecentWindow
Last24Hours RecentWindow
}
// Counts holds a number of events, of deliveries and of failed
// deliveries.
type Counts struct {
Events int64
Deliveries int64
Failures int64
}
// RecentWindow holds what happened in one recent window: the events
// received in it, and the deliveries that became delivered or failed in
// it.
type RecentWindow struct {
Events int64
Delivered int64
Failed int64
}
// TargetFinished is how many of one target's deliveries became
// delivered, and how many failed, in a recent window.
type TargetFinished struct {
TargetID string
Delivered int64
Failed int64
}
// FailurePercent is the share of the deliveries finished in the window
// that failed, or a dash when none finished. Deliveries still pending
// or retrying are not counted either way.
func (w RecentWindow) FailurePercent() string {
finished := w.Delivered + w.Failed
if finished == 0 {
return "—"
}
return fmt.Sprintf(
"%.1f%%", percent*float64(w.Failed)/float64(finished),
)
}
// loadWebhookStats gathers the figures for the statistics pane from the
// webhook's entrypoints and targets, as the page has already loaded
// them, and from its event database. It returns nil, and logs why, when
// the event database cannot be read.
func (h *Handlers) loadWebhookStats(
webhookID string,
entrypoints []database.Entrypoint,
targets []database.Target,
) *WebhookStats {
stats := &WebhookStats{
Entrypoints: len(entrypoints),
Targets: len(targets),
}
for i := range entrypoints {
if entrypoints[i].Active {
stats.ActiveEntrypoints++
}
}
for i := range targets {
if targets[i].Active {
stats.ActiveTargets++
}
}
// Opening an event database that does not exist would create it,
// and it would hold nothing to count.
if !h.dbMgr.DBExists(webhookID) {
return stats
}
webhookDB, err := h.dbMgr.GetDB(webhookID)
if err == nil {
err = readEventStats(webhookDB, time.Now(), stats)
}
if err != nil {
h.log.Error(
"failed to read webhook statistics",
"webhook_id", webhookID,
"error", err,
)
return nil
}
return stats
}
// readEventStats fills in the figures that come from the webhook's
// event database. None of them reads every stored row: the totals are
// one row for the events and one per target for the deliveries, and
// every other figure is read from an index, over only the rows it
// counts.
func readEventStats(
db *gorm.DB, now time.Time, stats *WebhookStats,
) error {
err := readTotals(db, stats)
if err != nil {
return err
}
err = db.Model(&database.Delivery{}).
Where("status IN ?", []database.DeliveryStatus{
database.DeliveryStatusPending,
database.DeliveryStatusRetrying,
}).
Count(&stats.InProgress).Error
if err != nil {
return fmt.Errorf("counting deliveries in progress: %w", err)
}
var newest []time.Time
err = db.Model(&database.Event{}).
Order("created_at DESC").
Limit(1).
Pluck("created_at", &newest).Error
if err != nil {
return fmt.Errorf("reading newest event time: %w", err)
}
if len(newest) > 0 {
stats.LastEventAt = &newest[0]
}
stats.Last10Minutes, err = readRecentWindow(
db, now.Add(-shortWindow),
)
if err != nil {
return err
}
stats.Last24Hours, err = readRecentWindow(
db, now.Add(-longWindow),
)
return err
}
// readTotals fills in the lifetime and within-retention figures from
// the running totals: the events' row, and the targets' rows summed.
func readTotals(db *gorm.DB, stats *WebhookStats) error {
var events database.EventTotals
err := db.Take(&events).Error
if err != nil {
return fmt.Errorf("reading event totals: %w", err)
}
var targets []database.TargetTotals
err = db.Find(&targets).Error
if err != nil {
return fmt.Errorf("reading target totals: %w", err)
}
stats.Lifetime.Events = events.Events
stats.WithinRetention.Events = events.Events - events.EventsRemoved
for _, t := range targets {
stats.Lifetime.Deliveries += t.Deliveries
stats.Lifetime.Failures += t.Failed
stats.WithinRetention.Deliveries += t.Deliveries - t.DeliveriesRemoved
stats.WithinRetention.Failures += t.Failed - t.FailedRemoved
}
return nil
}
// readRecentWindow counts the events received, and the deliveries that
// became delivered or failed, since the given time.
func readRecentWindow(
db *gorm.DB, since time.Time,
) (RecentWindow, error) {
var w RecentWindow
err := db.Model(&database.Event{}).
Where("created_at >= ?", since).
Count(&w.Events).Error
if err != nil {
return w, fmt.Errorf("counting recent events: %w", err)
}
byTarget, err := finishedByTarget(db, since)
if err != nil {
return w, err
}
for _, f := range byTarget {
w.Delivered += f.Delivered
w.Failed += f.Failed
}
return w, nil
}
// finishedByTarget counts, for each target, the deliveries that became
// delivered and those that failed since the given time, in one query
// over just that window of the deliveries' status index. A target with
// neither is left out.
func finishedByTarget(
db *gorm.DB, since time.Time,
) ([]TargetFinished, error) {
var byTarget []TargetFinished
err := db.Model(&database.Delivery{}).
Select("target_id, "+
"count(CASE WHEN status = ? THEN 1 END) AS delivered, "+
"count(CASE WHEN status = ? THEN 1 END) AS failed",
database.DeliveryStatusDelivered,
database.DeliveryStatusFailed).
Where("status IN ? AND finished_at >= ?",
[]database.DeliveryStatus{
database.DeliveryStatusDelivered,
database.DeliveryStatusFailed,
}, since).
Group("target_id").
Find(&byTarget).Error
if err != nil {
return nil, fmt.Errorf(
"counting deliveries finished by target: %w", err,
)
}
return byTarget, nil
}
-432
View File
@@ -1,432 +0,0 @@
package handlers_test
import (
"net/http"
"strings"
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.uber.org/fx/fxtest"
"gorm.io/gorm"
"gorm.io/gorm/clause"
"sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/session"
)
// statsEntrypoint adds an entrypoint to a webhook and returns its path.
func statsEntrypoint(
t *testing.T, db *database.Database, webhookID string, active bool,
) string {
t.Helper()
ep := &database.Entrypoint{
WebhookID: webhookID,
Path: uuid.New().String(),
}
require.NoError(t, db.DB().Omit(clause.Associations).Create(ep).Error)
require.NoError(t, db.DB().Model(ep).Update("active", active).Error)
return ep.Path
}
// statsDelivery returns an event's delivery to a target.
func statsDelivery(
t *testing.T, webhookDB *gorm.DB, eventID, targetID string,
) database.Delivery {
t.Helper()
var d database.Delivery
require.NoError(t, webhookDB.Where(
"event_id = ? AND target_id = ?", eventID, targetID,
).First(&d).Error)
return d
}
// statsFinish settles a delivery as the delivery engine does: its
// final status and the time it finished, and one more on its target's
// delivered or failed total, in one transaction.
func statsFinish(
t *testing.T,
webhookDB *gorm.DB,
d database.Delivery,
status database.DeliveryStatus,
at time.Time,
) {
t.Helper()
add := database.TargetTotals{TargetID: d.TargetID, Delivered: 1}
if status == database.DeliveryStatusFailed {
add = database.TargetTotals{TargetID: d.TargetID, Failed: 1}
}
require.NoError(t, webhookDB.Transaction(func(tx *gorm.DB) error {
err := tx.Model(&database.Delivery{}).
Where("id = ?", d.ID).
Updates(map[string]any{"status": status, "finished_at": at}).
Error
if err != nil {
return err
}
return database.AddTargetTotals(tx, add)
}))
}
// statsAge moves an event's arrival back to the given time.
func statsAge(
t *testing.T, webhookDB *gorm.DB, eventID string, at time.Time,
) {
t.Helper()
require.NoError(t, webhookDB.Model(&database.Event{}).
Where("id = ?", eventID).
Update("created_at", at).Error)
}
// statsTargetTotals reads a webhook database's target totals, keyed by
// target.
func statsTargetTotals(
t *testing.T, webhookDB *gorm.DB,
) map[string]database.TargetTotals {
t.Helper()
var rows []database.TargetTotals
require.NoError(t, webhookDB.Find(&rows).Error)
byTarget := make(map[string]database.TargetTotals, len(rows))
for _, row := range rows {
byTarget[row.TargetID] = row
}
return byTarget
}
// statsHistory is the webhook seedStatsHistory builds: its event
// database, its newest event, and its two active targets.
type statsHistory struct {
webhook *database.Webhook
webhookDB *gorm.DB
newest database.Event
first, second string
}
// seedStatsHistory builds the webhook the statistics test checks: one
// day of retention, two entrypoints (one inactive) and three targets
// (one inactive). Three events arrive through the receiver, and so
// each has a delivery to the two active targets. The oldest event is
// past retention, the middle one six hours old, the newest just in.
// Their deliveries are settled as the delivery engine would, and a
// replay adds a pending delivery to the oldest event.
func seedStatsHistory(
t *testing.T,
h *handlers.Handlers,
sess *session.Session,
db *database.Database,
dbMgr *database.WebhookDBManager,
) statsHistory {
t.Helper()
wh := &database.Webhook{
UserID: deleteTestUserID, Name: "stats", RetentionDays: 1,
}
require.NoError(t, db.DB().Omit(clause.Associations).Create(wh).Error)
path := statsEntrypoint(t, db, wh.ID, true)
statsEntrypoint(t, db, wh.ID, false)
first := seedConfiguredTarget(
t, db, wh.ID, database.TargetTypeHTTP,
`{"url":"`+replayTargetURL+`"}`,
)
second := seedTarget(t, db, wh.ID, database.TargetTypeLog)
inactive := seedTarget(t, db, wh.ID, database.TargetTypeLog)
require.NoError(t, db.DB().Model(inactive).
Update("active", false).Error)
router := receiverRouter(h)
for range 3 {
require.Equal(t, http.StatusOK, postReceiver(t, router, path))
}
webhookDB, err := dbMgr.GetDB(wh.ID)
require.NoError(t, err)
events := listEvents(t, webhookDB)
require.Len(t, events, 3)
oldest, middle, newest := events[0], events[1], events[2]
now := time.Now()
statsAge(t, webhookDB, oldest.ID, now.Add(-50*time.Hour))
statsAge(t, webhookDB, middle.ID, now.Add(-6*time.Hour))
oldestFailure := statsDelivery(t, webhookDB, oldest.ID, first.ID)
statsFinish(t, webhookDB, oldestFailure,
database.DeliveryStatusFailed, now.Add(-49*time.Hour))
statsFinish(t, webhookDB,
statsDelivery(t, webhookDB, oldest.ID, second.ID),
database.DeliveryStatusDelivered, now.Add(-49*time.Hour))
statsFinish(t, webhookDB,
statsDelivery(t, webhookDB, middle.ID, first.ID),
database.DeliveryStatusFailed, now.Add(-5*time.Hour))
statsFinish(t, webhookDB,
statsDelivery(t, webhookDB, middle.ID, second.ID),
database.DeliveryStatusFailed, now.Add(-time.Minute))
statsFinish(t, webhookDB,
statsDelivery(t, webhookDB, newest.ID, first.ID),
database.DeliveryStatusDelivered, now.Add(-2*time.Minute))
require.Equal(t, http.StatusSeeOther,
postReplay(t, h, sess, wh.ID, oldestFailure.ID).Code)
return statsHistory{
webhook: wh,
webhookDB: webhookDB,
newest: newest,
first: first.ID,
second: second.ID,
}
}
// statsPrune runs the real retention reaper until it has removed one
// event from the webhook's database, then stops it.
func statsPrune(
t *testing.T,
db *database.Database,
dbMgr *database.WebhookDBManager,
log *logger.Logger,
webhookDB *gorm.DB,
) {
t.Helper()
lc := fxtest.NewLifecycle(t)
database.NewRetentionReaper(lc, database.RetentionReaperParams{
Config: &config.Config{
RetentionSweepInterval: 10 * time.Millisecond,
},
Database: db,
DBManager: dbMgr,
Logger: log,
})
lc.RequireStart()
require.Eventually(t, func() bool {
var totals database.EventTotals
err := webhookDB.Take(&totals).Error
return err == nil && totals.EventsRemoved == 1
}, 10*time.Second, 10*time.Millisecond)
lc.RequireStop()
}
// statsPane returns the statistics pane from a rendered webhook page:
// everything from its heading to the next heading on the page.
func statsPane(t *testing.T, page string) string {
t.Helper()
_, pane, found := strings.Cut(page, ">Statistics</h2>")
require.True(t, found, "the page has no statistics pane")
pane, _, _ = strings.Cut(pane, "<h2")
return pane
}
// TestWebhookStats_EveryFigureAcrossRetentionPrune checks every figure
// the statistics pane shows for the history seedStatsHistory builds,
// and each target's totals and recent figures, before and after the
// real retention reaper removes the oldest event.
func TestWebhookStats_EveryFigureAcrossRetentionPrune(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
dbMgr *database.WebhookDBManager
log *logger.Logger
)
app := newTestApp(t, &h, &sess, &db, &dbMgr, &log)
app.RequireStart()
t.Cleanup(app.RequireStop)
hist := seedStatsHistory(t, h, sess, db, dbMgr)
first, second := hist.first, hist.second
stats := h.WebhookStatsForTest(hist.webhook.ID)
require.NotNil(t, stats)
assert.Equal(t, 2, stats.Entrypoints)
assert.Equal(t, 1, stats.ActiveEntrypoints)
assert.Equal(t, 3, stats.Targets)
assert.Equal(t, 2, stats.ActiveTargets)
assert.Equal(t, handlers.Counts{Events: 3, Deliveries: 7, Failures: 3},
stats.Lifetime)
assert.Equal(t, stats.Lifetime, stats.WithinRetention)
assert.Equal(t, int64(2), stats.InProgress)
require.NotNil(t, stats.LastEventAt)
assert.True(t, hist.newest.CreatedAt.Equal(*stats.LastEventAt))
assert.Equal(t, handlers.RecentWindow{
Events: 1, Delivered: 1, Failed: 1,
}, stats.Last10Minutes)
assert.Equal(t, handlers.RecentWindow{
Events: 2, Delivered: 1, Failed: 2,
}, stats.Last24Hours)
assert.Equal(t, "50.0%", stats.Last10Minutes.FailurePercent())
assert.Equal(t, "66.7%", stats.Last24Hours.FailurePercent())
// The first target has three deliveries and the replay, the second
// three; the inactive target has none and so no row.
assert.Equal(t, map[string]database.TargetTotals{
first: {TargetID: first, Deliveries: 4, Delivered: 1, Failed: 2},
second: {
TargetID: second, Deliveries: 3, Delivered: 1, Failed: 1,
},
}, statsTargetTotals(t, hist.webhookDB))
lastDay, err := handlers.FinishedByTargetForTest(
hist.webhookDB, time.Now().Add(-24*time.Hour),
)
require.NoError(t, err)
assert.ElementsMatch(t, []handlers.TargetFinished{
{TargetID: first, Delivered: 1, Failed: 1},
{TargetID: second, Failed: 1},
}, lastDay)
// Retention removes the oldest event with its three deliveries:
// the first target's failed one and the pending replay, and the
// second target's delivered one.
statsPrune(t, db, dbMgr, log, hist.webhookDB)
after := h.WebhookStatsForTest(hist.webhook.ID)
require.NotNil(t, after)
assert.Equal(t, stats.Lifetime, after.Lifetime)
assert.Equal(t, handlers.Counts{Events: 2, Deliveries: 4, Failures: 2},
after.WithinRetention)
assert.Equal(t, int64(1), after.InProgress)
assert.Equal(t, stats.LastEventAt, after.LastEventAt)
assert.Equal(t, stats.Last10Minutes, after.Last10Minutes)
assert.Equal(t, stats.Last24Hours, after.Last24Hours)
assert.Equal(t, map[string]database.TargetTotals{
first: {
TargetID: first, Deliveries: 4, Delivered: 1, Failed: 2,
DeliveriesRemoved: 2, FailedRemoved: 1,
},
second: {
TargetID: second, Deliveries: 3, Delivered: 1, Failed: 1,
DeliveriesRemoved: 1,
},
}, statsTargetTotals(t, hist.webhookDB))
pane := statsPane(t, renderSourceDetailPage(t, h, sess, hist.webhook.ID))
assert.Contains(t, pane, "Within retention")
assert.Contains(t, pane, "50.0%")
assert.Contains(t, pane, "66.7%")
}
// TestWebhookStats_PaneShowsRetentionPeriod checks that the statistics
// pane itself, not only the line at the foot of the page, shows the
// webhook's retention period, for a finite one and for forever.
func TestWebhookStats_PaneShowsRetentionPeriod(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
)
app := newTestApp(t, &h, &sess, &db)
app.RequireStart()
t.Cleanup(app.RequireStop)
tests := []struct {
retentionDays int
want string
}{
{30, "30 days"},
{database.RetentionForeverDays, "forever"},
}
for _, tt := range tests {
wh := &database.Webhook{
UserID: deleteTestUserID,
Name: "retention",
RetentionDays: tt.retentionDays,
}
require.NoError(t,
db.DB().Omit(clause.Associations).Create(wh).Error)
pane := statsPane(t, renderSourceDetailPage(t, h, sess, wh.ID))
assert.Contains(t, pane, "Retention", tt.want)
assert.Contains(t, pane, tt.want)
}
}
// TestWebhookStats_WebhookWithNoEvents covers a webhook whose event
// database has never been opened: every count is zero, the
// percentages are a dash, and showing the page does not create the
// database.
func TestWebhookStats_WebhookWithNoEvents(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
dbMgr *database.WebhookDBManager
)
app := newTestApp(t, &h, &sess, &db, &dbMgr)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
assert.Equal(t, &handlers.WebhookStats{}, h.WebhookStatsForTest(wh.ID))
assert.Equal(t, "—", handlers.RecentWindow{}.FailurePercent())
statsPane(t, renderSourceDetailPage(t, h, sess, wh.ID))
assert.False(t, dbMgr.DBExists(wh.ID))
}
// TestRecentWindow_FailurePercent pins the percentage: failed
// deliveries out of all that finished in the window.
func TestRecentWindow_FailurePercent(t *testing.T) {
t.Parallel()
tests := []struct {
window handlers.RecentWindow
want string
}{
{handlers.RecentWindow{}, "—"},
{handlers.RecentWindow{Events: 4}, "—"},
{handlers.RecentWindow{Delivered: 3, Failed: 1}, "25.0%"},
{handlers.RecentWindow{Failed: 2}, "100.0%"},
{handlers.RecentWindow{Delivered: 2}, "0.0%"},
}
for _, tt := range tests {
assert.Equal(t, tt.want, tt.window.FailurePercent(), tt.window)
}
}
+27 -20
View File
@@ -3,17 +3,18 @@
// deliveries are attempted, how they end, how long they take, how // deliveries are attempted, how they end, how long they take, how
// deep the queues are, and how many circuit breakers are open. // deep the queues are, and how many circuit breakers are open.
// //
// The inbound HTTP metrics come from the go-http-metrics recorder in // It also builds the registry the authenticated /metrics route
// internal/middleware and land on prometheus.DefaultRegisterer. These // serves. In production, these collectors, the inbound HTTP metrics
// collectors register there too, so both surfaces are gathered by the // recorded in internal/middleware, and the Go runtime and process
// one promhttp handler mounted on the authenticated /metrics route. // collectors all register on that one registry, never on Prometheus's
// global default.
package metrics package metrics
import ( import (
"sync"
"time" "time"
"github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/collectors"
"github.com/prometheus/client_golang/prometheus/promauto" "github.com/prometheus/client_golang/prometheus/promauto"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
) )
@@ -57,25 +58,31 @@ var knownTargetTypes = []database.TargetType{
database.TargetTypeSlack, database.TargetTypeSlack,
} }
// defaultSet is the process-wide metric set, registered on the same // NewRegistry returns the registry /metrics serves, carrying the Go
// registry the HTTP middleware and the /metrics handler already use. // runtime and process collectors that Prometheus's global default
// It is built on first use rather than in an init so that a test // registry carries, so the go_* and process_* series stay in the
// binary that never touches metrics never registers them. // scrape.
// //
//nolint:gochecknoglobals // one process-wide registration, by design // A registry of its own, rather than the global default, is what lets
var defaultSet = sync.OnceValue(func() *Set { // two dependency graphs in one process — two tests, say — each
return New(prometheus.DefaultRegisterer) // register their collectors without the second registration
}) // panicking.
func NewRegistry() *prometheus.Registry {
reg := prometheus.NewRegistry()
reg.MustRegister(
collectors.NewGoCollector(),
collectors.NewProcessCollector(
collectors.ProcessCollectorOpts{},
),
)
// Default returns the process-wide metric set. return reg
func Default() *Set {
return defaultSet()
} }
// Set is one registered group of webhooker's delivery collectors. // Set is one registered group of webhooker's delivery collectors.
// Production uses the single Default set; tests build their own // Production builds one on the registry /metrics serves; tests build
// against a private registry so assertions are not disturbed by // their own against a private registry so assertions are not
// deliveries other tests are making concurrently. // disturbed by deliveries other tests are making concurrently.
type Set struct { type Set struct {
eventsReceived prometheus.Counter eventsReceived prometheus.Counter
deliveryAttempts *prometheus.CounterVec deliveryAttempts *prometheus.CounterVec
@@ -93,7 +100,7 @@ type Set struct {
// New registers a full set of delivery collectors on reg and returns // New registers a full set of delivery collectors on reg and returns
// it. It panics if reg already holds them, which is the intended // it. It panics if reg already holds them, which is the intended
// behaviour for a duplicate registration. // behaviour for a duplicate registration.
func New(reg prometheus.Registerer) *Set { func New(reg *prometheus.Registry) *Set {
factory := promauto.With(reg) factory := promauto.With(reg)
s := &Set{ s := &Set{
+1 -2
View File
@@ -10,8 +10,7 @@ import (
// MetricsMiddlewareForTest builds the metrics recording middleware // MetricsMiddlewareForTest builds the metrics recording middleware
// against a caller-supplied recorder, so a test can gather from its // against a caller-supplied recorder, so a test can gather from its
// own Prometheus registry rather than the process-wide default one // own Prometheus registry without building a whole Middleware.
// that Middleware.Metrics uses.
func MetricsMiddlewareForTest( func MetricsMiddlewareForTest(
rec httpmetrics.Recorder, rec httpmetrics.Recorder,
) func(http.Handler) http.Handler { ) func(http.Handler) http.Handler {
+7 -8
View File
@@ -7,7 +7,6 @@ import (
"github.com/go-chi/chi" "github.com/go-chi/chi"
httpmetrics "github.com/slok/go-http-metrics/metrics" httpmetrics "github.com/slok/go-http-metrics/metrics"
prommetrics "github.com/slok/go-http-metrics/metrics/prometheus"
ghmm "github.com/slok/go-http-metrics/middleware" ghmm "github.com/slok/go-http-metrics/middleware"
"github.com/slok/go-http-metrics/middleware/std" "github.com/slok/go-http-metrics/middleware/std"
) )
@@ -151,17 +150,17 @@ func (r boundedLabelRecorder) AddInflightRequests(
var _ httpmetrics.Recorder = boundedLabelRecorder{} var _ httpmetrics.Recorder = boundedLabelRecorder{}
// Metrics returns middleware that records Prometheus HTTP metrics on // Metrics returns middleware that records Prometheus HTTP metrics
// the default registry, which is the one the /metrics route gathers. // with the Middleware's one recorder, which New builds on the registry
// the /metrics route serves and NewForTest on a registry of its own.
// Every call reuses that recorder, so any number of routers can
// install it.
func (s *Middleware) Metrics() func(http.Handler) http.Handler { func (s *Middleware) Metrics() func(http.Handler) http.Handler {
return metricsMiddleware( return metricsMiddleware(s.metricsRecorder)
prommetrics.NewRecorder(prommetrics.Config{}),
)
} }
// metricsMiddleware builds the recording middleware against a given // metricsMiddleware builds the recording middleware against a given
// recorder, so tests can gather from a registry of their own instead // recorder, so tests can gather from a registry of their own.
// of the process-wide default.
func metricsMiddleware( func metricsMiddleware(
rec httpmetrics.Recorder, rec httpmetrics.Recorder,
) func(http.Handler) http.Handler { ) func(http.Handler) http.Handler {
+28 -3
View File
@@ -57,9 +57,8 @@ const (
// Server.setupWebhookRoutes inside it. That ordering is the whole // Server.setupWebhookRoutes inside it. That ordering is the whole
// defect, so a test that flattens it would prove nothing. // defect, so a test that flattens it would prove nothing.
// //
// The recorder writes to a registry of the test's own rather than the // The recorder writes to a registry of the test's own, so each test
// process-wide default one, so each test observes only its own // observes only its own traffic.
// traffic.
func metricsTestRouter( func metricsTestRouter(
t *testing.T, t *testing.T,
receiverLimit int, receiverLimit int,
@@ -455,3 +454,29 @@ func TestMetrics_StatusAndSizeStillRecorded(t *testing.T) {
"the interceptor must still count written bytes", "the interceptor must still count written bytes",
) )
} }
// TestMetrics_WorksOnNewForTestMiddleware pins that a Middleware built
// by NewForTest has a recorder of its own: its Metrics() serves a
// request instead of panicking, and a second one does not collide
// with the first.
func TestMetrics_WorksOnNewForTestMiddleware(t *testing.T) {
t.Parallel()
log := slog.New(slog.DiscardHandler)
cfg := &config.Config{Environment: "prod"}
ok := http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
_, _ = w.Write([]byte(okBody))
})
for range 2 {
h := middleware.NewForTest(log, cfg, nil).Metrics()(ok)
req := httptest.NewRequestWithContext(
t.Context(), http.MethodGet, okRoute, nil,
)
w := httptest.NewRecorder()
h.ServeHTTP(w, req)
assert.Equal(t, http.StatusOK, w.Code)
}
}
+19 -4
View File
@@ -13,6 +13,9 @@ import (
"github.com/go-chi/chi" "github.com/go-chi/chi"
"github.com/go-chi/chi/middleware" "github.com/go-chi/chi/middleware"
"github.com/go-chi/cors" "github.com/go-chi/cors"
"github.com/prometheus/client_golang/prometheus"
httpmetrics "github.com/slok/go-http-metrics/metrics"
prommetrics "github.com/slok/go-http-metrics/metrics/prometheus"
"go.uber.org/fx" "go.uber.org/fx"
"sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/globals" "sneak.berlin/go/webhooker/internal/globals"
@@ -148,10 +151,11 @@ const (
type MiddlewareParams struct { type MiddlewareParams struct {
fx.In fx.In
Logger *logger.Logger Logger *logger.Logger
Globals *globals.Globals Globals *globals.Globals
Config *config.Config Config *config.Config
Session *session.Session Session *session.Session
Registry *prometheus.Registry
} }
// Middleware provides HTTP middleware for logging, CORS, auth, and // Middleware provides HTTP middleware for logging, CORS, auth, and
@@ -161,6 +165,14 @@ type Middleware struct {
params *MiddlewareParams params *MiddlewareParams
session *session.Session session *session.Session
// metricsRecorder records the inbound HTTP metrics. New builds
// it on the registry /metrics serves, NewForTest on a registry
// of its own. Either way it is built once per Middleware and
// Metrics reuses it, because building it registers its
// collectors, and a second registration on the same registry
// panics.
metricsRecorder httpmetrics.Recorder
// loginGuard counts failed credential verifications and bounds // loginGuard counts failed credential verifications and bounds
// concurrent password hashing. It is built on first use so that // concurrent password hashing. It is built on first use so that
// every construction path gets one; see guard(). // every construction path gets one; see guard().
@@ -179,6 +191,9 @@ func New(
s.params = &params s.params = &params
s.log = params.Logger.Get() s.log = params.Logger.Get()
s.session = params.Session s.session = params.Session
s.metricsRecorder = prommetrics.NewRecorder(
prommetrics.Config{Registry: params.Registry},
)
return s, nil return s, nil
} }
+8
View File
@@ -3,12 +3,17 @@ package middleware
import ( import (
"log/slog" "log/slog"
"github.com/prometheus/client_golang/prometheus"
prommetrics "github.com/slok/go-http-metrics/metrics/prometheus"
"sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/session" "sneak.berlin/go/webhooker/internal/session"
) )
// NewForTest creates a Middleware with the minimum dependencies // NewForTest creates a Middleware with the minimum dependencies
// needed for testing. This bypasses the fx lifecycle. // needed for testing. This bypasses the fx lifecycle.
//
// Its metrics recorder writes to a fresh registry of its own, so
// Metrics() works on it and two of them never collide.
func NewForTest( func NewForTest(
log *slog.Logger, log *slog.Logger,
cfg *config.Config, cfg *config.Config,
@@ -20,5 +25,8 @@ func NewForTest(
Config: cfg, Config: cfg,
}, },
session: sess, session: sess,
metricsRecorder: prommetrics.NewRecorder(
prommetrics.Config{Registry: prometheus.NewRegistry()},
),
} }
} }
+3
View File
@@ -24,6 +24,7 @@ import (
"sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck" "sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware" "sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/resetpw" "sneak.berlin/go/webhooker/internal/resetpw"
"sneak.berlin/go/webhooker/internal/session" "sneak.berlin/go/webhooker/internal/session"
@@ -163,6 +164,8 @@ func newServerApp(
session.New, session.New,
func() delivery.Notifier { return &noopNotifier{} }, func() delivery.Notifier { return &noopNotifier{} },
func() delivery.WebhookEvictor { return &noopEvictor{} }, func() delivery.WebhookEvictor { return &noopEvictor{} },
metrics.NewRegistry,
metrics.New,
middleware.New, middleware.New,
delivery.NewGuard, delivery.NewGuard,
handlers.New, handlers.New,
+1 -7
View File
@@ -7,7 +7,6 @@ import (
sentryhttp "github.com/getsentry/sentry-go/http" sentryhttp "github.com/getsentry/sentry-go/http"
"github.com/go-chi/chi" "github.com/go-chi/chi"
"github.com/go-chi/chi/middleware" "github.com/go-chi/chi/middleware"
"github.com/prometheus/client_golang/prometheus/promhttp"
"sneak.berlin/go/webhooker/static" "sneak.berlin/go/webhooker/static"
) )
@@ -130,12 +129,7 @@ func (s *Server) setupRoutes() {
if s.params.Config.MetricsAuthEnabled() { if s.params.Config.MetricsAuthEnabled() {
s.router.Group(func(r chi.Router) { s.router.Group(func(r chi.Router) {
r.Use(s.mw.MetricsAuth()) r.Use(s.mw.MetricsAuth())
r.Get( r.Get("/metrics", s.h.HandleMetrics())
"/metrics",
http.HandlerFunc(
promhttp.Handler().ServeHTTP,
),
)
}) })
} }
+46
View File
@@ -24,6 +24,7 @@ import (
"sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/healthcheck" "sneak.berlin/go/webhooker/internal/healthcheck"
"sneak.berlin/go/webhooker/internal/logger" "sneak.berlin/go/webhooker/internal/logger"
"sneak.berlin/go/webhooker/internal/metrics"
"sneak.berlin/go/webhooker/internal/middleware" "sneak.berlin/go/webhooker/internal/middleware"
"sneak.berlin/go/webhooker/internal/server" "sneak.berlin/go/webhooker/internal/server"
"sneak.berlin/go/webhooker/internal/session" "sneak.berlin/go/webhooker/internal/session"
@@ -113,6 +114,8 @@ func newTestEnvWithConfig(
session.New, session.New,
func() delivery.Notifier { return &noopNotifier{} }, func() delivery.Notifier { return &noopNotifier{} },
func() delivery.WebhookEvictor { return &noopEvictor{} }, func() delivery.WebhookEvictor { return &noopEvictor{} },
metrics.NewRegistry,
metrics.New,
middleware.New, middleware.New,
delivery.NewGuard, delivery.NewGuard,
handlers.New, handlers.New,
@@ -1027,3 +1030,46 @@ func TestMetricsRouteUnmountedOnHalfSetConfig(t *testing.T) {
}) })
} }
} }
// TestTwoMetricsRoutersInOneProcess pins
// https://git.eeqj.de/sneak/webhooker/issues/227: a second
// metrics-enabled router in one process used to panic, because the
// HTTP metrics registered on Prometheus's global default registry.
// Two routers are built over separate dependency graphs and a third
// over the first graph again, and each must still serve the HTTP,
// delivery, Go runtime and process series, and the series counting
// scrapes of /metrics itself.
func TestTwoMetricsRoutersInOneProcess(t *testing.T) {
t.Parallel()
first := newTestEnvWithConfig(
t, metricsConfig(t, metricsUser, metricsAuthValue),
)
second := newTestEnvWithConfig(
t, metricsConfig(t, metricsUser, metricsAuthValue),
)
third := &testEnv{
router: server.NewRouterForTest(
first.log.Get(), first.cfg, first.mw, first.hnd,
),
}
for _, env := range []*testEnv{first, second, third} {
env.get("/", nil)
scrape := env.metricsRequest(metricsUser, metricsAuthValue)
require.Equal(t, http.StatusOK, scrape.Code)
for _, series := range []string{
"http_request_duration_seconds",
"http_response_size_bytes",
"http_requests_inflight",
"webhooker_events_received_total",
"go_goroutines",
"process_start_time_seconds",
"promhttp_metric_handler_requests_total",
} {
assert.Contains(t, scrape.Body.String(), series)
}
}
}
-2
View File
@@ -24,8 +24,6 @@
</div> </div>
</div> </div>
{{template "webhook_stats" .}}
<div class="grid grid-cols-1 lg:grid-cols-2 gap-6"> <div class="grid grid-cols-1 lg:grid-cols-2 gap-6">
<!-- Entrypoints --> <!-- Entrypoints -->
<div class="card"> <div class="card">
-93
View File
@@ -1,93 +0,0 @@
{{define "webhook_stats"}}
<!-- Statistics pane at the top of the webhook page. -->
<div class="card mb-6">
<div class="p-4 border-b border-gray-200">
<h2 class="text-lg font-medium text-gray-900">Statistics</h2>
</div>
{{with .Stats}}
<div class="p-4 flex flex-wrap gap-6 text-sm border-b border-gray-200">
<div>
<span class="text-gray-500">Entrypoints</span>
<span class="font-medium text-gray-900">{{.Entrypoints}}</span>
<span class="text-gray-500">({{.ActiveEntrypoints}} active)</span>
</div>
<div>
<span class="text-gray-500">Targets</span>
<span class="font-medium text-gray-900">{{.Targets}}</span>
<span class="text-gray-500">({{.ActiveTargets}} active)</span>
</div>
<div>
<span class="text-gray-500">Deliveries in progress</span>
<span class="font-medium text-gray-900">{{.InProgress}}</span>
</div>
<div>
<span class="text-gray-500">Last event</span>
<span class="font-medium text-gray-900">{{with .LastEventAt}}{{.Format "2006-01-02 15:04:05 UTC"}}{{else}}none{{end}}</span>
</div>
<div>
<span class="text-gray-500">Retention</span>
<span class="font-medium text-gray-900">{{$.Webhook.RetentionLabel}}</span>
</div>
</div>
<div class="p-4 grid grid-cols-1 lg:grid-cols-2 gap-6 text-sm">
<table class="w-full text-center text-gray-900">
<thead>
<tr class="border-b border-gray-200 text-xs text-gray-500 uppercase tracking-wide">
<th></th>
<th class="py-2 font-medium">Lifetime</th>
<th class="py-2 font-medium">Within retention</th>
</tr>
</thead>
<tbody>
<tr>
<td class="py-2 text-left text-gray-600">Events</td>
<td class="py-2">{{.Lifetime.Events}}</td>
<td class="py-2">{{.WithinRetention.Events}}</td>
</tr>
<tr>
<td class="py-2 text-left text-gray-600">Deliveries</td>
<td class="py-2">{{.Lifetime.Deliveries}}</td>
<td class="py-2">{{.WithinRetention.Deliveries}}</td>
</tr>
<tr>
<td class="py-2 text-left text-gray-600">Failures</td>
<td class="py-2">{{.Lifetime.Failures}}</td>
<td class="py-2">{{.WithinRetention.Failures}}</td>
</tr>
</tbody>
</table>
<div>
<table class="w-full text-center text-gray-900">
<thead>
<tr class="border-b border-gray-200 text-xs text-gray-500 uppercase tracking-wide">
<th></th>
<th class="py-2 font-medium">Last 10 minutes</th>
<th class="py-2 font-medium">Last 24 hours</th>
</tr>
</thead>
<tbody>
<tr>
<td class="py-2 text-left text-gray-600">Events</td>
<td class="py-2">{{.Last10Minutes.Events}}</td>
<td class="py-2">{{.Last24Hours.Events}}</td>
</tr>
<tr>
<td class="py-2 text-left text-gray-600">Failures</td>
<td class="py-2">{{.Last10Minutes.Failed}}</td>
<td class="py-2">{{.Last24Hours.Failed}}</td>
</tr>
<tr>
<td class="py-2 text-left text-gray-600">Failure percentage</td>
<td class="py-2">{{.Last10Minutes.FailurePercent}}</td>
<td class="py-2">{{.Last24Hours.FailurePercent}}</td>
</tr>
</tbody>
</table>
<p class="mt-2 text-xs text-gray-500">Failure percentage is the failed deliveries out of all deliveries that finished in the window. Deliveries still pending or retrying are not counted.</p>
</div>
</div>
{{else}}
<div class="p-4 text-sm text-gray-500">The statistics could not be read.</div>
{{end}}
</div>
{{end}}