diff --git a/README.md b/README.md index 8173e92..17b712f 100644 --- a/README.md +++ b/README.md @@ -1075,7 +1075,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 @@ -1392,7 +1392,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). @@ -1421,6 +1421,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) │ +│ └──────────────┘ │ └─────────────────────────────────────────────────────────────┘ ``` @@ -1674,6 +1681,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. @@ -1741,33 +1749,66 @@ 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 a range. 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 the +statistics compare `created_at` with a range (`>=`) and SQLite narrows by a +range 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 | | ------------ | --------- | ----------- | @@ -1809,6 +1850,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 @@ -2788,6 +2831,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) diff --git a/internal/database/event_tier_indexes_test.go b/internal/database/event_tier_indexes_test.go index d25dca5..6eb4436 100644 --- a/internal/database/event_tier_indexes_test.go +++ b/internal/database/event_tier_indexes_test.go @@ -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= ?", + []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) } diff --git a/internal/database/export_test.go b/internal/database/export_test.go index 5c10271..26283e3 100644 --- a/internal/database/export_test.go +++ b/internal/database/export_test.go @@ -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) diff --git a/internal/database/model_delivery.go b/internal/database/model_delivery.go index 0ebd2bf..4424756 100644 --- a/internal/database/model_delivery.go +++ b/internal/database/model_delivery.go @@ -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"` diff --git a/internal/database/model_totals.go b/internal/database/model_totals.go new file mode 100644 index 0000000..e992597 --- /dev/null +++ b/internal/database/model_totals.go @@ -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 +} diff --git a/internal/database/models.go b/internal/database/models.go index 0857a74..7cbf924 100644 --- a/internal/database/models.go +++ b/internal/database/models.go @@ -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( diff --git a/internal/database/retention.go b/internal/database/retention.go index 13051af..3490055 100644 --- a/internal/database/retention.go +++ b/internal/database/retention.go @@ -18,6 +18,19 @@ 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 + +// reapBatchPause is how long retention waits after one batch before +// starting the next. A writer waiting for the write lock checks for it +// again after at most 100 ms, so a longer pause lets it in between two +// batches instead of only after the whole prune. +const reapBatchPause = 200 * time.Millisecond + // RetentionReaperParams holds the fx dependencies for the // RetentionReaper. type RetentionReaperParams struct { @@ -265,57 +278,99 @@ 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 -// deleted. +// reapExpired hard-deletes the events older than cutoff, with their +// deliveries and delivery results, reapBatchSize events per +// transaction with reapBatchPause between transactions, 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 + } + + time.Sleep(reapBatchPause) + } +} + +// 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}) } diff --git a/internal/database/totals_test.go b/internal/database/totals_test.go new file mode 100644 index 0000000..3139974 --- /dev/null +++ b/internal/database/totals_test.go @@ -0,0 +1,305 @@ +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)) +} + +// seedExpiredEvents stores count events created at the given time, +// each with a delivered delivery to one target and a failed delivery +// to the other, and one attempt for each delivery. +func seedExpiredEvents( + t *testing.T, + db *gorm.DB, + webhookID string, + count int, + createdAt time.Time, + delivered, failed string, +) { + t.Helper() + + events := make([]database.Event, count) + deliveries := make([]database.Delivery, 0, 2*count) + + 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 = createdAt + + 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) +} + +// 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) + + expired := database.ExportReapBatchSize + 1 + delivered, failed := uuid.New().String(), uuid.New().String() + seedExpiredEvents(t, db, webhookID, expired, + time.Now().Add(-40*24*time.Hour), delivered, failed) + + // 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)) +} + +// TestRetentionReaper_WriteDuringPruneSucceeds verifies that a prune +// of several batches lets other writers in between its batches: an +// event stored once the first batch is deleted is stored while expired +// events are still left, not only after the prune has finished. +func TestRetentionReaper_WriteDuringPruneSucceeds(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) + + // Three batches of expired events, with nothing else stored: only + // the number of batches matters here. + expired := make([]database.Event, 3*database.ExportReapBatchSize) + for i := range expired { + expired[i] = database.Event{ + WebhookID: webhookID, + EntrypointID: uuid.New().String(), + Method: http.MethodPost, + } + expired[i].CreatedAt = time.Now().Add(-40 * 24 * time.Hour) + } + + require.NoError(t, db.CreateInBatches(expired, 500).Error) + + cutoff := time.Now().Add(-30 * 24 * time.Hour) + countExpired := func() int64 { + var count int64 + + require.NoError(t, db.Model(&database.Event{}). + Where("created_at < ?", cutoff). + Count(&count).Error) + + return count + } + + pruned := make(chan struct{}) + + go func() { + defer close(pruned) + + env.reaper.ExportSweep(context.Background()) + }() + + t.Cleanup(func() { <-pruned }) + + // Every stored event is expired until the write below. + require.Eventually(t, func() bool { + var count int64 + + err := db.Model(&database.Event{}).Count(&count).Error + + return err == nil && count < int64(len(expired)) + }, 10*time.Second, 10*time.Millisecond) + + event := &database.Event{ + WebhookID: webhookID, + EntrypointID: uuid.New().String(), + Method: http.MethodPost, + } + require.NoError(t, db.Create(event).Error) + + assert.Positive(t, countExpired(), + "the event was stored only after the whole prune") + + <-pruned + + assert.Zero(t, countExpired()) + + var stored database.Event + + require.NoError(t, db.First(&stored, "id = ?", event.ID).Error) +} diff --git a/internal/database/webhook_db_manager.go b/internal/database/webhook_db_manager.go index a628b56..380dbe5 100644 --- a/internal/database/webhook_db_manager.go +++ b/internal/database/webhook_db_manager.go @@ -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, diff --git a/internal/delivery/delivery_totals_test.go b/internal/delivery/delivery_totals_test.go new file mode 100644 index 0000000..bfff90b --- /dev/null +++ b/internal/delivery/delivery_totals_test.go @@ -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)) +} diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index 4786f08..fae2a83 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -531,6 +531,11 @@ func (e *Engine) processRetryTask( return } + // Set before anything below can fail the delivery: the failure is + // added to this target's totals. + d.EventID = task.EventID + d.TargetID = task.TargetID + if d.Status != database.DeliveryStatusRetrying { e.log.Debug( "skipping retry for delivery "+ @@ -562,8 +567,6 @@ func (e *Engine) processRetryTask( } target := buildTargetFromTask(task) - d.EventID = task.EventID - d.TargetID = task.TargetID d.Event = event d.Target = target @@ -1554,8 +1557,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", @@ -1574,6 +1578,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 diff --git a/internal/delivery/engine_test.go b/internal/delivery/engine_test.go index 13d1625..abe668b 100644 --- a/internal/delivery/engine_test.go +++ b/internal/delivery/engine_test.go @@ -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 } diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index 9dd2531..5310291 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -150,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, diff --git a/internal/delivery/terminal_state_test.go b/internal/delivery/terminal_state_test.go index 8eda81a..2cf9f4d 100644 --- a/internal/delivery/terminal_state_test.go +++ b/internal/delivery/terminal_state_test.go @@ -418,6 +418,38 @@ func TestProcessRetryTask_TargetDeleted_MakesNoAttempt( assert.Zero(t, s.Engine.ExportInflightHeld()) } +// TestProcessRetryTask_TargetDeleted_CountsFailureOnTarget verifies +// that the failure of a retry abandoned because its target is gone is +// added to that target's own totals, not to a row with no target. +func TestProcessRetryTask_TargetDeleted_CountsFailureOnTarget( + t *testing.T, +) { + t.Parallel() + + s := newISetup(t) + + var hits atomic.Int64 + + task, targetID := tRetryChainSetup( + t, s, "gone-counted", &hits, + ) + + require.NoError(t, s.MainDB.Delete( + &database.Target{}, "id = ?", targetID, + ).Error) + + s.Engine.ExportProcessRetryTask( + context.Background(), &task, + ) + + var rows []database.TargetTotals + + require.NoError(t, s.WebhookDB.Find(&rows).Error) + assert.Equal(t, []database.TargetTotals{ + {TargetID: targetID, Failed: 1}, + }, rows) +} + // TestProcessRetryTask_TargetPresent_StillDelivers is the guard's // mutation check: a liveness check that refused every retry would pass // the test above and break every retry there is. diff --git a/internal/handlers/delivery_replay.go b/internal/handlers/delivery_replay.go index ee9425b..594087d 100644 --- a/internal/handlers/delivery_replay.go +++ b/internal/handlers/delivery_replay.go @@ -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 } diff --git a/internal/handlers/export_test.go b/internal/handlers/export_test.go index 4612e94..4c1ec27 100644 --- a/internal/handlers/export_test.go +++ b/internal/handlers/export_test.go @@ -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. diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index ad3bcd5..f64d51e 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -91,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...), ) } @@ -131,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"), diff --git a/internal/handlers/source_management.go b/internal/handlers/source_management.go index 9147c1f..e76769f 100644 --- a/internal/handlers/source_management.go +++ b/internal/handlers/source_management.go @@ -457,6 +457,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) diff --git a/internal/handlers/webhook.go b/internal/handlers/webhook.go index 1889ec5..5a2540a 100644 --- a/internal/handlers/webhook.go +++ b/internal/handlers/webhook.go @@ -253,11 +253,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. @@ -297,6 +298,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( @@ -355,8 +363,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, @@ -380,6 +389,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, diff --git a/internal/handlers/webhook_stats.go b/internal/handlers/webhook_stats.go new file mode 100644 index 0000000..b1c750d --- /dev/null +++ b/internal/handlers/webhook_stats.go @@ -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 +} diff --git a/internal/handlers/webhook_stats_test.go b/internal/handlers/webhook_stats_test.go new file mode 100644 index 0000000..2895e96 --- /dev/null +++ b/internal/handlers/webhook_stats_test.go @@ -0,0 +1,474 @@ +package handlers_test + +import ( + "net/http" + "regexp" + "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 text of the statistics pane in a rendered +// webhook page, everything from its heading to the next heading on the +// page, with the markup taken out and each run of space made one +// space. A table then reads header by header and row by row, each +// row's label followed by its figures in column order. +func statsPane(t *testing.T, page string) string { + t.Helper() + + _, pane, found := strings.Cut(page, ">Statistics") + require.True(t, found, "the page has no statistics pane") + + pane, _, _ = strings.Cut(pane, "]*>`).ReplaceAllString(pane, " ") + + return strings.Join(strings.Fields(pane), " ") +} + +// assertStatsTargets checks, for the history seedStatsHistory builds, +// each target's totals and its deliveries finished in the last 24 +// hours. The first target has three deliveries and the replay, the +// second three; the inactive target has none and so no row. +func assertStatsTargets(t *testing.T, hist statsHistory) { + t.Helper() + + first, second := hist.first, hist.second + + 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) +} + +// assertStatsPaneAfterPrune checks the rendered statistics pane for the +// history seedStatsHistory builds, once retention has removed the +// oldest event: each figure after its label, in its column. +func assertStatsPaneAfterPrune( + t *testing.T, + h *handlers.Handlers, + sess *session.Session, + hist statsHistory, +) { + t.Helper() + + pane := statsPane(t, renderSourceDetailPage(t, h, sess, hist.webhook.ID)) + lastEvent := hist.newest.CreatedAt.Format("2006-01-02 15:04:05 UTC") + + assert.Contains(t, pane, "Entrypoints 2 (1 active) "+ + "Targets 3 (2 active) "+ + "Deliveries in progress 1 "+ + "Last event "+lastEvent+" "+ + "Retention 1 day") + assert.Contains(t, pane, "Lifetime Within retention "+ + "Events 3 2 "+ + "Deliveries 7 4 "+ + "Failures 3 2") + assert.Contains(t, pane, "Last 10 minutes Last 24 hours "+ + "Events 1 2 "+ + "Failures 1 2 "+ + "Failure percentage 50.0% 66.7%") +} + +// 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()) + + assertStatsTargets(t, hist) + + // 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)) + + assertStatsPaneAfterPrune(t, h, sess, hist) +} + +// 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) + } +} + +// 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()) + + pane := statsPane(t, renderSourceDetailPage(t, h, sess, wh.ID)) + assert.Contains(t, pane, "Last event none") + assert.Contains(t, pane, "Failure percentage — —") + 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) + } +} diff --git a/templates/source_detail.html b/templates/source_detail.html index 1b60eb8..d5ae6e1 100644 --- a/templates/source_detail.html +++ b/templates/source_detail.html @@ -24,6 +24,8 @@ + {{template "webhook_stats" .}} +
diff --git a/templates/webhook_stats.html b/templates/webhook_stats.html new file mode 100644 index 0000000..e738d09 --- /dev/null +++ b/templates/webhook_stats.html @@ -0,0 +1,93 @@ +{{define "webhook_stats"}} + +
+
+

Statistics

+
+ {{with .Stats}} +
+
+ Entrypoints + {{.Entrypoints}} + ({{.ActiveEntrypoints}} active) +
+
+ Targets + {{.Targets}} + ({{.ActiveTargets}} active) +
+
+ Deliveries in progress + {{.InProgress}} +
+
+ Last event + {{with .LastEventAt}}{{.Format "2006-01-02 15:04:05 UTC"}}{{else}}none{{end}} +
+
+ Retention + {{$.Webhook.RetentionLabel}} +
+
+
+ + + + + + + + + + + + + + + + + + + + + + + + + +
LifetimeWithin retention
Events{{.Lifetime.Events}}{{.WithinRetention.Events}}
Deliveries{{.Lifetime.Deliveries}}{{.WithinRetention.Deliveries}}
Failures{{.Lifetime.Failures}}{{.WithinRetention.Failures}}
+
+ + + + + + + + + + + + + + + + + + + + + + + + + +
Last 10 minutesLast 24 hours
Events{{.Last10Minutes.Events}}{{.Last24Hours.Events}}
Failures{{.Last10Minutes.Failed}}{{.Last24Hours.Failed}}
Failure percentage{{.Last10Minutes.FailurePercent}}{{.Last24Hours.FailurePercent}}
+

Failure percentage is the failed deliveries out of all deliveries that finished in the window. Deliveries still pending or retrying are not counted.

+
+
+ {{else}} +
The statistics could not be read.
+ {{end}} +
+{{end}}