check / check (push) Successful in 3m23s
Each entrypoint in the webhook page's entrypoint list now shows when its last event arrived (relative, with the full UTC time on hover), or "never", and how many events arrived through it within the webhook's retention period. Both come from one query per page, grouped by entrypoint, over a new events index on entrypoint_id, deleted_at and created_at, so the page reads only the index entries it counts. Pre-1.0: the index goes into the schema in place. A test checks the database's plan for the statement as the code builds it; another shows two entrypoints with different traffic each with their own figures, an event older than retention left out, and an unused entrypoint reading "never". Model: opus-5-5
292 lines
9.6 KiB
Go
292 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 figures for each entrypoint, the events since the retention
|
|
// cutoff and when the newest arrived (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, "+
|
|
"max(created_at), created_at AS last_event_at").
|
|
Where("entrypoint_id IN ?", entrypoints).
|
|
Where("created_at >= ?", time.Now()).
|
|
Group("entrypoint_id").Find(&rows),
|
|
"COVERING INDEX idx_events_entrypoint_id "+
|
|
"(entrypoint_id=? AND deleted_at=? 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())
|
|
}
|
|
}
|