check / check (push) Waiting to run
An audit of every file webhooker reads configuration or required state from found cases that carried on silently. A zero-length webhooker.db, and a missing or zero-length per-webhook database, now log the "created a new, empty database" warning naming the file; restart recovery opens every live webhook's database, checking under the manager's lock that it still exists, so a missing one is reported at start. The main database's open errors name webhooker.db, for the server and webhooker resetpw; resetpw refuses a zero-length webhooker.db. A directory in place of any database file or its -wal or -shm is refused naming it. The README says how each case is treated. Also closes #459. Model: opus-5-5
415 lines
9.7 KiB
Go
415 lines
9.7 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",
|
|
)
|
|
|
|
// ErrEventDBNotRemoved is in DeleteDB's error when the event
|
|
// database file itself could not be removed: it is still on disk.
|
|
var ErrEventDBNotRemoved = errors.New(
|
|
"event database file not removed",
|
|
)
|
|
|
|
// ErrSidecarNotRemoved is in DeleteDB's error when the event
|
|
// database file was removed, so its events are gone, but its -wal
|
|
// or -shm sidecar could not be.
|
|
var ErrSidecarNotRemoved = errors.New(
|
|
"event database file removed, but a -wal or -shm sidecar was not",
|
|
)
|
|
|
|
// 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, opening it on
|
|
// first use.
|
|
//
|
|
// The file is made by CreateDB when the webhook is created. One that is
|
|
// missing or zero-length here means the webhook's events and pending
|
|
// deliveries are gone: an empty database is created in its place so
|
|
// the webhook keeps receiving, and that is logged as a warning naming
|
|
// the file, as a new main database is.
|
|
func (m *WebhookDBManager) GetDB(
|
|
webhookID string,
|
|
) (*gorm.DB, error) {
|
|
return m.getDB(webhookID, false)
|
|
}
|
|
|
|
// GetDBIf is GetDB, done only when check reports true. check runs under
|
|
// the lock DeleteDB holds while it removes the files, so a caller can
|
|
// confirm the webhook still exists and open its database with no delete
|
|
// in between. The handle is nil when check reports false. check must
|
|
// not call the manager.
|
|
func (m *WebhookDBManager) GetDBIf(
|
|
webhookID string, check func() (bool, error),
|
|
) (*gorm.DB, error) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
|
|
ok, err := check()
|
|
if err != nil || !ok {
|
|
return nil, err
|
|
}
|
|
|
|
return m.getDBLocked(webhookID, false)
|
|
}
|
|
|
|
// 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 creates a new webhook's database file and runs
|
|
// migrations.
|
|
func (m *WebhookDBManager) CreateDB(
|
|
webhookID string,
|
|
) error {
|
|
_, err := m.getDB(webhookID, true)
|
|
|
|
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, with its -wal and -shm sidecars. The files are
|
|
// permanently removed. Each file is tried even when another could
|
|
// not be removed, and the error wraps ErrEventDBNotRemoved or
|
|
// ErrSidecarNotRemoved to say which was left, naming each file.
|
|
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()
|
|
}
|
|
}
|
|
}
|
|
|
|
path := m.dbPath(webhookID)
|
|
|
|
dbErr := removeFile(path)
|
|
sidecarErr := errors.Join(
|
|
removeFile(path+"-wal"),
|
|
removeFile(path+"-shm"),
|
|
)
|
|
|
|
if dbErr != nil {
|
|
return fmt.Errorf(
|
|
"%w: %w",
|
|
ErrEventDBNotRemoved, errors.Join(dbErr, sidecarErr),
|
|
)
|
|
}
|
|
|
|
if sidecarErr != nil {
|
|
return fmt.Errorf("%w: %w", ErrSidecarNotRemoved, sidecarErr)
|
|
}
|
|
|
|
m.log.Info(
|
|
"deleted per-webhook database",
|
|
"webhook_id", webhookID,
|
|
)
|
|
|
|
return nil
|
|
}
|
|
|
|
// removeFile removes path. A file that is already gone counts as
|
|
// removed; the error from any other failure names the file.
|
|
func removeFile(path string) error {
|
|
err := os.Remove(path)
|
|
if errors.Is(err, os.ErrNotExist) {
|
|
return nil
|
|
}
|
|
|
|
return err
|
|
}
|
|
|
|
// 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)
|
|
}
|
|
|
|
// getDB is GetDB, and CreateDB when isNew is true: the webhook has just
|
|
// been created, so a missing file is expected rather than lost.
|
|
func (m *WebhookDBManager) getDB(
|
|
webhookID string, isNew bool,
|
|
) (*gorm.DB, error) {
|
|
// Fast path: already open
|
|
if val, ok := m.dbs.Load(webhookID); ok {
|
|
return asGormDB(val, webhookID)
|
|
}
|
|
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
|
|
return m.getDBLocked(webhookID, isNew)
|
|
}
|
|
|
|
// getDBLocked is getDB's slow path, run with m.mu held. It looks in the
|
|
// cache again first: a caller that raced another one to the lock then
|
|
// gets its handle instead of opening a second one.
|
|
func (m *WebhookDBManager) getDBLocked(
|
|
webhookID string, isNew bool,
|
|
) (*gorm.DB, error) {
|
|
if val, ok := m.dbs.Load(webhookID); ok {
|
|
return asGormDB(val, webhookID)
|
|
}
|
|
|
|
// Checked before opening, which creates the file. See GetDB.
|
|
path := m.dbPath(webhookID)
|
|
replaced := !isNew && missingOrEmpty(path)
|
|
|
|
db, err := m.openDB(webhookID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
if replaced {
|
|
m.log.Warn(
|
|
"created a new, empty database",
|
|
"webhook_id", webhookID,
|
|
"path", path,
|
|
)
|
|
}
|
|
|
|
m.dbs.Store(webhookID, db)
|
|
|
|
return db, nil
|
|
}
|
|
|
|
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
|
|
}
|