check / check (push) Waiting to run
Each webhook's event database keeps running totals: one row for its events, and one row per target for that target's deliveries, delivered and failed, each with what retention removed. Every write to them shares the transaction of the rows it counts. Deliveries get a finished_at column; it and target_id end the status index, so each target's deliveries finished in a window come from one index-range query grouped by target. Retention deletes 1000 expired events per transaction. The pane is its own template, its figures in tables. The schema changes in place with nothing back-filled, so an existing database must be recreated. Model: opus-5-5
330 lines
7.2 KiB
Go
330 lines
7.2 KiB
Go
package database
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"os"
|
|
"path/filepath"
|
|
"sync"
|
|
|
|
"go.uber.org/fx"
|
|
"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"
|
|
)
|
|
|
|
// WebhookDBManagerParams holds the fx dependencies for
|
|
// WebhookDBManager.
|
|
type WebhookDBManagerParams struct {
|
|
fx.In
|
|
|
|
Config *config.Config
|
|
Logger *logger.Logger
|
|
}
|
|
|
|
// errInvalidCachedDBType indicates a type assertion failure
|
|
// when retrieving a cached database connection.
|
|
var errInvalidCachedDBType = errors.New(
|
|
"invalid cached database type",
|
|
)
|
|
|
|
// WebhookDBManager manages per-webhook SQLite database files
|
|
// for event storage. Each webhook gets its own dedicated
|
|
// database containing Events, Deliveries, DeliveryResults and the
|
|
// running totals of them (EventTotals, TargetTotals).
|
|
// Database connections are opened lazily and cached.
|
|
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
|
|
// registers lifecycle hooks.
|
|
func NewWebhookDBManager(
|
|
lc fx.Lifecycle,
|
|
params WebhookDBManagerParams,
|
|
) (*WebhookDBManager, error) {
|
|
m := &WebhookDBManager{
|
|
dataDir: params.Config.DataDir,
|
|
log: params.Logger.Get(),
|
|
}
|
|
|
|
// 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",
|
|
m.dataDir,
|
|
err,
|
|
)
|
|
}
|
|
|
|
lc.Append(fx.Hook{
|
|
OnStop: func(_ context.Context) error {
|
|
return m.CloseAll()
|
|
},
|
|
})
|
|
|
|
m.log.Info(
|
|
"webhook database manager initialized",
|
|
"data_dir", m.dataDir,
|
|
)
|
|
|
|
return m, nil
|
|
}
|
|
|
|
// GetDB returns the database connection for a webhook,
|
|
// creating the database file lazily if it doesn't exist.
|
|
func (m *WebhookDBManager) GetDB(
|
|
webhookID string,
|
|
) (*gorm.DB, error) {
|
|
// Fast path: already open
|
|
if val, ok := m.dbs.Load(webhookID); ok {
|
|
return asGormDB(val, webhookID)
|
|
}
|
|
|
|
// 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)
|
|
}
|
|
|
|
db, err := m.openDB(webhookID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
m.dbs.Store(webhookID, db)
|
|
|
|
return db, 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
|
|
}
|
|
|
|
// CreateDB explicitly creates a new per-webhook database file
|
|
// and runs migrations.
|
|
func (m *WebhookDBManager) CreateDB(
|
|
webhookID string,
|
|
) error {
|
|
_, err := m.GetDB(webhookID)
|
|
|
|
return err
|
|
}
|
|
|
|
// DBExists checks if a per-webhook database file exists on
|
|
// disk.
|
|
func (m *WebhookDBManager) DBExists(
|
|
webhookID string,
|
|
) bool {
|
|
_, err := os.Stat(m.dbPath(webhookID))
|
|
|
|
return err == nil
|
|
}
|
|
|
|
// DeleteDB closes the connection and deletes the database file
|
|
// for a webhook. The file is permanently removed.
|
|
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 {
|
|
sqlDB, err := gormDB.DB()
|
|
if err == nil {
|
|
_ = sqlDB.Close()
|
|
}
|
|
}
|
|
}
|
|
|
|
// Delete the main DB file and WAL/SHM files
|
|
path := m.dbPath(webhookID)
|
|
for _, suffix := range []string{"", "-wal", "-shm"} {
|
|
err := os.Remove(path + suffix)
|
|
if err != nil && !os.IsNotExist(err) {
|
|
return fmt.Errorf(
|
|
"deleting webhook database file %s%s: %w",
|
|
path, suffix, err,
|
|
)
|
|
}
|
|
}
|
|
|
|
m.log.Info(
|
|
"deleted per-webhook database",
|
|
"webhook_id", webhookID,
|
|
)
|
|
|
|
return nil
|
|
}
|
|
|
|
// 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 {
|
|
if gormDB, castOK := value.(*gorm.DB); castOK {
|
|
sqlDB, err := gormDB.DB()
|
|
if err == nil {
|
|
closeErr := sqlDB.Close()
|
|
if closeErr != nil {
|
|
lastErr = closeErr
|
|
m.log.Error(
|
|
"failed to close webhook database",
|
|
"webhook_id", key,
|
|
"error", closeErr,
|
|
)
|
|
}
|
|
}
|
|
}
|
|
|
|
m.dbs.Delete(key)
|
|
|
|
return true
|
|
})
|
|
|
|
return lastErr
|
|
}
|
|
|
|
// DBPath returns the filesystem path for a webhook's database
|
|
// file.
|
|
func (m *WebhookDBManager) DBPath(
|
|
webhookID string,
|
|
) string {
|
|
return m.dbPath(webhookID)
|
|
}
|
|
|
|
func (m *WebhookDBManager) dbPath(
|
|
webhookID string,
|
|
) string {
|
|
return filepath.Join(
|
|
m.dataDir,
|
|
fmt.Sprintf("events-%s.db", webhookID),
|
|
)
|
|
}
|
|
|
|
// openDB opens (or creates) a per-webhook SQLite database and
|
|
// runs migrations.
|
|
func (m *WebhookDBManager) openDB(
|
|
webhookID string,
|
|
) (*gorm.DB, error) {
|
|
path := m.dbPath(webhookID)
|
|
|
|
// See sqlite_open.go: WAL, a busy timeout, immediate-transaction
|
|
// locking, and a bounded pool, all of which this file needs most —
|
|
// it is the one every delivery worker writes to concurrently.
|
|
sqlDB, err := OpenSQLite(path, SQLiteModeCreate)
|
|
if err != nil {
|
|
return nil, fmt.Errorf(
|
|
"opening webhook database %s: %w",
|
|
webhookID, err,
|
|
)
|
|
}
|
|
|
|
db, err := gorm.Open(sqlite.Dialector{
|
|
Conn: sqlDB,
|
|
}, &gorm.Config{
|
|
// Never leave this at GORM's default. See internal/gormlog.
|
|
Logger: gormlog.New(m.log),
|
|
})
|
|
if err != nil {
|
|
_ = sqlDB.Close()
|
|
|
|
return nil, fmt.Errorf(
|
|
"connecting to webhook database %s: %w",
|
|
webhookID, err,
|
|
)
|
|
}
|
|
|
|
// Keep main-database rows out of this file. See
|
|
// event_db_isolation.go.
|
|
err = omitAssociations(db)
|
|
if err != nil {
|
|
_ = sqlDB.Close()
|
|
|
|
return nil, fmt.Errorf(
|
|
"guarding webhook database %s: %w",
|
|
webhookID, err,
|
|
)
|
|
}
|
|
|
|
err = purgeTargetRows(db, m.log, webhookID)
|
|
if err != nil {
|
|
_ = sqlDB.Close()
|
|
|
|
return nil, err
|
|
}
|
|
|
|
// Run migrations for event-tier models only
|
|
err = db.AutoMigrate(
|
|
&Event{}, &Delivery{}, &DeliveryResult{},
|
|
&EventTotals{}, &TargetTotals{},
|
|
)
|
|
if err != nil {
|
|
_ = sqlDB.Close()
|
|
|
|
return nil, fmt.Errorf(
|
|
"migrating webhook database %s: %w",
|
|
webhookID, err,
|
|
)
|
|
}
|
|
|
|
// A new database gets its row of event totals, all zero. Target
|
|
// totals rows are created by the first delivery to each target.
|
|
err = db.FirstOrCreate(&EventTotals{}).Error
|
|
if err != nil {
|
|
_ = sqlDB.Close()
|
|
|
|
return nil, fmt.Errorf(
|
|
"creating event totals for webhook database %s: %w",
|
|
webhookID, err,
|
|
)
|
|
}
|
|
|
|
m.log.Info(
|
|
"opened per-webhook database",
|
|
"webhook_id", webhookID,
|
|
"path", path,
|
|
)
|
|
|
|
return db, nil
|
|
}
|