check / check (push) Waiting to run
The entrypoint list showed no sign of whether anything uses an entrypoint, so an operator with several could not tell which senders are live before deactivating or deleting one. Each entrypoint now shows when its last event arrived, relative with the UTC time on hover, or "never", and how many events arrived through it within the webhook's retention. The last-event time comes from a new entrypoint_totals row written in the transaction that stores the event and left by retention, so a sender quieter than the retention period does not read "never". The count is one grouped query over a new index. Resubmitted copies count in neither. Pre-1.0: schema changed in place. Model: opus-5-5
293 lines
9.6 KiB
Go
293 lines
9.6 KiB
Go
package database_test
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"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"
|
|
)
|
|
|
|
// TestWebhookDBManager_OpenAddsEventTierIndexes verifies that opening a
|
|
// per-webhook database that predates these indexes creates them. It
|
|
// stands in for an older database file by dropping the indexes
|
|
// AutoMigrate just created, then reopening the same file.
|
|
func TestWebhookDBManager_OpenAddsEventTierIndexes(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
indexes := []struct {
|
|
model any
|
|
name string
|
|
}{
|
|
{&database.Delivery{}, "idx_deliveries_status"},
|
|
{&database.Delivery{}, "idx_deliveries_event_id"},
|
|
{&database.DeliveryResult{}, "idx_delivery_results_delivery_id"},
|
|
{&database.Event{}, "idx_events_deleted_at_created_at"},
|
|
{&database.Event{}, "idx_events_created_at"},
|
|
}
|
|
|
|
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)
|
|
|
|
// A fresh database has them.
|
|
for _, ix := range indexes {
|
|
require.True(t, db.Migrator().HasIndex(ix.model, ix.name))
|
|
}
|
|
|
|
// Stand in for a database file created before the indexes existed.
|
|
for _, ix := range indexes {
|
|
require.NoError(t, db.Migrator().DropIndex(ix.model, ix.name))
|
|
require.False(t, db.Migrator().HasIndex(ix.model, ix.name))
|
|
}
|
|
|
|
// Drop the cached connection so the next open reopens the file and
|
|
// runs AutoMigrate against it, as a restart would.
|
|
require.NoError(t, mgr.CloseAll())
|
|
|
|
db, err = mgr.GetDB(webhookID)
|
|
require.NoError(t, err)
|
|
|
|
for _, ix := range indexes {
|
|
assert.True(t, db.Migrator().HasIndex(ix.model, ix.name),
|
|
"opening the existing database should create %s", ix.name)
|
|
}
|
|
}
|
|
|
|
// TestEventTierQueriesUseTheirIndexes verifies that the statements the
|
|
// indexes are for use them. GORM builds each statement in a dry run as
|
|
// the code named above it does, soft-delete condition included, and
|
|
// SQLite, which keeps no statistics on these tables, must plan to seek
|
|
// on each index listed by the columns in parentheses.
|
|
func TestEventTierQueriesUseTheirIndexes(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})
|
|
ids := []string{
|
|
uuid.New().String(), uuid.New().String(), uuid.New().String(),
|
|
}
|
|
cutoff := time.Now()
|
|
|
|
var (
|
|
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=?)"
|
|
|
|
// The delivery engine: recovery and the retry sweep, the sweep for
|
|
// stranded pending deliveries, and the queue depth count.
|
|
assertPlanUses(t, db, dry.Where(
|
|
"status = ?", database.DeliveryStatusRetrying,
|
|
).Find(&deliveries), byStatus)
|
|
assertPlanUses(t, db, dry.Where(
|
|
"status = ? AND updated_at < ?",
|
|
database.DeliveryStatusPending, cutoff,
|
|
).Limit(500).Find(&deliveries), byStatus)
|
|
assertPlanUses(t, db, dry.Model(&database.Delivery{}).
|
|
Select("target_id", "status", "count(*) as depth").
|
|
Where("status IN ?", []database.DeliveryStatus{
|
|
database.DeliveryStatusPending,
|
|
database.DeliveryStatusRetrying,
|
|
}).Group("target_id, status").Find(&depths), byStatus)
|
|
|
|
// The event log: each event's deliveries, then their attempts
|
|
// (loadEventsWithDeliveries, loadDeliveryResults).
|
|
assertPlanUses(t, db, dry.Where("event_id = ?", ids[0]).
|
|
Find(&deliveries), byEvent)
|
|
assertPlanUses(t, db, dry.Where("delivery_id IN ?", ids).
|
|
Order("attempt_num ASC").Find(&results),
|
|
"idx_delivery_results_delivery_id (delivery_id=? AND deleted_at=?)")
|
|
|
|
// 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.Unscoped().Model(&database.Delivery{}).
|
|
Select("id").Where("event_id IN ?", ids),
|
|
).Delete(&database.DeliveryResult{}),
|
|
"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, and events received since a time.
|
|
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
|
|
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>?)")
|
|
}
|
|
|
|
// TestResubmitCountUsesItsIndex does the same for the event log's count
|
|
// of the events resubmitted from each of a page's events (resubmitCounts
|
|
// in the handlers). It passes a full page of 25 ids: with an index on
|
|
// resubmitted_from_id alone, SQLite uses it for three ids and turns to
|
|
// the deleted_at index from five.
|
|
func TestResubmitCountUsesItsIndex(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})
|
|
|
|
page := make([]string, 25)
|
|
for i := range page {
|
|
page[i] = uuid.New().String()
|
|
}
|
|
|
|
var counts []struct{ Total int }
|
|
|
|
assertPlanUses(t, db, dry.Model(&database.Event{}).
|
|
Select("resubmitted_from_id, count(*) AS total").
|
|
Where("resubmitted_from_id IN ?", page).
|
|
Group("resubmitted_from_id").Find(&counts),
|
|
"idx_events_resubmitted_from_id "+
|
|
"(resubmitted_from_id=? AND deleted_at=?)")
|
|
}
|
|
|
|
// TestEntrypointEventsUseTheirIndex does the same for the webhook
|
|
// page's count, for each entrypoint, of the events that arrived on its
|
|
// URL since the retention cutoff (addEntrypointEvents in the
|
|
// handlers), which must come from the index alone. It passes 25
|
|
// entrypoints, as TestResubmitCountUsesItsIndex passes 25 events.
|
|
func TestEntrypointEventsUseTheirIndex(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})
|
|
|
|
entrypoints := make([]string, 25)
|
|
for i := range entrypoints {
|
|
entrypoints[i] = uuid.New().String()
|
|
}
|
|
|
|
var rows []struct{ Events int }
|
|
|
|
assertPlanUses(t, db, dry.Model(&database.Event{}).
|
|
Select("entrypoint_id, count(*) AS events").
|
|
Where("entrypoint_id IN ? AND resubmitted_from_id IS NULL",
|
|
entrypoints).
|
|
Where("created_at >= ?", time.Now()).
|
|
Group("entrypoint_id").Find(&rows),
|
|
"COVERING INDEX idx_events_entrypoint_id "+
|
|
"(entrypoint_id=? AND deleted_at=? AND "+
|
|
"resubmitted_from_id=? AND created_at>?)")
|
|
}
|
|
|
|
// assertPlanUses asserts that SQLite's plan for a statement GORM built
|
|
// in a dry run, run with the same SQL and arguments GORM would send,
|
|
// names each of the given indexes.
|
|
func assertPlanUses(
|
|
t *testing.T, db, built *gorm.DB, indexes ...string,
|
|
) {
|
|
t.Helper()
|
|
|
|
var plan []struct{ Detail string }
|
|
|
|
require.NoError(t, db.Raw(
|
|
"EXPLAIN QUERY PLAN "+built.Statement.SQL.String(),
|
|
built.Statement.Vars...,
|
|
).Scan(&plan).Error)
|
|
|
|
for _, index := range indexes {
|
|
assert.Contains(t, fmt.Sprint(plan), index,
|
|
built.Statement.SQL.String())
|
|
}
|
|
}
|