check / check (push) Waiting to run
The event log had no way to list only the events whose delivery failed, and once it showed only the 50 newest, an older failure could not be found at all. It now has All, Failed (N) and Pending (N) links, carried in a `show` query parameter, so they work without the page's script library. Each filtered list keeps the 50-row limit and newest-first order, and lists an event once. It finds matching deliveries through `idx_deliveries_status` and looks their events up by ID, so its cost follows the matches, not the webhook's size. Replay returns to the list it was pressed in. The heading line says what a filter counts. Model: opus-5-5
349 lines
12 KiB
Go
349 lines
12 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=?)")
|
|
}
|
|
|
|
// TestEventLogFiltersUseTheStatusIndex does the same for the event log's
|
|
// Failed and Pending lists, of the newest events with a delivery in
|
|
// given statuses, and for their counts (eventsWithStatus and
|
|
// countEventsWithStatus in the handlers). The lists must also reach
|
|
// the events table only by ID: from the matching deliveries, then from
|
|
// the newest of those events.
|
|
func TestEventLogFiltersUseTheStatusIndex(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)
|
|
|
|
dry := db.Session(&gorm.Session{DryRun: true})
|
|
byStatus := "idx_deliveries_status (status=? AND deleted_at=?)"
|
|
pending := []database.DeliveryStatus{
|
|
database.DeliveryStatusPending,
|
|
database.DeliveryStatusRetrying,
|
|
}
|
|
|
|
var (
|
|
rows []struct{ ID string }
|
|
count int64
|
|
)
|
|
|
|
matching := dry.Model(&database.Delivery{}).
|
|
Distinct("event_id").Where("status IN ?", pending)
|
|
newest := dry.Table("(?) AS matching", matching).
|
|
Joins("CROSS JOIN events ON events.id = matching.event_id").
|
|
Where(
|
|
"events.webhook_id = ? AND events.deleted_at IS NULL",
|
|
webhookID,
|
|
).
|
|
Order("events.created_at DESC").Limit(50).
|
|
Select("events.id AS event_id")
|
|
|
|
// Each step of the plan is printed in braces, so these name the
|
|
// lookup that follows each scan.
|
|
byID := "{SEARCH events USING INDEX sqlite_autoindex_events_1 (id=?)}"
|
|
|
|
assertPlanUses(t, db, dry.Table("(?) AS newest", newest).
|
|
Joins("CROSS JOIN events ON events.id = newest.event_id").
|
|
Select("id").Order("created_at DESC").Limit(50).Find(&rows),
|
|
byStatus, "{SCAN matching} "+byID, "{SCAN newest} "+byID)
|
|
assertPlanUses(t, db, dry.Model(&database.Delivery{}).
|
|
Distinct("event_id").Where("status IN ?", pending).Count(&count),
|
|
byStatus)
|
|
}
|
|
|
|
// 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())
|
|
}
|
|
}
|