From 190cabe0f2e92fc7c2059b2c16a5810eebd61e39 Mon Sep 17 00:00:00 2001 From: clawbot Date: Sun, 9 Aug 2026 02:25:33 +0000 Subject: [PATCH] Evict archive writers on deletion and sweep idle archives (closes #89) The per-webhook archiveWriter registry in the database delivery target was never evicted, so a deleted webhook's writer -- and any archive file handle open within its debounce window -- lingered for the process lifetime. Separately, expiry pruning ran only when an archive was (re)opened, and reopens only happen on writes, so an archive belonging to a webhook that stopped receiving events kept its expired rows forever. Eviction: a new one-method delivery.WebhookEvictor interface (kept separate from Notifier: archiving lifecycle is not notification) is implemented by the Engine and injected into the handlers. Deleting a webhook, or deleting its last database target, drops the writer from the registry and closes its handle under the writer's own mutex, so eviction can never race an in-flight write. An evicted writer refuses further writes rather than reopening a file nothing holds. The archive file is deliberately left on disk: it is long-term storage an operator may want to keep or move away, and destroying it as a side effect of deleting a webhook would be unrecoverable. Idle sweep: a new ArchiveSweeper, modelled on the event RetentionReaper (fx lifecycle hooks, cancellable context, WaitGroup, ticker loop), prunes archives whose database target declares a positive expiry. It reuses the existing RETENTION_SWEEP_INTERVAL rather than adding a config key. It never creates an archive -- a missing file is skipped, and the reopen uses SQLite mode=rw so the file cannot be conjured even if it disappears mid-sweep -- routes the prune through the per-webhook writer so its mutex orders the sweep against concurrent writes, and leaves the archive closed so the move-the-file-away workflow keeps working. A failure for one webhook is logged and the sweep continues. Archives with no expiry or the expiry "never" are untouched. --- README.md | 24 + TODO.md | 5 + cmd/webhooker/main.go | 7 + internal/delivery/archive_sweeper.go | 189 ++++++++ internal/delivery/archive_sweeper_test.go | 425 ++++++++++++++++++ internal/delivery/engine.go | 34 ++ internal/delivery/export_test.go | 84 ++++ internal/delivery/target.go | 5 +- internal/delivery/target_database.go | 94 +++- internal/delivery/target_database_archive.go | 106 ++++- .../delivery/target_database_evict_test.go | 130 ++++++ internal/handlers/handlers.go | 3 + internal/handlers/handlers_test.go | 33 ++ internal/handlers/source_delete_test.go | 305 +++++++++++++ internal/handlers/source_management.go | 81 +++- 15 files changed, 1514 insertions(+), 11 deletions(-) create mode 100644 internal/delivery/archive_sweeper.go create mode 100644 internal/delivery/archive_sweeper_test.go create mode 100644 internal/delivery/target_database_evict_test.go create mode 100644 internal/handlers/source_delete_test.go diff --git a/README.md b/README.md index 1514b61..b43bafc 100644 --- a/README.md +++ b/README.md @@ -531,6 +531,30 @@ older than the expiry are pruned each time the archive is (re)opened. An archive write failure is never silent success: the delivery records a failed attempt with the error and is marked failed. +Because reopens only happen on writes, an archive belonging to a webhook +that has stopped receiving events would never be pruned. A background +**archive sweeper** closes that gap: on the same interval as the event +retention reaper (`RETENTION_SWEEP_INTERVAL`) it prunes every archive +whose database target declares a positive expiry, whether or not the +webhook is still receiving traffic. The sweep never creates an archive — +a webhook whose archive file does not yet exist is skipped, not +initialised — it takes the same per-webhook lock the write path uses, so +it can never interleave with a write, and it leaves the archive closed +afterwards so the move-the-file-away workflow keeps working. Archives +with no expiry, or the expiry `never`, are not touched by the sweep at +all. + +Deleting a webhook releases its archive: the delivery engine's cached +archive writer is dropped and its file handle closed, so nothing lingers +after the webhook is gone. The archive **file itself is deliberately +left on disk**. Unlike the event database — per-webhook working storage +that is hard-deleted with the webhook — an archive is long-term storage +an operator may still want to keep or move away for offline retention, +and destroying it as a side effect of deleting a webhook would be +unrecoverable. Removing `archive-{webhookID}.db` is the operator's call. +Deleting a webhook's last `database` target releases the writer the same +way, and for the same reason leaves the file alone. + The **Slack target type** sends webhook events as formatted messages to any Slack-compatible incoming webhook URL (works with Slack, Mattermost, and other compatible services). Each message includes event metadata diff --git a/TODO.md b/TODO.md index c869649..ec3dce2 100644 --- a/TODO.md +++ b/TODO.md @@ -28,6 +28,11 @@ databases currently grow without bound. # Completed Steps +- 2026-08-09 Archive writer lifecycle (#89): deleting a webhook (or its + last `database` target) evicts the cached archive writer and closes + its handle while deliberately leaving `archive-{webhookID}.db` on + disk, and a new `ArchiveSweeper` prunes idle archives on the existing + `RETENTION_SWEEP_INTERVAL` without ever creating an archive file - 2026-08-07 Update golangci-lint to v2.12.2 (Docker image digest in `Dockerfile`, release-archive sha256 pins in `script/bootstrap`), adopt the canonical `.golangci.yml` (v2 `linters.settings` layout so diff --git a/cmd/webhooker/main.go b/cmd/webhooker/main.go index 59fc655..793f44c 100644 --- a/cmd/webhooker/main.go +++ b/cmd/webhooker/main.go @@ -40,9 +40,15 @@ func main() { handlers.New, middleware.New, delivery.New, + delivery.NewArchiveSweeper, // Wire *delivery.Engine as delivery.Notifier so the // webhook handler can notify the engine of new deliveries. func(e *delivery.Engine) delivery.Notifier { return e }, + // Wire *delivery.Engine as delivery.WebhookEvictor so + // deleting a webhook releases its archive writer. + func(e *delivery.Engine) delivery.WebhookEvictor { + return e + }, server.New, ), fx.Invoke( @@ -50,6 +56,7 @@ func main() { *server.Server, *delivery.Engine, *database.RetentionReaper, + *delivery.ArchiveSweeper, ) { }, ), diff --git a/internal/delivery/archive_sweeper.go b/internal/delivery/archive_sweeper.go new file mode 100644 index 0000000..09efd2d --- /dev/null +++ b/internal/delivery/archive_sweeper.go @@ -0,0 +1,189 @@ +package delivery + +import ( + "context" + "log/slog" + "sync" + "time" + + "go.uber.org/fx" + "sneak.berlin/go/webhooker/internal/config" + "sneak.berlin/go/webhooker/internal/database" + "sneak.berlin/go/webhooker/internal/logger" +) + +// ArchiveSweeperParams holds the fx dependencies for the +// ArchiveSweeper. +type ArchiveSweeperParams struct { + fx.In + + Config *config.Config + Database *database.Database + Engine *Engine + Logger *logger.Logger +} + +// ArchiveSweeper periodically prunes expired rows from +// per-webhook archive databases whose database target carries a +// positive expiry. +// +// Without it, pruning happens only when an archive is +// (re)opened, and archives are only ever reopened by writes: an +// archive belonging to a webhook that has stopped receiving +// events would keep its expired rows forever. The sweep closes +// that gap without changing anything for archives whose expiry +// is unset or "never". +// +// It reuses Config.RetentionSweepInterval rather than +// introducing a second interval: this is a retention sweep with +// the same semantics as the event retention reaper. +type ArchiveSweeper struct { + db *database.Database + eng *Engine + log *slog.Logger + interval time.Duration + cancel context.CancelFunc + wg sync.WaitGroup +} + +// NewArchiveSweeper creates the archive sweeper and registers +// its fx lifecycle hooks. The background sweep loop starts on +// OnStart and stops cleanly on OnStop via context cancellation. +func NewArchiveSweeper( + lc fx.Lifecycle, + params ArchiveSweeperParams, +) *ArchiveSweeper { + s := &ArchiveSweeper{ + db: params.Database, + eng: params.Engine, + log: params.Logger.Get(), + interval: params.Config.RetentionSweepInterval, + } + + lc.Append(fx.Hook{ + OnStart: func(ctx context.Context) error { + s.start(ctx) + + return nil + }, + OnStop: func(_ context.Context) error { + s.stop() + + return nil + }, + }) + + return s +} + +func (s *ArchiveSweeper) start(ctx context.Context) { + ctx, cancel := context.WithCancel(ctx) + s.cancel = cancel + + s.wg.Add(1) + + go s.run(ctx) + + s.log.Info( + "archive sweeper started", + "interval", s.interval.String(), + ) +} + +func (s *ArchiveSweeper) stop() { + s.log.Info("archive sweeper stopping") + + if s.cancel != nil { + s.cancel() + } + + s.wg.Wait() + s.log.Info("archive sweeper stopped") +} + +func (s *ArchiveSweeper) run(ctx context.Context) { + defer s.wg.Done() + + ticker := time.NewTicker(s.interval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + s.sweep(ctx) + } + } +} + +// sweep prunes every archive whose database target declares a +// positive expiry. Targets belonging to a deleted webhook are +// soft-deleted along with it, so GORM's default scope already +// excludes them. +// +// A failure for one webhook is logged and the sweep continues, +// matching how the write path already treats a prune error as +// non-fatal. +func (s *ArchiveSweeper) sweep(ctx context.Context) { + var targets []database.Target + + err := s.db.DB(). + Model(&database.Target{}). + Where("type = ?", database.TargetTypeDatabase). + Find(&targets).Error + if err != nil { + s.log.Error( + "archive sweep: failed to list database targets", + "error", err, + ) + + return + } + + for i := range targets { + select { + case <-ctx.Done(): + return + default: + } + + s.sweepTarget(&targets[i]) + } +} + +// sweepTarget prunes the archive of a single database target. +// A missing, empty, or "never" expiry parses as a zero duration +// and is skipped entirely, so those archives keep exactly the +// behaviour they had before the sweep existed. +func (s *ArchiveSweeper) sweepTarget(target *database.Target) { + expiry, err := parseArchiveExpiry(target.Config) + if err != nil { + s.log.Error( + "archive sweep: invalid database target config", + "webhook_id", target.WebhookID, + "target_id", target.ID, + "error", err, + ) + + return + } + + if expiry <= 0 { + return + } + + if s.eng == nil || s.eng.dbTarget == nil { + return + } + + err = s.eng.dbTarget.sweepWebhook(target.WebhookID, expiry) + if err != nil { + s.log.Error( + "archive sweep: failed to prune archive", + "webhook_id", target.WebhookID, + "target_id", target.ID, + "error", err, + ) + } +} diff --git a/internal/delivery/archive_sweeper_test.go b/internal/delivery/archive_sweeper_test.go new file mode 100644 index 0000000..ea9e87e --- /dev/null +++ b/internal/delivery/archive_sweeper_test.go @@ -0,0 +1,425 @@ +package delivery_test + +import ( + "context" + "database/sql" + "fmt" + "net/http" + "path/filepath" + "sync" + "testing" + "time" + + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/driver/sqlite" + "gorm.io/gorm" + "gorm.io/gorm/clause" + _ "modernc.org/sqlite" // Pure Go SQLite driver. + "sneak.berlin/go/webhooker/internal/database" + "sneak.berlin/go/webhooker/internal/delivery" +) + +const ( + // sweepRowOld and sweepRowNew are the event ids + // seedArchiveRows assigns to the first and second seeded + // rows. + sweepRowOld = "ev-0" + sweepRowNew = "ev-1" + + // sweepConcurrentWrites is how many deliveries the + // concurrent write-plus-sweep test races against the sweep. + sweepConcurrentWrites = 20 +) + +// sweeperEnv bundles the pieces an archive sweep test drives: +// a main configuration database holding webhooks and targets, a +// delivery engine owning the archive writer registry, and the +// data directory the archive files live in. +type sweeperEnv struct { + sweeper *delivery.ArchiveSweeper + eng *delivery.Engine + mainDB *database.Database + dataDir string +} + +func setupSweeperTest(t *testing.T) *sweeperEnv { + t.Helper() + + dataDir := t.TempDir() + log := archiveTestLogger() + + sqlDB, err := sql.Open( + "sqlite", + fmt.Sprintf( + "file:%s?mode=rwc", + filepath.Join(dataDir, "main.db"), + ), + ) + require.NoError(t, err) + + t.Cleanup(func() { _ = sqlDB.Close() }) + + gdb, err := gorm.Open( + sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, + ) + require.NoError(t, err) + + mainDB := database.NewTestDatabase(gdb) + require.NoError(t, mainDB.Migrate()) + + eng := delivery.NewTestEngineWithDB( + mainDB, + database.NewTestWebhookDBManager(dataDir), + log, + &http.Client{Timeout: 5 * time.Second}, + 1, + ) + + return &sweeperEnv{ + sweeper: delivery.NewTestArchiveSweeper( + mainDB, eng, log, + ), + eng: eng, + mainDB: mainDB, + dataDir: dataDir, + } +} + +// archivePath returns where the engine keeps a webhook's +// archive file. +func (env *sweeperEnv) archivePath(webhookID string) string { + return filepath.Join( + env.dataDir, fmt.Sprintf("archive-%s.db", webhookID), + ) +} + +// seedDatabaseTarget creates a webhook with one database target +// carrying the given target config JSON, and returns the +// webhook id. +func (env *sweeperEnv) seedDatabaseTarget( + t *testing.T, configJSON string, +) string { + t.Helper() + + wh := &database.Webhook{ + UserID: uuid.New().String(), + Name: "sweep-test", + } + require.NoError( + t, + env.mainDB.DB(). + Omit(clause.Associations). + Create(wh).Error, + ) + + tgt := &database.Target{ + WebhookID: wh.ID, + Name: "archive", + Type: database.TargetTypeDatabase, + Active: true, + Config: configJSON, + } + require.NoError( + t, + env.mainDB.DB(). + Omit(clause.Associations). + Create(tgt).Error, + ) + + return wh.ID +} + +// seedArchiveRows creates the archive file for a webhook and +// inserts one row per supplied archived-at timestamp, returning +// the archive path. The handle is closed before returning, so +// the archive is idle exactly as it would be with no traffic. +func (env *sweeperEnv) seedArchiveRows( + t *testing.T, webhookID string, archivedAt ...time.Time, +) string { + t.Helper() + + path := env.archivePath(webhookID) + + sqlDB, err := sql.Open( + "sqlite", fmt.Sprintf("file:%s?mode=rwc", path), + ) + require.NoError(t, err) + + gdb, err := gorm.Open( + sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, + ) + require.NoError(t, err) + + require.NoError( + t, gdb.AutoMigrate(&delivery.ExportArchivedEvent{}), + ) + + for i, at := range archivedAt { + row := delivery.ExportArchivedEvent{ + EventID: fmt.Sprintf("ev-%d", i), + WebhookID: webhookID, + Method: http.MethodPost, + Body: `{"seeded":true}`, + ArchivedAt: at, + } + require.NoError(t, gdb.Create(&row).Error) + } + + require.NoError(t, sqlDB.Close()) + + return path +} + +// archivedEventIDs returns the event ids currently stored in an +// archive file, read through a separate read-only handle. +func archivedEventIDs( + t *testing.T, path string, +) []string { + t.Helper() + + var rows []delivery.ExportArchivedEvent + + rdb := openArchiveDBForRead(t, path) + require.NoError(t, rdb.Order("event_id").Find(&rows).Error) + + ids := make([]string, 0, len(rows)) + for i := range rows { + ids = append(ids, rows[i].EventID) + } + + return ids +} + +// TestArchiveSweep_PrunesIdleArchive is the core regression +// test for this issue: an archive that receives no further +// writes must still lose its expired rows. Before the sweeper +// existed, pruning only ever ran on a write-triggered reopen, +// so an idle archive kept expired rows forever. +func TestArchiveSweep_PrunesIdleArchive(t *testing.T) { + t.Parallel() + + env := setupSweeperTest(t) + + webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + + now := time.Now() + path := env.seedArchiveRows( + t, webhookID, + now.Add(-48*time.Hour), + now.Add(-time.Minute), + ) + + require.Equal( + t, []string{sweepRowOld, sweepRowNew}, + archivedEventIDs(t, path), + ) + + env.sweeper.ExportSweep(context.Background()) + + assert.Equal( + t, []string{sweepRowNew}, archivedEventIDs(t, path), + "the sweep should prune rows older than the expiry "+ + "from an idle archive and keep the rest", + ) +} + +// TestArchiveSweep_LeavesArchiveClosed proves the sweep does +// not hold the archive open afterwards, so an operator can +// still move the file away for offline retention. +func TestArchiveSweep_LeavesArchiveClosed(t *testing.T) { + t.Parallel() + + env := setupSweeperTest(t) + + webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + env.seedArchiveRows( + t, webhookID, time.Now().Add(-48*time.Hour), + ) + + env.sweeper.ExportSweep(context.Background()) + + assert.False( + t, env.eng.ExportArchiveHandleOpen(webhookID), + "an idle archive must end the sweep closed", + ) +} + +// TestArchiveSweep_NeverExpiryUntouched proves the sweep is a +// no-op for the default retention policy, so archives with no +// expiry (or the literal "never") behave exactly as before. +func TestArchiveSweep_NeverExpiryUntouched(t *testing.T) { + t.Parallel() + + for _, configJSON := range []string{ + `{"expiry":"never"}`, + `{"expiry":""}`, + "", + } { + env := setupSweeperTest(t) + + webhookID := env.seedDatabaseTarget(t, configJSON) + path := env.seedArchiveRows( + t, webhookID, + time.Now().Add(-10000*time.Hour), + ) + + env.sweeper.ExportSweep(context.Background()) + + assert.Equal( + t, []string{sweepRowOld}, archivedEventIDs(t, path), + "config %q must keep rows forever", configJSON, + ) + } +} + +// TestArchiveSweep_DoesNotCreateArchiveFile proves the sweep +// never conjures an archive: a webhook with a database target +// that has never received an event must still have no archive +// file (nor SQLite sidecar) after a sweep. +func TestArchiveSweep_DoesNotCreateArchiveFile(t *testing.T) { + t.Parallel() + + env := setupSweeperTest(t) + + webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + path := env.archivePath(webhookID) + + require.NoFileExists(t, path) + + env.sweeper.ExportSweep(context.Background()) + + for _, suffix := range []string{"", "-wal", "-shm"} { + assert.NoFileExists( + t, path+suffix, + "the sweep must not create an archive file", + ) + } +} + +// TestArchiveSweep_DoesNotCreateAfterWriterExists covers the +// same guarantee once a writer is cached in the registry but +// the file itself is still absent (for instance because the +// operator moved the archive away). +func TestArchiveSweep_DoesNotCreateAfterWriterExists( + t *testing.T, +) { + t.Parallel() + + env := setupSweeperTest(t) + + webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + + path, err := env.eng.ExportEnsureArchiveWriter(webhookID) + require.NoError(t, err) + require.NoFileExists(t, path) + + env.sweeper.ExportSweep(context.Background()) + + assert.NoFileExists(t, path) +} + +// TestArchiveSweep_SkipsDeletedWebhookTargets proves that the +// sweep ignores targets soft-deleted along with their webhook, +// so a deleted webhook's archive is never reopened. +func TestArchiveSweep_SkipsDeletedWebhookTargets(t *testing.T) { + t.Parallel() + + env := setupSweeperTest(t) + + webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + path := env.seedArchiveRows( + t, webhookID, time.Now().Add(-48*time.Hour), + ) + + require.NoError( + t, + env.mainDB.DB(). + Where("webhook_id = ?", webhookID). + Delete(&database.Target{}).Error, + ) + + env.sweeper.ExportSweep(context.Background()) + + assert.Equal( + t, []string{sweepRowOld}, archivedEventIDs(t, path), + "a deleted target's archive must be left alone", + ) +} + +// TestArchiveSweep_ConcurrentWrites proves the sweep serialises +// against writes through the per-webhook writer mutex. Run +// under -race, an unsynchronised sweep would be caught here. +func TestArchiveSweep_ConcurrentWrites(t *testing.T) { + t.Parallel() + + env := setupSweeperTest(t) + + webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + + webhookDB := testWebhookDB(t) + + // The deliveries are seeded up front, on the test's own + // goroutine: the seed helpers assert, and testify assertions + // must not run off the test goroutine. + deliveries := make( + []*database.Delivery, 0, sweepConcurrentWrites, + ) + + for range sweepConcurrentWrites { + event := seedEvent(t, webhookDB, `{"n":1}`) + event.WebhookID = webhookID + + deliveries = append( + deliveries, + seedDatabaseTargetDelivery( + t, webhookDB, event, `{"expiry":"1h"}`, + ), + ) + } + + var wg sync.WaitGroup + + wg.Add(2) + + go func() { + defer wg.Done() + + for _, d := range deliveries { + env.eng.ExportDeliverDatabase(webhookDB, d) + } + }() + + go func() { + defer wg.Done() + + for range sweepConcurrentWrites { + env.sweeper.ExportSweep(context.Background()) + } + }() + + wg.Wait() + + assert.FileExists(t, env.archivePath(webhookID)) +} + +// TestArchiveSweeper_StopsCleanly proves the background loop +// exits on OnStop rather than leaking a goroutine. +func TestArchiveSweeper_StopsCleanly(t *testing.T) { + t.Parallel() + + env := setupSweeperTest(t) + + webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + env.seedArchiveRows( + t, webhookID, time.Now().Add(-48*time.Hour), + ) + + env.sweeper.ExportSetInterval(time.Millisecond) + env.sweeper.ExportStart(context.Background()) + + // stop blocks on the loop's WaitGroup, so returning at all + // proves the loop observed the cancellation and exited. + env.sweeper.ExportStop() +} diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index 6011558..f45ca3d 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -94,6 +94,23 @@ type Notifier interface { Notify(tasks []Task) } +// WebhookEvictor releases the delivery engine's per-webhook +// state for a webhook that no longer needs it — currently the +// cached archive writer of the database target, whose open +// file handle would otherwise outlive the webhook. +// +// It is deliberately separate from Notifier and deliberately +// one method wide: archiving lifecycle is not notification, and +// a single-method interface keeps the handlers package free of +// any dependency on the engine's internals while staying +// trivially fakeable in tests. +// +// EvictWebhook never deletes an archive file. It is idempotent +// and is a no-op for a webhook with no engine state. +type WebhookEvictor interface { + EvictWebhook(webhookID string) +} + // EngineParams are the fx dependencies for the delivery // engine. type EngineParams struct { @@ -127,6 +144,10 @@ type Engine struct { // httpTarget is retained so tests can reach the HTTP // target's shared client and circuit breakers. httpTarget *httpTarget + + // dbTarget is retained so the engine can reach the archive + // writer registry for webhook eviction and the idle sweep. + dbTarget *databaseTarget } // New creates and registers the delivery engine with the @@ -182,6 +203,19 @@ func (e *Engine) Notify(tasks []Task) { } } +// EvictWebhook implements WebhookEvictor. It releases the +// engine's per-webhook archiving state: the database target's +// cached archive writer is dropped from the registry and its +// file handle closed. The archive file itself is left on disk +// — it is long-term storage the operator owns. +func (e *Engine) EvictWebhook(webhookID string) { + if e.dbTarget == nil { + return + } + + e.dbTarget.evict(webhookID) +} + // ScheduleRetry schedules a task to be re-enqueued onto the // retry channel after delay. It implements the Scheduler // interface the targets use to own their durable retries. diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index 739eb03..606aeda 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -328,6 +328,90 @@ func (e *ExportArchiveWriter) DB() *gorm.DB { return e.w.db } +// ExportHasArchiveWriter reports whether the database target +// currently caches an archive writer for a webhook. +func (e *Engine) ExportHasArchiveWriter( + webhookID string, +) bool { + e.dbTarget.mu.Lock() + defer e.dbTarget.mu.Unlock() + + _, ok := e.dbTarget.writers[webhookID] + + return ok +} + +// ExportArchiveHandleOpen reports whether the cached archive +// writer for a webhook holds an open database handle. It +// returns false when no writer is cached. +func (e *Engine) ExportArchiveHandleOpen( + webhookID string, +) bool { + e.dbTarget.mu.Lock() + w, ok := e.dbTarget.writers[webhookID] + e.dbTarget.mu.Unlock() + + if !ok { + return false + } + + w.mu.Lock() + defer w.mu.Unlock() + + return w.db != nil +} + +// ExportEnsureArchiveWriter creates (if needed) and returns the +// archive file path of the cached writer for a webhook, so a +// test can prime the registry the way a delivery would. +func (e *Engine) ExportEnsureArchiveWriter( + webhookID string, +) (string, error) { + w, err := e.dbTarget.writerFor(webhookID) + if err != nil { + return "", err + } + + return w.path, nil +} + +// NewTestArchiveSweeper builds an ArchiveSweeper backed by the +// given main database and engine, without the fx lifecycle. +// Intended for tests. +func NewTestArchiveSweeper( + db *database.Database, + eng *Engine, + log *slog.Logger, +) *ArchiveSweeper { + return &ArchiveSweeper{ + db: db, + eng: eng, + log: log, + interval: time.Hour, + } +} + +// ExportSweep runs a single archive sweep synchronously for +// tests. +func (s *ArchiveSweeper) ExportSweep(ctx context.Context) { + s.sweep(ctx) +} + +// ExportStart starts the sweeper's background loop for tests. +func (s *ArchiveSweeper) ExportStart(ctx context.Context) { + s.start(ctx) +} + +// ExportStop stops the sweeper's background loop for tests. +func (s *ArchiveSweeper) ExportStop() { + s.stop() +} + +// ExportSetInterval overrides the sweep interval for tests. +func (s *ArchiveSweeper) ExportSetInterval(d time.Duration) { + s.interval = d +} + // ExportParseArchiveExpiry exposes parseArchiveExpiry. func ExportParseArchiveExpiry( configJSON string, diff --git a/internal/delivery/target.go b/internal/delivery/target.go index 264ce2f..9e131b2 100644 --- a/internal/delivery/target.go +++ b/internal/delivery/target.go @@ -90,12 +90,15 @@ func (e *Engine) initTargets(client *http.Client) { client: client, } + dbT := &databaseTarget{eng: e} + e.httpTarget = httpT + e.dbTarget = dbT e.targets = map[database.TargetType]Target{ database.TargetTypeHTTP: httpT, database.TargetTypeSlack: slackT, - database.TargetTypeDatabase: &databaseTarget{eng: e}, + database.TargetTypeDatabase: dbT, database.TargetTypeLog: &logTarget{eng: e}, } } diff --git a/internal/delivery/target_database.go b/internal/delivery/target_database.go index 0443aaa..06e7db0 100644 --- a/internal/delivery/target_database.go +++ b/internal/delivery/target_database.go @@ -5,6 +5,7 @@ import ( "fmt" "path/filepath" "sync" + "time" "gorm.io/gorm" "sneak.berlin/go/webhooker/internal/database" @@ -111,15 +112,11 @@ func (t *databaseTarget) archive(d *database.Delivery) error { func (t *databaseTarget) writerFor( webhookID string, ) (*archiveWriter, error) { - if t.eng.dbManager == nil { - return nil, errArchiveNoDataDir + path, err := t.archivePath(webhookID) + if err != nil { + return nil, err } - dir := filepath.Dir(t.eng.dbManager.DBPath(webhookID)) - path := filepath.Join( - dir, fmt.Sprintf("archive-%s.db", webhookID), - ) - t.mu.Lock() defer t.mu.Unlock() @@ -135,3 +132,86 @@ func (t *databaseTarget) writerFor( return w, nil } + +// archivePath returns the archive file path for a webhook: it +// lives beside the per-webhook event database in the data +// directory. It does not touch the filesystem. +func (t *databaseTarget) archivePath( + webhookID string, +) (string, error) { + if t.eng.dbManager == nil { + return "", errArchiveNoDataDir + } + + dir := filepath.Dir(t.eng.dbManager.DBPath(webhookID)) + + return filepath.Join( + dir, fmt.Sprintf("archive-%s.db", webhookID), + ), nil +} + +// evict drops a webhook's archive writer from the registry and +// closes its handle, so a deleted webhook does not leave a +// writer (and an open archive handle within its debounce +// window) alive for the process lifetime. +// +// The map entry is removed under the registry lock, which is +// then released before the handle is closed under the writer's +// own lock: that ordering keeps the registry available to other +// webhooks while an in-flight write on this one drains, and +// closing under the writer's lock means eviction can never race +// a write. +// +// Eviction is idempotent and silent for a webhook with no +// writer, which is the common case: a webhook with no database +// target never creates one. It never deletes the archive file. +func (t *databaseTarget) evict(webhookID string) { + t.mu.Lock() + + w, ok := t.writers[webhookID] + if ok { + delete(t.writers, webhookID) + } + + t.mu.Unlock() + + if !ok { + return + } + + w.evict() + + t.eng.log.Info( + "evicted archive writer", + "webhook_id", webhookID, + "path", w.path, + ) +} + +// sweepWebhook prunes one webhook's archive of rows older than +// expiry, without requiring a write. It returns nil (nothing to +// do) when the archive file does not exist, so a sweep never +// creates an archive for a webhook that has a database target +// but has never received an event. +func (t *databaseTarget) sweepWebhook( + webhookID string, expiry time.Duration, +) error { + path, err := t.archivePath(webhookID) + if err != nil { + return err + } + + // Check before taking a writer at all: a webhook whose + // archive has never been created gets no writer, no handle, + // and no file. + if !fileExists(path) { + return nil + } + + w, err := t.writerFor(webhookID) + if err != nil { + return err + } + + return w.sweepExpired(expiry) +} diff --git a/internal/delivery/target_database_archive.go b/internal/delivery/target_database_archive.go index 547a0af..1b8d056 100644 --- a/internal/delivery/target_database_archive.go +++ b/internal/delivery/target_database_archive.go @@ -24,6 +24,20 @@ const archiveExpiryNever = "never" // offline archiving, but never more than once per this window. const archiveReopenDebounce = time.Second +const ( + // archiveModeCreate is the SQLite URI mode used by the write + // path: open the archive file, creating it if missing, so a + // first write (or a write after the operator moved the file + // away) recreates it. + archiveModeCreate = "rwc" + + // archiveModeExisting is the SQLite URI mode used by the idle + // sweep: open read-write but never create. A sweep must never + // conjure an empty archive file for a webhook that has a + // database target but has never received an event. + archiveModeExisting = "rw" +) + var ( // errArchiveMissingWebhookID is returned when an event to // archive has no webhook id to key its archive file on. @@ -44,6 +58,15 @@ var ( errArchiveExpiryNotPositive = errors.New( "expiry must be a positive duration or \"never\"", ) + + // errArchiveWriterEvicted is returned when a writer that has + // been evicted (its webhook was deleted, or its last database + // target was removed) is used again. An evicted writer is no + // longer in the registry, so reopening its file would leak a + // handle nothing owns. + errArchiveWriterEvicted = errors.New( + "archive writer has been evicted", + ) ) // databaseTargetConfig is the optional per-target JSON config @@ -161,6 +184,12 @@ type archiveWriter struct { db *gorm.DB lastReopen time.Time reopens int + + // evicted marks a writer that has been removed from the + // per-webhook registry. Its handle is closed and it must + // never open the file again: nothing holds it any more, so a + // reopen would leak the handle for the process lifetime. + evicted bool } // newArchiveWriter builds an archiveWriter for a file path with @@ -185,6 +214,12 @@ func (w *archiveWriter) write( w.mu.Lock() defer w.mu.Unlock() + if w.evicted { + return fmt.Errorf( + "%w: %s", errArchiveWriterEvicted, w.path, + ) + } + if w.db == nil || !fileExists(w.path) { err := w.reopen(expiry) if err != nil { @@ -212,7 +247,19 @@ func (w *archiveWriter) write( // its schema, records the reopen time, and prunes expired rows // when expiry is positive. func (w *archiveWriter) open(expiry time.Duration) error { - dbURL := fmt.Sprintf("file:%s?mode=rwc", w.path) + return w.openMode(archiveModeCreate, expiry) +} + +// openMode opens the archive file with the given SQLite URI +// mode, migrates its schema, records the reopen time, and +// prunes expired rows when expiry is positive. The write path +// passes archiveModeCreate so a missing file is recreated; the +// idle sweep passes archiveModeExisting so a missing file is an +// error rather than a newly conjured empty archive. +func (w *archiveWriter) openMode( + mode string, expiry time.Duration, +) error { + dbURL := fmt.Sprintf("file:%s?mode=%s", w.path, mode) sqlDB, err := sql.Open("sqlite", dbURL) if err != nil { @@ -275,6 +322,63 @@ func (w *archiveWriter) close() { w.db = nil } +// sweepExpired prunes an archive that may have gone idle, with +// no write to trigger the usual on-reopen prune. It takes the +// writer's own mutex for the whole operation, so a sweep is +// ordered against concurrent writes rather than reaching around +// them to the file. +// +// It never creates the archive file: a missing file is skipped, +// and the reopen uses archiveModeExisting so SQLite itself +// refuses to create one if the file disappears between the +// check and the open. +// +// The archive is left CLOSED afterwards. An idle archive holding +// no handle is what keeps the operator's move-the-file-away +// workflow working; the next write reopens (and recreates) the +// file as it always has. +func (w *archiveWriter) sweepExpired(expiry time.Duration) error { + w.mu.Lock() + defer w.mu.Unlock() + + if w.evicted { + return fmt.Errorf( + "%w: %s", errArchiveWriterEvicted, w.path, + ) + } + + if !fileExists(w.path) { + return nil + } + + // Drop any live handle first so the prune runs against a + // freshly opened file, matching the write path's semantics. + w.close() + + err := w.openMode(archiveModeExisting, expiry) + if err != nil { + return err + } + + w.close() + + return nil +} + +// evict closes the writer's handle and marks it unusable. It is +// called when the writer leaves the registry, either because the +// webhook was deleted or because its last database target was +// removed. The archive FILE is deliberately left on disk: it is +// long-term storage an operator may still want. +func (w *archiveWriter) evict() { + w.mu.Lock() + defer w.mu.Unlock() + + w.evicted = true + + w.close() +} + // prune deletes archived rows older than expiry, measured from // each row's archived time. It runs on every (re)open, and // because the file is reopened after writes this keeps the diff --git a/internal/delivery/target_database_evict_test.go b/internal/delivery/target_database_evict_test.go new file mode 100644 index 0000000..d45cfb6 --- /dev/null +++ b/internal/delivery/target_database_evict_test.go @@ -0,0 +1,130 @@ +package delivery_test + +import ( + "fmt" + "net/http" + "path/filepath" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "sneak.berlin/go/webhooker/internal/database" + "sneak.berlin/go/webhooker/internal/delivery" +) + +// evictTestEngine builds an engine backed by a temporary data +// directory and returns it along with that directory. +func evictTestEngine(t *testing.T) (*delivery.Engine, string) { + t.Helper() + + dataDir := t.TempDir() + + eng := delivery.NewTestEngineWithDB( + nil, + database.NewTestWebhookDBManager(dataDir), + archiveTestLogger(), + &http.Client{Timeout: 5 * time.Second}, + 1, + ) + + return eng, dataDir +} + +// TestEvictWebhook_ClosesAndRemovesWriter proves that evicting +// a webhook drops its archive writer from the registry and +// closes the open archive handle, rather than leaving both +// alive for the process lifetime. +func TestEvictWebhook_ClosesAndRemovesWriter(t *testing.T) { + t.Parallel() + + eng, dataDir := evictTestEngine(t) + + webhookDB := testWebhookDB(t) + event := seedEvent(t, webhookDB, `{"archived":true}`) + d := seedDatabaseTargetDelivery(t, webhookDB, event, "") + + eng.ExportDeliverDatabase(webhookDB, d) + + webhookID := event.WebhookID + + require.True( + t, eng.ExportHasArchiveWriter(webhookID), + "a delivery should have cached an archive writer", + ) + require.True( + t, eng.ExportArchiveHandleOpen(webhookID), + "the writer should hold an open handle after a write", + ) + + eng.EvictWebhook(webhookID) + + assert.False( + t, eng.ExportHasArchiveWriter(webhookID), + "eviction should remove the registry entry", + ) + assert.False( + t, eng.ExportArchiveHandleOpen(webhookID), + "eviction should close the archive handle", + ) + + archivePath := filepath.Join( + dataDir, fmt.Sprintf("archive-%s.db", webhookID), + ) + assert.FileExists( + t, archivePath, + "eviction must not delete the archive file", + ) +} + +// TestEvictWebhook_UnknownWebhookIsNoOp proves eviction is safe +// for the common case of a webhook that never had a database +// target, and that repeating it does not panic. +func TestEvictWebhook_UnknownWebhookIsNoOp(t *testing.T) { + t.Parallel() + + eng, _ := evictTestEngine(t) + + assert.NotPanics(t, func() { + eng.EvictWebhook("no-such-webhook") + eng.EvictWebhook("no-such-webhook") + }) + + assert.False( + t, eng.ExportHasArchiveWriter("no-such-webhook"), + "eviction must not create a writer", + ) +} + +// TestEvictWebhook_EvictedWriterDoesNotReopen proves an evicted +// writer refuses further writes instead of silently reopening +// the archive file: nothing holds it any more, so a reopened +// handle would leak. +func TestEvictWebhook_EvictedWriterDoesNotReopen(t *testing.T) { + t.Parallel() + + eng, _ := evictTestEngine(t) + + webhookDB := testWebhookDB(t) + event := seedEvent(t, webhookDB, `{"archived":true}`) + d := seedDatabaseTargetDelivery(t, webhookDB, event, "") + + eng.ExportDeliverDatabase(webhookDB, d) + require.True( + t, eng.ExportHasArchiveWriter(event.WebhookID), + ) + + eng.EvictWebhook(event.WebhookID) + + // A fresh delivery for the same webhook gets a brand new + // writer from the registry, so archiving keeps working. + second := seedDatabaseTargetDelivery( + t, webhookDB, event, "", + ) + eng.ExportDeliverDatabase(webhookDB, second) + + assert.True( + t, eng.ExportHasArchiveWriter(event.WebhookID), + "a later delivery should recreate the writer", + ) +} diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index 193856c..c98a360 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -51,6 +51,7 @@ type HandlersParams struct { Healthcheck *healthcheck.Healthcheck Session *session.Session Notifier delivery.Notifier + Evictor delivery.WebhookEvictor } // Handlers provides HTTP handler methods for all application @@ -63,6 +64,7 @@ type Handlers struct { dbMgr *database.WebhookDBManager session *session.Session notifier delivery.Notifier + evictor delivery.WebhookEvictor templates map[string]*template.Template } @@ -97,6 +99,7 @@ func New( s.dbMgr = params.WebhookDBMgr s.session = params.Session s.notifier = params.Notifier + s.evictor = params.Evictor // Parse all page templates once at startup s.templates = map[string]*template.Template{ diff --git a/internal/handlers/handlers_test.go b/internal/handlers/handlers_test.go index 8aae54a..ad927c5 100644 --- a/internal/handlers/handlers_test.go +++ b/internal/handlers/handlers_test.go @@ -4,6 +4,7 @@ import ( "context" "net/http" "net/http/httptest" + "sync" "testing" "github.com/stretchr/testify/assert" @@ -24,6 +25,32 @@ type noopNotifier struct{} func (n *noopNotifier) Notify([]delivery.Task) {} +// recordingEvictor is a delivery.WebhookEvictor that records +// the webhook ids it was asked to evict, so a test can prove +// that a deletion path reached the delivery engine. +type recordingEvictor struct { + mu sync.Mutex + evicted []string +} + +func (r *recordingEvictor) EvictWebhook(webhookID string) { + r.mu.Lock() + defer r.mu.Unlock() + + r.evicted = append(r.evicted, webhookID) +} + +// Evicted returns a copy of the recorded webhook ids. +func (r *recordingEvictor) Evicted() []string { + r.mu.Lock() + defer r.mu.Unlock() + + out := make([]string, len(r.evicted)) + copy(out, r.evicted) + + return out +} + func newTestApp( t *testing.T, targets ...any, @@ -47,6 +74,12 @@ func newTestApp( func() delivery.Notifier { return &noopNotifier{} }, + func() *recordingEvictor { + return &recordingEvictor{} + }, + func(r *recordingEvictor) delivery.WebhookEvictor { + return r + }, handlers.New, ), fx.Populate(targets...), diff --git a/internal/handlers/source_delete_test.go b/internal/handlers/source_delete_test.go new file mode 100644 index 0000000..4bb8469 --- /dev/null +++ b/internal/handlers/source_delete_test.go @@ -0,0 +1,305 @@ +package handlers_test + +import ( + "context" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "testing" + + "github.com/go-chi/chi" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/gorm/clause" + "sneak.berlin/go/webhooker/internal/database" + "sneak.berlin/go/webhooker/internal/handlers" + "sneak.berlin/go/webhooker/internal/session" +) + +const ( + deleteTestUserID = "test-user-id" + deleteTestUsername = "testuser" + + // paramSourceID and paramTargetID are the chi URL parameter + // names the deletion handlers read. + paramSourceID = "sourceID" + paramTargetID = "targetID" +) + +// seedWebhook inserts a webhook owned by the test user and +// returns it. +func seedWebhook( + t *testing.T, + db *database.Database, +) *database.Webhook { + t.Helper() + + wh := &database.Webhook{ + UserID: deleteTestUserID, + Name: "delete-me", + } + + require.NoError( + t, + db.DB().Omit(clause.Associations).Create(wh).Error, + ) + + return wh +} + +// seedTarget inserts a target of the given type for a webhook +// and returns it. +func seedTarget( + t *testing.T, + db *database.Database, + webhookID string, + targetType database.TargetType, +) *database.Target { + t.Helper() + + tgt := &database.Target{ + WebhookID: webhookID, + Name: "t-" + string(targetType), + Type: targetType, + Active: true, + } + + require.NoError( + t, + db.DB().Omit(clause.Associations).Create(tgt).Error, + ) + + return tgt +} + +// archivePathFor returns the archive database path the +// delivery engine would use for a webhook: beside the webhook's +// event database in the data directory. +func archivePathFor( + t *testing.T, + mgr *database.WebhookDBManager, + webhookID string, +) string { + t.Helper() + + return filepath.Join( + filepath.Dir(mgr.DBPath(webhookID)), + "archive-"+webhookID+".db", + ) +} + +// writeArchivePlaceholder creates a stand-in archive file so a +// test can assert the file survives webhook deletion. +func writeArchivePlaceholder(path string) error { + return os.WriteFile(path, []byte("archive"), 0o600) +} + +// postRequest builds an authenticated POST request carrying the +// given chi URL parameters. +func postRequest( + path string, + cookies []*http.Cookie, + params map[string]string, +) *http.Request { + req := httptest.NewRequestWithContext( + context.Background(), http.MethodPost, path, nil, + ) + + for _, c := range cookies { + req.AddCookie(c) + } + + rctx := chi.NewRouteContext() + for k, v := range params { + rctx.URLParams.Add(k, v) + } + + return req.WithContext( + context.WithValue(req.Context(), chi.RouteCtxKey, rctx), + ) +} + +// TestHandleSourceDelete_EvictsArchiveWriter proves that +// deleting a webhook reaches the delivery engine and releases +// the webhook's archive writer, exercised through the real +// deletion handler rather than by calling the evictor directly. +func TestHandleSourceDelete_EvictsArchiveWriter(t *testing.T) { + t.Parallel() + + var ( + h *handlers.Handlers + sess *session.Session + db *database.Database + ev *recordingEvictor + ) + + app := newTestApp(t, &h, &sess, &db, &ev) + app.RequireStart() + + t.Cleanup(app.RequireStop) + + wh := seedWebhook(t, db) + seedTarget(t, db, wh.ID, database.TargetTypeDatabase) + + cookies := authenticatedCookies( + t, sess, deleteTestUserID, deleteTestUsername, + ) + + req := postRequest( + "/source/"+wh.ID+"/delete", + cookies, + map[string]string{paramSourceID: wh.ID}, + ) + w := httptest.NewRecorder() + + h.HandleSourceDelete().ServeHTTP(w, req) + + require.Equal(t, http.StatusSeeOther, w.Code) + assert.Equal( + t, []string{wh.ID}, ev.Evicted(), + "deleting a webhook should evict its archive writer", + ) +} + +// TestHandleSourceDelete_KeepsArchiveFile proves that deleting +// a webhook does not remove its archive database file: the +// archive is long-term storage the operator owns. +func TestHandleSourceDelete_KeepsArchiveFile(t *testing.T) { + t.Parallel() + + var ( + h *handlers.Handlers + sess *session.Session + db *database.Database + mgr *database.WebhookDBManager + ) + + app := newTestApp(t, &h, &sess, &db, &mgr) + app.RequireStart() + + t.Cleanup(app.RequireStop) + + wh := seedWebhook(t, db) + + // Place an archive file where the delivery engine would. + archivePath := archivePathFor(t, mgr, wh.ID) + require.NoError( + t, + writeArchivePlaceholder(archivePath), + ) + + cookies := authenticatedCookies( + t, sess, deleteTestUserID, deleteTestUsername, + ) + + req := postRequest( + "/source/"+wh.ID+"/delete", + cookies, + map[string]string{paramSourceID: wh.ID}, + ) + w := httptest.NewRecorder() + + h.HandleSourceDelete().ServeHTTP(w, req) + + require.Equal(t, http.StatusSeeOther, w.Code) + assert.FileExists( + t, archivePath, + "webhook deletion must not destroy the archive file", + ) +} + +// TestHandleTargetDelete_EvictsWhenLastDatabaseTargetGone +// proves that removing the last database target releases the +// archive writer. +func TestHandleTargetDelete_EvictsWhenLastDatabaseTargetGone( + t *testing.T, +) { + t.Parallel() + + var ( + h *handlers.Handlers + sess *session.Session + db *database.Database + ev *recordingEvictor + ) + + app := newTestApp(t, &h, &sess, &db, &ev) + app.RequireStart() + + t.Cleanup(app.RequireStop) + + wh := seedWebhook(t, db) + tgt := seedTarget( + t, db, wh.ID, database.TargetTypeDatabase, + ) + + cookies := authenticatedCookies( + t, sess, deleteTestUserID, deleteTestUsername, + ) + + req := postRequest( + "/source/"+wh.ID+"/targets/"+tgt.ID+"/delete", + cookies, + map[string]string{ + paramSourceID: wh.ID, + paramTargetID: tgt.ID, + }, + ) + w := httptest.NewRecorder() + + h.HandleTargetDelete().ServeHTTP(w, req) + + require.Equal(t, http.StatusSeeOther, w.Code) + assert.Equal( + t, []string{wh.ID}, ev.Evicted(), + "removing the last database target should evict", + ) +} + +// TestHandleTargetDelete_KeepsWriterWhileDatabaseTargetRemains +// proves that deleting an unrelated target, or one of several +// database targets, leaves a still-needed archive writer alone. +func TestHandleTargetDelete_KeepsWriterWhileDatabaseTargetRemains( + t *testing.T, +) { + t.Parallel() + + var ( + h *handlers.Handlers + sess *session.Session + db *database.Database + ev *recordingEvictor + ) + + app := newTestApp(t, &h, &sess, &db, &ev) + app.RequireStart() + + t.Cleanup(app.RequireStop) + + wh := seedWebhook(t, db) + seedTarget(t, db, wh.ID, database.TargetTypeDatabase) + other := seedTarget(t, db, wh.ID, database.TargetTypeLog) + + cookies := authenticatedCookies( + t, sess, deleteTestUserID, deleteTestUsername, + ) + + req := postRequest( + "/source/"+wh.ID+"/targets/"+other.ID+"/delete", + cookies, + map[string]string{ + paramSourceID: wh.ID, + paramTargetID: other.ID, + }, + ) + w := httptest.NewRecorder() + + h.HandleTargetDelete().ServeHTTP(w, req) + + require.Equal(t, http.StatusSeeOther, w.Code) + assert.Empty( + t, ev.Evicted(), + "a surviving database target must keep its writer", + ) +} diff --git a/internal/handlers/source_management.go b/internal/handlers/source_management.go index eaaf7e9..4f295dc 100644 --- a/internal/handlers/source_management.go +++ b/internal/handlers/source_management.go @@ -533,6 +533,13 @@ func (h *Handlers) deleteWebhookResources( return } + // Release the delivery engine's per-webhook archiving state + // so a deleted webhook's archive writer (and any handle open + // within its debounce window) does not linger for the + // process lifetime. The archive file itself is deliberately + // left on disk; see evictArchiveWriter. + h.evictArchiveWriter(webhook.ID) + err = h.dbMgr.DeleteDB(webhook.ID) if err != nil { h.log.Error( @@ -551,6 +558,64 @@ func (h *Handlers) deleteWebhookResources( http.Redirect(w, r, "/sources", http.StatusSeeOther) } +// evictArchiveWriter asks the delivery engine to drop its +// cached archive writer for a webhook, closing the archive file +// handle. +// +// The archive database file is NOT deleted. Unlike the event +// database — which is per-webhook working storage and is +// hard-deleted with the webhook — an archive is explicitly +// long-term storage that an operator may want to keep or move +// away for offline retention. Destroying it as a side effect of +// deleting a webhook would be a surprising and unrecoverable +// data loss, so the file is left for the operator to handle. +func (h *Handlers) evictArchiveWriter(webhookID string) { + if h.evictor == nil { + return + } + + h.evictor.EvictWebhook(webhookID) +} + +// evictArchiveWriterIfUnused releases a webhook's archive +// writer once the webhook has no database target left to feed +// it. +// +// It is called after any child resource of a webhook is +// deleted, and is correct without knowing which kind was: it +// evicts only when no database target remains, so deleting one +// of several database targets — or deleting an unrelated +// target type — leaves a still-needed writer alone. When no +// database target ever existed there is no writer and eviction +// is a no-op. Soft-deleted targets are excluded by GORM's +// default scope, so the row just deleted is not counted. +func (h *Handlers) evictArchiveWriterIfUnused(webhookID string) { + var remaining int64 + + err := h.db.DB(). + Model(&database.Target{}). + Where( + "webhook_id = ? AND type = ?", + webhookID, database.TargetTypeDatabase, + ). + Count(&remaining).Error + if err != nil { + h.log.Error( + "failed to count remaining database targets", + "webhook_id", webhookID, + "error", err, + ) + + return + } + + if remaining > 0 { + return + } + + h.evictArchiveWriter(webhookID) +} + // HandleSourceLogs shows the request/response logs for a // webhook. func (h *Handlers) HandleSourceLogs() http.HandlerFunc { @@ -1024,23 +1089,31 @@ func (h *Handlers) HandleEntrypointDelete() http.HandlerFunc { return h.deleteChildResource( "entrypointID", &database.Entrypoint{}, "failed to delete entrypoint", + nil, ) } -// HandleTargetDelete handles deleting a target. +// HandleTargetDelete handles deleting a target. Deleting the +// last database target of a webhook leaves its archive writer +// with nothing to write, so the writer is evicted and its +// handle closed; the archive file is left on disk. func (h *Handlers) HandleTargetDelete() http.HandlerFunc { return h.deleteChildResource( "targetID", &database.Target{}, "failed to delete target", + h.evictArchiveWriterIfUnused, ) } // deleteChildResource returns a handler that deletes a child -// resource (entrypoint or target) belonging to a webhook. +// resource (entrypoint or target) belonging to a webhook. The +// optional afterDelete hook runs with the webhook's id once the +// delete has succeeded, before the redirect. func (h *Handlers) deleteChildResource( idParam string, model any, errMsg string, + afterDelete func(webhookID string), ) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { userID, ok := h.getUserID(r) @@ -1080,6 +1153,10 @@ func (h *Handlers) deleteChildResource( return } + if afterDelete != nil { + afterDelete(webhook.ID) + } + http.Redirect( w, r, "/source/"+webhook.ID,