Brings `prod`, which upaas deploys, up to `main` at `9cf9cdd`, the merge of #321. `prod` was cut from `main` at `251cb3d` (1.0.0b1). What it deploys is everything listed in #321. For running it: - With `WEBHOOKER_ENVIRONMENT` unset, the instance runs as `prod` and sends no `Access-Control-Allow-Origin: *`. - Each event database gains its new indexes the first time it is opened after the upgrade. - `webhooker_delivery_retries_total` no longer counts a circuit breaker holding back a delivery that is already `retrying`. Not in this PR yet: #340, in which the container sets its own data directory owner and mode before start. It is in progress on `next`. Once it reaches `main`, this PR carries it, because the PR follows `main`. Model: opus-5-5 Co-authored-by: Jeffrey Paul <1+sneak@noreply.example.org> Reviewed-on: #343
This commit is contained in:
@@ -17,12 +17,12 @@ import (
|
||||
"gorm.io/gorm"
|
||||
"sneak.berlin/go/webhooker/internal/banner"
|
||||
"sneak.berlin/go/webhooker/internal/config"
|
||||
"sneak.berlin/go/webhooker/internal/datadir"
|
||||
"sneak.berlin/go/webhooker/internal/gormlog"
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
)
|
||||
|
||||
const (
|
||||
dataDirPerm = 0750
|
||||
randomPasswordLen = 16
|
||||
sessionKeyLen = 32
|
||||
)
|
||||
@@ -185,7 +185,9 @@ func (d *Database) connect() error {
|
||||
// caller's decision.
|
||||
func (d *Database) connectTo(dataDir string) error {
|
||||
// Ensure the data directory exists before opening the database.
|
||||
err := os.MkdirAll(dataDir, dataDirPerm)
|
||||
// datadir.DirPerm is the single source of the directory mode; this
|
||||
// package creates the directory too, since either may run first.
|
||||
err := os.MkdirAll(dataDir, datadir.DirPerm)
|
||||
if err != nil {
|
||||
return fmt.Errorf(
|
||||
"creating data directory %s: %w",
|
||||
|
||||
@@ -0,0 +1,166 @@
|
||||
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 }
|
||||
)
|
||||
|
||||
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<?)"
|
||||
|
||||
// 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's three deletes (reapExpired), whose subqueries are built
|
||||
// afresh for each statement as it builds them.
|
||||
expiredEventIDs := func() *gorm.DB {
|
||||
return dry.Model(&database.Event{}).Select("id").
|
||||
Where("created_at < ?", cutoff)
|
||||
}
|
||||
|
||||
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||
"delivery_id IN (?)", dry.Model(&database.Delivery{}).
|
||||
Select("id").Where("event_id IN (?)", expiredEventIDs()),
|
||||
).Delete(&database.DeliveryResult{}),
|
||||
"idx_delivery_results_delivery_id (delivery_id=?)", byEvent, byAge)
|
||||
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||
"event_id IN (?)", expiredEventIDs(),
|
||||
).Delete(&database.Delivery{}),
|
||||
"idx_deliveries_event_id (event_id=?)", byAge)
|
||||
assertPlanUses(t, db, dry.Unscoped().Where(
|
||||
"created_at < ?", cutoff,
|
||||
).Delete(&database.Event{}), "idx_events_created_at (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())
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,7 @@
|
||||
package database
|
||||
|
||||
import "gorm.io/gorm"
|
||||
|
||||
// DeliveryStatus represents the status of a delivery
|
||||
type DeliveryStatus string
|
||||
|
||||
@@ -29,12 +31,19 @@ func (s DeliveryStatus) Terminal() bool {
|
||||
}
|
||||
|
||||
// Delivery represents a delivery attempt for an event to a target
|
||||
//
|
||||
//nolint:lll // a struct tag cannot wrap
|
||||
type Delivery struct {
|
||||
BaseModel
|
||||
|
||||
EventID string `gorm:"type:uuid;not null" json:"eventId"`
|
||||
TargetID string `gorm:"type:uuid;not null" json:"targetId"`
|
||||
Status DeliveryStatus `gorm:"not null;default:'pending'" json:"status"`
|
||||
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"`
|
||||
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
|
||||
// of the event_id and status indexes, for the reason DeliveryResult
|
||||
// gives.
|
||||
DeletedAt gorm.DeletedAt `gorm:"index:idx_deliveries_event_id,priority:2;index:idx_deliveries_status,priority:2" json:"deletedAt,omitzero"`
|
||||
|
||||
// Relations
|
||||
Event Event `json:"event,omitzero"`
|
||||
|
||||
@@ -1,10 +1,21 @@
|
||||
package database
|
||||
|
||||
import "gorm.io/gorm"
|
||||
|
||||
// DeliveryResult represents the result of a delivery attempt
|
||||
//
|
||||
//nolint:lll // a struct tag cannot wrap
|
||||
type DeliveryResult struct {
|
||||
BaseModel
|
||||
|
||||
DeliveryID string `gorm:"type:uuid;not null" json:"deliveryId"`
|
||||
// DeliveryID and DeletedAt make up one index, in that order.
|
||||
// DeletedAt repeats the BaseModel field only to join it: GORM adds
|
||||
// "deleted_at IS NULL" to almost every query, and where a column is
|
||||
// matched against several values SQLite otherwise reads through the
|
||||
// deleted_at index, which every live row matches.
|
||||
DeliveryID string `gorm:"type:uuid;not null;index:idx_delivery_results_delivery_id,priority:1" json:"deliveryId"`
|
||||
DeletedAt gorm.DeletedAt `gorm:"index:idx_delivery_results_delivery_id,priority:2" json:"deletedAt,omitzero"`
|
||||
|
||||
AttemptNum int `gorm:"not null" json:"attemptNum"`
|
||||
Success bool `json:"success"`
|
||||
StatusCode int `json:"statusCode,omitempty"`
|
||||
|
||||
@@ -1,9 +1,27 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// Event represents a captured webhook event
|
||||
//
|
||||
//nolint:lll // a struct tag cannot wrap
|
||||
type Event struct {
|
||||
BaseModel
|
||||
|
||||
// CreatedAt and DeletedAt repeat the BaseModel fields only to index
|
||||
// them for retention, which finds events by age. Its lookups carry
|
||||
// GORM's "deleted_at IS NULL" (see DeliveryResult) and compare
|
||||
// created_at with <, so their index has deleted_at first: SQLite
|
||||
// narrows by a < only on the last column it uses. Its final delete
|
||||
// has no deleted_at condition and uses the index on created_at
|
||||
// alone. The other tables keep the unindexed BaseModel created_at.
|
||||
CreatedAt time.Time `gorm:"index;index:idx_events_deleted_at_created_at,priority:2" json:"createdAt"`
|
||||
DeletedAt gorm.DeletedAt `gorm:"index:idx_events_deleted_at_created_at,priority:1" json:"deletedAt,omitzero"`
|
||||
|
||||
WebhookID string `gorm:"type:uuid;not null" json:"webhookId"`
|
||||
EntrypointID string `gorm:"type:uuid;not null" json:"entrypointId"`
|
||||
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
"sneak.berlin/go/webhooker/internal/config"
|
||||
"sneak.berlin/go/webhooker/internal/datadir"
|
||||
"sneak.berlin/go/webhooker/internal/gormlog"
|
||||
"sneak.berlin/go/webhooker/internal/logger"
|
||||
)
|
||||
@@ -40,6 +41,11 @@ type WebhookDBManager struct {
|
||||
dataDir string
|
||||
dbs sync.Map // map[webhookID]*gorm.DB
|
||||
log *slog.Logger
|
||||
|
||||
// mu is held while a database is opened, deleted, or closed, so
|
||||
// each file has at most one open handle. Reading an already cached
|
||||
// handle does not take it.
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
// NewWebhookDBManager creates a new WebhookDBManager and
|
||||
@@ -53,8 +59,9 @@ func NewWebhookDBManager(
|
||||
log: params.Logger.Get(),
|
||||
}
|
||||
|
||||
// Create data directory if it doesn't exist
|
||||
err := os.MkdirAll(m.dataDir, dataDirPerm)
|
||||
// Create data directory if it doesn't exist. datadir.DirPerm is the
|
||||
// single source of the directory mode; either package may run first.
|
||||
err := os.MkdirAll(m.dataDir, datadir.DirPerm)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf(
|
||||
"creating data directory %s: %w",
|
||||
@@ -84,43 +91,39 @@ func (m *WebhookDBManager) GetDB(
|
||||
) (*gorm.DB, error) {
|
||||
// Fast path: already open
|
||||
if val, ok := m.dbs.Load(webhookID); ok {
|
||||
cachedDB, castOK := val.(*gorm.DB)
|
||||
if !castOK {
|
||||
return nil, fmt.Errorf(
|
||||
"%w for webhook %s",
|
||||
errInvalidCachedDBType,
|
||||
webhookID,
|
||||
)
|
||||
}
|
||||
return asGormDB(val, webhookID)
|
||||
}
|
||||
|
||||
return cachedDB, nil
|
||||
// Slow path: open the database under the lock, looking in the
|
||||
// cache again first. A caller that raced another one here then
|
||||
// waits for its handle instead of opening a second one.
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
|
||||
if val, ok := m.dbs.Load(webhookID); ok {
|
||||
return asGormDB(val, webhookID)
|
||||
}
|
||||
|
||||
// Slow path: open/create the database
|
||||
db, err := m.openDB(webhookID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Store it; if another goroutine beat us, close ours
|
||||
actual, loaded := m.dbs.LoadOrStore(webhookID, db)
|
||||
if loaded {
|
||||
// Another goroutine created it first; close our duplicate
|
||||
sqlDB, closeErr := db.DB()
|
||||
if closeErr == nil {
|
||||
_ = sqlDB.Close()
|
||||
}
|
||||
m.dbs.Store(webhookID, db)
|
||||
|
||||
existingDB, castOK := actual.(*gorm.DB)
|
||||
if !castOK {
|
||||
return nil, fmt.Errorf(
|
||||
"%w for webhook %s",
|
||||
errInvalidCachedDBType,
|
||||
webhookID,
|
||||
)
|
||||
}
|
||||
return db, nil
|
||||
}
|
||||
|
||||
return existingDB, nil
|
||||
// asGormDB returns a value read from the cache as the database
|
||||
// handle it is.
|
||||
func asGormDB(val any, webhookID string) (*gorm.DB, error) {
|
||||
db, ok := val.(*gorm.DB)
|
||||
if !ok {
|
||||
return nil, fmt.Errorf(
|
||||
"%w for webhook %s",
|
||||
errInvalidCachedDBType,
|
||||
webhookID,
|
||||
)
|
||||
}
|
||||
|
||||
return db, nil
|
||||
@@ -151,6 +154,11 @@ func (m *WebhookDBManager) DBExists(
|
||||
func (m *WebhookDBManager) DeleteDB(
|
||||
webhookID string,
|
||||
) error {
|
||||
// Held until the files are gone, so GetDB cannot open the file
|
||||
// again between the close and the removal.
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
|
||||
// Close and remove from cache
|
||||
if val, ok := m.dbs.LoadAndDelete(webhookID); ok {
|
||||
if gormDB, castOK := val.(*gorm.DB); castOK {
|
||||
@@ -184,6 +192,11 @@ func (m *WebhookDBManager) DeleteDB(
|
||||
// CloseAll closes all open per-webhook database connections.
|
||||
// Called during application shutdown.
|
||||
func (m *WebhookDBManager) CloseAll() error {
|
||||
// An open already under way finishes and is cached first, so it
|
||||
// is closed here rather than cached after this loop has passed.
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
|
||||
var lastErr error
|
||||
|
||||
m.dbs.Range(func(key, value any) bool {
|
||||
|
||||
@@ -1,10 +1,14 @@
|
||||
package database_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
@@ -104,6 +108,54 @@ func TestWebhookDBManager_CreateAndGetDB(t *testing.T) {
|
||||
assert.Equal(t, `{"test": true}`, readEvent.Body)
|
||||
}
|
||||
|
||||
// Many callers ask for one webhook's database at the same moment,
|
||||
// before it is cached. Only one of them may open the file; the others
|
||||
// must wait for its handle. openDB logs one "opened per-webhook
|
||||
// database" line per open, and those lines are what is counted.
|
||||
func TestWebhookDBManager_ConcurrentFirstTouchOpensOnce(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var logs bytes.Buffer
|
||||
|
||||
mgr := database.NewTestWebhookDBManagerWithLogger(
|
||||
t.TempDir(),
|
||||
slog.New(slog.NewTextHandler(&logs, nil)),
|
||||
)
|
||||
|
||||
t.Cleanup(func() { assert.NoError(t, mgr.CloseAll()) })
|
||||
|
||||
webhookID := uuid.New().String()
|
||||
|
||||
const callers = 16
|
||||
|
||||
start := make(chan struct{})
|
||||
handles := make([]*gorm.DB, callers)
|
||||
errs := make([]error, callers)
|
||||
|
||||
var wg sync.WaitGroup
|
||||
|
||||
for i := range callers {
|
||||
wg.Go(func() {
|
||||
<-start
|
||||
|
||||
handles[i], errs[i] = mgr.GetDB(webhookID)
|
||||
})
|
||||
}
|
||||
|
||||
close(start)
|
||||
wg.Wait()
|
||||
|
||||
for i := range callers {
|
||||
require.NoError(t, errs[i])
|
||||
assert.Same(t, handles[0], handles[i])
|
||||
}
|
||||
|
||||
assert.Equal(
|
||||
t, 1,
|
||||
strings.Count(logs.String(), "opened per-webhook database"),
|
||||
)
|
||||
}
|
||||
|
||||
func TestWebhookDBManager_DeleteDB(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user