Close archive writers when the delivery engine stops (closes #280)
check / check (push) Successful in 3m36s

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
This commit is contained in:
2026-09-29 05:00:30 +00:00
parent 51580a2bc6
commit fa67d883e9
5 changed files with 173 additions and 20 deletions
+13 -20
View File
@@ -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
+11
View File
@@ -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
@@ -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",
)
}
+18
View File
@@ -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
@@ -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",
)
}