Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ff0018cf43 |
@@ -1065,7 +1065,7 @@ unconditionally against whatever files it finds:
|
||||
- the main database on connect — `Setting`, `User`, `APIKey`, `Webhook`,
|
||||
`Entrypoint`, `Target`
|
||||
- each event database when it is lazily opened — `Event`, `Delivery`,
|
||||
`DeliveryResult`
|
||||
`DeliveryResult`, `EventTotals`, `TargetTotals`
|
||||
- each archive database on every open and reopen
|
||||
|
||||
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
|
||||
|
||||
webhooker's data model has nine entities organized into two tiers: the
|
||||
webhooker's data model has eleven entities organized into two tiers: the
|
||||
**application tier** (user and webhook configuration) and the **event
|
||||
tier** (event ingestion, delivery, and logging).
|
||||
|
||||
@@ -1410,6 +1410,13 @@ tier** (event ingestion, delivery, and logging).
|
||||
│ ┌──────────┐ ┌──────────┐ ┌─────────────────┐ │
|
||||
│ │ Event │──1:N──│ Delivery │──1:N──│ DeliveryResult │ │
|
||||
│ └──────────┘ └──────────┘ └─────────────────┘ │
|
||||
│ │
|
||||
│ ┌──────────────┐ (one row: running counts of events) │
|
||||
│ │ EventTotals │ │
|
||||
│ └──────────────┘ │
|
||||
│ ┌──────────────┐ (one row per target: running counts │
|
||||
│ │ TargetTotals │ of its deliveries) │
|
||||
│ └──────────────┘ │
|
||||
└─────────────────────────────────────────────────────────────┘
|
||||
```
|
||||
|
||||
@@ -1662,6 +1669,7 @@ status across potentially multiple attempts.
|
||||
| `event_id` | UUID | Foreign key → Event |
|
||||
| `target_id`| UUID | Foreign key → Target |
|
||||
| `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
|
||||
DeliveryResults.
|
||||
@@ -1729,33 +1737,65 @@ retries) is individually logged for full observability.
|
||||
|
||||
**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
|
||||
|
||||
These indexes on the per-webhook event databases are declared in the model
|
||||
tags, so `AutoMigrate` creates them on a fresh and on an existing database:
|
||||
tags, so `AutoMigrate` creates them on a fresh database:
|
||||
|
||||
| Table | Columns | Serves |
|
||||
| ------------------ | --------------------------- | ------ |
|
||||
| `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 selects and deletes the deliveries of expired events |
|
||||
| `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` | `event_id`, `deleted_at` | The event log, which loads each event's deliveries, and retention, which counts 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 |
|
||||
| `events` | `deleted_at`, `created_at` | Retention, which selects expired events by age |
|
||||
| `events` | `created_at` | Retention's delete of the expired events themselves |
|
||||
| `events` | `deleted_at`, `created_at` | The webhook page's statistics, which count recent events and find the newest |
|
||||
| `events` | `created_at` | Retention, which selects expired events by age |
|
||||
|
||||
GORM's soft delete adds `deleted_at IS NULL` to these queries; retention's
|
||||
deletes leave it out, but their lookups of expired rows keep it. SQLite keeps
|
||||
no statistics on these tables, and without them it rates the `deleted_at`
|
||||
index, which every live row matches, above an index on a column matched
|
||||
against several values or compared with `<`. So every index but the last also
|
||||
covers `deleted_at`. It comes second, so that retention's deletes can use the
|
||||
index without it, except in `events`, where `created_at` is compared with `<`
|
||||
and SQLite narrows by a `<` only on the last column it uses.
|
||||
GORM's soft delete adds `deleted_at IS NULL` to these queries; retention
|
||||
leaves it out. SQLite keeps no statistics on these tables, and without them it
|
||||
rates the `deleted_at` index, which every live row matches, above an index on
|
||||
a column matched against several values or compared with `<`. So every index
|
||||
but the last also covers `deleted_at`. It comes second, so that retention can
|
||||
use the 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
|
||||
|
||||
Every entity except `Setting` includes these fields from `BaseModel`.
|
||||
`Setting` is a bare key-value row with no `id`, no timestamps and no
|
||||
soft delete:
|
||||
Every entity except `Setting`, `EventTotals` and `TargetTotals` includes
|
||||
these fields from `BaseModel`. `Setting` is a bare key-value row with no
|
||||
`id`, no timestamps and no soft delete, and the two totals tables hold
|
||||
only counts, keyed by a numeric `id` and by `target_id`:
|
||||
|
||||
| Field | Type | Description |
|
||||
| ------------ | --------- | ----------- |
|
||||
@@ -1797,6 +1837,8 @@ encryption key is generated and stored, and an `admin` user is created.
|
||||
- **Events** — captured incoming webhook payloads
|
||||
- **Deliveries** — event-to-target pairings and their status
|
||||
- **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
|
||||
created (and lazily on first access for webhooks that predate this
|
||||
@@ -2777,6 +2819,7 @@ webhooker/
|
||||
│ │ ├── model_event.go # Event entity (per-webhook DB)
|
||||
│ │ ├── model_delivery.go # Delivery 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
|
||||
│ │ ├── password.go # Argon2id hashing and verification
|
||||
│ │ ├── retention.go # Retention reaper (per-webhook event expiry)
|
||||
|
||||
@@ -16,7 +16,6 @@ import (
|
||||
"sneak.berlin/go/webhooker/internal/handlers"
|
||||
"sneak.berlin/go/webhooker/internal/healthcheck"
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
"sneak.berlin/go/webhooker/internal/metrics"
|
||||
"sneak.berlin/go/webhooker/internal/middleware"
|
||||
"sneak.berlin/go/webhooker/internal/resetpw"
|
||||
"sneak.berlin/go/webhooker/internal/server"
|
||||
@@ -178,10 +177,6 @@ func newApp() *fx.App {
|
||||
healthcheck.New,
|
||||
session.New,
|
||||
handlers.New,
|
||||
// The registry /metrics serves, and the delivery
|
||||
// collectors registered on it.
|
||||
metrics.NewRegistry,
|
||||
metrics.New,
|
||||
middleware.New,
|
||||
// The one SSRF guard both target-creation validation
|
||||
// and the delivery dialer consult, so they cannot
|
||||
|
||||
@@ -93,11 +93,11 @@ func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
|
||||
deliveries []database.Delivery
|
||||
results []database.DeliveryResult
|
||||
depths []struct{ Depth int }
|
||||
removed []database.TargetTotals
|
||||
)
|
||||
|
||||
byStatus := "idx_deliveries_status (status=? 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
|
||||
// stranded pending deliveries, and the queue depth count.
|
||||
@@ -123,25 +123,89 @@ func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
|
||||
Order("attempt_num ASC").Find(&results),
|
||||
"idx_delivery_results_delivery_id (delivery_id=? AND deleted_at=?)")
|
||||
|
||||
// Retention's three deletes (reapExpired), whose subqueries are built
|
||||
// afresh for each statement as it builds them.
|
||||
expiredEventIDs := func() *gorm.DB {
|
||||
return dry.Model(&database.Event{}).Select("id").
|
||||
Where("created_at < ?", cutoff)
|
||||
}
|
||||
// Retention (reapExpired, deleteEvents): one batch of expired
|
||||
// events, then their attempts, deliveries and the events.
|
||||
var expired []string
|
||||
|
||||
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(
|
||||
"delivery_id IN (?)", dry.Model(&database.Delivery{}).
|
||||
Select("id").Where("event_id IN (?)", expiredEventIDs()),
|
||||
"delivery_id IN (?)", dry.Unscoped().Model(&database.Delivery{}).
|
||||
Select("id").Where("event_id IN ?", ids),
|
||||
).Delete(&database.DeliveryResult{}),
|
||||
"idx_delivery_results_delivery_id (delivery_id=?)", byEvent, byAge)
|
||||
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||
"event_id IN (?)", expiredEventIDs(),
|
||||
).Delete(&database.Delivery{}),
|
||||
"idx_deliveries_event_id (event_id=?)", byAge)
|
||||
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||
"created_at < ?", cutoff,
|
||||
).Delete(&database.Event{}), "idx_events_created_at (created_at<?)")
|
||||
"idx_delivery_results_delivery_id (delivery_id=?)",
|
||||
"idx_deliveries_event_id (event_id=?)")
|
||||
assertPlanUses(t, db, dry.Unscoped().Model(&database.Delivery{}).
|
||||
Select("target_id, count(*) AS deliveries_removed, "+
|
||||
"count(CASE WHEN status = ? THEN 1 END) AS failed_removed",
|
||||
database.DeliveryStatusFailed).
|
||||
Where("event_id IN ?", ids).Group("target_id").Find(&removed),
|
||||
"idx_deliveries_event_id (event_id=?)")
|
||||
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
|
||||
@@ -152,6 +216,18 @@ func assertPlanUses(
|
||||
) {
|
||||
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 }
|
||||
|
||||
require.NoError(t, db.Raw(
|
||||
@@ -159,8 +235,5 @@ func assertPlanUses(
|
||||
built.Statement.Vars...,
|
||||
).Scan(&plan).Error)
|
||||
|
||||
for _, index := range indexes {
|
||||
assert.Contains(t, fmt.Sprint(plan), index,
|
||||
built.Statement.SQL.String())
|
||||
}
|
||||
return fmt.Sprint(plan)
|
||||
}
|
||||
|
||||
@@ -28,6 +28,10 @@ func NewTestRetentionReaper(
|
||||
}
|
||||
}
|
||||
|
||||
// ExportReapBatchSize exposes how many expired events one retention
|
||||
// transaction deletes.
|
||||
const ExportReapBatchSize = reapBatchSize
|
||||
|
||||
// ExportSweep runs a single retention sweep synchronously for tests.
|
||||
func (r *RetentionReaper) ExportSweep(ctx context.Context) {
|
||||
r.sweep(ctx)
|
||||
|
||||
@@ -1,6 +1,10 @@
|
||||
package database
|
||||
|
||||
import "gorm.io/gorm"
|
||||
import (
|
||||
"time"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// DeliveryStatus represents the status of a delivery
|
||||
type DeliveryStatus string
|
||||
@@ -37,7 +41,7 @@ type Delivery struct {
|
||||
BaseModel
|
||||
|
||||
EventID string `gorm:"type:uuid;not null;index:idx_deliveries_event_id,priority:1" json:"eventId"`
|
||||
TargetID string `gorm:"type:uuid;not null" json:"targetId"`
|
||||
TargetID string `gorm:"type:uuid;not null;index:idx_deliveries_status,priority:4" json:"targetId"`
|
||||
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
|
||||
@@ -45,6 +49,13 @@ type Delivery struct {
|
||||
// gives.
|
||||
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
|
||||
Event Event `json:"event,omitzero"`
|
||||
Target Target `json:"target,omitzero"`
|
||||
|
||||
@@ -0,0 +1,91 @@
|
||||
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
|
||||
}
|
||||
@@ -2,7 +2,8 @@ package database
|
||||
|
||||
// Migrate runs database migrations for the main application database.
|
||||
// Only configuration-tier models are stored in the main database.
|
||||
// Event-tier models (Event, Delivery, DeliveryResult) live in
|
||||
// Event-tier models (Event, Delivery, DeliveryResult, EventTotals,
|
||||
// TargetTotals) live in
|
||||
// per-webhook dedicated databases managed by WebhookDBManager.
|
||||
func (d *Database) Migrate() error {
|
||||
return d.db.AutoMigrate(
|
||||
|
||||
@@ -18,6 +18,13 @@ import (
|
||||
// computation.
|
||||
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
|
||||
// RetentionReaper.
|
||||
type RetentionReaperParams struct {
|
||||
@@ -265,57 +272,97 @@ func retentionCutoff(
|
||||
), true
|
||||
}
|
||||
|
||||
// reapExpired hard-deletes, in foreign-key-safe order, the delivery
|
||||
// results, deliveries, and events associated with events older than
|
||||
// cutoff. Deletes are unscoped so rows are physically removed rather
|
||||
// than soft-deleted, reclaiming disk. It returns the number of events
|
||||
// reapExpired hard-deletes the events older than cutoff, with their
|
||||
// deliveries and delivery results, reapBatchSize events per
|
||||
// transaction until none is left. It returns the number of events
|
||||
// deleted.
|
||||
func reapExpired(db *gorm.DB, cutoff time.Time) (int64, error) {
|
||||
// Fresh subqueries are built per statement to avoid reusing a
|
||||
// mutated builder across executions.
|
||||
expiredEventIDs := func() *gorm.DB {
|
||||
return db.Model(&Event{}).
|
||||
Select("id").
|
||||
Where("created_at < ?", cutoff)
|
||||
}
|
||||
expiredDeliveryIDs := func() *gorm.DB {
|
||||
return db.Model(&Delivery{}).
|
||||
Select("id").
|
||||
Where("event_id IN (?)", expiredEventIDs())
|
||||
}
|
||||
var total int64
|
||||
|
||||
// 1. Delivery results whose delivery belongs to an expired event.
|
||||
res := db.Unscoped().
|
||||
Where("delivery_id IN (?)", expiredDeliveryIDs()).
|
||||
Delete(&DeliveryResult{})
|
||||
if res.Error != nil {
|
||||
return 0, fmt.Errorf(
|
||||
"deleting expired delivery results: %w",
|
||||
res.Error,
|
||||
)
|
||||
}
|
||||
for {
|
||||
var eventIDs []string
|
||||
|
||||
// 2. Deliveries belonging to an expired event.
|
||||
del := db.Unscoped().
|
||||
Where("event_id IN (?)", expiredEventIDs()).
|
||||
Delete(&Delivery{})
|
||||
if del.Error != nil {
|
||||
return 0, fmt.Errorf(
|
||||
"deleting expired deliveries: %w",
|
||||
del.Error,
|
||||
)
|
||||
}
|
||||
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)
|
||||
}
|
||||
|
||||
// 3. The expired events themselves.
|
||||
ev := db.Unscoped().
|
||||
Where("created_at < ?", cutoff).
|
||||
Delete(&Event{})
|
||||
if ev.Error != nil {
|
||||
return 0, fmt.Errorf(
|
||||
"deleting expired events: %w",
|
||||
ev.Error,
|
||||
)
|
||||
}
|
||||
if len(eventIDs) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
return ev.RowsAffected, 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").
|
||||
Where("event_id IN ?", eventIDs)).
|
||||
Delete(&DeliveryResult{}).Error
|
||||
if err != nil {
|
||||
return fmt.Errorf("deleting expired delivery results: %w", err)
|
||||
}
|
||||
|
||||
// 2. The events' deliveries, after counting them, and the failed
|
||||
// ones among them, per target. The status is tested in the select
|
||||
// list rather than the WHERE clause: there, SQLite would read every
|
||||
// failed delivery the webhook has through the status index,
|
||||
// instead of only these through the event_id index.
|
||||
var removed []TargetTotals
|
||||
|
||||
err = tx.Unscoped().Model(&Delivery{}).
|
||||
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().
|
||||
Where("event_id IN ?", eventIDs).
|
||||
Delete(&Delivery{}).Error
|
||||
if err != nil {
|
||||
return fmt.Errorf("deleting expired deliveries: %w", err)
|
||||
}
|
||||
|
||||
// 3. The events themselves.
|
||||
ev := tx.Unscoped().Where("id IN ?", eventIDs).Delete(&Event{})
|
||||
if ev.Error != nil {
|
||||
return fmt.Errorf("deleting expired events: %w", ev.Error)
|
||||
}
|
||||
|
||||
for i := range removed {
|
||||
err = AddTargetTotals(tx, removed[i])
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return AddEventTotals(tx, EventTotals{EventsRemoved: ev.RowsAffected})
|
||||
}
|
||||
|
||||
@@ -0,0 +1,215 @@
|
||||
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))
|
||||
}
|
||||
@@ -35,7 +35,8 @@ var errInvalidCachedDBType = errors.New(
|
||||
|
||||
// WebhookDBManager manages per-webhook SQLite database files
|
||||
// for event storage. Each webhook gets its own dedicated
|
||||
// database containing Events, Deliveries, and DeliveryResults.
|
||||
// database containing Events, Deliveries, DeliveryResults and the
|
||||
// running totals of them (EventTotals, TargetTotals).
|
||||
// Database connections are opened lazily and cached.
|
||||
type WebhookDBManager struct {
|
||||
dataDir string
|
||||
@@ -295,6 +296,7 @@ func (m *WebhookDBManager) openDB(
|
||||
// Run migrations for event-tier models only
|
||||
err = db.AutoMigrate(
|
||||
&Event{}, &Delivery{}, &DeliveryResult{},
|
||||
&EventTotals{}, &TargetTotals{},
|
||||
)
|
||||
if err != nil {
|
||||
_ = sqlDB.Close()
|
||||
@@ -305,6 +307,18 @@ 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(
|
||||
"opened per-webhook database",
|
||||
"webhook_id", webhookID,
|
||||
|
||||
@@ -0,0 +1,116 @@
|
||||
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))
|
||||
}
|
||||
@@ -148,7 +148,6 @@ type EngineParams struct {
|
||||
DBManager *database.WebhookDBManager
|
||||
Logger *logger.Logger
|
||||
SSRFGuard *Guard
|
||||
Metrics *metrics.Set
|
||||
}
|
||||
|
||||
// Engine processes queued deliveries in the background
|
||||
@@ -168,10 +167,10 @@ type Engine struct {
|
||||
retryCh chan Task
|
||||
workers int
|
||||
|
||||
// mtr is the delivery metric set. Production wires the one
|
||||
// registered on the registry /metrics serves; a test can
|
||||
// substitute a set registered on a registry it holds, so it can
|
||||
// gather what its own deliveries recorded.
|
||||
// mtr is the delivery metric set. Production wires the
|
||||
// process-wide one; a test can substitute a set registered on
|
||||
// a private registry so its assertions are not disturbed by
|
||||
// deliveries other tests are making at the same time.
|
||||
mtr *metrics.Set
|
||||
|
||||
// targets maps each target type to its implementation.
|
||||
@@ -205,7 +204,7 @@ func New(
|
||||
deliveryCh: make(chan Task, deliveryChannelSize),
|
||||
retryCh: make(chan Task, retryChannelSize),
|
||||
workers: defaultWorkers,
|
||||
mtr: params.Metrics,
|
||||
mtr: metrics.Default(),
|
||||
}
|
||||
|
||||
e.initTargets(&http.Client{
|
||||
@@ -1555,8 +1554,9 @@ func (e *Engine) updateDeliveryStatus(
|
||||
targetType database.TargetType,
|
||||
status database.DeliveryStatus,
|
||||
) error {
|
||||
err := webhookDB.Model(d).
|
||||
Update("status", status).Error
|
||||
err := webhookDB.Transaction(func(tx *gorm.DB) error {
|
||||
return writeDeliveryStatus(tx, d, status)
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf(
|
||||
"updating delivery %s to status %s: %w",
|
||||
@@ -1575,6 +1575,36 @@ func (e *Engine) updateDeliveryStatus(
|
||||
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
|
||||
// failed write through bookkeepingFailed, which leaves the row
|
||||
// recoverable. It exists so the target call sites read as one
|
||||
|
||||
@@ -57,7 +57,10 @@ func testWebhookDB(t *testing.T) *gorm.DB {
|
||||
&database.Event{},
|
||||
&database.Delivery{},
|
||||
&database.DeliveryResult{},
|
||||
&database.EventTotals{},
|
||||
&database.TargetTotals{},
|
||||
))
|
||||
require.NoError(t, db.Create(&database.EventTotals{}).Error)
|
||||
|
||||
return db
|
||||
}
|
||||
|
||||
@@ -9,7 +9,6 @@ import (
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"go.uber.org/fx"
|
||||
"gorm.io/gorm"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
@@ -151,6 +150,16 @@ 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.
|
||||
func (e *Engine) ExportProcessNewTask(
|
||||
ctx context.Context, task *Task,
|
||||
@@ -390,7 +399,7 @@ func NewTestEngine(
|
||||
deliveryCh: make(chan Task, deliveryChannelSize),
|
||||
retryCh: make(chan Task, retryChannelSize),
|
||||
workers: workers,
|
||||
mtr: metrics.New(prometheus.NewRegistry()),
|
||||
mtr: metrics.Default(),
|
||||
}
|
||||
e.initTargets(client)
|
||||
|
||||
@@ -405,7 +414,7 @@ func NewTestEngineSmallRetry(
|
||||
e := &Engine{
|
||||
log: log,
|
||||
retryCh: make(chan Task, 1),
|
||||
mtr: metrics.New(prometheus.NewRegistry()),
|
||||
mtr: metrics.Default(),
|
||||
}
|
||||
e.initTargets(nil)
|
||||
|
||||
@@ -428,7 +437,7 @@ func NewTestEngineWithDB(
|
||||
deliveryCh: make(chan Task, deliveryChannelSize),
|
||||
retryCh: make(chan Task, retryChannelSize),
|
||||
workers: workers,
|
||||
mtr: metrics.New(prometheus.NewRegistry()),
|
||||
mtr: metrics.Default(),
|
||||
}
|
||||
e.initTargets(client)
|
||||
|
||||
@@ -436,7 +445,8 @@ func NewTestEngineWithDB(
|
||||
}
|
||||
|
||||
// ExportSetMetrics substitutes the engine's metric set, so a test can
|
||||
// assert on collectors registered on a registry it holds.
|
||||
// assert on collectors registered on a private registry instead of
|
||||
// the process-wide ones every other test is also moving.
|
||||
func (e *Engine) ExportSetMetrics(mtr *metrics.Set) {
|
||||
e.mtr = mtr
|
||||
}
|
||||
|
||||
@@ -35,8 +35,9 @@ const (
|
||||
)
|
||||
|
||||
// mIsolate gives the setup's engine a metric set registered on a
|
||||
// registry this test holds, so its exact assertions can gather from
|
||||
// it.
|
||||
// private registry. The process-wide collectors are moved by every
|
||||
// other delivery test running in parallel, so exact assertions are
|
||||
// only possible against a registry this test owns.
|
||||
func mIsolate(
|
||||
t *testing.T, s iSetup,
|
||||
) *prometheus.Registry {
|
||||
|
||||
@@ -299,8 +299,9 @@ func countInFlightDeliveries(
|
||||
return count, err
|
||||
}
|
||||
|
||||
// createReplayDelivery writes the new pending delivery row and returns
|
||||
// the task that carries it to the delivery engine.
|
||||
// createReplayDelivery writes the new pending delivery row, adds it to
|
||||
// its target's totals in the same transaction, and returns the task
|
||||
// that carries it to the delivery engine.
|
||||
//
|
||||
// The row is written with associations omitted, and neither Event nor
|
||||
// Target is populated on it: GORM's SaveBeforeAssociations would
|
||||
@@ -319,7 +320,16 @@ func createReplayDelivery(
|
||||
Status: database.DeliveryStatusPending,
|
||||
}
|
||||
|
||||
err := webhookDB.Omit(clause.Associations).Create(dlv).Error
|
||||
err := webhookDB.Transaction(func(tx *gorm.DB) 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 {
|
||||
return delivery.Task{}, err
|
||||
}
|
||||
|
||||
@@ -4,7 +4,9 @@ import (
|
||||
"html/template"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"gorm.io/gorm"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
)
|
||||
|
||||
@@ -69,6 +71,29 @@ func (s *Handlers) LoadEventLogViewsForTest(
|
||||
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
|
||||
// the handlers_test package can drive the render path with a
|
||||
// template of its own.
|
||||
|
||||
@@ -12,7 +12,6 @@ import (
|
||||
"net/http"
|
||||
"sync/atomic"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"go.uber.org/fx"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
"sneak.berlin/go/webhooker/internal/delivery"
|
||||
@@ -62,8 +61,6 @@ type HandlersParams struct {
|
||||
Notifier delivery.Notifier
|
||||
Evictor delivery.WebhookEvictor
|
||||
SSRFGuard *delivery.Guard
|
||||
Metrics *metrics.Set
|
||||
Registry *prometheus.Registry
|
||||
}
|
||||
|
||||
// Handlers provides HTTP handler methods for all application
|
||||
@@ -94,18 +91,22 @@ type Handlers struct {
|
||||
|
||||
// parsePageTemplate parses a page-specific template set from the
|
||||
// embedded FS. Each page template is combined with the shared
|
||||
// base, htmlheader, and navbar templates. The page file must be
|
||||
// listed first so that its root action ({{template "base" .}})
|
||||
// becomes the template set's entry point.
|
||||
func parsePageTemplate(pageFile string) *template.Template {
|
||||
// base, htmlheader, and navbar templates, and with any further files
|
||||
// the page includes. The page file must be listed first so that its
|
||||
// root action ({{template "base" .}}) becomes the template set's entry
|
||||
// point.
|
||||
func parsePageTemplate(
|
||||
pageFile string, included ...string,
|
||||
) *template.Template {
|
||||
files := append([]string{
|
||||
pageFile,
|
||||
"base.html",
|
||||
"htmlheader.html",
|
||||
"navbar.html",
|
||||
}, included...)
|
||||
|
||||
return template.Must(
|
||||
template.ParseFS(
|
||||
templates.Templates,
|
||||
pageFile,
|
||||
"base.html",
|
||||
"htmlheader.html",
|
||||
"navbar.html",
|
||||
),
|
||||
template.ParseFS(templates.Templates, files...),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -125,7 +126,7 @@ func New(
|
||||
s.mw = params.Middleware
|
||||
s.notifier = params.Notifier
|
||||
s.evictor = params.Evictor
|
||||
s.mtr = params.Metrics
|
||||
s.mtr = metrics.Default()
|
||||
s.ssrf = params.SSRFGuard
|
||||
|
||||
// Parse all page templates once at startup
|
||||
@@ -134,7 +135,7 @@ func New(
|
||||
"profile.html": parsePageTemplate("profile.html"),
|
||||
"sources_list.html": parsePageTemplate("sources_list.html"),
|
||||
"sources_new.html": parsePageTemplate("sources_new.html"),
|
||||
"source_detail.html": parsePageTemplate("source_detail.html"),
|
||||
"source_detail.html": parsePageTemplate("source_detail.html", "webhook_stats.html"),
|
||||
"source_edit.html": parsePageTemplate("source_edit.html"),
|
||||
"source_logs.html": parsePageTemplate("source_logs.html"),
|
||||
"target_edit.html": parsePageTemplate("target_edit.html"),
|
||||
|
||||
@@ -20,7 +20,6 @@ import (
|
||||
"sneak.berlin/go/webhooker/internal/handlers"
|
||||
"sneak.berlin/go/webhooker/internal/healthcheck"
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
"sneak.berlin/go/webhooker/internal/metrics"
|
||||
"sneak.berlin/go/webhooker/internal/middleware"
|
||||
"sneak.berlin/go/webhooker/internal/session"
|
||||
)
|
||||
@@ -110,8 +109,6 @@ func newTestApp(
|
||||
func(r *recordingEvictor) delivery.WebhookEvictor {
|
||||
return r
|
||||
},
|
||||
metrics.NewRegistry,
|
||||
metrics.New,
|
||||
middleware.New,
|
||||
delivery.NewGuard,
|
||||
handlers.New,
|
||||
|
||||
@@ -1,21 +0,0 @@
|
||||
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
|
||||
}
|
||||
@@ -450,6 +450,7 @@ func (h *Handlers) renderSourceDetail(
|
||||
"Targets": delivery.NewTargetViews(targets),
|
||||
"Events": events,
|
||||
"BaseURL": baseURL,
|
||||
"Stats": h.loadWebhookStats(webhook.ID, entrypoints, targets),
|
||||
}
|
||||
|
||||
h.renderTemplate(w, r, "source_detail.html", data)
|
||||
|
||||
@@ -252,11 +252,12 @@ func requestEventSource(
|
||||
}
|
||||
}
|
||||
|
||||
// createAndFanOut writes the event and one pending delivery per target
|
||||
// in a single transaction, then hands the tasks to the delivery
|
||||
// engine. It is the only path by which an event and its deliveries are
|
||||
// created, so a resubmitted event is retried, SSRF-guarded and
|
||||
// circuit-broken exactly as a received one is.
|
||||
// createAndFanOut writes the event and one pending delivery per target,
|
||||
// and adds them to the webhook's running totals, in a single
|
||||
// transaction, then hands the tasks to the delivery engine. It is the
|
||||
// only path by which an event and its deliveries are created, so a
|
||||
// resubmitted event is retried, SSRF-guarded and circuit-broken
|
||||
// exactly as a received one is.
|
||||
//
|
||||
// The tasks are returned as well as queued, so a caller can report how
|
||||
// many targets the event went to.
|
||||
@@ -296,6 +297,13 @@ func (h *Handlers) createAndFanOut(
|
||||
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
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf(
|
||||
@@ -354,8 +362,9 @@ func (h *Handlers) finishWebhookResponse(
|
||||
}
|
||||
|
||||
// buildDeliveryTasks creates one pending delivery per target in the
|
||||
// transaction and returns the tasks for the delivery engine. The
|
||||
// caller owns the transaction and rolls it back on error.
|
||||
// transaction, adds each to its target's totals, and returns the tasks
|
||||
// for the delivery engine. The caller owns the transaction and rolls
|
||||
// it back on error.
|
||||
func buildDeliveryTasks(
|
||||
tx *gorm.DB,
|
||||
event *database.Event,
|
||||
@@ -379,6 +388,13 @@ 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{
|
||||
DeliveryID: dlv.ID,
|
||||
EventID: event.ID,
|
||||
|
||||
@@ -0,0 +1,271 @@
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,432 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
+20
-27
@@ -3,18 +3,17 @@
|
||||
// deliveries are attempted, how they end, how long they take, how
|
||||
// deep the queues are, and how many circuit breakers are open.
|
||||
//
|
||||
// It also builds the registry the authenticated /metrics route
|
||||
// serves. In production, these collectors, the inbound HTTP metrics
|
||||
// recorded in internal/middleware, and the Go runtime and process
|
||||
// collectors all register on that one registry, never on Prometheus's
|
||||
// global default.
|
||||
// The inbound HTTP metrics come from the go-http-metrics recorder in
|
||||
// internal/middleware and land on prometheus.DefaultRegisterer. These
|
||||
// collectors register there too, so both surfaces are gathered by the
|
||||
// one promhttp handler mounted on the authenticated /metrics route.
|
||||
package metrics
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"github.com/prometheus/client_golang/prometheus/collectors"
|
||||
"github.com/prometheus/client_golang/prometheus/promauto"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
)
|
||||
@@ -58,31 +57,25 @@ var knownTargetTypes = []database.TargetType{
|
||||
database.TargetTypeSlack,
|
||||
}
|
||||
|
||||
// NewRegistry returns the registry /metrics serves, carrying the Go
|
||||
// runtime and process collectors that Prometheus's global default
|
||||
// registry carries, so the go_* and process_* series stay in the
|
||||
// scrape.
|
||||
// defaultSet is the process-wide metric set, registered on the same
|
||||
// registry the HTTP middleware and the /metrics handler already use.
|
||||
// It is built on first use rather than in an init so that a test
|
||||
// binary that never touches metrics never registers them.
|
||||
//
|
||||
// A registry of its own, rather than the global default, is what lets
|
||||
// two dependency graphs in one process — two tests, say — each
|
||||
// register their collectors without the second registration
|
||||
// panicking.
|
||||
func NewRegistry() *prometheus.Registry {
|
||||
reg := prometheus.NewRegistry()
|
||||
reg.MustRegister(
|
||||
collectors.NewGoCollector(),
|
||||
collectors.NewProcessCollector(
|
||||
collectors.ProcessCollectorOpts{},
|
||||
),
|
||||
)
|
||||
//nolint:gochecknoglobals // one process-wide registration, by design
|
||||
var defaultSet = sync.OnceValue(func() *Set {
|
||||
return New(prometheus.DefaultRegisterer)
|
||||
})
|
||||
|
||||
return reg
|
||||
// Default returns the process-wide metric set.
|
||||
func Default() *Set {
|
||||
return defaultSet()
|
||||
}
|
||||
|
||||
// Set is one registered group of webhooker's delivery collectors.
|
||||
// Production builds one on the registry /metrics serves; tests build
|
||||
// their own against a private registry so assertions are not
|
||||
// disturbed by deliveries other tests are making concurrently.
|
||||
// Production uses the single Default set; tests build their own
|
||||
// against a private registry so assertions are not disturbed by
|
||||
// deliveries other tests are making concurrently.
|
||||
type Set struct {
|
||||
eventsReceived prometheus.Counter
|
||||
deliveryAttempts *prometheus.CounterVec
|
||||
@@ -100,7 +93,7 @@ type Set struct {
|
||||
// New registers a full set of delivery collectors on reg and returns
|
||||
// it. It panics if reg already holds them, which is the intended
|
||||
// behaviour for a duplicate registration.
|
||||
func New(reg *prometheus.Registry) *Set {
|
||||
func New(reg prometheus.Registerer) *Set {
|
||||
factory := promauto.With(reg)
|
||||
|
||||
s := &Set{
|
||||
|
||||
@@ -10,7 +10,8 @@ import (
|
||||
|
||||
// MetricsMiddlewareForTest builds the metrics recording middleware
|
||||
// against a caller-supplied recorder, so a test can gather from its
|
||||
// own Prometheus registry without building a whole Middleware.
|
||||
// own Prometheus registry rather than the process-wide default one
|
||||
// that Middleware.Metrics uses.
|
||||
func MetricsMiddlewareForTest(
|
||||
rec httpmetrics.Recorder,
|
||||
) func(http.Handler) http.Handler {
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
|
||||
"github.com/go-chi/chi"
|
||||
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"
|
||||
"github.com/slok/go-http-metrics/middleware/std"
|
||||
)
|
||||
@@ -150,17 +151,17 @@ func (r boundedLabelRecorder) AddInflightRequests(
|
||||
|
||||
var _ httpmetrics.Recorder = boundedLabelRecorder{}
|
||||
|
||||
// Metrics returns middleware that records Prometheus HTTP metrics
|
||||
// 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.
|
||||
// Metrics returns middleware that records Prometheus HTTP metrics on
|
||||
// the default registry, which is the one the /metrics route gathers.
|
||||
func (s *Middleware) Metrics() func(http.Handler) http.Handler {
|
||||
return metricsMiddleware(s.metricsRecorder)
|
||||
return metricsMiddleware(
|
||||
prommetrics.NewRecorder(prommetrics.Config{}),
|
||||
)
|
||||
}
|
||||
|
||||
// metricsMiddleware builds the recording middleware against a given
|
||||
// recorder, so tests can gather from a registry of their own.
|
||||
// recorder, so tests can gather from a registry of their own instead
|
||||
// of the process-wide default.
|
||||
func metricsMiddleware(
|
||||
rec httpmetrics.Recorder,
|
||||
) func(http.Handler) http.Handler {
|
||||
|
||||
@@ -57,8 +57,9 @@ const (
|
||||
// Server.setupWebhookRoutes inside it. That ordering is the whole
|
||||
// defect, so a test that flattens it would prove nothing.
|
||||
//
|
||||
// The recorder writes to a registry of the test's own, so each test
|
||||
// observes only its own traffic.
|
||||
// The recorder writes to a registry of the test's own rather than the
|
||||
// process-wide default one, so each test observes only its own
|
||||
// traffic.
|
||||
func metricsTestRouter(
|
||||
t *testing.T,
|
||||
receiverLimit int,
|
||||
@@ -454,29 +455,3 @@ func TestMetrics_StatusAndSizeStillRecorded(t *testing.T) {
|
||||
"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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,9 +13,6 @@ import (
|
||||
"github.com/go-chi/chi"
|
||||
"github.com/go-chi/chi/middleware"
|
||||
"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"
|
||||
"sneak.berlin/go/webhooker/internal/config"
|
||||
"sneak.berlin/go/webhooker/internal/globals"
|
||||
@@ -151,11 +148,10 @@ const (
|
||||
type MiddlewareParams struct {
|
||||
fx.In
|
||||
|
||||
Logger *logger.Logger
|
||||
Globals *globals.Globals
|
||||
Config *config.Config
|
||||
Session *session.Session
|
||||
Registry *prometheus.Registry
|
||||
Logger *logger.Logger
|
||||
Globals *globals.Globals
|
||||
Config *config.Config
|
||||
Session *session.Session
|
||||
}
|
||||
|
||||
// Middleware provides HTTP middleware for logging, CORS, auth, and
|
||||
@@ -165,14 +161,6 @@ type Middleware struct {
|
||||
params *MiddlewareParams
|
||||
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
|
||||
// concurrent password hashing. It is built on first use so that
|
||||
// every construction path gets one; see guard().
|
||||
@@ -191,9 +179,6 @@ func New(
|
||||
s.params = ¶ms
|
||||
s.log = params.Logger.Get()
|
||||
s.session = params.Session
|
||||
s.metricsRecorder = prommetrics.NewRecorder(
|
||||
prommetrics.Config{Registry: params.Registry},
|
||||
)
|
||||
|
||||
return s, nil
|
||||
}
|
||||
|
||||
@@ -3,17 +3,12 @@ package middleware
|
||||
import (
|
||||
"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/session"
|
||||
)
|
||||
|
||||
// NewForTest creates a Middleware with the minimum dependencies
|
||||
// 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(
|
||||
log *slog.Logger,
|
||||
cfg *config.Config,
|
||||
@@ -25,8 +20,5 @@ func NewForTest(
|
||||
Config: cfg,
|
||||
},
|
||||
session: sess,
|
||||
metricsRecorder: prommetrics.NewRecorder(
|
||||
prommetrics.Config{Registry: prometheus.NewRegistry()},
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -24,7 +24,6 @@ import (
|
||||
"sneak.berlin/go/webhooker/internal/handlers"
|
||||
"sneak.berlin/go/webhooker/internal/healthcheck"
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
"sneak.berlin/go/webhooker/internal/metrics"
|
||||
"sneak.berlin/go/webhooker/internal/middleware"
|
||||
"sneak.berlin/go/webhooker/internal/resetpw"
|
||||
"sneak.berlin/go/webhooker/internal/session"
|
||||
@@ -164,8 +163,6 @@ func newServerApp(
|
||||
session.New,
|
||||
func() delivery.Notifier { return &noopNotifier{} },
|
||||
func() delivery.WebhookEvictor { return &noopEvictor{} },
|
||||
metrics.NewRegistry,
|
||||
metrics.New,
|
||||
middleware.New,
|
||||
delivery.NewGuard,
|
||||
handlers.New,
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
sentryhttp "github.com/getsentry/sentry-go/http"
|
||||
"github.com/go-chi/chi"
|
||||
"github.com/go-chi/chi/middleware"
|
||||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||
"sneak.berlin/go/webhooker/static"
|
||||
)
|
||||
|
||||
@@ -129,7 +130,12 @@ func (s *Server) setupRoutes() {
|
||||
if s.params.Config.MetricsAuthEnabled() {
|
||||
s.router.Group(func(r chi.Router) {
|
||||
r.Use(s.mw.MetricsAuth())
|
||||
r.Get("/metrics", s.h.HandleMetrics())
|
||||
r.Get(
|
||||
"/metrics",
|
||||
http.HandlerFunc(
|
||||
promhttp.Handler().ServeHTTP,
|
||||
),
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -24,7 +24,6 @@ import (
|
||||
"sneak.berlin/go/webhooker/internal/handlers"
|
||||
"sneak.berlin/go/webhooker/internal/healthcheck"
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
"sneak.berlin/go/webhooker/internal/metrics"
|
||||
"sneak.berlin/go/webhooker/internal/middleware"
|
||||
"sneak.berlin/go/webhooker/internal/server"
|
||||
"sneak.berlin/go/webhooker/internal/session"
|
||||
@@ -114,8 +113,6 @@ func newTestEnvWithConfig(
|
||||
session.New,
|
||||
func() delivery.Notifier { return &noopNotifier{} },
|
||||
func() delivery.WebhookEvictor { return &noopEvictor{} },
|
||||
metrics.NewRegistry,
|
||||
metrics.New,
|
||||
middleware.New,
|
||||
delivery.NewGuard,
|
||||
handlers.New,
|
||||
@@ -1030,46 +1027,3 @@ 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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -24,6 +24,8 @@
|
||||
</div>
|
||||
</div>
|
||||
|
||||
{{template "webhook_stats" .}}
|
||||
|
||||
<div class="grid grid-cols-1 lg:grid-cols-2 gap-6">
|
||||
<!-- Entrypoints -->
|
||||
<div class="card">
|
||||
|
||||
@@ -0,0 +1,93 @@
|
||||
{{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}}
|
||||
Reference in New Issue
Block a user