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
369 lines
9.7 KiB
Go
369 lines
9.7 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
|
|
|
|
// 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]
|
|
|
|
// Skip retain-forever webhooks before building any query.
|
|
// RetainsForever covers both the RetentionForeverDays
|
|
// sentinel and the non-positive values that predate it: the
|
|
// sentinel is a positive number, so without this the reaper
|
|
// would compute a cutoff a thousand years in the past and
|
|
// issue a DELETE matching nothing on every single sweep.
|
|
if wh.RetainsForever() {
|
|
continue
|
|
}
|
|
|
|
// Nothing to reap if the per-webhook database has never
|
|
// been created.
|
|
if !r.dbManager.DBExists(wh.ID) {
|
|
continue
|
|
}
|
|
|
|
r.reapWebhook(wh.ID, wh.RetentionDays)
|
|
}
|
|
}
|
|
|
|
// reapWebhook removes every expired event (and its dependents) from a
|
|
// single webhook's database.
|
|
func (r *RetentionReaper) reapWebhook(
|
|
webhookID string,
|
|
retentionDays int,
|
|
) {
|
|
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
|
|
}
|
|
|
|
cutoff, ok := retentionCutoff(time.Now(), retentionDays)
|
|
if !ok {
|
|
return
|
|
}
|
|
|
|
deleted, err := reapExpired(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 until none is left. It returns the number of events
|
|
// deleted.
|
|
func reapExpired(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
|
|
}
|
|
}
|
|
}
|
|
|
|
// 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})
|
|
}
|