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
379 lines
10 KiB
Go
379 lines
10 KiB
Go
package database
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log/slog"
|
|
"sync"
|
|
"time"
|
|
|
|
"go.uber.org/fx"
|
|
"gorm.io/gorm"
|
|
"sneak.berlin/go/webhooker/internal/config"
|
|
"sneak.berlin/go/webhooker/internal/lifecycle"
|
|
"sneak.berlin/go/webhooker/internal/logger"
|
|
)
|
|
|
|
// hoursPerDay converts a RetentionDays count into hours for cutoff
|
|
// computation.
|
|
const hoursPerDay = 24
|
|
|
|
// reapBatchSize is how many expired events one retention transaction
|
|
// deletes. A transaction holds the event database's write lock, which
|
|
// the receiver and the delivery workers wait for, so a large prune is
|
|
// split into transactions each short enough to finish well inside the
|
|
// busy timeout.
|
|
const reapBatchSize = 1000
|
|
|
|
// reapBatchPause is how long retention waits after one batch before
|
|
// starting the next. A writer waiting for the write lock checks for it
|
|
// again after at most 100 ms, so a longer pause lets it in between two
|
|
// batches instead of only after the whole prune.
|
|
const reapBatchPause = 200 * time.Millisecond
|
|
|
|
// RetentionReaperParams holds the fx dependencies for the
|
|
// RetentionReaper.
|
|
type RetentionReaperParams struct {
|
|
fx.In
|
|
|
|
Config *config.Config
|
|
Database *Database
|
|
DBManager *WebhookDBManager
|
|
Logger *logger.Logger
|
|
}
|
|
|
|
// RetentionReaper periodically deletes expired events (and their
|
|
// dependent deliveries and delivery results) from each per-webhook
|
|
// database, enforcing every webhook's RetentionDays. Rows are removed
|
|
// permanently so that per-webhook SQLite files do not grow without
|
|
// bound.
|
|
type RetentionReaper struct {
|
|
db *Database
|
|
dbManager *WebhookDBManager
|
|
log *slog.Logger
|
|
interval time.Duration
|
|
cancel context.CancelFunc
|
|
wg sync.WaitGroup
|
|
}
|
|
|
|
// NewRetentionReaper creates the retention reaper and registers its
|
|
// fx lifecycle hooks. The background sweep loop starts on OnStart and
|
|
// stops cleanly on OnStop via context cancellation.
|
|
func NewRetentionReaper(
|
|
lc fx.Lifecycle,
|
|
params RetentionReaperParams,
|
|
) *RetentionReaper {
|
|
r := &RetentionReaper{
|
|
db: params.Database,
|
|
dbManager: params.DBManager,
|
|
log: params.Logger.Get(),
|
|
interval: params.Config.RetentionSweepInterval,
|
|
}
|
|
|
|
r.registerHooks(lc)
|
|
|
|
return r
|
|
}
|
|
|
|
// registerHooks wires the reaper's start and stop into the fx
|
|
// lifecycle. The start hook's context is deliberately ignored (see
|
|
// start for why the sweep loop must not inherit it); the stop hook's
|
|
// context is honoured (see stop).
|
|
func (r *RetentionReaper) registerHooks(lc fx.Lifecycle) {
|
|
lc.Append(fx.Hook{
|
|
//nolint:contextcheck // Not inheriting the hook context is
|
|
// the point: see start.
|
|
OnStart: func(_ context.Context) error {
|
|
r.start()
|
|
|
|
return nil
|
|
},
|
|
OnStop: func(ctx context.Context) error {
|
|
return r.stop(ctx)
|
|
},
|
|
})
|
|
}
|
|
|
|
// start launches the background sweep loop.
|
|
//
|
|
// The loop's context is derived from context.Background(), NOT from
|
|
// the fx OnStart hook context. The hook context carries fx's start
|
|
// timeout (15s by default) and is cancelled once the start phase
|
|
// completes, so a loop derived from it dies 45 minutes before its
|
|
// first tick under the default one-hour sweep interval, leaving a
|
|
// reaper that never reaps. A long-lived goroutine must outlive the
|
|
// startup phase, so its lifetime is bounded by OnStop instead: stop
|
|
// cancels this context and waits on the WaitGroup.
|
|
func (r *RetentionReaper) start() {
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
r.cancel = cancel
|
|
|
|
r.wg.Add(1)
|
|
|
|
go r.run(ctx)
|
|
|
|
r.log.Info(
|
|
"retention reaper started",
|
|
"interval", r.interval.String(),
|
|
)
|
|
}
|
|
|
|
// stop cancels the sweep loop's context and waits for it to
|
|
// exit, bounded by the stop hook's context: a sweep wedged on a
|
|
// locked database must not hang the process past fx's stop
|
|
// timeout.
|
|
func (r *RetentionReaper) stop(ctx context.Context) error {
|
|
r.log.Info("retention reaper stopping")
|
|
|
|
if r.cancel != nil {
|
|
r.cancel()
|
|
}
|
|
|
|
err := lifecycle.WaitForShutdown(
|
|
ctx, r.log, "retention reaper", &r.wg,
|
|
)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
r.log.Info("retention reaper stopped")
|
|
|
|
return nil
|
|
}
|
|
|
|
func (r *RetentionReaper) run(ctx context.Context) {
|
|
defer r.wg.Done()
|
|
|
|
ticker := time.NewTicker(r.interval)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
r.sweep(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
// sweep lists every webhook from the main database and reaps expired
|
|
// rows from each per-webhook database that has a finite retention
|
|
// policy. Webhooks set to retain forever are skipped entirely.
|
|
func (r *RetentionReaper) sweep(ctx context.Context) {
|
|
var webhooks []Webhook
|
|
|
|
err := r.db.DB().
|
|
Model(&Webhook{}).
|
|
Find(&webhooks).Error
|
|
if err != nil {
|
|
r.log.Error(
|
|
"retention sweep: failed to list webhooks",
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
for i := range webhooks {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
default:
|
|
}
|
|
|
|
wh := webhooks[i]
|
|
|
|
// A missing database has nothing to reap. Restart recovery
|
|
// reports a lost one (see WebhookDBManager.GetDB).
|
|
if !r.dbManager.DBExists(wh.ID) {
|
|
continue
|
|
}
|
|
|
|
r.reapWebhook(ctx, wh.ID, wh.RetentionDays)
|
|
}
|
|
}
|
|
|
|
// reapWebhook removes every expired event (and its dependents) from a
|
|
// single webhook's database, or as many as it reaches before ctx is
|
|
// cancelled.
|
|
func (r *RetentionReaper) reapWebhook(
|
|
ctx context.Context,
|
|
webhookID string,
|
|
retentionDays int,
|
|
) {
|
|
// A retain-forever webhook has no cutoff, so its database is not
|
|
// even opened.
|
|
cutoff, ok := retentionCutoff(time.Now(), retentionDays)
|
|
if !ok {
|
|
return
|
|
}
|
|
|
|
db, err := r.dbManager.GetDB(webhookID)
|
|
if err != nil {
|
|
r.log.Error(
|
|
"retention sweep: failed to open webhook database",
|
|
"webhook_id", webhookID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
deleted, err := reapExpired(ctx, db, cutoff)
|
|
if err != nil {
|
|
r.log.Error(
|
|
"retention sweep: failed to reap expired events",
|
|
"webhook_id", webhookID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
if deleted > 0 {
|
|
r.log.Info(
|
|
"retention sweep: reaped expired events",
|
|
"webhook_id", webhookID,
|
|
"retention_days", retentionDays,
|
|
"events_deleted", deleted,
|
|
)
|
|
}
|
|
}
|
|
|
|
// retentionCutoff returns the timestamp before which a webhook's
|
|
// events have expired, and whether any cutoff applies at all. It
|
|
// reports false for a retain-forever policy, so no DELETE is issued.
|
|
//
|
|
// The day count is clamped to MaxFiniteRetentionDays first. This is
|
|
// defense in depth rather than decoration: a time.Duration is an int64
|
|
// nanosecond count, so an unclamped multiplication overflows above
|
|
// that ceiling and wraps the span negative. Subtracting a negative
|
|
// span moves the cutoff into the far future, where it matches every
|
|
// row in the database: the sweep then deletes every event, delivery,
|
|
// and delivery result, including ones created seconds ago. Rejecting
|
|
// out-of-range input at the form is the primary guard; saturating here
|
|
// means an old row, a migration, or a future call site cannot turn a
|
|
// too-large retention into total data loss.
|
|
func retentionCutoff(
|
|
now time.Time,
|
|
retentionDays int,
|
|
) (time.Time, bool) {
|
|
if retainsForever(retentionDays) {
|
|
return time.Time{}, false
|
|
}
|
|
|
|
if retentionDays > MaxFiniteRetentionDays {
|
|
retentionDays = MaxFiniteRetentionDays
|
|
}
|
|
|
|
return now.Add(
|
|
-time.Duration(retentionDays*hoursPerDay) * time.Hour,
|
|
), true
|
|
}
|
|
|
|
// reapExpired hard-deletes the events older than cutoff, with their
|
|
// deliveries and delivery results, reapBatchSize events per
|
|
// transaction with reapBatchPause between transactions, until none is
|
|
// left. Once ctx is cancelled it returns after the batch in hand,
|
|
// leaving the rest to the next sweep, so stopping the app does not
|
|
// wait for a long prune. It returns the number of events deleted.
|
|
func reapExpired(
|
|
ctx context.Context, db *gorm.DB, cutoff time.Time,
|
|
) (int64, error) {
|
|
var total int64
|
|
|
|
for {
|
|
var eventIDs []string
|
|
|
|
err := db.Transaction(func(tx *gorm.DB) error {
|
|
err := tx.Unscoped().Model(&Event{}).
|
|
Where("created_at < ?", cutoff).
|
|
Limit(reapBatchSize).
|
|
Pluck("id", &eventIDs).Error
|
|
if err != nil {
|
|
return fmt.Errorf("selecting expired events: %w", err)
|
|
}
|
|
|
|
if len(eventIDs) == 0 {
|
|
return nil
|
|
}
|
|
|
|
return deleteEvents(tx, eventIDs)
|
|
})
|
|
if err != nil {
|
|
return total, err
|
|
}
|
|
|
|
total += int64(len(eventIDs))
|
|
|
|
if len(eventIDs) < reapBatchSize {
|
|
return total, nil
|
|
}
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
return total, nil
|
|
case <-time.After(reapBatchPause):
|
|
}
|
|
}
|
|
}
|
|
|
|
// deleteEvents hard-deletes the given events and, in foreign-key-safe
|
|
// order before them, their delivery results and deliveries, then adds
|
|
// what it deleted to the running totals. It runs on reapExpired's
|
|
// transaction, so the totals change exactly when the rows do. Deletes
|
|
// are unscoped so rows are physically removed rather than
|
|
// soft-deleted, reclaiming disk.
|
|
func deleteEvents(tx *gorm.DB, eventIDs []string) error {
|
|
// 1. The delivery results of the events' deliveries.
|
|
err := tx.Unscoped().
|
|
Where("delivery_id IN (?)", tx.Unscoped().Model(&Delivery{}).
|
|
Select("id").
|
|
Where("event_id IN ?", eventIDs)).
|
|
Delete(&DeliveryResult{}).Error
|
|
if err != nil {
|
|
return fmt.Errorf("deleting expired delivery results: %w", err)
|
|
}
|
|
|
|
// 2. The events' deliveries, after counting them, and the failed
|
|
// ones among them, per target. The status is tested in the select
|
|
// list rather than the WHERE clause: there, SQLite would read every
|
|
// failed delivery the webhook has through the status index,
|
|
// instead of only these through the event_id index.
|
|
var removed []TargetTotals
|
|
|
|
err = tx.Unscoped().Model(&Delivery{}).
|
|
Select("target_id, count(*) AS deliveries_removed, "+
|
|
"count(CASE WHEN status = ? THEN 1 END) AS failed_removed",
|
|
DeliveryStatusFailed).
|
|
Where("event_id IN ?", eventIDs).
|
|
Group("target_id").
|
|
Find(&removed).Error
|
|
if err != nil {
|
|
return fmt.Errorf("counting expired deliveries: %w", err)
|
|
}
|
|
|
|
err = tx.Unscoped().
|
|
Where("event_id IN ?", eventIDs).
|
|
Delete(&Delivery{}).Error
|
|
if err != nil {
|
|
return fmt.Errorf("deleting expired deliveries: %w", err)
|
|
}
|
|
|
|
// 3. The events themselves.
|
|
ev := tx.Unscoped().Where("id IN ?", eventIDs).Delete(&Event{})
|
|
if ev.Error != nil {
|
|
return fmt.Errorf("deleting expired events: %w", ev.Error)
|
|
}
|
|
|
|
for i := range removed {
|
|
err = AddTargetTotals(tx, removed[i])
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return AddEventTotals(tx, EventTotals{EventsRemoved: ev.RowsAffected})
|
|
}
|