From 17e6be34ab7b5c4d7c39216f269c926a0e43b804 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Fri, 2 Oct 2026 05:10:03 +0000 Subject: [PATCH] Name each database target's archive for its webhook and target (closes #376) Each database target now has its own archive file, archive-WEBHOOKNAME-TARGETNAME-TARGETID.db, instead of one archive-WEBHOOKID.db per webhook. delivery.ArchiveFileName builds the name: each name is lowercased, keeps ASCII letters and digits, turns every other run of characters into one dash, and is cut to 40 characters. A change of webhook or target name renames its archive files under the archive writer's lock, before the new name is saved, and back again if the save fails. A rename never replaces a file: if one already has the new name, the edit is refused. Deleting a target evicts only that target's writer. Archive files are never deleted, and nothing looks for files under the old name. Model: opus-5-5 --- README.md | 116 ++++--- cmd/webhooker/main.go | 9 +- internal/delivery/archive_sweeper.go | 26 +- internal/delivery/archive_sweeper_test.go | 276 ++++++++-------- internal/delivery/engine.go | 71 ++-- internal/delivery/engine_lifecycle_test.go | 33 +- internal/delivery/engine_test.go | 45 +-- internal/delivery/export_test.go | 46 +-- internal/delivery/target_database.go | 293 +++++++++++------ internal/delivery/target_database_archive.go | 98 +++++- .../delivery/target_database_evict_test.go | 190 +++++------ internal/delivery/target_database_test.go | 304 +++++++++++++++--- internal/handlers/handlers.go | 6 +- internal/handlers/handlers_test.go | 95 +++++- internal/handlers/source_delete_test.go | 113 ++----- internal/handlers/source_management.go | 130 +++++--- internal/handlers/source_management_test.go | 106 +++++- internal/handlers/target_edit.go | 49 ++- internal/handlers/target_edit_test.go | 65 ++++ internal/resetpw/resetpw_test.go | 12 +- internal/server/routes_test.go | 18 +- 21 files changed, 1438 insertions(+), 663 deletions(-) diff --git a/README.md b/README.md index dc4a10d..a4066e7 100644 --- a/README.md +++ b/README.md @@ -698,7 +698,8 @@ The app runs as a non-root user (`webhooker`, UID 1000), exposes port The `/var/lib/webhooker` volume holds all SQLite databases: the main application database (`webhooker.db`), the per-webhook event databases (`events-{uuid}.db`), and any archive databases written by `database` -targets (`archive-{uuid}.db`). Mount this as a persistent volume to +targets (`archive-{webhook_name}-{target_name}-{target_uuid}.db`). Mount +this as a persistent volume to preserve data across container restarts. **The container sets its data directory's owner and mode itself @@ -937,13 +938,13 @@ is both the simplest and the only complete rule: encryption key), users, API keys, webhooks, entrypoints, targets. - `events-{webhook_uuid}.db` — **one per webhook**. Events, deliveries, delivery results. -- `archive-{webhook_uuid}.db` — **one per webhook that has a `database` - target**. Archived events. Keyed on the webhook UUID, not the target - UUID: a webhook with several `database` targets still has exactly one - archive file. +- `archive-{webhook_name}-{target_name}-{target_uuid}.db` — **one per + `database` target**. Archived events. The two names are made safe for + a file name, and the file is renamed when the webhook or the target is + (see [Database Architecture](#database-architecture)). -`{webhook_uuid}` is the webhook's UUID primary key in its canonical -36-character hyphenated form, so a real filename looks like +`{webhook_uuid}` and `{target_uuid}` are UUID primary keys in their +canonical 36-character hyphenated form, so a real filename looks like `events-3f2a1c9e-....db`. The only other file is `webhooker.lock`, the always-empty [single-instance lock](#single-instance-lock); it holds no state and is not part of the backup set — a copied one is stale and @@ -1017,8 +1018,8 @@ stopped copy. Archive databases are the one exception the service is built for: the archive writer closes and reopens its handle around writes (debounced -to at most one reopen per second), so an operator can move -`archive-{uuid}.db` away for offline retention while the service runs, +to at most one reopen per second), so an operator can move an +`archive-….db` away for offline retention while the service runs, and it is recreated on the next write. See [Database Architecture](#database-architecture). That is a move-the-file-away workflow, not a substitute for the backup procedures @@ -1036,7 +1037,7 @@ 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. A -clean stop closes it too. So either move `archive-{uuid}.db` together +clean stop closes it too. So either move the `archive-….db` together with any `-wal`/`-shm` beside it, or wait until there are none. ### Restore @@ -1171,7 +1172,7 @@ commit still produce a byte-identical binary. Treat a backup with the same care as the credentials inside it. Encrypt backups at rest and restrict who can read them. -- `events-{uuid}.db` and `archive-{uuid}.db` hold the **full payload +- `events-{uuid}.db` and `archive-….db` hold the **full payload body and headers** of every event as received, including whatever the sending service put in them — tokens, signatures, personal data. - Event databases written before @@ -1589,8 +1590,9 @@ events should be forwarded. is built on the same HTTP core as `http` and honours `max_retries` identically, circuit breaker included. See the Slack target section under "Per-Webhook Event Databases" for the message format. -- **`database`** — Archive the full event as a row into a separate - per-webhook archive database (`archive-{webhookID}.db`) for long-term +- **`database`** — Archive the full event as a row into the target's + own archive database + (`archive-{webhook_name}-{target_name}-{target_uuid}.db`) for long-term retention, with an optional creation-validated expiry (default: keep forever). No external delivery and no retries; an archive write failure fails the delivery. See the database target section under @@ -1904,9 +1906,35 @@ The **database target type** builds on this architecture to provide long-term archiving, separate from the per-webhook event database (which may prune events under its own retention). Delivering to a database target writes the full event — body, headers, method, content type, and -webhook/entrypoint/event identifiers — as a row into a dedicated archive -database, `archive-{webhookID}.db`, stored under the data directory -beside the event database. After each write the archive handle is closed +webhook/entrypoint/event identifiers — as a row into the target's own +archive database, `archive-{webhook_name}-{target_name}-{target_uuid}.db`, +stored under the data directory beside the event database. Each +`database` target has its own archive file, even when one webhook has +several. + +Both names are made safe for a file name the same way: lowercased, ASCII +letters and digits kept, every other run of characters turned into a +single `-`, no `-` at either end, cut to 40 characters, and `unnamed` +when nothing is left. The target UUID keeps the file name unique. A +webhook named `Orders (EU)` with a target named `Long-term archive` +archives into `archive-orders-eu-long-term-archive-{target_uuid}.db`. +Renaming the webhook or the target renames the file, under the same +lock the archive writes and the archive sweeper take, so the name on +disk matches the UI. A rename never replaces a file: if one already has +the new name, the edit is refused with an error naming that file, and +the stored name stays. If the archive is not there (the operator moved +it away), the rename is not an error, and the next write creates the +file under the new name. + +The file is moved just before the new name is saved. If the process +stops between the two, the archive is left under the new name while the +UI still shows the old one, and the next delivery starts a second +archive under the name shown. To bring them back together, move the +file under the new name back to the name shown; if a second archive is +already there, move the older file out of the data directory instead +and keep it as you would any archive moved away. + +After each write the archive handle is closed and reopened, debounced to at most once per second, so an operator can move the archive file away for offline archiving without stopping the service; a moved or removed archive file is recreated automatically on @@ -1917,35 +1945,33 @@ 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. +Because reopens only happen on writes, an archive whose target 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 target +is still receiving traffic. The sweep never creates an archive — a +target whose archive file does not yet exist is skipped, not initialised +— it takes the same per-target 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. -Note that a webhook has one archive file but may carry more than one -`database` target, each with its own `expiry`. The shortest expiry -configured on any of them therefore governs the whole archive, and the -sweep applies it whether or not the webhook is still receiving events. -Configure a single `database` target per webhook unless you intend that. +Because each `database` target has its own archive file, a target's +`expiry` governs only its own archive. Two `database` targets on one +webhook with different expiries keep two archives, each pruned on its +own schedule. -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. +Deleting a webhook releases its archives: the delivery engine's cached +archive writers are dropped and their file handles closed, so nothing +lingers after the webhook is gone. The archive **files themselves are +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 an `archive-….db` is the +operator's call. Deleting a `database` target releases its writer the +same way, and for the same reason leaves its file alone. The **Slack target type** sends webhook events as formatted messages to any Slack-compatible incoming webhook URL (works with Slack, Mattermost, @@ -2955,8 +2981,8 @@ Components are wired via Uber fx in this order: 11. `delivery.New` — Event-driven delivery engine 12. `delivery.NewArchiveSweeper` — Periodic pruning of idle archives 13. `delivery.Engine` → `delivery.Notifier` — interface bridge -14. `delivery.Engine` → `delivery.WebhookEvictor` — interface bridge so - deleting a webhook releases its archive writer +14. `delivery.Engine` → `delivery.Archives` — interface bridge so + deleting or renaming a webhook or target reaches its archive files 15. `server.New` — HTTP server and router The server starts via `fx.Invoke(func(*server.Server, *delivery.Engine, diff --git a/cmd/webhooker/main.go b/cmd/webhooker/main.go index 0da8fa0..b42f88a 100644 --- a/cmd/webhooker/main.go +++ b/cmd/webhooker/main.go @@ -187,11 +187,10 @@ func newApp() *fx.App { // 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 - }, + // Wire *delivery.Engine as delivery.Archives so deleting + // or renaming a webhook or target reaches its archive + // files. + func(e *delivery.Engine) delivery.Archives { return e }, server.New, ), fx.Invoke( diff --git a/internal/delivery/archive_sweeper.go b/internal/delivery/archive_sweeper.go index 2cd2934..646b1e9 100644 --- a/internal/delivery/archive_sweeper.go +++ b/internal/delivery/archive_sweeper.go @@ -8,6 +8,7 @@ import ( "time" "go.uber.org/fx" + "gorm.io/gorm" "sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/lifecycle" @@ -25,14 +26,14 @@ type ArchiveSweeperParams struct { Logger *logger.Logger } -// ArchiveSweeper periodically prunes expired rows from -// per-webhook archive databases whose database target carries a -// positive expiry. +// ArchiveSweeper periodically prunes expired rows from the +// archive databases of database targets that carry 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 +// archive whose target 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". // @@ -155,7 +156,7 @@ func (s *ArchiveSweeper) run(ctx context.Context) { // soft-deleted along with it, so GORM's default scope already // excludes them. // -// A failure for one webhook is logged and the sweep continues, +// A failure for one target 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) { @@ -210,19 +211,20 @@ func (s *ArchiveSweeper) sweepTarget(target *database.Target) { return } - err = s.eng.dbTarget.sweepWebhook(target.WebhookID, expiry) + err = s.eng.dbTarget.sweepArchive(target.ID, expiry) if err == nil { return } - // A writer evicted underneath the sweep means the operator - // deleted the webhook (or its last database target) while the - // sweep was walking the target list. That is an ordinary + // A writer evicted, or a target row gone, underneath the sweep + // means the operator deleted the target or its webhook while + // the sweep was walking the target list. That is an ordinary // interleaving, not a failure, so it must not produce an // error line. - if errors.Is(err, errArchiveWriterEvicted) { + if errors.Is(err, errArchiveWriterEvicted) || + errors.Is(err, gorm.ErrRecordNotFound) { s.log.Debug( - "archive sweep: writer evicted mid-sweep", + "archive sweep: target deleted mid-sweep", "webhook_id", target.WebhookID, "target_id", target.ID, ) diff --git a/internal/delivery/archive_sweeper_test.go b/internal/delivery/archive_sweeper_test.go index cdfab09..590c598 100644 --- a/internal/delivery/archive_sweeper_test.go +++ b/internal/delivery/archive_sweeper_test.go @@ -34,18 +34,23 @@ const ( 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 { +// archiveTestWebhookName is the name of every webhook +// seedDatabaseTarget creates. It is not safe in a file name as it +// stands, so every archive test goes through archiveNamePart. +const archiveTestWebhookName = "Sweep Test!" + +// archiveEnv bundles the pieces an archive test drives: a main +// configuration database holding webhooks and targets, a delivery +// engine owning the archive writer registry, the archive sweeper, +// and the data directory the archive files live in. +type archiveEnv struct { sweeper *delivery.ArchiveSweeper eng *delivery.Engine mainDB *database.Database dataDir string } -func setupSweeperTest(t *testing.T) *sweeperEnv { +func setupArchiveTest(t *testing.T) *archiveEnv { t.Helper() dataDir := t.TempDir() @@ -78,7 +83,7 @@ func setupSweeperTest(t *testing.T) *sweeperEnv { 1, ) - return &sweeperEnv{ + return &archiveEnv{ sweeper: delivery.NewTestArchiveSweeper( mainDB, eng, log, ), @@ -88,25 +93,27 @@ func setupSweeperTest(t *testing.T) *sweeperEnv { } } -// archivePath returns where the engine keeps a webhook's -// archive file. -func (env *sweeperEnv) archivePath(webhookID string) string { +// archivePath returns where the engine keeps a database target's +// archive file, for the names seedDatabaseTarget gave it. +func (env *archiveEnv) archivePath(tgt *database.Target) string { return filepath.Join( - env.dataDir, fmt.Sprintf("archive-%s.db", webhookID), + env.dataDir, + delivery.ArchiveFileName( + archiveTestWebhookName, tgt.Name, tgt.ID, + ), ) } // seedDatabaseTarget creates a webhook with one database target -// carrying the given target config JSON, and returns the -// webhook id. -func (env *sweeperEnv) seedDatabaseTarget( +// carrying the given target config JSON, and returns the target. +func (env *archiveEnv) seedDatabaseTarget( t *testing.T, configJSON string, -) string { +) *database.Target { t.Helper() wh := &database.Webhook{ UserID: uuid.New().String(), - Name: "sweep-test", + Name: archiveTestWebhookName, } require.NoError( t, @@ -115,9 +122,19 @@ func (env *sweeperEnv) seedDatabaseTarget( Create(wh).Error, ) + return env.addDatabaseTarget(t, wh.ID, configJSON) +} + +// addDatabaseTarget creates one more database target on an +// existing webhook and returns it. +func (env *archiveEnv) addDatabaseTarget( + t *testing.T, webhookID, configJSON string, +) *database.Target { + t.Helper() + tgt := &database.Target{ - WebhookID: wh.ID, - Name: "archive", + WebhookID: webhookID, + Name: "Archive", Type: database.TargetTypeDatabase, Active: true, Config: configJSON, @@ -129,19 +146,19 @@ func (env *sweeperEnv) seedDatabaseTarget( Create(tgt).Error, ) - return wh.ID + return tgt } -// seedArchiveRows creates the archive file for a webhook and +// seedArchiveRows creates the archive file for a target 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, +func (env *archiveEnv) seedArchiveRows( + t *testing.T, tgt *database.Target, archivedAt ...time.Time, ) string { t.Helper() - path := env.archivePath(webhookID) + path := env.archivePath(tgt) sqlDB, err := sql.Open( "sqlite", fmt.Sprintf("file:%s?mode=rwc", path), @@ -160,7 +177,7 @@ func (env *sweeperEnv) seedArchiveRows( for i, at := range archivedAt { row := delivery.ExportArchivedEvent{ EventID: fmt.Sprintf("ev-%d", i), - WebhookID: webhookID, + WebhookID: tgt.WebhookID, Method: http.MethodPost, Body: `{"seeded":true}`, ArchivedAt: at, @@ -243,13 +260,13 @@ func TestArchiveSweeper_LoopOutlivesStartHookContext( ) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) now := time.Now() path := env.seedArchiveRows( - t, webhookID, + t, tgt, now.Add(-48*time.Hour), now.Add(-time.Minute), ) @@ -287,60 +304,60 @@ func TestArchiveSweeper_LoopOutlivesStartHookContext( } // TestArchiveSweep_DoesNotResurrectEvictedWriter covers the -// interleaving where a sweep tick has already listed a webhook's -// target when the webhook is deleted and its writer evicted. The -// sweep must not put a writer back into the registry: nothing -// would ever evict it again, which is precisely the leak this -// change exists to close. +// interleaving where a sweep tick has already listed a target +// when the target is deleted and its writer evicted. The sweep +// must not put a writer back into the registry: nothing would +// ever evict it again, which is precisely the leak this change +// exists to close. func TestArchiveSweep_DoesNotResurrectEvictedWriter( t *testing.T, ) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) env.seedArchiveRows( - t, webhookID, time.Now().Add(-48*time.Hour), + t, tgt, time.Now().Add(-48*time.Hour), ) // Prime the registry the way a delivery would, then evict as // the deletion path does. The target row is deliberately left - // in place: this is the tick that listed the webhook before + // in place: this is the tick that listed the target before // the deletion committed. - _, err := env.eng.ExportEnsureArchiveWriter(webhookID) + _, err := env.eng.ExportEnsureArchiveWriter(tgt.ID) require.NoError(t, err) - env.eng.EvictWebhook(webhookID) - require.False(t, env.eng.ExportHasArchiveWriter(webhookID)) + env.eng.EvictTarget(tgt.ID) + require.False(t, env.eng.ExportHasArchiveWriter(tgt.ID)) env.sweeper.ExportSweep(context.Background()) assert.False( - t, env.eng.ExportHasArchiveWriter(webhookID), - "a sweep must never re-register a writer for a webhook "+ + t, env.eng.ExportHasArchiveWriter(tgt.ID), + "a sweep must never re-register a writer for a target "+ "whose registry entry has already been released", ) } // TestArchiveSweep_LeavesNoRegistryEntry states the same // invariant in its general form: sweeping an archive whose -// webhook has no cached writer must not leave one behind, so the +// target has no cached writer must not leave one behind, so the // registry keeps holding only writers a delivery created and an // eviction can reach. func TestArchiveSweep_LeavesNoRegistryEntry(t *testing.T) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) path := env.seedArchiveRows( - t, webhookID, + t, tgt, time.Now().Add(-48*time.Hour), time.Now().Add(-time.Minute), ) - require.False(t, env.eng.ExportHasArchiveWriter(webhookID)) + require.False(t, env.eng.ExportHasArchiveWriter(tgt.ID)) env.sweeper.ExportSweep(context.Background()) @@ -349,7 +366,7 @@ func TestArchiveSweep_LeavesNoRegistryEntry(t *testing.T) { "the sweep must still prune an idle archive", ) assert.False( - t, env.eng.ExportHasArchiveWriter(webhookID), + t, env.eng.ExportHasArchiveWriter(tgt.ID), "the sweep must release the registry entry it created", ) } @@ -364,34 +381,31 @@ func TestArchiveSweep_KeepsWriterAdoptedByDelivery( ) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) env.seedArchiveRows( - t, webhookID, time.Now().Add(-48*time.Hour), + t, tgt, time.Now().Add(-48*time.Hour), ) webhookDB := testWebhookDB(t) event := seedEvent(t, webhookDB, `{"n":1}`) - event.WebhookID = webhookID - d := seedDatabaseTargetDelivery( - t, webhookDB, event, `{"expiry":"1h"}`, - ) + d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt) env.sweeper.ExportSweep(context.Background()) - require.False(t, env.eng.ExportHasArchiveWriter(webhookID)) + require.False(t, env.eng.ExportHasArchiveWriter(tgt.ID)) env.eng.ExportDeliverDatabase(webhookDB, d) assert.True( - t, env.eng.ExportHasArchiveWriter(webhookID), + t, env.eng.ExportHasArchiveWriter(tgt.ID), "a delivery's writer must stay registered", ) env.sweeper.ExportSweep(context.Background()) assert.True( - t, env.eng.ExportHasArchiveWriter(webhookID), + t, env.eng.ExportHasArchiveWriter(tgt.ID), "a sweep must not drop a writer a delivery owns", ) } @@ -423,15 +437,15 @@ func TestArchiveSweep_KeepsWriterAdoptedDuringSweep( ) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) env.seedArchiveRows( - t, webhookID, time.Now().Add(-48*time.Hour), + t, tgt, time.Now().Add(-48*time.Hour), ) sweepWriter, created, err := env.eng.ExportSweepWriterFor( - webhookID, + tgt.ID, ) require.NoError(t, err) require.True( @@ -442,37 +456,34 @@ func TestArchiveSweep_KeepsWriterAdoptedDuringSweep( // The delivery lands mid-sweep and adopts the entry. webhookDB := testWebhookDB(t) event := seedEvent(t, webhookDB, `{"n":1}`) - event.WebhookID = webhookID - d := seedDatabaseTargetDelivery( - t, webhookDB, event, `{"expiry":"1h"}`, - ) + d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt) env.eng.ExportDeliverDatabase(webhookDB, d) - adopted := env.eng.ExportArchiveWriterFor(webhookID) + adopted := env.eng.ExportArchiveWriterFor(tgt.ID) require.NotNil(t, adopted) require.True( t, sweepWriter.Same(adopted), "the delivery must have adopted the sweep's writer", ) require.True( - t, env.eng.ExportArchiveHandleOpen(webhookID), + t, env.eng.ExportArchiveHandleOpen(tgt.ID), "the delivery leaves the archive handle open", ) // The sweep finishes. - env.eng.ExportReleaseSweepWriter(webhookID, sweepWriter) + env.eng.ExportReleaseSweepWriter(tgt.ID, sweepWriter) require.True( - t, env.eng.ExportHasArchiveWriter(webhookID), + t, env.eng.ExportHasArchiveWriter(tgt.ID), "a writer adopted by a delivery during a sweep must "+ "stay registered, or its open handle is unreachable", ) - env.eng.EvictWebhook(webhookID) + env.eng.EvictTarget(tgt.ID) assert.False( - t, env.eng.ExportHasArchiveWriter(webhookID), + t, env.eng.ExportHasArchiveWriter(tgt.ID), "the adopted writer must still be evictable", ) assert.False( @@ -481,34 +492,34 @@ func TestArchiveSweep_KeepsWriterAdoptedDuringSweep( ) } -// TestArchiveSweep_ContinuesAfterPerWebhookFailure proves a -// failure for one webhook does not abort the sweep for the +// TestArchiveSweep_ContinuesAfterPerTargetFailure proves a +// failure for one target does not abort the sweep for the // others: an unparseable expiry and an unreadable archive both // have to be logged and stepped over. -func TestArchiveSweep_ContinuesAfterPerWebhookFailure( +func TestArchiveSweep_ContinuesAfterPerTargetFailure( t *testing.T, ) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) // Seeded first so the sweep reaches them before the healthy - // webhook: targets come back in insertion order. - badConfigID := env.seedDatabaseTarget(t, `{"expiry":"!!!"}`) + // target: targets come back in insertion order. + badConfig := env.seedDatabaseTarget(t, `{"expiry":"!!!"}`) env.seedArchiveRows( - t, badConfigID, time.Now().Add(-48*time.Hour), + t, badConfig, time.Now().Add(-48*time.Hour), ) - corruptID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + corrupt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) require.NoError(t, os.WriteFile( - env.archivePath(corruptID), + env.archivePath(corrupt), []byte("this is not a sqlite database"), 0o600, )) - healthyID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + healthy := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) healthyPath := env.seedArchiveRows( - t, healthyID, + t, healthy, time.Now().Add(-48*time.Hour), time.Now().Add(-time.Minute), ) @@ -518,14 +529,14 @@ func TestArchiveSweep_ContinuesAfterPerWebhookFailure( assert.Equal( t, []string{sweepRowNew}, archivedEventIDs(t, healthyPath), - "a failure for an earlier webhook must not stop the "+ + "a failure for an earlier target must not stop the "+ "sweep from pruning the ones after it", ) } // TestArchiveSweep_OpenExistingDoesNotCreateFile pins the second // of the two no-create guards. The first is the stat in -// sweepWebhook; this one is the SQLite open mode, which is what +// sweepExpired; this one is the SQLite open mode, which is what // protects the window between that stat and the open. Flipping // the sweep's mode to create-if-missing makes this fail. func TestArchiveSweep_OpenExistingDoesNotCreateFile( @@ -561,13 +572,13 @@ func TestArchiveSweep_OpenExistingDoesNotCreateFile( func TestArchiveSweep_PrunesIdleArchive(t *testing.T) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) now := time.Now() path := env.seedArchiveRows( - t, webhookID, + t, tgt, now.Add(-48*time.Hour), now.Add(-time.Minute), ) @@ -600,11 +611,11 @@ func TestArchiveSweep_PrunesIdleArchive(t *testing.T) { func TestArchiveSweep_LeavesArchiveClosed(t *testing.T) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) path := env.seedArchiveRows( - t, webhookID, time.Now().Add(-48*time.Hour), + t, tgt, time.Now().Add(-48*time.Hour), ) w := delivery.NewExportArchiveWriter( @@ -640,35 +651,32 @@ func TestArchiveSweep_ClosesHandleOfRegisteredWriter( ) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) env.seedArchiveRows( - t, webhookID, time.Now().Add(-48*time.Hour), + t, tgt, time.Now().Add(-48*time.Hour), ) webhookDB := testWebhookDB(t) event := seedEvent(t, webhookDB, `{"n":1}`) - event.WebhookID = webhookID - d := seedDatabaseTargetDelivery( - t, webhookDB, event, `{"expiry":"1h"}`, - ) + d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt) env.eng.ExportDeliverDatabase(webhookDB, d) require.True( - t, env.eng.ExportArchiveHandleOpen(webhookID), + t, env.eng.ExportArchiveHandleOpen(tgt.ID), "the delivery must leave the archive handle open", ) env.sweeper.ExportSweep(context.Background()) require.True( - t, env.eng.ExportHasArchiveWriter(webhookID), + t, env.eng.ExportHasArchiveWriter(tgt.ID), "the delivery's registry entry must survive the sweep", ) assert.False( - t, env.eng.ExportArchiveHandleOpen(webhookID), + t, env.eng.ExportArchiveHandleOpen(tgt.ID), "the sweep must leave the archive closed", ) } @@ -684,11 +692,11 @@ func TestArchiveSweep_NeverExpiryUntouched(t *testing.T) { `{"expiry":""}`, "", } { - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, configJSON) + tgt := env.seedDatabaseTarget(t, configJSON) path := env.seedArchiveRows( - t, webhookID, + t, tgt, time.Now().Add(-10000*time.Hour), ) @@ -699,7 +707,7 @@ func TestArchiveSweep_NeverExpiryUntouched(t *testing.T) { "config %q must keep rows forever", configJSON, ) assert.False( - t, env.eng.ExportHasArchiveWriter(webhookID), + t, env.eng.ExportHasArchiveWriter(tgt.ID), "config %q must leave no registry entry behind", configJSON, ) @@ -722,10 +730,10 @@ func TestArchiveSweep_NeverExpirySkipsBeforeOpening( ) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"never"}`) - path := env.archivePath(webhookID) + tgt := env.seedDatabaseTarget(t, `{"expiry":"never"}`) + path := env.archivePath(tgt) seedUnmigratedArchive(t, path) require.False(t, archiveTableExists(t, path)) @@ -768,16 +776,16 @@ func archiveTableExists(t *testing.T, path string) bool { } // 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. +// never conjures an archive: a database target that has never +// received an event must still have no archive file (nor SQLite +// sidecar) after a sweep, and no registry entry either. func TestArchiveSweep_DoesNotCreateArchiveFile(t *testing.T) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) - path := env.archivePath(webhookID) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + path := env.archivePath(tgt) require.NoFileExists(t, path) @@ -789,6 +797,11 @@ func TestArchiveSweep_DoesNotCreateArchiveFile(t *testing.T) { "the sweep must not create an archive file", ) } + + assert.False( + t, env.eng.ExportHasArchiveWriter(tgt.ID), + "the sweep must leave no registry entry behind", + ) } // TestArchiveSweep_DoesNotCreateAfterWriterExists covers the @@ -800,11 +813,11 @@ func TestArchiveSweep_DoesNotCreateAfterWriterExists( ) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) - path, err := env.eng.ExportEnsureArchiveWriter(webhookID) + path, err := env.eng.ExportEnsureArchiveWriter(tgt.ID) require.NoError(t, err) require.NoFileExists(t, path) @@ -819,17 +832,17 @@ func TestArchiveSweep_DoesNotCreateAfterWriterExists( func TestArchiveSweep_SkipsDeletedWebhookTargets(t *testing.T) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) path := env.seedArchiveRows( - t, webhookID, time.Now().Add(-48*time.Hour), + t, tgt, time.Now().Add(-48*time.Hour), ) require.NoError( t, env.mainDB.DB(). - Where("webhook_id = ?", webhookID). + Where("webhook_id = ?", tgt.WebhookID). Delete(&database.Target{}).Error, ) @@ -842,14 +855,14 @@ func TestArchiveSweep_SkipsDeletedWebhookTargets(t *testing.T) { } // TestArchiveSweep_ConcurrentWrites proves the sweep serialises -// against writes through the per-webhook writer mutex. Run -// under -race, an unsynchronised sweep would be caught here. +// against writes through the target's writer mutex. Run under +// -race, an unsynchronised sweep would be caught here. func TestArchiveSweep_ConcurrentWrites(t *testing.T) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) webhookDB := testWebhookDB(t) @@ -862,13 +875,10 @@ func TestArchiveSweep_ConcurrentWrites(t *testing.T) { for range sweepConcurrentWrites { event := seedEvent(t, webhookDB, `{"n":1}`) - event.WebhookID = webhookID deliveries = append( deliveries, - seedDatabaseTargetDelivery( - t, webhookDB, event, `{"expiry":"1h"}`, - ), + seedDatabaseTargetDelivery(t, webhookDB, event, tgt), ) } @@ -894,7 +904,7 @@ func TestArchiveSweep_ConcurrentWrites(t *testing.T) { wg.Wait() - assert.FileExists(t, env.archivePath(webhookID)) + assert.FileExists(t, env.archivePath(tgt)) } // TestArchiveSweeper_StopsCleanly proves the background loop @@ -902,11 +912,11 @@ func TestArchiveSweep_ConcurrentWrites(t *testing.T) { func TestArchiveSweeper_StopsCleanly(t *testing.T) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) - webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) + tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) env.seedArchiveRows( - t, webhookID, time.Now().Add(-48*time.Hour), + t, tgt, time.Now().Add(-48*time.Hour), ) env.sweeper.ExportSetInterval(time.Millisecond) @@ -930,7 +940,7 @@ func TestArchiveSweeper_StopHookHonoursStopTimeout( ) { t.Parallel() - env := setupSweeperTest(t) + env := setupArchiveTest(t) lc := &recordingLifecycle{} env.sweeper.ExportRegisterHooks(lc) diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index 0265ffd..45050cf 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -122,21 +122,24 @@ 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. +// Archives is how the handlers keep the database targets' archive +// files in step with the configuration. Deleting a webhook or a +// target releases the cached archive writers, whose open file +// handles would otherwise outlive them; renaming one renames the +// archive files, which are named for the webhook and the target +// (see ArchiveFileName). // -// 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. +// It is deliberately separate from Notifier: archiving lifecycle +// is not notification, and a small 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 { +// Neither eviction deletes an archive file. Both are idempotent +// and are no-ops for a webhook or target with no engine state. +type Archives interface { EvictWebhook(webhookID string) + EvictTarget(targetID string) + Rename(targetID, webhookName, targetName string) error } // EngineParams are the fx dependencies for the delivery @@ -181,7 +184,7 @@ type Engine struct { httpTarget *httpTarget // dbTarget is retained so the engine can reach the archive - // writer registry for webhook eviction and the idle sweep. + // writer registry for eviction, renames and the idle sweep. dbTarget *databaseTarget // inflight is the set of deliveries this engine currently owns. @@ -249,17 +252,44 @@ 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. +// EvictWebhook implements Archives. The cached archive writer of +// every database target of the webhook is dropped from the +// registry and its file handle closed. The archive files +// themselves are left on disk — they are long-term storage the +// operator owns. func (e *Engine) EvictWebhook(webhookID string) { if e.dbTarget == nil { return } - e.dbTarget.evict(webhookID) + e.dbTarget.evictWebhook(webhookID) +} + +// EvictTarget implements Archives. It is EvictWebhook for a single +// database target, and leaves the archive file on disk the same +// way. +func (e *Engine) EvictTarget(targetID string) { + if e.dbTarget == nil { + return + } + + e.dbTarget.evict(targetID) +} + +// Rename implements Archives. It renames a database target's +// archive file to ArchiveFileName(webhookName, targetName, +// targetID), under the lock the target's archive writes and the +// idle sweep take. It never replaces a file: if one already has the +// new name, the error is ErrArchiveNameTaken. The caller renames +// before it saves the new name: see databaseTarget.rename. +func (e *Engine) Rename( + targetID, webhookName, targetName string, +) error { + if e.dbTarget == nil { + return nil + } + + return e.dbTarget.rename(targetID, webhookName, targetName) } // ScheduleRetry schedules a task to be re-enqueued onto the @@ -366,7 +396,8 @@ func (e *Engine) start() { // 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 +// engine, and deleting or renaming a webhook or target only closes +// or moves one. If the pool did // not drain in time, the writers are left open, as a kill would // leave them. Closing them would wait for any write in progress, // and a worker still running would then open new writers that diff --git a/internal/delivery/engine_lifecycle_test.go b/internal/delivery/engine_lifecycle_test.go index 9fadf9d..69e0d43 100644 --- a/internal/delivery/engine_lifecycle_test.go +++ b/internal/delivery/engine_lifecycle_test.go @@ -2,7 +2,6 @@ package delivery_test import ( "context" - "fmt" "path/filepath" "testing" "time" @@ -10,6 +9,7 @@ import ( "github.com/google/uuid" "github.com/stretchr/testify/require" "go.uber.org/fx" + "gorm.io/gorm/clause" "sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/delivery" ) @@ -272,22 +272,35 @@ 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 { +// deliverToArchive gives the setup's webhook a database target, +// runs one delivery to it through the running engine, and returns +// the target's ID and archive file path. The archive writer holds +// the file open afterwards. +func deliverToArchive(t *testing.T, s iSetup) (string, string) { t.Helper() + iCreateWebhook(t, s.MainDB, s.WebhookID, "hook") + + tgt := &database.Target{ + WebhookID: s.WebhookID, + Name: "archive", + Type: database.TargetTypeDatabase, + } + require.NoError( + t, s.MainDB.Omit(clause.Associations).Create(tgt).Error, + ) + deliveryID, task := seedLogTask(t, s) + task.TargetID = tgt.ID task.TargetType = database.TargetTypeDatabase s.Engine.Notify([]delivery.Task{task}) iWaitForDelivered(t, s.WebhookDB, deliveryID) - return filepath.Join( + return tgt.ID, filepath.Join( filepath.Dir(s.DBMgr.DBPath(s.WebhookID)), - fmt.Sprintf("archive-%s.db", s.WebhookID), + "archive-hook-archive-"+tgt.ID+".db", ) } @@ -304,7 +317,7 @@ func TestEngine_StopHookClosesArchives(t *testing.T) { lc := startEngineViaHook(t, s.Engine) - path := deliverToArchive(t, s) + _, path := deliverToArchive(t, s) require.FileExists( t, path+"-wal", "an open archive should have a -wal for the stop to remove", @@ -338,7 +351,7 @@ func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) { lc := startEngineViaHook(t, s.Engine) - deliverToArchive(t, s) + targetID, _ := deliverToArchive(t, s) release := make(chan struct{}) @@ -352,7 +365,7 @@ func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) { requireStopHookExpires(t, lc.hooks[0], "delivery engine") require.True( - t, s.Engine.ExportArchiveHandleOpen(s.WebhookID), + t, s.Engine.ExportArchiveHandleOpen(targetID), "a stop that timed out must not close archive writers", ) } diff --git a/internal/delivery/engine_test.go b/internal/delivery/engine_test.go index abe668b..15733cf 100644 --- a/internal/delivery/engine_test.go +++ b/internal/delivery/engine_test.go @@ -351,23 +351,15 @@ func TestDeliverDatabase_ImmediateSuccess( db := testWebhookDB(t) - // The database target archives for real now, so the engine - // needs a webhook DB manager to locate the data directory. - e := delivery.NewTestEngineWithDB( - nil, - database.NewTestWebhookDBManager(t.TempDir()), - slog.New(slog.NewTextHandler( - os.Stderr, - &slog.HandlerOptions{Level: slog.LevelDebug}, - )), - &http.Client{Timeout: 5 * time.Second}, - 1, - ) + // The database target archives for real, so the engine needs + // the target in the main database and a data directory. + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, "") event := seedEvent(t, db, `{"db":"target"}`) - d := seedDatabaseTargetDelivery(t, db, event, "") + d := seedDatabaseTargetDelivery(t, db, event, tgt) - e.ExportDeliverDatabase(db, d) + env.eng.ExportDeliverDatabase(db, d) var updated database.Delivery @@ -1336,32 +1328,27 @@ func TestProcessDelivery_RoutesToCorrectHandler( db := testWebhookDB(t) - // The database target archives for real now, so the engine - // needs a webhook DB manager to locate the data directory. - e := delivery.NewTestEngineWithDB( - nil, - database.NewTestWebhookDBManager(t.TempDir()), - slog.New(slog.NewTextHandler( - os.Stderr, - &slog.HandlerOptions{Level: slog.LevelDebug}, - )), - &http.Client{Timeout: 5 * time.Second}, - 1, - ) + // The database target archives for real, so the engine needs + // the target in the main database and a data directory. + env := setupArchiveTest(t) + archive := env.seedDatabaseTarget(t, "") tests := []struct { name string targetType database.TargetType + targetID string wantStatus database.DeliveryStatus }{ { "database target", database.TargetTypeDatabase, + archive.ID, database.DeliveryStatusDelivered, }, { "log target", database.TargetTypeLog, + uuid.New().String(), database.DeliveryStatusDelivered, }, } @@ -1371,7 +1358,7 @@ func TestProcessDelivery_RoutesToCorrectHandler( t.Parallel() runRoutingSubtest( - t, db, e, tt.targetType, + t, db, env.eng, tt.targetType, tt.targetID, tt.wantStatus, ) }) @@ -1383,6 +1370,7 @@ func runRoutingSubtest( db *gorm.DB, e *delivery.Engine, targetType database.TargetType, + targetID string, wantStatus database.DeliveryStatus, ) { t.Helper() @@ -1390,8 +1378,7 @@ func runRoutingSubtest( event := seedEvent(t, db, `{"routing":"test"}`) dlv := seedDelivery( - t, db, event.ID, - uuid.New().String(), + t, db, event.ID, targetID, database.DeliveryStatusPending, ) diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index 3a803cf..76e4df5 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -473,7 +473,7 @@ func NewTestCircuitBreaker( type ExportArchivedEvent = archivedEvent // ExportArchiveWriter wraps an archiveWriter so black-box tests -// can exercise the per-webhook archive file mechanics. +// can exercise the archive file mechanics. type ExportArchiveWriter struct { w *archiveWriter } @@ -548,6 +548,12 @@ func (e *ExportArchiveWriter) Evict() { e.w.evict() } +// Rename gives the archive file a new name in the same directory, +// as a rename of the webhook or target does. +func (e *ExportArchiveWriter) Rename(name string) error { + return e.w.rename(name) +} + // HandleOpen reports whether the writer currently holds an open // archive handle. func (e *ExportArchiveWriter) HandleOpen() bool { @@ -567,16 +573,16 @@ func (e *ExportArchiveWriter) Same( } // ExportArchiveWriterFor returns the archive writer the registry -// currently caches for a webhook, or nil when none is cached. It -// never creates one, so a test can hold a reference to the very -// writer an eviction is about to detach. +// currently caches for a database target, or nil when none is +// cached. It never creates one, so a test can hold a reference to +// the very writer an eviction is about to detach. func (e *Engine) ExportArchiveWriterFor( - webhookID string, + targetID string, ) *ExportArchiveWriter { e.dbTarget.mu.Lock() defer e.dbTarget.mu.Unlock() - w, ok := e.dbTarget.writers[webhookID] + w, ok := e.dbTarget.writers[targetID] if !ok { return nil } @@ -585,26 +591,26 @@ func (e *Engine) ExportArchiveWriterFor( } // ExportHasArchiveWriter reports whether the database target -// currently caches an archive writer for a webhook. +// type currently caches an archive writer for a target. func (e *Engine) ExportHasArchiveWriter( - webhookID string, + targetID string, ) bool { e.dbTarget.mu.Lock() defer e.dbTarget.mu.Unlock() - _, ok := e.dbTarget.writers[webhookID] + _, ok := e.dbTarget.writers[targetID] return ok } // ExportArchiveHandleOpen reports whether the cached archive -// writer for a webhook holds an open database handle. It +// writer for a target holds an open database handle. It // returns false when no writer is cached. func (e *Engine) ExportArchiveHandleOpen( - webhookID string, + targetID string, ) bool { e.dbTarget.mu.Lock() - w, ok := e.dbTarget.writers[webhookID] + w, ok := e.dbTarget.writers[targetID] e.dbTarget.mu.Unlock() if !ok { @@ -618,12 +624,12 @@ func (e *Engine) ExportArchiveHandleOpen( } // ExportEnsureArchiveWriter creates (if needed) and returns the -// archive file path of the cached writer for a webhook, so a +// archive file path of the cached writer for a target, so a // test can prime the registry the way a delivery would. func (e *Engine) ExportEnsureArchiveWriter( - webhookID string, + targetID string, ) (string, error) { - w, err := e.dbTarget.writerFor(webhookID) + w, err := e.dbTarget.writerFor(targetID) if err != nil { return "", err } @@ -631,14 +637,14 @@ func (e *Engine) ExportEnsureArchiveWriter( return w.path, nil } -// ExportSweepWriterFor takes a webhook's registry writer exactly +// ExportSweepWriterFor takes a target's registry writer exactly // as the idle sweep does, reporting whether the sweep had to // create the entry. It lets a test drive the registry through the // sweep's own entry point instead of choreographing goroutines. func (e *Engine) ExportSweepWriterFor( - webhookID string, + targetID string, ) (*ExportArchiveWriter, bool, error) { - w, created, err := e.dbTarget.sweepWriterFor(webhookID) + w, created, err := e.dbTarget.sweepWriterFor(targetID) if err != nil { return nil, false, err } @@ -649,9 +655,9 @@ func (e *Engine) ExportSweepWriterFor( // ExportReleaseSweepWriter releases a sweep-created registry entry // exactly as a finished sweep does. func (e *Engine) ExportReleaseSweepWriter( - webhookID string, w *ExportArchiveWriter, + targetID string, w *ExportArchiveWriter, ) { - e.dbTarget.releaseSweepWriter(webhookID, w.w) + e.dbTarget.releaseSweepWriter(targetID, w.w) } // NewTestArchiveSweeper builds an ArchiveSweeper backed by the diff --git a/internal/delivery/target_database.go b/internal/delivery/target_database.go index 7173b71..3fb52db 100644 --- a/internal/delivery/target_database.go +++ b/internal/delivery/target_database.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "path/filepath" + "strings" "sync" "time" @@ -11,22 +12,75 @@ import ( "sneak.berlin/go/webhooker/internal/database" ) -// databaseTarget is a no-retry target that archives the -// full inbound event into a per-webhook archive SQLite file, -// separate from the per-webhook event database. The event is -// already persisted in the per-webhook event DB by the time -// delivery runs; the database target additionally writes a -// durable long-term copy into archive-{webhookID}.db and then -// records a single attempt whose outcome reflects whether the -// archive write succeeded. See archiveWriter for the -// close/reopen, auto-recreate, and expiry semantics. +// archiveNameMaxLen is how many characters of a webhook or target +// name an archive file name keeps. +const archiveNameMaxLen = 40 + +// databaseTarget is a no-retry target that archives the full +// inbound event into the target's own archive SQLite file, separate +// from the per-webhook event database. The event is already +// persisted in the per-webhook event DB by the time delivery runs; +// the database target additionally writes a durable long-term copy +// into the file ArchiveFileName names and then records a single +// attempt whose outcome reflects whether the archive write +// succeeded. See archiveWriter for the close/reopen, auto-recreate, +// and expiry semantics. type databaseTarget struct { eng *Engine + // writers holds one archive writer per database target, keyed + // by target ID. mu sync.Mutex writers map[string]*archiveWriter } +// ArchiveFileName returns the file name of a database target's +// archive: archive-WEBHOOKNAME-TARGETNAME-TARGETID.db, with both +// names passed through archiveNamePart. The target ID keeps the +// name unique when two targets' names come out the same. +func ArchiveFileName(webhookName, targetName, targetID string) string { + return "archive-" + archiveNamePart(webhookName) + "-" + + archiveNamePart(targetName) + "-" + targetID + ".db" +} + +// archiveNamePart makes a webhook or target name safe to put in a +// file name. It is lowercased; ASCII letters and digits are kept, +// every other run of characters becomes a single "-", and no "-" is +// left at either end. It is cut to archiveNameMaxLen characters, and +// a name with nothing left is "unnamed". +func archiveNamePart(name string) string { + var b strings.Builder + + dash := false + + for _, r := range strings.ToLower(name) { + if (r < 'a' || r > 'z') && (r < '0' || r > '9') { + dash = b.Len() > 0 + + continue + } + + if dash { + b.WriteByte('-') + + dash = false + } + + b.WriteRune(r) + } + + part := b.String() + if len(part) > archiveNameMaxLen { + part = strings.TrimRight(part[:archiveNameMaxLen], "-") + } + + if part == "" { + return "unnamed" + } + + return part +} + // Deliver implements Target. It archives the event, then // records one successful attempt and marks the delivery // delivered. An archiving error fails the delivery: the @@ -92,7 +146,7 @@ func (t *databaseTarget) Deliver( ) } -// archive writes the full event as a row into the webhook's +// archive writes the full event as a row into the target's // archive database, honouring the optional per-target expiry // parsed from the target config JSON. func (t *databaseTarget) archive(d *database.Delivery) error { @@ -106,7 +160,7 @@ func (t *databaseTarget) archive(d *database.Delivery) error { return err } - w, err := t.writerFor(webhookID) + w, err := t.writerFor(d.TargetID) if err != nil { return err } @@ -124,30 +178,31 @@ func (t *databaseTarget) archive(d *database.Delivery) error { return w.write(row, expiry) } -// writerFor returns the archiveWriter for a webhook, creating -// and caching it on first use. Each webhook has one writer so -// its close/reopen debounce state is shared across concurrent -// deliveries. The archive file lives beside the per-webhook -// event database in the data directory. +// writerFor returns the archive writer for a database target, +// creating and caching it on first use. Each target has one writer +// so its close/reopen debounce state is shared across concurrent +// deliveries, and so a rename and the idle sweep take the same lock +// as its writes. func (t *databaseTarget) writerFor( - webhookID string, + targetID string, ) (*archiveWriter, error) { - path, err := t.archivePath(webhookID) - if err != nil { - return nil, err - } - t.mu.Lock() defer t.mu.Unlock() - if t.writers == nil { - t.writers = make(map[string]*archiveWriter) - } - - w, ok := t.writers[webhookID] + w, ok := t.writers[targetID] if !ok { - w = newArchiveWriter(path, t.eng.log) - t.writers[webhookID] = w + var err error + + w, err = t.newWriter(targetID) + if err != nil { + return nil, err + } + + if t.writers == nil { + t.writers = make(map[string]*archiveWriter) + } + + t.writers[targetID] = w } // A delivery claims the entry: even if the idle sweep created @@ -159,40 +214,39 @@ func (t *databaseTarget) writerFor( } // sweepWriterFor returns the archive writer the idle sweep should -// prune a webhook through, together with whether the sweep itself -// created the registry entry. +// prune a target's archive through, together with whether the sweep +// itself created the registry entry. // // The sweep must route its prune through the registered writer so // the writer's mutex orders it against concurrent writes, but it // must never leave a registry entry behind: a sweep that ran -// concurrently with the webhook's deletion would otherwise +// concurrently with the target's deletion would otherwise // re-create an entry that nothing will ever evict again, which is // exactly the leak eviction exists to prevent. An entry the sweep // creates is therefore marked sweep-owned and handed back to // releaseSweepWriter when the sweep is done. func (t *databaseTarget) sweepWriterFor( - webhookID string, + targetID string, ) (*archiveWriter, bool, error) { - path, err := t.archivePath(webhookID) + t.mu.Lock() + defer t.mu.Unlock() + + w, ok := t.writers[targetID] + if ok { + return w, false, nil + } + + w, err := t.newWriter(targetID) if err != nil { return nil, false, err } - t.mu.Lock() - defer t.mu.Unlock() - if t.writers == nil { t.writers = make(map[string]*archiveWriter) } - w, ok := t.writers[webhookID] - if ok { - return w, false, nil - } - - w = newArchiveWriter(path, t.eng.log) w.sweepOwned = true - t.writers[webhookID] = w + t.writers[targetID] = w return w, true, nil } @@ -209,57 +263,95 @@ func (t *databaseTarget) sweepWriterFor( // delivery that adopted the writer keeps a registered, evictable // one. func (t *databaseTarget) releaseSweepWriter( - webhookID string, w *archiveWriter, + targetID string, w *archiveWriter, ) { t.mu.Lock() defer t.mu.Unlock() - cur, ok := t.writers[webhookID] + cur, ok := t.writers[targetID] if !ok || cur != w || !cur.sweepOwned { return } - delete(t.writers, webhookID) + delete(t.writers, targetID) } -// 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) { +// newWriter builds the writer for a database target's archive. The +// file lives beside the webhook's event database in the data +// directory and is named for the webhook and the target as the main +// database has them now; from then on only rename changes the name +// the writer uses. It does not touch the archive file. +func (t *databaseTarget) newWriter( + targetID string, +) (*archiveWriter, error) { if t.eng.dbManager == nil { - return "", errArchiveNoDataDir + return nil, errArchiveNoDataDir } - dir := filepath.Dir(t.eng.dbManager.DBPath(webhookID)) + var target database.Target - return filepath.Join( - dir, fmt.Sprintf("archive-%s.db", webhookID), - ), nil + err := t.eng.database.DB(). + Preload("Webhook"). + First(&target, "id = ?", targetID).Error + if err != nil { + return nil, fmt.Errorf( + "loading database target %s: %w", targetID, err, + ) + } + + dir := filepath.Dir(t.eng.dbManager.DBPath(target.WebhookID)) + name := ArchiveFileName( + target.Webhook.Name, target.Name, target.ID, + ) + + w := newArchiveWriter(filepath.Join(dir, name), t.eng.log) + w.webhookID = target.WebhookID + + return w, 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. +// rename moves a database target's archive file to the name for +// webhookName and targetName. It goes through the target's writer, +// so the move holds the lock that writes and the idle sweep take, +// and later writes use the new name. +// +// The writer is created if there is none, and it stays cached. The +// handlers rename before they save the new name, so until the save +// the main database still has the old one; a delivery in that window +// must find this writer rather than build one from the old name. +func (t *databaseTarget) rename( + targetID, webhookName, targetName string, +) error { + w, err := t.writerFor(targetID) + if err != nil { + return err + } + + return w.rename(ArchiveFileName(webhookName, targetName, targetID)) +} + +// evict drops a database target's archive writer from the registry +// and closes its handle, so a deleted target 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 +// targets 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) { +// Eviction is idempotent and silent for a target with no writer, +// which is the common case: only a database target that has +// received an event or been renamed has one. It never deletes the +// archive file. +func (t *databaseTarget) evict(targetID string) { t.mu.Lock() - w, ok := t.writers[webhookID] + w, ok := t.writers[targetID] if ok { - delete(t.writers, webhookID) + delete(t.writers, targetID) } t.mu.Unlock() @@ -272,13 +364,41 @@ func (t *databaseTarget) evict(webhookID string) { t.eng.log.Info( "evicted archive writer", - "webhook_id", webhookID, + "target_id", targetID, "path", w.path, ) } +// evictWebhook evicts, exactly as evict does, the writer of every +// database target of a webhook. +func (t *databaseTarget) evictWebhook(webhookID string) { + t.mu.Lock() + + var gone []*archiveWriter + + for targetID, w := range t.writers { + if w.webhookID == webhookID { + delete(t.writers, targetID) + + gone = append(gone, w) + } + } + + t.mu.Unlock() + + for _, w := range gone { + w.evict() + + t.eng.log.Info( + "evicted archive writer", + "webhook_id", webhookID, + "path", w.path, + ) + } +} + // evictAll evicts every cached archive writer, exactly as evict -// does for one webhook. The engine calls it at shutdown, once its +// does for one target. 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. @@ -295,38 +415,25 @@ func (t *databaseTarget) evictAll() { } } -// 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. +// sweepArchive prunes one database target's archive of rows older +// than expiry, without requiring a write. A missing archive file is +// left missing (see sweepExpired), so a sweep never creates an +// archive for a target that has never received an event. // // It also never leaves a registry entry behind: an entry it had // to create to reach the writer's mutex is released again once -// the prune is done, so a sweep racing a webhook deletion cannot +// the prune is done, so a sweep racing a target deletion cannot // resurrect the writer the eviction just dropped. -func (t *databaseTarget) sweepWebhook( - webhookID string, expiry time.Duration, +func (t *databaseTarget) sweepArchive( + targetID 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, created, err := t.sweepWriterFor(webhookID) + w, created, err := t.sweepWriterFor(targetID) if err != nil { return err } if created { - defer t.releaseSweepWriter(webhookID, w) + defer t.releaseSweepWriter(targetID, w) } return w.sweepExpired(expiry) diff --git a/internal/delivery/target_database_archive.go b/internal/delivery/target_database_archive.go index 0937854..a1ee122 100644 --- a/internal/delivery/target_database_archive.go +++ b/internal/delivery/target_database_archive.go @@ -4,8 +4,10 @@ import ( "encoding/json" "errors" "fmt" + "io/fs" "log/slog" "os" + "path/filepath" "sync" "time" @@ -41,7 +43,7 @@ const ( var ( // errArchiveMissingWebhookID is returned when an event to - // archive has no webhook id to key its archive file on. + // archive has no webhook id to record in its archive row. errArchiveMissingWebhookID = errors.New( "cannot archive event without a webhook id", ) @@ -61,13 +63,19 @@ var ( ) // 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. + // been evicted (its target or its webhook was deleted) 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", ) + + // ErrArchiveNameTaken is returned when an archive cannot be + // renamed because a file already has the new name. That file may + // be an archive with rows of its own, so it is never replaced. + ErrArchiveNameTaken = errors.New( + "a file already has the archive's new name", + ) ) // databaseTargetConfig is the optional per-target JSON config @@ -80,7 +88,7 @@ type databaseTargetConfig struct { } // archivedEvent is one fully captured webhook event stored in a -// per-webhook archive database for long-term retention. It is a +// database target's archive for long-term retention. It is a // self-contained copy — independent of the per-webhook event // database, which may prune events under its own retention. type archivedEvent struct { @@ -170,8 +178,8 @@ func ValidateArchiveExpiry(expiry string) error { return nil } -// archiveWriter owns one per-webhook archive SQLite file. It -// serialises writes, and after each write closes and reopens +// archiveWriter owns one database target's archive SQLite file. +// It serialises writes, and after each write closes and reopens // the file (debounced to at most once per debounce window) so // an operator can move the file away for offline archiving. The // next write recreates a moved or removed file, because the @@ -187,16 +195,21 @@ type archiveWriter struct { 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. + // 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 + // webhookID is the webhook the archive's target belongs to, + // so deleting the webhook can find its writers. It is set + // when the writer is created and never changes. + webhookID string + // sweepOwned marks a registry entry that the idle sweep - // created because no writer was cached for the webhook. The + // created because no writer was cached for the target. The // sweep removes such an entry again when it is done, so a // sweep can never leave — or resurrect — a registry entry - // for a webhook that has been deleted. A delivery that adopts + // for a target that has been deleted. A delivery that adopts // the writer clears the flag, handing the entry to the // registry proper. // @@ -385,11 +398,62 @@ func (w *archiveWriter) sweepExpired(expiry time.Duration) error { return nil } +// rename gives the archive file a new name in the same directory, +// and the writer uses the file under that name from now on. The +// handle is closed first, which folds the -wal into the .db; any +// -wal or -shm still beside the file (left by a crash) is moved with +// it, because SQLite finds them by name. A missing file is not an +// error: the operator may have moved it away, and the next write +// creates it under the new name. +// +// If a file already has the new name, nothing is moved and the +// error is ErrArchiveNameTaken. +func (w *archiveWriter) rename(name string) error { + w.mu.Lock() + defer w.mu.Unlock() + + if w.evicted { + return fmt.Errorf( + "%w: %s", errArchiveWriterEvicted, w.path, + ) + } + + path := filepath.Join(filepath.Dir(w.path), name) + if path == w.path { + return nil + } + + suffixes := []string{"", "-wal", "-shm"} + + for _, suffix := range suffixes { + if fileExists(path + suffix) { + return fmt.Errorf( + "%w: %s", ErrArchiveNameTaken, name+suffix, + ) + } + } + + w.close() + + for _, suffix := range suffixes { + err := os.Rename(w.path+suffix, path+suffix) + if err != nil && !errors.Is(err, fs.ErrNotExist) { + return fmt.Errorf( + "renaming archive %s to %s: %w", w.path, path, err, + ) + } + } + + w.path = path + + 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. +// called when the writer leaves the registry, because its target +// or its webhook was deleted, or at shutdown. 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() diff --git a/internal/delivery/target_database_evict_test.go b/internal/delivery/target_database_evict_test.go index 095823d..5022ba4 100644 --- a/internal/delivery/target_database_evict_test.go +++ b/internal/delivery/target_database_evict_test.go @@ -17,85 +17,109 @@ import ( "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) { +// deliverTo archives one event to a database target, leaving the +// target's writer cached with its handle open. +func deliverTo( + t *testing.T, env *archiveEnv, tgt *database.Target, +) { 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", + env.eng.ExportDeliverDatabase( + webhookDB, seedDatabaseTargetDelivery(t, webhookDB, event, tgt), ) +} - eng.EvictWebhook(webhookID) +// TestEvictWebhook_ClosesAndRemovesWriter proves that evicting +// a webhook drops the archive writers of its database targets +// from the registry and closes their open handles, rather than +// leaving them alive for the process lifetime, and leaves another +// webhook's writer alone. +func TestEvictWebhook_ClosesAndRemovesWriter(t *testing.T) { + t.Parallel() - assert.False( - t, eng.ExportHasArchiveWriter(webhookID), - "eviction should remove the registry entry", - ) - assert.False( - t, eng.ExportArchiveHandleOpen(webhookID), - "eviction should close the archive handle", - ) + env := setupArchiveTest(t) + first := env.seedDatabaseTarget(t, "") + second := env.addDatabaseTarget(t, first.WebhookID, "") + other := env.seedDatabaseTarget(t, "") - archivePath := filepath.Join( - dataDir, fmt.Sprintf("archive-%s.db", webhookID), + for _, tgt := range []*database.Target{first, second, other} { + deliverTo(t, env, tgt) + + require.True( + t, env.eng.ExportArchiveHandleOpen(tgt.ID), + "the writer should hold an open handle after a write", + ) + } + + env.eng.EvictWebhook(first.WebhookID) + + for _, tgt := range []*database.Target{first, second} { + assert.False( + t, env.eng.ExportHasArchiveWriter(tgt.ID), + "eviction should remove the registry entry", + ) + assert.False( + t, env.eng.ExportArchiveHandleOpen(tgt.ID), + "eviction should close the archive handle", + ) + assert.FileExists( + t, env.archivePath(tgt), + "eviction must not delete the archive file", + ) + } + + assert.True( + t, env.eng.ExportArchiveHandleOpen(other.ID), + "another webhook's writer must be left alone", ) +} + +// TestEvictTarget_LeavesOtherTargets proves that evicting one +// database target leaves the writer of another target of the same +// webhook in place. +func TestEvictTarget_LeavesOtherTargets(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + doomed := env.seedDatabaseTarget(t, "") + kept := env.addDatabaseTarget(t, doomed.WebhookID, "") + + deliverTo(t, env, doomed) + deliverTo(t, env, kept) + + env.eng.EvictTarget(doomed.ID) + + assert.False(t, env.eng.ExportHasArchiveWriter(doomed.ID)) assert.FileExists( - t, archivePath, + t, env.archivePath(doomed), "eviction must not delete the archive file", ) + assert.True( + t, env.eng.ExportArchiveHandleOpen(kept.ID), + "the other target's writer must be left alone", + ) } // 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. +// for the common case of a webhook or target that never had an +// archive writer, and that repeating it does not panic. func TestEvictWebhook_UnknownWebhookIsNoOp(t *testing.T) { t.Parallel() - eng, _ := evictTestEngine(t) + env := setupArchiveTest(t) assert.NotPanics(t, func() { - eng.EvictWebhook("no-such-webhook") - eng.EvictWebhook("no-such-webhook") + env.eng.EvictWebhook("no-such-webhook") + env.eng.EvictWebhook("no-such-webhook") + env.eng.EvictTarget("no-such-target") + env.eng.EvictTarget("no-such-target") }) assert.False( - t, eng.ExportHasArchiveWriter("no-such-webhook"), + t, env.eng.ExportHasArchiveWriter("no-such-target"), "eviction must not create a writer", ) } @@ -289,17 +313,14 @@ func TestEvictWebhook_RacingWriteDoesNotReopenHandle( ) { t.Parallel() - eng, _ := evictTestEngine(t) - - webhookDB := testWebhookDB(t) - event := seedEvent(t, webhookDB, `{"archived":true}`) - d := seedDatabaseTargetDelivery(t, webhookDB, event, "") + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, "") // Prime the registry so the test can hold the very writer the // eviction is about to detach. - eng.ExportDeliverDatabase(webhookDB, d) + deliverTo(t, env, tgt) - w := eng.ExportArchiveWriterFor(event.WebhookID) + w := env.eng.ExportArchiveWriterFor(tgt.ID) require.NotNil(t, w) require.True(t, w.HandleOpen()) @@ -309,7 +330,7 @@ func TestEvictWebhook_RacingWriteDoesNotReopenHandle( // eviction has to contend for the writer's mutex. race.awaitFirstWrite() - eng.EvictWebhook(event.WebhookID) + env.eng.EvictWebhook(tgt.WebhookID) sawEvicted, otherErr := race.wait() @@ -324,41 +345,33 @@ func TestEvictWebhook_RacingWriteDoesNotReopenHandle( "been evicted", ) assert.False( - t, eng.ExportHasArchiveWriter(event.WebhookID), + t, env.eng.ExportHasArchiveWriter(tgt.ID), "the registry entry must stay gone", ) } // TestEvictWebhook_LaterDeliveryRecreatesWriter proves eviction -// does not break archiving for a webhook that is still alive: a +// does not break archiving for a target that is still alive: a // subsequent delivery gets a brand new writer from the registry. // It says nothing about the evicted writer itself — that is what // TestEvictedWriter_WriteDoesNotReopenFile covers. func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) { t.Parallel() - eng, _ := evictTestEngine(t) + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, "") - webhookDB := testWebhookDB(t) - event := seedEvent(t, webhookDB, `{"archived":true}`) - d := seedDatabaseTargetDelivery(t, webhookDB, event, "") + deliverTo(t, env, tgt) + require.True(t, env.eng.ExportHasArchiveWriter(tgt.ID)) - eng.ExportDeliverDatabase(webhookDB, d) - require.True( - t, eng.ExportHasArchiveWriter(event.WebhookID), - ) + env.eng.EvictWebhook(tgt.WebhookID) - eng.EvictWebhook(event.WebhookID) - - // A fresh delivery for the same webhook gets a brand new + // A fresh delivery for the same target gets a brand new // writer from the registry, so archiving keeps working. - second := seedDatabaseTargetDelivery( - t, webhookDB, event, "", - ) - eng.ExportDeliverDatabase(webhookDB, second) + deliverTo(t, env, tgt) assert.True( - t, eng.ExportHasArchiveWriter(event.WebhookID), + t, env.eng.ExportHasArchiveWriter(tgt.ID), "a later delivery should recreate the writer", ) } @@ -370,19 +383,16 @@ func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) { func TestEngineStop_WriteAfterStopIsRefused(t *testing.T) { t.Parallel() - eng, _ := evictTestEngine(t) + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, "") - webhookDB := testWebhookDB(t) - event := seedEvent(t, webhookDB, `{"archived":true}`) - d := seedDatabaseTargetDelivery(t, webhookDB, event, "") + deliverTo(t, env, tgt) - eng.ExportDeliverDatabase(webhookDB, d) - - w := eng.ExportArchiveWriterFor(event.WebhookID) + w := env.eng.ExportArchiveWriterFor(tgt.ID) require.NotNil(t, w) require.True(t, w.HandleOpen()) - require.NoError(t, eng.ExportStop(context.Background())) + require.NoError(t, env.eng.ExportStop(context.Background())) err := w.Write(evictTestRow("ev-after-stop"), 0) @@ -395,7 +405,7 @@ func TestEngineStop_WriteAfterStopIsRefused(t *testing.T) { "a refused write must not reopen the archive", ) assert.False( - t, eng.ExportHasArchiveWriter(event.WebhookID), + t, env.eng.ExportHasArchiveWriter(tgt.ID), "the stop should empty the registry", ) diff --git a/internal/delivery/target_database_test.go b/internal/delivery/target_database_test.go index dc2e65c..ba4842d 100644 --- a/internal/delivery/target_database_test.go +++ b/internal/delivery/target_database_test.go @@ -4,13 +4,12 @@ import ( "database/sql" "fmt" "log/slog" - "net/http" "os" "path/filepath" + "strings" "testing" "time" - "github.com/google/uuid" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "gorm.io/driver/sqlite" @@ -74,25 +73,18 @@ func removeArchiveFiles(t *testing.T, path string) { // TestDeliverDatabase_ArchivesEvent verifies that delivering to // a database target marks the delivery delivered and archives -// the full event into a separate per-webhook archive file. +// the full event into the target's own archive file. func TestDeliverDatabase_ArchivesEvent(t *testing.T) { t.Parallel() - dataDir := t.TempDir() - dbMgr := database.NewTestWebhookDBManager(dataDir) - - e := delivery.NewTestEngineWithDB( - nil, dbMgr, - archiveTestLogger(), - &http.Client{Timeout: 5 * time.Second}, - 1, - ) + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, "") webhookDB := testWebhookDB(t) event := seedEvent(t, webhookDB, `{"archived":true}`) - d := seedDatabaseTargetDelivery(t, webhookDB, event, "") + d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt) - e.ExportDeliverDatabase(webhookDB, d) + env.eng.ExportDeliverDatabase(webhookDB, d) var updated database.Delivery @@ -105,8 +97,7 @@ func TestDeliverDatabase_ArchivesEvent(t *testing.T) { ) archivePath := filepath.Join( - dataDir, - fmt.Sprintf("archive-%s.db", event.WebhookID), + env.dataDir, "archive-sweep-test-archive-"+tgt.ID+".db", ) assert.FileExists(t, archivePath) @@ -288,31 +279,31 @@ func TestParseArchiveExpiry(t *testing.T) { } } -// seedDatabaseTargetDelivery seeds a pending delivery for a -// database target with the given config JSON and returns the -// in-memory delivery the target handler is invoked with. +// seedDatabaseTargetDelivery seeds a pending delivery of an event +// to a database target and returns the in-memory delivery the +// target handler is invoked with. func seedDatabaseTargetDelivery( t *testing.T, webhookDB *gorm.DB, event database.Event, - config string, + tgt *database.Target, ) *database.Delivery { t.Helper() dlv := seedDelivery( - t, webhookDB, event.ID, uuid.New().String(), + t, webhookDB, event.ID, tgt.ID, database.DeliveryStatusPending, ) d := &database.Delivery{ EventID: event.ID, - TargetID: dlv.TargetID, + TargetID: tgt.ID, Status: database.DeliveryStatusPending, Event: event, Target: database.Target{ - Name: "test-db", + Name: tgt.Name, Type: database.TargetTypeDatabase, - Config: config, + Config: tgt.Config, }, } d.ID = dlv.ID @@ -330,22 +321,14 @@ func TestDeliverDatabase_ArchiveFailureFailsDelivery( ) { t.Parallel() - dataDir := t.TempDir() - - e := delivery.NewTestEngineWithDB( - nil, database.NewTestWebhookDBManager(dataDir), - archiveTestLogger(), - &http.Client{Timeout: 5 * time.Second}, - 1, - ) + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, `{"expiry":"nonsense"}`) webhookDB := testWebhookDB(t) event := seedEvent(t, webhookDB, `{"archived":false}`) - d := seedDatabaseTargetDelivery( - t, webhookDB, event, `{"expiry":"nonsense"}`, - ) + d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt) - e.ExportDeliverDatabase(webhookDB, d) + env.eng.ExportDeliverDatabase(webhookDB, d) var updated database.Delivery @@ -373,10 +356,7 @@ func TestDeliverDatabase_ArchiveFailureFailsDelivery( ) assert.NoFileExists(t, - filepath.Join( - dataDir, - fmt.Sprintf("archive-%s.db", event.WebhookID), - ), + env.archivePath(tgt), "no archive file should exist for a failed config", ) } @@ -400,3 +380,247 @@ func TestValidateArchiveExpiry(t *testing.T) { ) } } + +// TestArchiveFileName pins the archive file name and the rules +// that make a webhook or target name safe to put in it. +func TestArchiveFileName(t *testing.T) { + t.Parallel() + + const id = "3f2a1c9e-8d4b-4c1a-9e2f-0a1b2c3d4e5f" + + cases := []struct { + name string + webhook string + target string + want string + }{ + { + "plain names", "orders", "archive", + "archive-orders-archive-" + id + ".db", + }, + { + "lowercased", "Orders", "Main Archive", + "archive-orders-main-archive-" + id + ".db", + }, + { + "a run of other characters is one dash", + `a /\..b`, "c__--d", + "archive-a-b-c-d-" + id + ".db", + }, + { + "no dash at either end", " --orders!! ", "(archive)", + "archive-orders-archive-" + id + ".db", + }, + { + "path separators", "../../etc/passwd", "a/b", + "archive-etc-passwd-a-b-" + id + ".db", + }, + { + "letters outside ASCII are dropped", + "Bestellungen Größe", "café", + "archive-bestellungen-gr-e-caf-" + id + ".db", + }, + { + "nothing left is unnamed", "", "!!!", + "archive-unnamed-unnamed-" + id + ".db", + }, + { + "cut to 40 characters", strings.Repeat("a", 50), "x", + "archive-" + strings.Repeat("a", 40) + "-x-" + id + ".db", + }, + { + "no dash left by the cut", + strings.Repeat("a", 39) + " b", "x", + "archive-" + strings.Repeat("a", 39) + "-x-" + id + ".db", + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + assert.Equal( + t, tc.want, + delivery.ArchiveFileName(tc.webhook, tc.target, id), + ) + }) + } +} + +// TestDeliverDatabase_EachTargetHasItsOwnArchive proves two +// database targets of one webhook archive into separate files. +func TestDeliverDatabase_EachTargetHasItsOwnArchive(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + first := env.seedDatabaseTarget(t, "") + second := env.addDatabaseTarget(t, first.WebhookID, "") + + webhookDB := testWebhookDB(t) + event := seedEvent(t, webhookDB, `{"n":1}`) + + for _, tgt := range []*database.Target{first, second} { + env.eng.ExportDeliverDatabase( + webhookDB, + seedDatabaseTargetDelivery(t, webhookDB, event, tgt), + ) + } + + require.NotEqual( + t, env.archivePath(first), env.archivePath(second), + ) + assert.Equal( + t, []string{event.ID}, + archivedEventIDs(t, env.archivePath(first)), + ) + assert.Equal( + t, []string{event.ID}, + archivedEventIDs(t, env.archivePath(second)), + ) +} + +// TestRename_MovesTheFile proves a rename moves the archive, rows +// and all, and that later writes go to the new name. +func TestRename_MovesTheFile(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, "") + oldPath := env.archivePath(tgt) + + webhookDB := testWebhookDB(t) + first := seedEvent(t, webhookDB, `{"n":1}`) + env.eng.ExportDeliverDatabase( + webhookDB, seedDatabaseTargetDelivery(t, webhookDB, first, tgt), + ) + require.FileExists(t, oldPath) + + require.NoError( + t, env.eng.Rename(tgt.ID, "Orders", "Long Term"), + ) + + newPath := filepath.Join( + env.dataDir, "archive-orders-long-term-"+tgt.ID+".db", + ) + + assert.NoFileExists(t, oldPath) + assert.Equal(t, []string{first.ID}, archivedEventIDs(t, newPath)) + + second := seedEvent(t, webhookDB, `{"n":2}`) + env.eng.ExportDeliverDatabase( + webhookDB, + seedDatabaseTargetDelivery(t, webhookDB, second, tgt), + ) + + assert.ElementsMatch( + t, []string{first.ID, second.ID}, + archivedEventIDs(t, newPath), + ) + assert.NoFileExists( + t, oldPath, "a write after the rename must use the new name", + ) +} + +// TestRename_NeverReplacesAFile plants a file at the new name and +// proves the rename is refused, the planted file survives, and the +// archive keeps its name and its rows. +func TestRename_NeverReplacesAFile(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, "") + oldPath := env.archivePath(tgt) + + webhookDB := testWebhookDB(t) + first := seedEvent(t, webhookDB, `{"n":1}`) + env.eng.ExportDeliverDatabase( + webhookDB, seedDatabaseTargetDelivery(t, webhookDB, first, tgt), + ) + + newPath := filepath.Join( + env.dataDir, "archive-orders-long-term-"+tgt.ID+".db", + ) + require.NoError(t, os.WriteFile(newPath, []byte("planted"), 0o600)) + + require.ErrorIs( + t, env.eng.Rename(tgt.ID, "Orders", "Long Term"), + delivery.ErrArchiveNameTaken, + ) + + //nolint:gosec // reads the file the test planted under t.TempDir() + planted, err := os.ReadFile(newPath) + require.NoError(t, err) + assert.Equal(t, "planted", string(planted)) + + second := seedEvent(t, webhookDB, `{"n":2}`) + env.eng.ExportDeliverDatabase( + webhookDB, + seedDatabaseTargetDelivery(t, webhookDB, second, tgt), + ) + + assert.ElementsMatch( + t, []string{first.ID, second.ID}, + archivedEventIDs(t, oldPath), + ) +} + +// TestRename_BeforeTheNameIsSaved covers the order the handlers +// use: they rename before they save the new name, so a delivery in +// between must write under the new name although the main database +// still has the old one. It also shows that renaming an archive that +// does not exist yet is not an error. +func TestRename_BeforeTheNameIsSaved(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, "") + + require.NoError( + t, env.eng.Rename(tgt.ID, "Orders", "Archive"), + ) + + webhookDB := testWebhookDB(t) + event := seedEvent(t, webhookDB, `{"n":1}`) + env.eng.ExportDeliverDatabase( + webhookDB, seedDatabaseTargetDelivery(t, webhookDB, event, tgt), + ) + + assert.FileExists( + t, + filepath.Join( + env.dataDir, "archive-orders-archive-"+tgt.ID+".db", + ), + ) + assert.NoFileExists(t, env.archivePath(tgt)) +} + +// TestArchiveWriter_RenameMovesSidecars proves a rename carries +// the -wal and -shm a crash can leave beside an archive no handle +// has opened since. SQLite finds them by name, so a -wal left +// behind would lose the transactions it holds. +func TestArchiveWriter_RenameMovesSidecars(t *testing.T) { + t.Parallel() + + dir := t.TempDir() + oldPath := filepath.Join(dir, "archive-old.db") + newPath := filepath.Join(dir, "archive-new.db") + + for _, suffix := range archiveFileSuffixes() { + require.NoError( + t, os.WriteFile(oldPath+suffix, []byte(suffix), 0o600), + ) + } + + w := delivery.NewExportArchiveWriter( + oldPath, archiveTestLogger(), 0, + ) + + require.NoError(t, w.Rename("archive-new.db")) + + for _, suffix := range archiveFileSuffixes() { + assert.NoFileExists(t, oldPath+suffix) + assert.FileExists(t, newPath+suffix) + } + + assert.Equal(t, newPath, w.Path()) +} diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index 6eca35f..934297c 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -62,7 +62,7 @@ type HandlersParams struct { Session *session.Session Middleware *middleware.Middleware Notifier delivery.Notifier - Evictor delivery.WebhookEvictor + Archives delivery.Archives SSRFGuard *delivery.Guard } @@ -77,7 +77,7 @@ type Handlers struct { session *session.Session mw *middleware.Middleware notifier delivery.Notifier - evictor delivery.WebhookEvictor + archives delivery.Archives mtr *metrics.Set templates map[string]*template.Template @@ -128,7 +128,7 @@ func New( s.session = params.Session s.mw = params.Middleware s.notifier = params.Notifier - s.evictor = params.Evictor + s.archives = params.Archives s.mtr = metrics.Default() s.ssrf = params.SSRFGuard diff --git a/internal/handlers/handlers_test.go b/internal/handlers/handlers_test.go index f6a7273..96dd7d8 100644 --- a/internal/handlers/handlers_test.go +++ b/internal/handlers/handlers_test.go @@ -3,6 +3,7 @@ package handlers_test import ( "context" "errors" + "fmt" "html/template" "net/http" "net/http/httptest" @@ -51,23 +52,73 @@ func (n *recordingNotifier) Tasks() []delivery.Task { return out } -// 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 +// recordingArchives is a delivery.Archives that records what it +// was asked to do, so a test can prove that a deletion or rename +// path reached the delivery engine. After FailRenames, every +// rename fails with the given error. +type recordingArchives struct { + mu sync.Mutex + evicted []string + evictedTargets []string + renames []archiveRename + renameErr error } -func (r *recordingEvictor) EvictWebhook(webhookID string) { +// errInjectedRename is the failure a test hands FailRenames. +var errInjectedRename = errors.New("injected rename failure") + +// errNameTaken is what the delivery engine returns when a file +// already has an archive's new name, here archive-taken.db. +var errNameTaken = fmt.Errorf( + "%w: archive-taken.db", delivery.ErrArchiveNameTaken, +) + +// archiveRename is one recorded Rename call. +type archiveRename struct { + TargetID string + WebhookName string + TargetName string +} + +func (r *recordingArchives) EvictWebhook(webhookID string) { r.mu.Lock() defer r.mu.Unlock() r.evicted = append(r.evicted, webhookID) } +func (r *recordingArchives) EvictTarget(targetID string) { + r.mu.Lock() + defer r.mu.Unlock() + + r.evictedTargets = append(r.evictedTargets, targetID) +} + +func (r *recordingArchives) Rename( + targetID, webhookName, targetName string, +) error { + r.mu.Lock() + defer r.mu.Unlock() + + r.renames = append(r.renames, archiveRename{ + TargetID: targetID, + WebhookName: webhookName, + TargetName: targetName, + }) + + return r.renameErr +} + +// FailRenames makes every later rename fail with err. +func (r *recordingArchives) FailRenames(err error) { + r.mu.Lock() + defer r.mu.Unlock() + + r.renameErr = err +} + // Evicted returns a copy of the recorded webhook ids. -func (r *recordingEvictor) Evicted() []string { +func (r *recordingArchives) Evicted() []string { r.mu.Lock() defer r.mu.Unlock() @@ -77,6 +128,28 @@ func (r *recordingEvictor) Evicted() []string { return out } +// EvictedTargets returns a copy of the recorded target ids. +func (r *recordingArchives) EvictedTargets() []string { + r.mu.Lock() + defer r.mu.Unlock() + + out := make([]string, len(r.evictedTargets)) + copy(out, r.evictedTargets) + + return out +} + +// Renames returns a copy of the recorded renames. +func (r *recordingArchives) Renames() []archiveRename { + r.mu.Lock() + defer r.mu.Unlock() + + out := make([]archiveRename, len(r.renames)) + copy(out, r.renames) + + return out +} + func newTestApp( t *testing.T, targets ...any, @@ -103,10 +176,10 @@ func newTestApp( func(n *recordingNotifier) delivery.Notifier { return n }, - func() *recordingEvictor { - return &recordingEvictor{} + func() *recordingArchives { + return &recordingArchives{} }, - func(r *recordingEvictor) delivery.WebhookEvictor { + func(r *recordingArchives) delivery.Archives { return r }, middleware.New, diff --git a/internal/handlers/source_delete_test.go b/internal/handlers/source_delete_test.go index 58db5c5..5e623c2 100644 --- a/internal/handlers/source_delete_test.go +++ b/internal/handlers/source_delete_test.go @@ -15,6 +15,7 @@ import ( "gorm.io/gorm" "gorm.io/gorm/clause" "sneak.berlin/go/webhooker/internal/database" + "sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/session" ) @@ -147,18 +148,19 @@ func failDeleteOnTable( } // archivePathFor returns the archive database path the -// delivery engine would use for a webhook: beside the webhook's -// event database in the data directory. +// delivery engine would use for a database target: beside the +// webhook's event database in the data directory. func archivePathFor( t *testing.T, mgr *database.WebhookDBManager, - webhookID string, + wh *database.Webhook, + tgt *database.Target, ) string { t.Helper() return filepath.Join( - filepath.Dir(mgr.DBPath(webhookID)), - "archive-"+webhookID+".db", + filepath.Dir(mgr.DBPath(wh.ID)), + delivery.ArchiveFileName(wh.Name, tgt.Name, tgt.ID), ) } @@ -195,8 +197,8 @@ func postRequest( // 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. +// the webhook's archive writers, exercised through the real +// deletion handler rather than by calling the engine directly. func TestHandleSourceDelete_EvictsArchiveWriter(t *testing.T) { t.Parallel() @@ -204,7 +206,7 @@ func TestHandleSourceDelete_EvictsArchiveWriter(t *testing.T) { h *handlers.Handlers sess *session.Session db *database.Database - ev *recordingEvictor + ev *recordingArchives ) app := newTestApp(t, &h, &sess, &db, &ev) @@ -254,9 +256,10 @@ func TestHandleSourceDelete_KeepsArchiveFile(t *testing.T) { t.Cleanup(app.RequireStop) wh := seedWebhook(t, db) + tgt := seedTarget(t, db, wh.ID, database.TargetTypeDatabase) // Place an archive file where the delivery engine would. - archivePath := archivePathFor(t, mgr, wh.ID) + archivePath := archivePathFor(t, mgr, wh, tgt) require.NoError( t, writeArchivePlaceholder(archivePath), @@ -435,68 +438,17 @@ func TestHandleSourceDelete_RemovesConfigAndEventDatabase( ) } -// TestHandleTargetDelete_EvictsWhenLastDatabaseTargetGone -// proves that removing the last database target releases the -// archive writer. -func TestHandleTargetDelete_EvictsWhenLastDatabaseTargetGone( - t *testing.T, -) { +// TestHandleTargetDelete_EvictsThatTarget proves that deleting a +// database target releases that target's archive writer and no +// other: the webhook's other database target keeps its own. +func TestHandleTargetDelete_EvictsThatTarget(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( - "/hook/"+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_KeepsWriterWhenDatabaseTargetRemains -// proves that deleting one of several database targets leaves -// the still-needed archive writer alone: the surviving target -// keeps archiving to the same file, so the writer must stay. -func TestHandleTargetDelete_KeepsWriterWhenDatabaseTargetRemains( - t *testing.T, -) { - t.Parallel() - - var ( - h *handlers.Handlers - sess *session.Session - db *database.Database - ev *recordingEvictor + ev *recordingArchives ) app := newTestApp(t, &h, &sess, &db, &ev) @@ -527,17 +479,17 @@ func TestHandleTargetDelete_KeepsWriterWhenDatabaseTargetRemains( h.HandleTargetDelete().ServeHTTP(w, req) require.Equal(t, http.StatusSeeOther, w.Code) - assert.Empty( - t, ev.Evicted(), - "a second database target still needs the writer", + assert.Equal( + t, []string{doomed.ID}, ev.EvictedTargets(), + "deleting a database target should evict its writer", ) + assert.Empty(t, ev.Evicted(), "the webhook is not deleted") } -// TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted proves -// that deleting a target of an unrelated type leaves a -// still-needed archive writer alone: the webhook's database -// target is untouched, so its writer must stay. -func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted( +// TestHandleTargetDelete_IgnoresAnotherWebhooksTarget proves that +// a target id from the URL that is not a target of the webhook +// deletes nothing and so evicts nothing. +func TestHandleTargetDelete_IgnoresAnotherWebhooksTarget( t *testing.T, ) { t.Parallel() @@ -546,7 +498,7 @@ func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted( h *handlers.Handlers sess *session.Session db *database.Database - ev *recordingEvictor + ev *recordingArchives ) app := newTestApp(t, &h, &sess, &db, &ev) @@ -555,19 +507,20 @@ func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted( t.Cleanup(app.RequireStop) wh := seedWebhook(t, db) - seedTarget(t, db, wh.ID, database.TargetTypeDatabase) - other := seedTarget(t, db, wh.ID, database.TargetTypeLog) + elsewhere := seedTarget( + t, db, seedWebhook(t, db).ID, database.TargetTypeDatabase, + ) cookies := authenticatedCookies( t, sess, deleteTestUserID, deleteTestUsername, ) req := postRequest( - "/hook/"+wh.ID+"/targets/"+other.ID+"/delete", + "/hook/"+wh.ID+"/targets/"+elsewhere.ID+"/delete", cookies, map[string]string{ paramSourceID: wh.ID, - paramTargetID: other.ID, + paramTargetID: elsewhere.ID, }, ) w := httptest.NewRecorder() @@ -576,7 +529,7 @@ func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted( require.Equal(t, http.StatusSeeOther, w.Code) assert.Empty( - t, ev.Evicted(), - "a surviving database target must keep its writer", + t, ev.EvictedTargets(), + "another webhook's target must not be evicted", ) } diff --git a/internal/handlers/source_management.go b/internal/handlers/source_management.go index 330e787..6a2b881 100644 --- a/internal/handlers/source_management.go +++ b/internal/handlers/source_management.go @@ -549,6 +549,7 @@ func (h *Handlers) applyWebhookEdit( return } + oldName := webhook.Name webhook.Name = name webhook.Description = r.PostFormValue("description") @@ -571,8 +572,40 @@ func (h *Handlers) applyWebhookEdit( webhook.RetentionDays = retentionDays - err := h.db.DB().Save(webhook).Error + // A new name renames the archive files before it is saved (see + // delivery.Engine.Rename). If either step fails, they go back to + // the name that is still stored. + err := h.renameWebhookArchives(webhook.ID, oldName, webhook.Name) + if err == nil { + err = h.db.DB().Save(webhook).Error + } + if err != nil { + restoreErr := h.renameWebhookArchives( + webhook.ID, webhook.Name, oldName, + ) + if restoreErr != nil { + h.log.Error( + "failed to rename archives back", + "webhook_id", webhook.ID, + "error", restoreErr, + ) + } + + if errors.Is(err, delivery.ErrArchiveNameTaken) { + data := map[string]any{ + tmplKeyWebhook: webhook, + tmplKeyError: "Not saved: " + err.Error() + + ". Move that file out of the data directory, " + + "then save again.", + } + + w.WriteHeader(http.StatusConflict) + h.renderTemplate(w, r, "source_edit.html", data) + + return + } + h.serverError(w, r, "failed to update webhook", err) return @@ -707,11 +740,11 @@ func (h *Handlers) commitWebhookDeletion( return tx.Commit().Error } -// evictArchiveWriter asks the delivery engine to drop its -// cached archive writer for a webhook, closing the archive file -// handle. +// evictArchiveWriter asks the delivery engine to drop the cached +// archive writers of a webhook's database targets, closing their +// archive file handles. // -// The archive database file is NOT deleted. Unlike the event +// The archive database files are 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 @@ -719,50 +752,58 @@ func (h *Handlers) commitWebhookDeletion( // 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 { + if h.archives == nil { return } - h.evictor.EvictWebhook(webhookID) + h.archives.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 +// evictTargetArchiveWriter is evictArchiveWriter for one deleted +// target, and leaves its archive file on disk for the same reason. +// A target that is not a database target has no writer, and +// evicting it does nothing. +func (h *Handlers) evictTargetArchiveWriter(targetID string) { + if h.archives == nil { + return + } + + h.archives.EvictTarget(targetID) +} + +// renameWebhookArchives renames the archive file of every database +// target of a webhook from the webhook name oldName to newName, +// keeping each target's own name. It does nothing when the name is +// unchanged, and stops at the first failure. +func (h *Handlers) renameWebhookArchives( + webhookID, oldName, newName string, +) error { + if h.archives == nil || oldName == newName { + return nil + } + + var targets []database.Target err := h.db.DB(). - Model(&database.Target{}). Where( "webhook_id = ? AND type = ?", webhookID, database.TargetTypeDatabase, ). - Count(&remaining).Error + Find(&targets).Error if err != nil { - h.log.Error( - "failed to count remaining database targets", - "webhook_id", webhookID, - "error", err, + return err + } + + for i := range targets { + err = h.archives.Rename( + targets[i].ID, newName, targets[i].Name, ) - - return + if err != nil { + return err + } } - if remaining > 0 { - return - } - - h.evictArchiveWriter(webhookID) + return nil } // ownedWebhook resolves the request's sourceID parameter to a @@ -1646,27 +1687,26 @@ func (h *Handlers) HandleEntrypointDelete() http.HandlerFunc { ) } -// 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. +// HandleTargetDelete handles deleting a target. A deleted +// database target's archive 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, + h.evictTargetArchiveWriter, ) } // deleteChildResource returns a handler that deletes a child // 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. +// optional afterDelete hook runs with the child's id once the +// delete has removed it, before the redirect. func (h *Handlers) deleteChildResource( idParam string, model any, errMsg string, - afterDelete func(webhookID string), + afterDelete func(childID string), ) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { userID, ok := h.getUserID(r) @@ -1702,8 +1742,10 @@ func (h *Handlers) deleteChildResource( return } - if afterDelete != nil { - afterDelete(webhook.ID) + // Only for a row this webhook really had: the id came from + // the URL and may name another webhook's child. + if afterDelete != nil && result.RowsAffected > 0 { + afterDelete(childID) } http.Redirect( diff --git a/internal/handlers/source_management_test.go b/internal/handlers/source_management_test.go index 9967a85..9071f8f 100644 --- a/internal/handlers/source_management_test.go +++ b/internal/handlers/source_management_test.go @@ -187,6 +187,7 @@ func storedRetentionDays( type sourceTestEnv struct { handlers *handlers.Handlers db *database.Database + archives *recordingArchives cookies []*http.Cookie } @@ -199,7 +200,9 @@ func setupSourceTest(t *testing.T) *sourceTestEnv { var db *database.Database - app := newTestApp(t, &h, &sess, &db) + var archives *recordingArchives + + app := newTestApp(t, &h, &sess, &db, &archives) app.RequireStart() t.Cleanup(app.RequireStop) @@ -207,6 +210,7 @@ func setupSourceTest(t *testing.T) *sourceTestEnv { return &sourceTestEnv{ handlers: h, db: db, + archives: archives, cookies: authenticatedCookies( t, sess, sourceTestUserID, "sourceuser", ), @@ -498,6 +502,106 @@ func TestHandleSourceEditSubmit_EmptyRetentionLeavesValueUnchanged( assert.Equal(t, 7, storedRetentionDays(t, env.db, wh.ID)) } +// renamedWebhookName is the name the rename tests give a webhook. +const renamedWebhookName = "Renamed" + +// TestHandleSourceEditSubmit_RenamesArchives proves that a save +// that keeps the webhook's name renames nothing, and that renaming a +// webhook renames the archive of each of its database targets and +// asks nothing of its other targets. +func TestHandleSourceEditSubmit_RenamesArchives(t *testing.T) { + t.Parallel() + + env := setupSourceTest(t) + wh := seedWebhookWithRetention(t, env.db, 7) + first := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase) + second := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase) + seedTarget(t, env.db, wh.ID, database.TargetTypeLog) + + w := submitEdit(t, env, wh, "") + require.Equal(t, http.StatusSeeOther, w.Code) + assert.Empty(t, env.archives.Renames()) + + wh.Name = renamedWebhookName + + w = submitEdit(t, env, wh, "") + require.Equal(t, http.StatusSeeOther, w.Code) + + assert.ElementsMatch( + t, + []archiveRename{ + {first.ID, renamedWebhookName, first.Name}, + {second.ID, renamedWebhookName, second.Name}, + }, + env.archives.Renames(), + ) +} + +// TestHandleSourceEditSubmit_FailedRenameKeepsTheName proves that a +// webhook whose archive cannot be renamed keeps its stored name, so +// the name on disk and the name in the UI do not part, and that the +// handler puts back what it may already have moved. +func TestHandleSourceEditSubmit_FailedRenameKeepsTheName( + t *testing.T, +) { + t.Parallel() + + env := setupSourceTest(t) + wh := seedWebhookWithRetention(t, env.db, 7) + tgt := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase) + + env.archives.FailRenames(errInjectedRename) + + oldName := wh.Name + wh.Name = renamedWebhookName + + w := submitEdit(t, env, wh, "") + require.Equal(t, http.StatusInternalServerError, w.Code) + + var stored database.Webhook + + require.NoError( + t, env.db.DB().First(&stored, "id = ?", wh.ID).Error, + ) + assert.Equal(t, oldName, stored.Name) + + assert.Equal( + t, + []archiveRename{ + {tgt.ID, renamedWebhookName, tgt.Name}, + {tgt.ID, oldName, tgt.Name}, + }, + env.archives.Renames(), + ) +} + +// TestHandleSourceEditSubmit_ArchiveNameTaken proves that when a file +// already has an archive's new name, the edit is refused with an +// error naming that file, and the webhook keeps its stored name. +func TestHandleSourceEditSubmit_ArchiveNameTaken(t *testing.T) { + t.Parallel() + + env := setupSourceTest(t) + wh := seedWebhookWithRetention(t, env.db, 7) + seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase) + + env.archives.FailRenames(errNameTaken) + + oldName := wh.Name + wh.Name = renamedWebhookName + + w := submitEdit(t, env, wh, "") + require.Equal(t, http.StatusConflict, w.Code) + assert.Contains(t, w.Body.String(), "archive-taken.db") + + var stored database.Webhook + + require.NoError( + t, env.db.DB().First(&stored, "id = ?", wh.ID).Error, + ) + assert.Equal(t, oldName, stored.Name) +} + // TestSourceEditForm_ForeverWebhookRoundTrips walks the exact path that // the removed max="365" cap used to break: render the edit form for a // retain-forever webhook, confirm the pre-filled sentinel is not capped diff --git a/internal/handlers/target_edit.go b/internal/handlers/target_edit.go index 810b801..051ed9c 100644 --- a/internal/handlers/target_edit.go +++ b/internal/handlers/target_edit.go @@ -1,6 +1,7 @@ package handlers import ( + "errors" "net/http" "github.com/go-chi/chi" @@ -150,11 +151,42 @@ func (h *Handlers) applyTargetEdit( target.MaxRetries = retries } + oldName := target.Name target.Name = name target.Config = configJSON - err = h.db.DB().Save(target).Error + // A new name renames the archive file before it is saved (see + // delivery.Engine.Rename). If either step fails, it goes back to + // the name that is still stored. + err = h.renameTargetArchive(target, webhook.Name, oldName, name) + if err == nil { + err = h.db.DB().Save(target).Error + } + if err != nil { + restoreErr := h.renameTargetArchive( + target, webhook.Name, name, oldName, + ) + if restoreErr != nil { + h.log.Error( + "failed to rename archive back", + "target_id", target.ID, + "error", restoreErr, + ) + } + + if errors.Is(err, delivery.ErrArchiveNameTaken) { + http.Error( + w, + "Not saved: "+err.Error()+ + ". Move that file out of the data directory, "+ + "then save again.", + http.StatusConflict, + ) + + return + } + h.serverError(w, r, "failed to update target", err) return @@ -165,6 +197,21 @@ func (h *Handlers) applyTargetEdit( ) } +// renameTargetArchive renames a database target's archive file from +// the target name oldName to newName. It does nothing when the name +// is unchanged; other target types have no archive. +func (h *Handlers) renameTargetArchive( + target *database.Target, + webhookName, oldName, newName string, +) error { + if h.archives == nil || oldName == newName || + target.Type != database.TargetTypeDatabase { + return nil + } + + return h.archives.Rename(target.ID, webhookName, newName) +} + // renderTargetEdit renders the target edit page with an optional // error message. func (h *Handlers) renderTargetEdit( diff --git a/internal/handlers/target_edit_test.go b/internal/handlers/target_edit_test.go index 40eed2c..14f8803 100644 --- a/internal/handlers/target_edit_test.go +++ b/internal/handlers/target_edit_test.go @@ -635,3 +635,68 @@ func assertWebhookOfAnotherUser404s( assert.Equal(t, http.StatusNotFound, w.Code) } + +// renamedTargetName is the name the rename tests give a target. +const renamedTargetName = "Long Term" + +// TestHandleTargetEditSubmit_RenamesArchive proves that renaming a +// database target renames its archive, that a save that keeps the +// name renames nothing, that a target of another type has no archive +// to rename, and that a target whose archive cannot be renamed keeps +// its stored name. When a file already has the archive's new name, +// the edit is refused with an error naming that file. +func TestHandleTargetEditSubmit_RenamesArchive(t *testing.T) { + t.Parallel() + + env := setupSourceTest(t) + wh := seedWebhookWithRetention(t, env.db, 7) + archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase) + rename := url.Values{"name": {renamedTargetName}} + + w := submitTargetEdit(env, wh.ID, archive.ID, rename) + require.Equal(t, http.StatusSeeOther, w.Code, w.Body.String()) + assert.Equal( + t, + []archiveRename{{archive.ID, wh.Name, renamedTargetName}}, + env.archives.Renames(), + ) + + w = submitTargetEdit(env, wh.ID, archive.ID, rename) + require.Equal(t, http.StatusSeeOther, w.Code, w.Body.String()) + assert.Len( + t, env.archives.Renames(), 1, + "a save that keeps the name renames nothing", + ) + + httpWebhook, httpTarget := seedHTTPTarget(t, env, "", "") + + w = submitTargetEdit( + env, httpWebhook.ID, httpTarget.ID, + editForm(editOriginalURL, "", ""), + ) + require.Equal(t, http.StatusSeeOther, w.Code, w.Body.String()) + assert.Len( + t, env.archives.Renames(), 1, + "an HTTP target has no archive to rename", + ) + + again := url.Values{"name": {"Again"}} + + env.archives.FailRenames(errInjectedRename) + + w = submitTargetEdit(env, wh.ID, archive.ID, again) + require.Equal(t, http.StatusInternalServerError, w.Code) + assert.Equal( + t, renamedTargetName, storedTarget(t, env, archive.ID).Name, + "a target whose archive was not renamed keeps its name", + ) + + env.archives.FailRenames(errNameTaken) + + w = submitTargetEdit(env, wh.ID, archive.ID, again) + require.Equal(t, http.StatusConflict, w.Code) + assert.Contains(t, w.Body.String(), "archive-taken.db") + assert.Equal( + t, renamedTargetName, storedTarget(t, env, archive.ID).Name, + ) +} diff --git a/internal/resetpw/resetpw_test.go b/internal/resetpw/resetpw_test.go index 495dfc4..e348a16 100644 --- a/internal/resetpw/resetpw_test.go +++ b/internal/resetpw/resetpw_test.go @@ -131,9 +131,15 @@ type noopNotifier struct{} func (n *noopNotifier) Notify([]delivery.Task) {} -type noopEvictor struct{} +type noopArchives struct{} -func (n *noopEvictor) EvictWebhook(string) {} +func (n *noopArchives) EvictWebhook(string) {} + +func (n *noopArchives) EvictTarget(string) {} + +func (n *noopArchives) Rename(_, _, _ string) error { + return nil +} // newServerApp starts the real login path against dir: the handlers, // the middleware that bounds password verification, the session store @@ -162,7 +168,7 @@ func newServerApp( healthcheck.New, session.New, func() delivery.Notifier { return &noopNotifier{} }, - func() delivery.WebhookEvictor { return &noopEvictor{} }, + func() delivery.Archives { return &noopArchives{} }, middleware.New, delivery.NewGuard, handlers.New, diff --git a/internal/server/routes_test.go b/internal/server/routes_test.go index cb344e7..f5355ab 100644 --- a/internal/server/routes_test.go +++ b/internal/server/routes_test.go @@ -46,12 +46,18 @@ type noopNotifier struct{} func (n *noopNotifier) Notify([]delivery.Task) {} -// noopEvictor satisfies handlers.New's delivery.WebhookEvictor -// dependency. No test here checks what gets evicted, so it records -// nothing. -type noopEvictor struct{} +// noopArchives satisfies handlers.New's delivery.Archives +// dependency. No test here checks what gets evicted or renamed, so +// it records nothing. +type noopArchives struct{} -func (e *noopEvictor) EvictWebhook(string) {} +func (e *noopArchives) EvictWebhook(string) {} + +func (e *noopArchives) EvictTarget(string) {} + +func (e *noopArchives) Rename(_, _, _ string) error { + return nil +} // testEnv is the real router from routes.go plus the collaborators // tests need to seed users and forge sessions. @@ -112,7 +118,7 @@ func newTestEnvWithConfig( healthcheck.New, session.New, func() delivery.Notifier { return &noopNotifier{} }, - func() delivery.WebhookEvictor { return &noopEvictor{} }, + func() delivery.Archives { return &noopArchives{} }, middleware.New, delivery.NewGuard, handlers.New,