Add a statistics pane to the webhook page (closes #368)
check / check (push) Failing after 1m4s
check / check (push) Failing after 1m4s
Each webhook's event database keeps one row of running totals: events, deliveries and failures, and how many of each retention removed. Storing an event, creating a delivery, a delivery becoming failed and the retention sweep each update it in the transaction that writes or deletes the rows it counts. Deliveries get a finished_at column, the last column of the status index, so the last-10-minutes and last-24-hours figures are index-range counts. The pane is its own template, included at the top of the page. The schema changes in place with nothing back-filled, so an existing database must be recreated. Model: opus-5-5
This commit is contained in:
@@ -123,7 +123,7 @@ 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
|
||||
// Retention's deletes (deleteExpired), whose subqueries are built
|
||||
// afresh for each statement as it builds them.
|
||||
expiredEventIDs := func() *gorm.DB {
|
||||
return dry.Model(&database.Event{}).Select("id").
|
||||
@@ -135,6 +135,14 @@ func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
|
||||
Select("id").Where("event_id IN (?)", expiredEventIDs()),
|
||||
).Delete(&database.DeliveryResult{}),
|
||||
"idx_delivery_results_delivery_id (delivery_id=?)", byEvent, byAge)
|
||||
var failed struct{ Count int64 }
|
||||
|
||||
assertPlanUses(t, db, dry.Unscoped().Model(&database.Delivery{}).
|
||||
Select("count(CASE WHEN status = ? THEN 1 END) AS count",
|
||||
database.DeliveryStatusFailed).
|
||||
Where("event_id IN (?)", expiredEventIDs()).
|
||||
Take(&failed),
|
||||
"idx_deliveries_event_id (event_id=?)", byAge)
|
||||
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||
"event_id IN (?)", expiredEventIDs(),
|
||||
).Delete(&database.Delivery{}),
|
||||
@@ -142,6 +150,36 @@ func TestEventTierQueriesUseTheirIndexes(t *testing.T) {
|
||||
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||
"created_at < ?", cutoff,
|
||||
).Delete(&database.Event{}), "idx_events_created_at (created_at<?)")
|
||||
|
||||
// The webhook page's statistics (readEventStats in the handlers):
|
||||
// deliveries in progress, deliveries finished and events received
|
||||
// since a time, and the newest event, which must come straight
|
||||
// off an index rather than from sorting every event.
|
||||
var (
|
||||
count int64
|
||||
newest []time.Time
|
||||
)
|
||||
|
||||
assertPlanUses(t, db, dry.Model(&database.Delivery{}).
|
||||
Where("status IN ?", []database.DeliveryStatus{
|
||||
database.DeliveryStatusPending,
|
||||
database.DeliveryStatusRetrying,
|
||||
}).Count(&count), byStatus)
|
||||
assertPlanUses(t, db, dry.Model(&database.Delivery{}).
|
||||
Where("status = ? AND finished_at >= ?",
|
||||
database.DeliveryStatusFailed, cutoff).Count(&count),
|
||||
"idx_deliveries_status "+
|
||||
"(status=? AND deleted_at=? AND finished_at>?)")
|
||||
assertPlanUses(t, db, dry.Model(&database.Event{}).
|
||||
Where("created_at >= ?", cutoff).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 +190,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 +209,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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
@@ -45,6 +49,12 @@ 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 ends the status index,
|
||||
// so the webhook page counts the deliveries that finished in a
|
||||
// recent window by reading 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,73 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// Totals is the single row of running totals in a webhook's event
|
||||
// database. It is what keeps the webhook page's lifetime figures right
|
||||
// after retention has removed the rows they count, and what lets the
|
||||
// page show them without counting every row.
|
||||
//
|
||||
// Storing an event, creating a delivery and failing a delivery each
|
||||
// add one, and retention adds what it deletes to the Removed columns.
|
||||
// Every addition goes through AddTotals, in the transaction that
|
||||
// writes or deletes the rows it counts.
|
||||
type Totals struct {
|
||||
ID int64 `gorm:"primaryKey"`
|
||||
|
||||
Events int64 `gorm:"not null"`
|
||||
Deliveries int64 `gorm:"not null"`
|
||||
Failures int64 `gorm:"not null"`
|
||||
|
||||
EventsRemoved int64 `gorm:"not null"`
|
||||
DeliveriesRemoved int64 `gorm:"not null"`
|
||||
FailuresRemoved int64 `gorm:"not null"`
|
||||
}
|
||||
|
||||
// TableName names the table AddTotals updates.
|
||||
func (Totals) TableName() string {
|
||||
return "totals"
|
||||
}
|
||||
|
||||
// EventsWithinRetention is how many of the webhook's events are still
|
||||
// stored.
|
||||
func (t Totals) EventsWithinRetention() int64 {
|
||||
return t.Events - t.EventsRemoved
|
||||
}
|
||||
|
||||
// DeliveriesWithinRetention is how many of the webhook's deliveries
|
||||
// are still stored.
|
||||
func (t Totals) DeliveriesWithinRetention() int64 {
|
||||
return t.Deliveries - t.DeliveriesRemoved
|
||||
}
|
||||
|
||||
// FailuresWithinRetention is how many of the webhook's failed
|
||||
// deliveries are still stored.
|
||||
func (t Totals) FailuresWithinRetention() int64 {
|
||||
return t.Failures - t.FailuresRemoved
|
||||
}
|
||||
|
||||
// AddTotals adds each count in add to the webhook's running totals.
|
||||
// Call it on the transaction that writes or deletes the rows it
|
||||
// counts, so the totals change exactly when those rows do.
|
||||
func AddTotals(tx *gorm.DB, add Totals) error {
|
||||
err := tx.Exec(
|
||||
`UPDATE totals SET
|
||||
events = events + ?,
|
||||
deliveries = deliveries + ?,
|
||||
failures = failures + ?,
|
||||
events_removed = events_removed + ?,
|
||||
deliveries_removed = deliveries_removed + ?,
|
||||
failures_removed = failures_removed + ?`,
|
||||
add.Events, add.Deliveries, add.Failures,
|
||||
add.EventsRemoved, add.DeliveriesRemoved, add.FailuresRemoved,
|
||||
).Error
|
||||
if err != nil {
|
||||
return fmt.Errorf("adding to running totals: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -2,7 +2,7 @@ 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, Totals) live in
|
||||
// per-webhook dedicated databases managed by WebhookDBManager.
|
||||
func (d *Database) Migrate() error {
|
||||
return d.db.AutoMigrate(
|
||||
|
||||
@@ -267,55 +267,101 @@ func retentionCutoff(
|
||||
|
||||
// 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.
|
||||
// cutoff, and adds what it deleted to the running totals, all in one
|
||||
// transaction. Deletes are unscoped so rows are physically removed
|
||||
// rather than soft-deleted, reclaiming disk. It returns the number of
|
||||
// events deleted.
|
||||
func reapExpired(db *gorm.DB, cutoff time.Time) (int64, error) {
|
||||
var removed Totals
|
||||
|
||||
err := db.Transaction(func(tx *gorm.DB) error {
|
||||
var err error
|
||||
|
||||
removed, err = deleteExpired(tx, cutoff)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return AddTotals(tx, removed)
|
||||
})
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
return removed.EventsRemoved, nil
|
||||
}
|
||||
|
||||
// deleteExpired runs reapExpired's deletes and returns how many
|
||||
// events, deliveries and failed deliveries they removed.
|
||||
func deleteExpired(tx *gorm.DB, cutoff time.Time) (Totals, error) {
|
||||
var removed Totals
|
||||
|
||||
// Fresh subqueries are built per statement to avoid reusing a
|
||||
// mutated builder across executions.
|
||||
expiredEventIDs := func() *gorm.DB {
|
||||
return db.Model(&Event{}).
|
||||
return tx.Model(&Event{}).
|
||||
Select("id").
|
||||
Where("created_at < ?", cutoff)
|
||||
}
|
||||
expiredDeliveryIDs := func() *gorm.DB {
|
||||
return db.Model(&Delivery{}).
|
||||
return tx.Model(&Delivery{}).
|
||||
Select("id").
|
||||
Where("event_id IN (?)", expiredEventIDs())
|
||||
}
|
||||
|
||||
// 1. Delivery results whose delivery belongs to an expired event.
|
||||
res := db.Unscoped().
|
||||
res := tx.Unscoped().
|
||||
Where("delivery_id IN (?)", expiredDeliveryIDs()).
|
||||
Delete(&DeliveryResult{})
|
||||
if res.Error != nil {
|
||||
return 0, fmt.Errorf(
|
||||
return removed, fmt.Errorf(
|
||||
"deleting expired delivery results: %w",
|
||||
res.Error,
|
||||
)
|
||||
}
|
||||
|
||||
// 2. Deliveries belonging to an expired event.
|
||||
del := db.Unscoped().
|
||||
// 2. Deliveries belonging to an expired event, after counting the
|
||||
// failed ones among them. 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 the expired ones through the event_id index.
|
||||
var failed struct{ Count int64 }
|
||||
|
||||
err := tx.Unscoped().Model(&Delivery{}).
|
||||
Select("count(CASE WHEN status = ? THEN 1 END) AS count",
|
||||
DeliveryStatusFailed).
|
||||
Where("event_id IN (?)", expiredEventIDs()).
|
||||
Take(&failed).Error
|
||||
if err != nil {
|
||||
return removed, fmt.Errorf(
|
||||
"counting expired failed deliveries: %w", err,
|
||||
)
|
||||
}
|
||||
|
||||
del := tx.Unscoped().
|
||||
Where("event_id IN (?)", expiredEventIDs()).
|
||||
Delete(&Delivery{})
|
||||
if del.Error != nil {
|
||||
return 0, fmt.Errorf(
|
||||
return removed, fmt.Errorf(
|
||||
"deleting expired deliveries: %w",
|
||||
del.Error,
|
||||
)
|
||||
}
|
||||
|
||||
// 3. The expired events themselves.
|
||||
ev := db.Unscoped().
|
||||
ev := tx.Unscoped().
|
||||
Where("created_at < ?", cutoff).
|
||||
Delete(&Event{})
|
||||
if ev.Error != nil {
|
||||
return 0, fmt.Errorf(
|
||||
return removed, fmt.Errorf(
|
||||
"deleting expired events: %w",
|
||||
ev.Error,
|
||||
)
|
||||
}
|
||||
|
||||
return ev.RowsAffected, nil
|
||||
removed.EventsRemoved = ev.RowsAffected
|
||||
removed.DeliveriesRemoved = del.RowsAffected
|
||||
removed.FailuresRemoved = failed.Count
|
||||
|
||||
return removed, nil
|
||||
}
|
||||
|
||||
@@ -0,0 +1,126 @@
|
||||
package database_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"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"
|
||||
)
|
||||
|
||||
// readTotals reads a webhook database's row of running totals,
|
||||
// asserting that it has exactly one.
|
||||
func readTotals(t *testing.T, db *gorm.DB) database.Totals {
|
||||
t.Helper()
|
||||
|
||||
var rows []database.Totals
|
||||
|
||||
require.NoError(t, db.Find(&rows).Error)
|
||||
require.Len(t, rows, 1)
|
||||
|
||||
return rows[0]
|
||||
}
|
||||
|
||||
// TestWebhookDBManager_TotalsRowSurvivesReopen verifies that a new
|
||||
// event database starts with one row of zero totals, and that opening
|
||||
// it again keeps that row and what was added to it.
|
||||
func TestWebhookDBManager_TotalsRowSurvivesReopen(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 := readTotals(t, db)
|
||||
assert.Equal(t, database.Totals{ID: fresh.ID}, fresh)
|
||||
|
||||
require.NoError(t, database.AddTotals(db, database.Totals{
|
||||
Events: 2, Deliveries: 3, Failures: 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.Totals{
|
||||
ID: fresh.ID, Events: 2, Deliveries: 3, Failures: 1,
|
||||
}, readTotals(t, db))
|
||||
}
|
||||
|
||||
// TestRetentionReaper_AddsWhatItRemovesToTotals verifies that a sweep
|
||||
// leaves the lifetime totals alone and adds the events, deliveries and
|
||||
// failed deliveries it deletes to the removed totals, so the totals
|
||||
// within retention match the rows still stored.
|
||||
func TestRetentionReaper_AddsWhatItRemovesToTotals(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)
|
||||
|
||||
now := time.Now()
|
||||
expired := now.Add(-40 * 24 * time.Hour)
|
||||
|
||||
seedEventChain(t, db, webhookID, expired)
|
||||
expiredFailure := seedEventChain(t, db, webhookID, expired)
|
||||
recentFailure := seedEventChain(
|
||||
t, db, webhookID, now.Add(-24*time.Hour),
|
||||
)
|
||||
|
||||
for _, id := range []string{
|
||||
expiredFailure.deliveryID, recentFailure.deliveryID,
|
||||
} {
|
||||
require.NoError(t, db.Model(&database.Delivery{}).
|
||||
Where("id = ?", id).
|
||||
Update("status", database.DeliveryStatusFailed).Error)
|
||||
}
|
||||
|
||||
// The totals storing those rows would have left.
|
||||
require.NoError(t, database.AddTotals(db, database.Totals{
|
||||
Events: 3, Deliveries: 3, Failures: 2,
|
||||
}))
|
||||
|
||||
env.reaper.ExportSweep(context.Background())
|
||||
|
||||
totals := readTotals(t, db)
|
||||
assert.Equal(t, database.Totals{
|
||||
ID: totals.ID,
|
||||
Events: 3, Deliveries: 3, Failures: 2,
|
||||
EventsRemoved: 2, DeliveriesRemoved: 2, FailuresRemoved: 1,
|
||||
}, totals)
|
||||
|
||||
var events, deliveries, failures int64
|
||||
|
||||
require.NoError(t, db.Model(&database.Event{}).Count(&events).Error)
|
||||
require.NoError(t, db.Model(&database.Delivery{}).
|
||||
Count(&deliveries).Error)
|
||||
require.NoError(t, db.Model(&database.Delivery{}).
|
||||
Where("status = ?", database.DeliveryStatusFailed).
|
||||
Count(&failures).Error)
|
||||
|
||||
assert.Equal(t, events, totals.EventsWithinRetention())
|
||||
assert.Equal(t, deliveries, totals.DeliveriesWithinRetention())
|
||||
assert.Equal(t, failures, totals.FailuresWithinRetention())
|
||||
|
||||
// A sweep with nothing left to remove changes nothing.
|
||||
env.reaper.ExportSweep(context.Background())
|
||||
|
||||
assert.Equal(t, totals, readTotals(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.
|
||||
// Database connections are opened lazily and cached.
|
||||
type WebhookDBManager struct {
|
||||
dataDir string
|
||||
@@ -294,7 +295,7 @@ func (m *WebhookDBManager) openDB(
|
||||
|
||||
// Run migrations for event-tier models only
|
||||
err = db.AutoMigrate(
|
||||
&Event{}, &Delivery{}, &DeliveryResult{},
|
||||
&Event{}, &Delivery{}, &DeliveryResult{}, &Totals{},
|
||||
)
|
||||
if err != nil {
|
||||
_ = sqlDB.Close()
|
||||
@@ -305,6 +306,17 @@ func (m *WebhookDBManager) openDB(
|
||||
)
|
||||
}
|
||||
|
||||
// A new database gets its row of running totals, all zero.
|
||||
err = db.FirstOrCreate(&Totals{}).Error
|
||||
if err != nil {
|
||||
_ = sqlDB.Close()
|
||||
|
||||
return nil, fmt.Errorf(
|
||||
"creating running totals for webhook database %s: %w",
|
||||
webhookID, err,
|
||||
)
|
||||
}
|
||||
|
||||
m.log.Info(
|
||||
"opened per-webhook database",
|
||||
"webhook_id", webhookID,
|
||||
|
||||
Reference in New Issue
Block a user