Some checks failed
check / check (push) Has been cancelled
Rewrites retention_days=0 to the RetentionForeverDays sentinel (365 * 1000) in Webhook.BeforeSave, so the GORM column default cannot win the race. The reaper skips retain-forever webhooks before building any query. Also bounds the reaper's cutoff arithmetic: a time.Duration is int64 nanoseconds, so day counts above MaxFiniteRetentionDays (106751) overflowed and wrapped the cutoff into the future, where created_at < cutoff matched every row and the sweep deleted everything. parseRetentionDays now rejects finite values above the ceiling, and retentionCutoff saturates so rows written by older versions cannot reach it either. Views render RetentionLabel() rather than the raw sentinel.
310 lines
7.7 KiB
Go
310 lines
7.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/logger"
|
|
)
|
|
|
|
// hoursPerDay converts a RetentionDays count into hours for cutoff
|
|
// computation.
|
|
const hoursPerDay = 24
|
|
|
|
// 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.
|
|
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(_ context.Context) error {
|
|
r.stop()
|
|
|
|
return nil
|
|
},
|
|
})
|
|
}
|
|
|
|
// 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(),
|
|
)
|
|
}
|
|
|
|
func (r *RetentionReaper) stop() {
|
|
r.log.Info("retention reaper stopping")
|
|
|
|
if r.cancel != nil {
|
|
r.cancel()
|
|
}
|
|
|
|
r.wg.Wait()
|
|
r.log.Info("retention reaper stopped")
|
|
}
|
|
|
|
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, in foreign-key-safe order, the delivery
|
|
// results, deliveries, and events associated with events older than
|
|
// cutoff. Deletes are unscoped so rows are physically removed rather
|
|
// than soft-deleted, reclaiming disk. It returns the number of events
|
|
// deleted.
|
|
func reapExpired(db *gorm.DB, cutoff time.Time) (int64, error) {
|
|
// Fresh subqueries are built per statement to avoid reusing a
|
|
// mutated builder across executions.
|
|
expiredEventIDs := func() *gorm.DB {
|
|
return db.Model(&Event{}).
|
|
Select("id").
|
|
Where("created_at < ?", cutoff)
|
|
}
|
|
expiredDeliveryIDs := func() *gorm.DB {
|
|
return db.Model(&Delivery{}).
|
|
Select("id").
|
|
Where("event_id IN (?)", expiredEventIDs())
|
|
}
|
|
|
|
// 1. Delivery results whose delivery belongs to an expired event.
|
|
res := db.Unscoped().
|
|
Where("delivery_id IN (?)", expiredDeliveryIDs()).
|
|
Delete(&DeliveryResult{})
|
|
if res.Error != nil {
|
|
return 0, fmt.Errorf(
|
|
"deleting expired delivery results: %w",
|
|
res.Error,
|
|
)
|
|
}
|
|
|
|
// 2. Deliveries belonging to an expired event.
|
|
del := db.Unscoped().
|
|
Where("event_id IN (?)", expiredEventIDs()).
|
|
Delete(&Delivery{})
|
|
if del.Error != nil {
|
|
return 0, fmt.Errorf(
|
|
"deleting expired deliveries: %w",
|
|
del.Error,
|
|
)
|
|
}
|
|
|
|
// 3. The expired events themselves.
|
|
ev := db.Unscoped().
|
|
Where("created_at < ?", cutoff).
|
|
Delete(&Event{})
|
|
if ev.Error != nil {
|
|
return 0, fmt.Errorf(
|
|
"deleting expired events: %w",
|
|
ev.Error,
|
|
)
|
|
}
|
|
|
|
return ev.RowsAffected, nil
|
|
}
|