From 4452ef71fba739b26b558b25cd19461bcd0d8100 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 05:00:30 +0000 Subject: [PATCH 1/2] Close archive writers when the delivery engine stops (closes #280) The engine cached archive writers and never closed them at shutdown, so after a clean stop an archive's rows could sit in its -wal while the .db held no table. The engine's stop hook now evicts every cached writer once its workers have returned, the same way deleting a webhook does, so a clean stop leaves each archive as one file and a late write is refused. If the workers do not return within the stop budget, the writers are left open as a kill would leave them. The README no longer says archives keep their sidecars across a clean stop. Model: opus-5-5 --- README.md | 33 +++---- internal/delivery/engine.go | 11 +++ internal/delivery/engine_lifecycle_test.go | 87 +++++++++++++++++++ internal/delivery/target_database.go | 18 ++++ .../delivery/target_database_evict_test.go | 44 ++++++++++ 5 files changed, 173 insertions(+), 20 deletions(-) diff --git a/README.md b/README.md index b3afd54..48a1dd0 100644 --- a/README.md +++ b/README.md @@ -968,15 +968,10 @@ scratch file**: it holds committed transactions that are not yet in the have no readable schema at all. `-shm` is regenerable, but there is no reason to separate the two — copy the directory and you have them. -A clean shutdown closes `webhooker.db` and every `events-*.db`, which -checkpoints and removes their sidecars; a killed or crashed instance -leaves them, and they must be carried with the `.db`. **Archive -databases are different**: their handle is not closed at shutdown, so -`archive-*.db-wal` and `-shm` normally survive a clean stop and the -`-wal` can hold every row the archive has. Measured on a stopped -instance: `archive-….db` 4096 bytes with no table, its `-wal` 157 KB -holding all 8 archived events. Copying `DATA_DIR` in full is what makes -this a non-issue; copying `.db` files out of it by name is not. +A clean shutdown closes every database, which checkpoints and removes +its sidecars; a killed or crashed instance leaves them, and they must be +carried with the `.db`. An archive the service has not opened since a +crash keeps that crash's sidecars, even across a later clean stop. Configuration is **not** in `DATA_DIR` — it comes from the environment and from a `.env` file read out of the process working directory. Back @@ -1051,10 +1046,9 @@ The file becomes self-contained again when the handle closes, which happens on the next write past the debounce window, when the connection pool retires the idle connection (about a minute after the last write), or at the idle archive sweep — measured, the same file was a complete -20 KB `.db` with no sidecars about a minute after its last write. -Shutdown is **not** on that list: the archive handle is not closed when -the service stops. So either move `archive-{uuid}.db` together with any -`-wal`/`-shm` beside it, or wait until there are none. +20 KB `.db` with no sidecars about a minute after its last write. A +clean stop closes it too. So either move `archive-{uuid}.db` together +with any `-wal`/`-shm` beside it, or wait until there are none. ### Restore @@ -1073,12 +1067,10 @@ the service stops. So either move `archive-{uuid}.db` together with any They are part of the database, and dropping a `-wal` silently discards every transaction it still holds. An `.backup` set will not contain any: it writes a single consolidated file per database. A - stop-and-copy set has none for `webhooker.db` or the `events-*.db`, - because a clean stop closes those and checkpoints their sidecars - away — but it will normally have them for `archive-*.db`, whose - handle stays open across shutdown, and those carry the archive's - rows. A copy salvaged from a crashed instance has them for - everything, and needs all of them. + stop-and-copy set normally has none, because a clean stop closes + every database and checkpoints its sidecars away; the exception is an + archive not opened since a crash. A copy salvaged from a crashed + instance has them for everything, and needs all of them. 4. **Fix ownership.** The container runs as the non-root `webhooker` user, UID 1000 / GID 1000. Restored files must be owned by (or @@ -3088,7 +3080,8 @@ each hook. The order, read off the fx stop-hook log: 3. `server` — the HTTP drain, bounded separately by `server.ShutdownTimeout` (**3 seconds**), then a Sentry flush if `SENTRY_DSN` is set -4. `delivery.Engine` +4. `delivery.Engine` — waits for its workers, then closes the archive + databases 5. `healthcheck` 6. `WebhookDBManager` 7. the database close diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index c96dda7..a9d06f2 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -362,6 +362,15 @@ func (e *Engine) start() { // stop cancels the worker pool's context and waits for the pool // to drain, bounded by the stop hook's context: a wedged worker // must not hang the process past fx's stop timeout. +// +// Once the pool has drained it closes the archive writers, so a +// clean stop leaves no archive -wal behind. Nothing else holds a +// writer for long by then: the archive sweeper stops before the +// engine, and deleting a webhook only closes one. If the pool did +// not drain in time, the writers are left open, as a kill would +// leave them: a worker still running may be mid-write, and +// closing its writer would wait on that write and then fail the +// next delivery the worker archives. func (e *Engine) stop(ctx context.Context) error { e.log.Info("delivery engine stopping") @@ -376,6 +385,8 @@ func (e *Engine) stop(ctx context.Context) error { return err } + e.dbTarget.evictAll() + e.log.Info("delivery engine stopped") return nil diff --git a/internal/delivery/engine_lifecycle_test.go b/internal/delivery/engine_lifecycle_test.go index 6e7ef2e..ae1afa5 100644 --- a/internal/delivery/engine_lifecycle_test.go +++ b/internal/delivery/engine_lifecycle_test.go @@ -2,6 +2,8 @@ package delivery_test import ( "context" + "fmt" + "path/filepath" "testing" "time" @@ -269,3 +271,88 @@ func TestEngine_StopHookHonoursStopTimeout(t *testing.T) { requireStopHookExpires(t, lc.hooks[0], "delivery engine") } + +// deliverToArchive runs one delivery to a database target through +// the running engine and returns the webhook's archive file path. +// The archive writer holds the file open afterwards. +func deliverToArchive(t *testing.T, s iSetup) string { + t.Helper() + + deliveryID, task := seedLogTask(t, s) + task.TargetType = database.TargetTypeDatabase + + s.Engine.Notify([]delivery.Task{task}) + + iWaitForDelivered(t, s.WebhookDB, deliveryID) + + return filepath.Join( + filepath.Dir(s.DBMgr.DBPath(s.WebhookID)), + fmt.Sprintf("archive-%s.db", s.WebhookID), + ) +} + +// TestEngine_StopHookClosesArchives is the regression test for an +// archive split across two files by a clean stop. The engine never +// closed its archive writers, so after a stop the archived rows +// could sit in archive-{id}.db-wal while archive-{id}.db held no +// table at all, and copying the .db on its own gave an empty +// database. +func TestEngine_StopHookClosesArchives(t *testing.T) { + t.Parallel() + + s := newISetup(t) + + lc := startEngineViaHook(t, s.Engine) + + path := deliverToArchive(t, s) + require.FileExists( + t, path+"-wal", + "an open archive should have a -wal for the stop to remove", + ) + + require.NoError(t, lc.hooks[0].OnStop(context.Background())) + + wals, err := filepath.Glob( + filepath.Join(filepath.Dir(path), "archive-*.db-wal"), + ) + require.NoError(t, err) + require.Empty( + t, wals, "a clean stop must leave no archive -wal behind", + ) + + // With no -wal beside it, the row can only be in the .db. + count, err := countArchivedRows(path) + require.NoError(t, err) + require.Equal(t, int64(1), count) +} + +// TestEngine_StopHookTimeoutLeavesArchivesOpen covers a stop whose +// budget runs out while a worker is still running. That worker may +// be in the middle of an archive write, so the archive writers are +// left open, as a kill would leave them, rather than closed +// underneath it. +func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) { + t.Parallel() + + s := newISetup(t) + + lc := startEngineViaHook(t, s.Engine) + + deliverToArchive(t, s) + + release := make(chan struct{}) + + t.Cleanup(func() { + close(release) + s.Engine.EvictWebhook(s.WebhookID) + }) + + s.Engine.ExportWedgeWorker(release) + + requireStopHookExpires(t, lc.hooks[0], "delivery engine") + + require.True( + t, s.Engine.ExportArchiveHandleOpen(s.WebhookID), + "a stop that timed out must not close archive writers", + ) +} diff --git a/internal/delivery/target_database.go b/internal/delivery/target_database.go index d319551..7173b71 100644 --- a/internal/delivery/target_database.go +++ b/internal/delivery/target_database.go @@ -277,6 +277,24 @@ func (t *databaseTarget) evict(webhookID string) { ) } +// evictAll evicts every cached archive writer, exactly as evict +// does for one webhook. The engine calls it at shutdown, once its +// workers have returned. Closing the last handle on an archive +// moves the contents of its -wal into the .db and removes the +// -wal, so a clean stop leaves each archive as a single file. +func (t *databaseTarget) evictAll() { + t.mu.Lock() + + writers := t.writers + t.writers = nil + + t.mu.Unlock() + + for _, w := range writers { + w.evict() + } +} + // 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 diff --git a/internal/delivery/target_database_evict_test.go b/internal/delivery/target_database_evict_test.go index 14e7945..095823d 100644 --- a/internal/delivery/target_database_evict_test.go +++ b/internal/delivery/target_database_evict_test.go @@ -1,6 +1,7 @@ package delivery_test import ( + "context" "errors" "fmt" "net/http" @@ -361,3 +362,46 @@ func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) { "a later delivery should recreate the writer", ) } + +// TestEngineStop_WriteAfterStopIsRefused proves the engine's stop +// closes each archive writer the way deleting its webhook does: a +// write that reaches a writer after the stop is refused, reopens +// nothing and adds no row. +func TestEngineStop_WriteAfterStopIsRefused(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) + + w := eng.ExportArchiveWriterFor(event.WebhookID) + require.NotNil(t, w) + require.True(t, w.HandleOpen()) + + require.NoError(t, eng.ExportStop(context.Background())) + + err := w.Write(evictTestRow("ev-after-stop"), 0) + + require.ErrorIs( + t, err, delivery.ErrExportArchiveWriterEvicted, + "a write after the stop must be refused", + ) + assert.False( + t, w.HandleOpen(), + "a refused write must not reopen the archive", + ) + assert.False( + t, eng.ExportHasArchiveWriter(event.WebhookID), + "the stop should empty the registry", + ) + + count, err := countArchivedRows(w.Path()) + require.NoError(t, err) + assert.Equal( + t, int64(1), count, "the refused row must not be written", + ) +} -- 2.54.0 From d03bb0ed08d79327502a149e805ddc6dec52813c Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 05:48:54 +0000 Subject: [PATCH 2/2] Correct the stop timeout comments on the archive writers The comments on the engine's stop and on the timeout test gave false reasons for leaving archive writers open when the stop budget runs out. Closing them would wait for any write in progress, and a worker still running would then open new writers that nothing closes, so closing gains nothing over a kill. Model: opus-5-5 --- internal/delivery/engine.go | 6 +++--- internal/delivery/engine_lifecycle_test.go | 8 ++++---- 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index a9d06f2..4786f08 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -368,9 +368,9 @@ func (e *Engine) start() { // writer for long by then: the archive sweeper stops before the // engine, and deleting a webhook only closes one. If the pool did // not drain in time, the writers are left open, as a kill would -// leave them: a worker still running may be mid-write, and -// closing its writer would wait on that write and then fail the -// next delivery the worker archives. +// leave them. Closing them would wait for any write in progress, +// and a worker still running would then open new writers that +// nothing closes, so it gains nothing over a kill. func (e *Engine) stop(ctx context.Context) error { e.log.Info("delivery engine stopping") diff --git a/internal/delivery/engine_lifecycle_test.go b/internal/delivery/engine_lifecycle_test.go index ae1afa5..9fadf9d 100644 --- a/internal/delivery/engine_lifecycle_test.go +++ b/internal/delivery/engine_lifecycle_test.go @@ -327,10 +327,10 @@ func TestEngine_StopHookClosesArchives(t *testing.T) { } // TestEngine_StopHookTimeoutLeavesArchivesOpen covers a stop whose -// budget runs out while a worker is still running. That worker may -// be in the middle of an archive write, so the archive writers are -// left open, as a kill would leave them, rather than closed -// underneath it. +// budget runs out while a worker is still running. The archive +// writers are left open, as a kill would leave them: closing them +// would wait for any write in progress, and that worker would then +// open new writers that nothing closes. func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) { t.Parallel() -- 2.54.0