1 Commits
Author SHA1 Message Date
clawbot 6257c6ec23 Name each database target's archive for its webhook and target (closes #376)
check / check (push) Waiting to run
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
2026-10-02 07:54:37 +00:00
21 changed files with 1545 additions and 663 deletions
+73 -45
View File
@@ -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 The `/var/lib/webhooker` volume holds all SQLite databases: the main
application database (`webhooker.db`), the per-webhook event databases application database (`webhooker.db`), the per-webhook event databases
(`events-{uuid}.db`), and any archive databases written by `database` (`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. preserve data across container restarts.
**The container sets its data directory's owner and mode itself **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. encryption key), users, API keys, webhooks, entrypoints, targets.
- `events-{webhook_uuid}.db` — **one per webhook**. Events, deliveries, - `events-{webhook_uuid}.db` — **one per webhook**. Events, deliveries,
delivery results. delivery results.
- `archive-{webhook_uuid}.db` — **one per webhook that has a `database` - `archive-{webhook_name}-{target_name}-{target_uuid}.db` — **one per
target**. Archived events. Keyed on the webhook UUID, not the target `database` target**. Archived events. The two names are made safe for
UUID: a webhook with several `database` targets still has exactly one a file name, and the file is renamed when the webhook or the target is
archive file. (see [Database Architecture](#database-architecture)).
`{webhook_uuid}` is the webhook's UUID primary key in its canonical `{webhook_uuid}` and `{target_uuid}` are UUID primary keys in their
36-character hyphenated form, so a real filename looks like canonical 36-character hyphenated form, so a real filename looks like
`events-3f2a1c9e-....db`. The only other file is `webhooker.lock`, the `events-3f2a1c9e-....db`. The only other file is `webhooker.lock`, the
always-empty [single-instance lock](#single-instance-lock); it holds no 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 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 databases are the one exception the service is built for: the
archive writer closes and reopens its handle around writes (debounced archive writer closes and reopens its handle around writes (debounced
to at most one reopen per second), so an operator can move to at most one reopen per second), so an operator can move an
`archive-{uuid}.db` away for offline retention while the service runs, `archive-….db` away for offline retention while the service runs,
and it is recreated on the next write. See and it is recreated on the next write. See
[Database Architecture](#database-architecture). That is a [Database Architecture](#database-architecture). That is a
move-the-file-away workflow, not a substitute for the backup procedures 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), 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 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 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. with any `-wal`/`-shm` beside it, or wait until there are none.
### Restore ### 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 Treat a backup with the same care as the credentials inside it. Encrypt
backups at rest and restrict who can read them. 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 body and headers** of every event as received, including whatever the
sending service put in them — tokens, signatures, personal data. sending service put in them — tokens, signatures, personal data.
- Event databases written before - 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` is built on the same HTTP core as `http` and honours `max_retries`
identically, circuit breaker included. See the Slack target section identically, circuit breaker included. See the Slack target section
under "Per-Webhook Event Databases" for the message format. under "Per-Webhook Event Databases" for the message format.
- **`database`** — Archive the full event as a row into a separate - **`database`** — Archive the full event as a row into the target's
per-webhook archive database (`archive-{webhookID}.db`) for long-term own archive database
(`archive-{webhook_name}-{target_name}-{target_uuid}.db`) for long-term
retention, with an optional creation-validated expiry (default: keep retention, with an optional creation-validated expiry (default: keep
forever). No external delivery and no retries; an archive write forever). No external delivery and no retries; an archive write
failure fails the delivery. See the database target section under failure fails the delivery. See the database target section under
@@ -1904,9 +1906,37 @@ The **database target type** builds on this architecture to provide
long-term archiving, separate from the per-webhook event database (which long-term archiving, separate from the per-webhook event database (which
may prune events under its own retention). Delivering to a database may prune events under its own retention). Delivering to a database
target writes the full event — body, headers, method, content type, and target writes the full event — body, headers, method, content type, and
webhook/entrypoint/event identifiers — as a row into a dedicated archive webhook/entrypoint/event identifiers — as a row into the target's own
database, `archive-{webhookID}.db`, stored under the data directory archive database, `archive-{webhook_name}-{target_name}-{target_uuid}.db`,
beside the event database. After each write the archive handle is closed 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, stop the
service before moving either file. If no file has the name shown, move
the file under the new name back to it. If a second archive already has
the name shown, move the file under the new name out of the data
directory instead and keep it as you would any archive moved away. Then
start the service again.
After each write the archive handle is closed
and reopened, debounced to at most once per second, so an operator can and reopened, debounced to at most once per second, so an operator can
move the archive file away for offline archiving without stopping the move the archive file away for offline archiving without stopping the
service; a moved or removed archive file is recreated automatically on service; a moved or removed archive file is recreated automatically on
@@ -1917,35 +1947,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 archive write failure is never silent success: the delivery records a
failed attempt with the error and is marked failed. failed attempt with the error and is marked failed.
Because reopens only happen on writes, an archive belonging to a webhook Because reopens only happen on writes, an archive whose target has
that has stopped receiving events would never be pruned. A background stopped receiving events would never be pruned. A background **archive
**archive sweeper** closes that gap: on the same interval as the event sweeper** closes that gap: on the same interval as the event retention
retention reaper (`RETENTION_SWEEP_INTERVAL`) it prunes every archive reaper (`RETENTION_SWEEP_INTERVAL`) it prunes every archive whose
whose database target declares a positive expiry, whether or not the database target declares a positive expiry, whether or not the target
webhook is still receiving traffic. The sweep never creates an archive — is still receiving traffic. The sweep never creates an archive — a
a webhook whose archive file does not yet exist is skipped, not target whose archive file does not yet exist is skipped, not initialised
initialised — it takes the same per-webhook lock the write path uses, so — it takes the same per-target lock the write path uses, so it can never
it can never interleave with a write, and it leaves the archive closed interleave with a write, and it leaves the archive closed afterwards so
afterwards so the move-the-file-away workflow keeps working. Archives the move-the-file-away workflow keeps working. Archives with no expiry,
with no expiry, or the expiry `never`, are not touched by the sweep at or the expiry `never`, are not touched by the sweep at all.
all.
Note that a webhook has one archive file but may carry more than one Because each `database` target has its own archive file, a target's
`database` target, each with its own `expiry`. The shortest expiry `expiry` governs only its own archive. Two `database` targets on one
configured on any of them therefore governs the whole archive, and the webhook with different expiries keep two archives, each pruned on its
sweep applies it whether or not the webhook is still receiving events. own schedule.
Configure a single `database` target per webhook unless you intend that.
Deleting a webhook releases its archive: the delivery engine's cached Deleting a webhook releases its archives: the delivery engine's cached
archive writer is dropped and its file handle closed, so nothing lingers archive writers are dropped and their file handles closed, so nothing
after the webhook is gone. The archive **file itself is deliberately lingers after the webhook is gone. The archive **files themselves are
left on disk**. Unlike the event database — per-webhook working storage deliberately left on disk**. Unlike the event database — per-webhook
that is hard-deleted with the webhook — an archive is long-term storage working storage that is hard-deleted with the webhook — an archive is
an operator may still want to keep or move away for offline retention, long-term storage an operator may still want to keep or move away for
and destroying it as a side effect of deleting a webhook would be offline retention, and destroying it as a side effect of deleting a
unrecoverable. Removing `archive-{webhookID}.db` is the operator's call. webhook would be unrecoverable. Removing an `archive-….db` is the
Deleting a webhook's last `database` target releases the writer the same operator's call. Deleting a `database` target releases its writer the
way, and for the same reason leaves the file alone. same way, and for the same reason leaves its file alone.
The **Slack target type** sends webhook events as formatted messages to The **Slack target type** sends webhook events as formatted messages to
any Slack-compatible incoming webhook URL (works with Slack, Mattermost, any Slack-compatible incoming webhook URL (works with Slack, Mattermost,
@@ -2957,8 +2985,8 @@ Components are wired via Uber fx in this order:
13. `delivery.New` — Event-driven delivery engine 13. `delivery.New` — Event-driven delivery engine
14. `delivery.NewArchiveSweeper` — Periodic pruning of idle archives 14. `delivery.NewArchiveSweeper` — Periodic pruning of idle archives
15. `delivery.Engine` → `delivery.Notifier` — interface bridge 15. `delivery.Engine` → `delivery.Notifier` — interface bridge
16. `delivery.Engine` → `delivery.WebhookEvictor` — interface bridge so 16. `delivery.Engine` → `delivery.Archives` — interface bridge so
deleting a webhook releases its archive writer deleting or renaming a webhook or target reaches its archive files
17. `server.New` — HTTP server and router 17. `server.New` — HTTP server and router
The server starts via `fx.Invoke(func(*server.Server, *delivery.Engine, The server starts via `fx.Invoke(func(*server.Server, *delivery.Engine,
+4 -5
View File
@@ -192,11 +192,10 @@ func newApp() *fx.App {
// Wire *delivery.Engine as delivery.Notifier so the // Wire *delivery.Engine as delivery.Notifier so the
// webhook handler can notify the engine of new deliveries. // webhook handler can notify the engine of new deliveries.
func(e *delivery.Engine) delivery.Notifier { return e }, func(e *delivery.Engine) delivery.Notifier { return e },
// Wire *delivery.Engine as delivery.WebhookEvictor so // Wire *delivery.Engine as delivery.Archives so deleting
// deleting a webhook releases its archive writer. // or renaming a webhook or target reaches its archive
func(e *delivery.Engine) delivery.WebhookEvictor { // files.
return e func(e *delivery.Engine) delivery.Archives { return e },
},
server.New, server.New,
), ),
fx.Invoke( fx.Invoke(
+14 -12
View File
@@ -8,6 +8,7 @@ import (
"time" "time"
"go.uber.org/fx" "go.uber.org/fx"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/config" "sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/lifecycle" "sneak.berlin/go/webhooker/internal/lifecycle"
@@ -25,14 +26,14 @@ type ArchiveSweeperParams struct {
Logger *logger.Logger Logger *logger.Logger
} }
// ArchiveSweeper periodically prunes expired rows from // ArchiveSweeper periodically prunes expired rows from the
// per-webhook archive databases whose database target carries a // archive databases of database targets that carry a positive
// positive expiry. // expiry.
// //
// Without it, pruning happens only when an archive is // Without it, pruning happens only when an archive is
// (re)opened, and archives are only ever reopened by writes: an // (re)opened, and archives are only ever reopened by writes: an
// archive belonging to a webhook that has stopped receiving // archive whose target has stopped receiving events would keep
// events would keep its expired rows forever. The sweep closes // its expired rows forever. The sweep closes
// that gap without changing anything for archives whose expiry // that gap without changing anything for archives whose expiry
// is unset or "never". // 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 // soft-deleted along with it, so GORM's default scope already
// excludes them. // 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 // matching how the write path already treats a prune error as
// non-fatal. // non-fatal.
func (s *ArchiveSweeper) sweep(ctx context.Context) { func (s *ArchiveSweeper) sweep(ctx context.Context) {
@@ -210,19 +211,20 @@ func (s *ArchiveSweeper) sweepTarget(target *database.Target) {
return return
} }
err = s.eng.dbTarget.sweepWebhook(target.WebhookID, expiry) err = s.eng.dbTarget.sweepArchive(target.ID, expiry)
if err == nil { if err == nil {
return return
} }
// A writer evicted underneath the sweep means the operator // A writer evicted, or a target row gone, underneath the sweep
// deleted the webhook (or its last database target) while the // means the operator deleted the target or its webhook while
// sweep was walking the target list. That is an ordinary // the sweep was walking the target list. That is an ordinary
// interleaving, not a failure, so it must not produce an // interleaving, not a failure, so it must not produce an
// error line. // error line.
if errors.Is(err, errArchiveWriterEvicted) { if errors.Is(err, errArchiveWriterEvicted) ||
errors.Is(err, gorm.ErrRecordNotFound) {
s.log.Debug( s.log.Debug(
"archive sweep: writer evicted mid-sweep", "archive sweep: target deleted mid-sweep",
"webhook_id", target.WebhookID, "webhook_id", target.WebhookID,
"target_id", target.ID, "target_id", target.ID,
) )
+143 -133
View File
@@ -34,18 +34,23 @@ const (
sweepConcurrentWrites = 20 sweepConcurrentWrites = 20
) )
// sweeperEnv bundles the pieces an archive sweep test drives: // archiveTestWebhookName is the name of every webhook
// a main configuration database holding webhooks and targets, a // seedDatabaseTarget creates. It is not safe in a file name as it
// delivery engine owning the archive writer registry, and the // stands, so every archive test goes through archiveNamePart.
// data directory the archive files live in. const archiveTestWebhookName = "Sweep Test!"
type sweeperEnv struct {
// 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 sweeper *delivery.ArchiveSweeper
eng *delivery.Engine eng *delivery.Engine
mainDB *database.Database mainDB *database.Database
dataDir string dataDir string
} }
func setupSweeperTest(t *testing.T) *sweeperEnv { func setupArchiveTest(t *testing.T) *archiveEnv {
t.Helper() t.Helper()
dataDir := t.TempDir() dataDir := t.TempDir()
@@ -78,7 +83,7 @@ func setupSweeperTest(t *testing.T) *sweeperEnv {
1, 1,
) )
return &sweeperEnv{ return &archiveEnv{
sweeper: delivery.NewTestArchiveSweeper( sweeper: delivery.NewTestArchiveSweeper(
mainDB, eng, log, mainDB, eng, log,
), ),
@@ -88,25 +93,27 @@ func setupSweeperTest(t *testing.T) *sweeperEnv {
} }
} }
// archivePath returns where the engine keeps a webhook's // archivePath returns where the engine keeps a database target's
// archive file. // archive file, for the names seedDatabaseTarget gave it.
func (env *sweeperEnv) archivePath(webhookID string) string { func (env *archiveEnv) archivePath(tgt *database.Target) string {
return filepath.Join( 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 // seedDatabaseTarget creates a webhook with one database target
// carrying the given target config JSON, and returns the // carrying the given target config JSON, and returns the target.
// webhook id. func (env *archiveEnv) seedDatabaseTarget(
func (env *sweeperEnv) seedDatabaseTarget(
t *testing.T, configJSON string, t *testing.T, configJSON string,
) string { ) *database.Target {
t.Helper() t.Helper()
wh := &database.Webhook{ wh := &database.Webhook{
UserID: uuid.New().String(), UserID: uuid.New().String(),
Name: "sweep-test", Name: archiveTestWebhookName,
} }
require.NoError( require.NoError(
t, t,
@@ -115,9 +122,19 @@ func (env *sweeperEnv) seedDatabaseTarget(
Create(wh).Error, 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{ tgt := &database.Target{
WebhookID: wh.ID, WebhookID: webhookID,
Name: "archive", Name: "Archive",
Type: database.TargetTypeDatabase, Type: database.TargetTypeDatabase,
Active: true, Active: true,
Config: configJSON, Config: configJSON,
@@ -129,19 +146,19 @@ func (env *sweeperEnv) seedDatabaseTarget(
Create(tgt).Error, 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 // inserts one row per supplied archived-at timestamp, returning
// the archive path. The handle is closed before returning, so // the archive path. The handle is closed before returning, so
// the archive is idle exactly as it would be with no traffic. // the archive is idle exactly as it would be with no traffic.
func (env *sweeperEnv) seedArchiveRows( func (env *archiveEnv) seedArchiveRows(
t *testing.T, webhookID string, archivedAt ...time.Time, t *testing.T, tgt *database.Target, archivedAt ...time.Time,
) string { ) string {
t.Helper() t.Helper()
path := env.archivePath(webhookID) path := env.archivePath(tgt)
sqlDB, err := sql.Open( sqlDB, err := sql.Open(
"sqlite", fmt.Sprintf("file:%s?mode=rwc", path), "sqlite", fmt.Sprintf("file:%s?mode=rwc", path),
@@ -160,7 +177,7 @@ func (env *sweeperEnv) seedArchiveRows(
for i, at := range archivedAt { for i, at := range archivedAt {
row := delivery.ExportArchivedEvent{ row := delivery.ExportArchivedEvent{
EventID: fmt.Sprintf("ev-%d", i), EventID: fmt.Sprintf("ev-%d", i),
WebhookID: webhookID, WebhookID: tgt.WebhookID,
Method: http.MethodPost, Method: http.MethodPost,
Body: `{"seeded":true}`, Body: `{"seeded":true}`,
ArchivedAt: at, ArchivedAt: at,
@@ -243,13 +260,13 @@ func TestArchiveSweeper_LoopOutlivesStartHookContext(
) { ) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
now := time.Now() now := time.Now()
path := env.seedArchiveRows( path := env.seedArchiveRows(
t, webhookID, t, tgt,
now.Add(-48*time.Hour), now.Add(-48*time.Hour),
now.Add(-time.Minute), now.Add(-time.Minute),
) )
@@ -287,60 +304,60 @@ func TestArchiveSweeper_LoopOutlivesStartHookContext(
} }
// TestArchiveSweep_DoesNotResurrectEvictedWriter covers the // TestArchiveSweep_DoesNotResurrectEvictedWriter covers the
// interleaving where a sweep tick has already listed a webhook's // interleaving where a sweep tick has already listed a target
// target when the webhook is deleted and its writer evicted. The // when the target is deleted and its writer evicted. The sweep
// sweep must not put a writer back into the registry: nothing // must not put a writer back into the registry: nothing would
// would ever evict it again, which is precisely the leak this // ever evict it again, which is precisely the leak this change
// change exists to close. // exists to close.
func TestArchiveSweep_DoesNotResurrectEvictedWriter( func TestArchiveSweep_DoesNotResurrectEvictedWriter(
t *testing.T, t *testing.T,
) { ) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows( 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 // Prime the registry the way a delivery would, then evict as
// the deletion path does. The target row is deliberately left // 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. // the deletion committed.
_, err := env.eng.ExportEnsureArchiveWriter(webhookID) _, err := env.eng.ExportEnsureArchiveWriter(tgt.ID)
require.NoError(t, err) require.NoError(t, err)
env.eng.EvictWebhook(webhookID) env.eng.EvictTarget(tgt.ID)
require.False(t, env.eng.ExportHasArchiveWriter(webhookID)) require.False(t, env.eng.ExportHasArchiveWriter(tgt.ID))
env.sweeper.ExportSweep(context.Background()) env.sweeper.ExportSweep(context.Background())
assert.False( assert.False(
t, env.eng.ExportHasArchiveWriter(webhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"a sweep must never re-register a writer for a webhook "+ "a sweep must never re-register a writer for a target "+
"whose registry entry has already been released", "whose registry entry has already been released",
) )
} }
// TestArchiveSweep_LeavesNoRegistryEntry states the same // TestArchiveSweep_LeavesNoRegistryEntry states the same
// invariant in its general form: sweeping an archive whose // 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 // registry keeps holding only writers a delivery created and an
// eviction can reach. // eviction can reach.
func TestArchiveSweep_LeavesNoRegistryEntry(t *testing.T) { func TestArchiveSweep_LeavesNoRegistryEntry(t *testing.T) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.seedArchiveRows( path := env.seedArchiveRows(
t, webhookID, t, tgt,
time.Now().Add(-48*time.Hour), time.Now().Add(-48*time.Hour),
time.Now().Add(-time.Minute), 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()) env.sweeper.ExportSweep(context.Background())
@@ -349,7 +366,7 @@ func TestArchiveSweep_LeavesNoRegistryEntry(t *testing.T) {
"the sweep must still prune an idle archive", "the sweep must still prune an idle archive",
) )
assert.False( assert.False(
t, env.eng.ExportHasArchiveWriter(webhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the sweep must release the registry entry it created", "the sweep must release the registry entry it created",
) )
} }
@@ -364,34 +381,31 @@ func TestArchiveSweep_KeepsWriterAdoptedByDelivery(
) { ) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows( env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour), t, tgt, time.Now().Add(-48*time.Hour),
) )
webhookDB := testWebhookDB(t) webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"n":1}`) event := seedEvent(t, webhookDB, `{"n":1}`)
event.WebhookID = webhookID d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt)
d := seedDatabaseTargetDelivery(
t, webhookDB, event, `{"expiry":"1h"}`,
)
env.sweeper.ExportSweep(context.Background()) 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) env.eng.ExportDeliverDatabase(webhookDB, d)
assert.True( assert.True(
t, env.eng.ExportHasArchiveWriter(webhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"a delivery's writer must stay registered", "a delivery's writer must stay registered",
) )
env.sweeper.ExportSweep(context.Background()) env.sweeper.ExportSweep(context.Background())
assert.True( assert.True(
t, env.eng.ExportHasArchiveWriter(webhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"a sweep must not drop a writer a delivery owns", "a sweep must not drop a writer a delivery owns",
) )
} }
@@ -423,15 +437,15 @@ func TestArchiveSweep_KeepsWriterAdoptedDuringSweep(
) { ) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows( env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour), t, tgt, time.Now().Add(-48*time.Hour),
) )
sweepWriter, created, err := env.eng.ExportSweepWriterFor( sweepWriter, created, err := env.eng.ExportSweepWriterFor(
webhookID, tgt.ID,
) )
require.NoError(t, err) require.NoError(t, err)
require.True( require.True(
@@ -442,37 +456,34 @@ func TestArchiveSweep_KeepsWriterAdoptedDuringSweep(
// The delivery lands mid-sweep and adopts the entry. // The delivery lands mid-sweep and adopts the entry.
webhookDB := testWebhookDB(t) webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"n":1}`) event := seedEvent(t, webhookDB, `{"n":1}`)
event.WebhookID = webhookID d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt)
d := seedDatabaseTargetDelivery(
t, webhookDB, event, `{"expiry":"1h"}`,
)
env.eng.ExportDeliverDatabase(webhookDB, d) env.eng.ExportDeliverDatabase(webhookDB, d)
adopted := env.eng.ExportArchiveWriterFor(webhookID) adopted := env.eng.ExportArchiveWriterFor(tgt.ID)
require.NotNil(t, adopted) require.NotNil(t, adopted)
require.True( require.True(
t, sweepWriter.Same(adopted), t, sweepWriter.Same(adopted),
"the delivery must have adopted the sweep's writer", "the delivery must have adopted the sweep's writer",
) )
require.True( require.True(
t, env.eng.ExportArchiveHandleOpen(webhookID), t, env.eng.ExportArchiveHandleOpen(tgt.ID),
"the delivery leaves the archive handle open", "the delivery leaves the archive handle open",
) )
// The sweep finishes. // The sweep finishes.
env.eng.ExportReleaseSweepWriter(webhookID, sweepWriter) env.eng.ExportReleaseSweepWriter(tgt.ID, sweepWriter)
require.True( require.True(
t, env.eng.ExportHasArchiveWriter(webhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"a writer adopted by a delivery during a sweep must "+ "a writer adopted by a delivery during a sweep must "+
"stay registered, or its open handle is unreachable", "stay registered, or its open handle is unreachable",
) )
env.eng.EvictWebhook(webhookID) env.eng.EvictTarget(tgt.ID)
assert.False( assert.False(
t, env.eng.ExportHasArchiveWriter(webhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the adopted writer must still be evictable", "the adopted writer must still be evictable",
) )
assert.False( assert.False(
@@ -481,34 +492,34 @@ func TestArchiveSweep_KeepsWriterAdoptedDuringSweep(
) )
} }
// TestArchiveSweep_ContinuesAfterPerWebhookFailure proves a // TestArchiveSweep_ContinuesAfterPerTargetFailure proves a
// failure for one webhook does not abort the sweep for the // failure for one target does not abort the sweep for the
// others: an unparseable expiry and an unreadable archive both // others: an unparseable expiry and an unreadable archive both
// have to be logged and stepped over. // have to be logged and stepped over.
func TestArchiveSweep_ContinuesAfterPerWebhookFailure( func TestArchiveSweep_ContinuesAfterPerTargetFailure(
t *testing.T, t *testing.T,
) { ) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
// Seeded first so the sweep reaches them before the healthy // Seeded first so the sweep reaches them before the healthy
// webhook: targets come back in insertion order. // target: targets come back in insertion order.
badConfigID := env.seedDatabaseTarget(t, `{"expiry":"!!!"}`) badConfig := env.seedDatabaseTarget(t, `{"expiry":"!!!"}`)
env.seedArchiveRows( 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( require.NoError(t, os.WriteFile(
env.archivePath(corruptID), env.archivePath(corrupt),
[]byte("this is not a sqlite database"), []byte("this is not a sqlite database"),
0o600, 0o600,
)) ))
healthyID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) healthy := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
healthyPath := env.seedArchiveRows( healthyPath := env.seedArchiveRows(
t, healthyID, t, healthy,
time.Now().Add(-48*time.Hour), time.Now().Add(-48*time.Hour),
time.Now().Add(-time.Minute), time.Now().Add(-time.Minute),
) )
@@ -518,14 +529,14 @@ func TestArchiveSweep_ContinuesAfterPerWebhookFailure(
assert.Equal( assert.Equal(
t, []string{sweepRowNew}, t, []string{sweepRowNew},
archivedEventIDs(t, healthyPath), 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", "sweep from pruning the ones after it",
) )
} }
// TestArchiveSweep_OpenExistingDoesNotCreateFile pins the second // TestArchiveSweep_OpenExistingDoesNotCreateFile pins the second
// of the two no-create guards. The first is the stat in // 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 // protects the window between that stat and the open. Flipping
// the sweep's mode to create-if-missing makes this fail. // the sweep's mode to create-if-missing makes this fail.
func TestArchiveSweep_OpenExistingDoesNotCreateFile( func TestArchiveSweep_OpenExistingDoesNotCreateFile(
@@ -561,13 +572,13 @@ func TestArchiveSweep_OpenExistingDoesNotCreateFile(
func TestArchiveSweep_PrunesIdleArchive(t *testing.T) { func TestArchiveSweep_PrunesIdleArchive(t *testing.T) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
now := time.Now() now := time.Now()
path := env.seedArchiveRows( path := env.seedArchiveRows(
t, webhookID, t, tgt,
now.Add(-48*time.Hour), now.Add(-48*time.Hour),
now.Add(-time.Minute), now.Add(-time.Minute),
) )
@@ -600,11 +611,11 @@ func TestArchiveSweep_PrunesIdleArchive(t *testing.T) {
func TestArchiveSweep_LeavesArchiveClosed(t *testing.T) { func TestArchiveSweep_LeavesArchiveClosed(t *testing.T) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.seedArchiveRows( path := env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour), t, tgt, time.Now().Add(-48*time.Hour),
) )
w := delivery.NewExportArchiveWriter( w := delivery.NewExportArchiveWriter(
@@ -640,35 +651,32 @@ func TestArchiveSweep_ClosesHandleOfRegisteredWriter(
) { ) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows( env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour), t, tgt, time.Now().Add(-48*time.Hour),
) )
webhookDB := testWebhookDB(t) webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"n":1}`) event := seedEvent(t, webhookDB, `{"n":1}`)
event.WebhookID = webhookID d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt)
d := seedDatabaseTargetDelivery(
t, webhookDB, event, `{"expiry":"1h"}`,
)
env.eng.ExportDeliverDatabase(webhookDB, d) env.eng.ExportDeliverDatabase(webhookDB, d)
require.True( require.True(
t, env.eng.ExportArchiveHandleOpen(webhookID), t, env.eng.ExportArchiveHandleOpen(tgt.ID),
"the delivery must leave the archive handle open", "the delivery must leave the archive handle open",
) )
env.sweeper.ExportSweep(context.Background()) env.sweeper.ExportSweep(context.Background())
require.True( require.True(
t, env.eng.ExportHasArchiveWriter(webhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the delivery's registry entry must survive the sweep", "the delivery's registry entry must survive the sweep",
) )
assert.False( assert.False(
t, env.eng.ExportArchiveHandleOpen(webhookID), t, env.eng.ExportArchiveHandleOpen(tgt.ID),
"the sweep must leave the archive closed", "the sweep must leave the archive closed",
) )
} }
@@ -684,11 +692,11 @@ func TestArchiveSweep_NeverExpiryUntouched(t *testing.T) {
`{"expiry":""}`, `{"expiry":""}`,
"", "",
} { } {
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, configJSON) tgt := env.seedDatabaseTarget(t, configJSON)
path := env.seedArchiveRows( path := env.seedArchiveRows(
t, webhookID, t, tgt,
time.Now().Add(-10000*time.Hour), time.Now().Add(-10000*time.Hour),
) )
@@ -699,7 +707,7 @@ func TestArchiveSweep_NeverExpiryUntouched(t *testing.T) {
"config %q must keep rows forever", configJSON, "config %q must keep rows forever", configJSON,
) )
assert.False( assert.False(
t, env.eng.ExportHasArchiveWriter(webhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"config %q must leave no registry entry behind", "config %q must leave no registry entry behind",
configJSON, configJSON,
) )
@@ -722,10 +730,10 @@ func TestArchiveSweep_NeverExpirySkipsBeforeOpening(
) { ) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"never"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"never"}`)
path := env.archivePath(webhookID) path := env.archivePath(tgt)
seedUnmigratedArchive(t, path) seedUnmigratedArchive(t, path)
require.False(t, archiveTableExists(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 // TestArchiveSweep_DoesNotCreateArchiveFile proves the sweep
// never conjures an archive: a webhook with a database target // never conjures an archive: a database target that has never
// that has never received an event must still have no archive // received an event must still have no archive file (nor SQLite
// file (nor SQLite sidecar) after a sweep. // sidecar) after a sweep, and no registry entry either.
func TestArchiveSweep_DoesNotCreateArchiveFile(t *testing.T) { func TestArchiveSweep_DoesNotCreateArchiveFile(t *testing.T) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.archivePath(webhookID) path := env.archivePath(tgt)
require.NoFileExists(t, path) require.NoFileExists(t, path)
@@ -789,6 +797,11 @@ func TestArchiveSweep_DoesNotCreateArchiveFile(t *testing.T) {
"the sweep must not create an archive file", "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 // TestArchiveSweep_DoesNotCreateAfterWriterExists covers the
@@ -800,11 +813,11 @@ func TestArchiveSweep_DoesNotCreateAfterWriterExists(
) { ) {
t.Parallel() 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.NoError(t, err)
require.NoFileExists(t, path) require.NoFileExists(t, path)
@@ -819,17 +832,17 @@ func TestArchiveSweep_DoesNotCreateAfterWriterExists(
func TestArchiveSweep_SkipsDeletedWebhookTargets(t *testing.T) { func TestArchiveSweep_SkipsDeletedWebhookTargets(t *testing.T) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.seedArchiveRows( path := env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour), t, tgt, time.Now().Add(-48*time.Hour),
) )
require.NoError( require.NoError(
t, t,
env.mainDB.DB(). env.mainDB.DB().
Where("webhook_id = ?", webhookID). Where("webhook_id = ?", tgt.WebhookID).
Delete(&database.Target{}).Error, Delete(&database.Target{}).Error,
) )
@@ -842,14 +855,14 @@ func TestArchiveSweep_SkipsDeletedWebhookTargets(t *testing.T) {
} }
// TestArchiveSweep_ConcurrentWrites proves the sweep serialises // TestArchiveSweep_ConcurrentWrites proves the sweep serialises
// against writes through the per-webhook writer mutex. Run // against writes through the target's writer mutex. Run under
// under -race, an unsynchronised sweep would be caught here. // -race, an unsynchronised sweep would be caught here.
func TestArchiveSweep_ConcurrentWrites(t *testing.T) { func TestArchiveSweep_ConcurrentWrites(t *testing.T) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
webhookDB := testWebhookDB(t) webhookDB := testWebhookDB(t)
@@ -862,13 +875,10 @@ func TestArchiveSweep_ConcurrentWrites(t *testing.T) {
for range sweepConcurrentWrites { for range sweepConcurrentWrites {
event := seedEvent(t, webhookDB, `{"n":1}`) event := seedEvent(t, webhookDB, `{"n":1}`)
event.WebhookID = webhookID
deliveries = append( deliveries = append(
deliveries, deliveries,
seedDatabaseTargetDelivery( seedDatabaseTargetDelivery(t, webhookDB, event, tgt),
t, webhookDB, event, `{"expiry":"1h"}`,
),
) )
} }
@@ -894,7 +904,7 @@ func TestArchiveSweep_ConcurrentWrites(t *testing.T) {
wg.Wait() wg.Wait()
assert.FileExists(t, env.archivePath(webhookID)) assert.FileExists(t, env.archivePath(tgt))
} }
// TestArchiveSweeper_StopsCleanly proves the background loop // TestArchiveSweeper_StopsCleanly proves the background loop
@@ -902,11 +912,11 @@ func TestArchiveSweep_ConcurrentWrites(t *testing.T) {
func TestArchiveSweeper_StopsCleanly(t *testing.T) { func TestArchiveSweeper_StopsCleanly(t *testing.T) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows( env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour), t, tgt, time.Now().Add(-48*time.Hour),
) )
env.sweeper.ExportSetInterval(time.Millisecond) env.sweeper.ExportSetInterval(time.Millisecond)
@@ -930,7 +940,7 @@ func TestArchiveSweeper_StopHookHonoursStopTimeout(
) { ) {
t.Parallel() t.Parallel()
env := setupSweeperTest(t) env := setupArchiveTest(t)
lc := &recordingLifecycle{} lc := &recordingLifecycle{}
env.sweeper.ExportRegisterHooks(lc) env.sweeper.ExportRegisterHooks(lc)
+51 -20
View File
@@ -123,21 +123,24 @@ type Notifier interface {
Notify(tasks []Task) Notify(tasks []Task)
} }
// WebhookEvictor releases the delivery engine's per-webhook // Archives is how the handlers keep the database targets' archive
// state for a webhook that no longer needs it — currently the // files in step with the configuration. Deleting a webhook or a
// cached archive writer of the database target, whose open // target releases the cached archive writers, whose open file
// file handle would otherwise outlive the webhook. // 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 // It is deliberately separate from Notifier: archiving lifecycle
// one method wide: archiving lifecycle is not notification, and // is not notification, and a small interface keeps the handlers
// a single-method interface keeps the handlers package free of // package free of any dependency on the engine's internals while
// any dependency on the engine's internals while staying // staying trivially fakeable in tests.
// trivially fakeable in tests.
// //
// EvictWebhook never deletes an archive file. It is idempotent // Neither eviction deletes an archive file. Both are idempotent
// and is a no-op for a webhook with no engine state. // and are no-ops for a webhook or target with no engine state.
type WebhookEvictor interface { type Archives interface {
EvictWebhook(webhookID string) EvictWebhook(webhookID string)
EvictTarget(targetID string)
Rename(targetID, webhookName, targetName string) error
} }
// EngineParams are the fx dependencies for the delivery // EngineParams are the fx dependencies for the delivery
@@ -188,7 +191,7 @@ type Engine struct {
httpTarget *httpTarget httpTarget *httpTarget
// dbTarget is retained so the engine can reach the archive // 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 dbTarget *databaseTarget
// inflight is the set of deliveries this engine currently owns. // inflight is the set of deliveries this engine currently owns.
@@ -257,17 +260,44 @@ func (e *Engine) Notify(tasks []Task) {
} }
} }
// EvictWebhook implements WebhookEvictor. It releases the // EvictWebhook implements Archives. The cached archive writer of
// engine's per-webhook archiving state: the database target's // every database target of the webhook is dropped from the
// cached archive writer is dropped from the registry and its // registry and its file handle closed. The archive files
// file handle closed. The archive file itself is left on disk // themselves are left on disk — they are long-term storage the
// — it is long-term storage the operator owns. // operator owns.
func (e *Engine) EvictWebhook(webhookID string) { func (e *Engine) EvictWebhook(webhookID string) {
if e.dbTarget == nil { if e.dbTarget == nil {
return 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 // ScheduleRetry schedules a task to be re-enqueued onto the
@@ -381,7 +411,8 @@ func (e *Engine) start() {
// Once the pool has drained it closes the archive writers, so a // Once the pool has drained it closes the archive writers, so a
// clean stop leaves no archive -wal behind. Nothing else holds a // clean stop leaves no archive -wal behind. Nothing else holds a
// writer for long by then: the archive sweeper stops before the // 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 // not drain in time, the writers are left open, as a kill would
// leave them. Closing them would wait for any write in progress, // leave them. Closing them would wait for any write in progress,
// and a worker still running would then open new writers that // and a worker still running would then open new writers that
+23 -10
View File
@@ -2,7 +2,6 @@ package delivery_test
import ( import (
"context" "context"
"fmt"
"path/filepath" "path/filepath"
"testing" "testing"
"time" "time"
@@ -10,6 +9,7 @@ import (
"github.com/google/uuid" "github.com/google/uuid"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
"go.uber.org/fx" "go.uber.org/fx"
"gorm.io/gorm/clause"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
) )
@@ -272,22 +272,35 @@ func TestEngine_StopHookHonoursStopTimeout(t *testing.T) {
requireStopHookExpires(t, lc.hooks[0], "delivery engine") requireStopHookExpires(t, lc.hooks[0], "delivery engine")
} }
// deliverToArchive runs one delivery to a database target through // deliverToArchive gives the setup's webhook a database target,
// the running engine and returns the webhook's archive file path. // runs one delivery to it through the running engine, and returns
// The archive writer holds the file open afterwards. // the target's ID and archive file path. The archive writer holds
func deliverToArchive(t *testing.T, s iSetup) string { // the file open afterwards.
func deliverToArchive(t *testing.T, s iSetup) (string, string) {
t.Helper() 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) deliveryID, task := seedLogTask(t, s)
task.TargetID = tgt.ID
task.TargetType = database.TargetTypeDatabase task.TargetType = database.TargetTypeDatabase
s.Engine.Notify([]delivery.Task{task}) s.Engine.Notify([]delivery.Task{task})
iWaitForDelivered(t, s.WebhookDB, deliveryID) iWaitForDelivered(t, s.WebhookDB, deliveryID)
return filepath.Join( return tgt.ID, filepath.Join(
filepath.Dir(s.DBMgr.DBPath(s.WebhookID)), 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) lc := startEngineViaHook(t, s.Engine)
path := deliverToArchive(t, s) _, path := deliverToArchive(t, s)
require.FileExists( require.FileExists(
t, path+"-wal", t, path+"-wal",
"an open archive should have a -wal for the stop to remove", "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) lc := startEngineViaHook(t, s.Engine)
deliverToArchive(t, s) targetID, _ := deliverToArchive(t, s)
release := make(chan struct{}) release := make(chan struct{})
@@ -352,7 +365,7 @@ func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) {
requireStopHookExpires(t, lc.hooks[0], "delivery engine") requireStopHookExpires(t, lc.hooks[0], "delivery engine")
require.True( require.True(
t, s.Engine.ExportArchiveHandleOpen(s.WebhookID), t, s.Engine.ExportArchiveHandleOpen(targetID),
"a stop that timed out must not close archive writers", "a stop that timed out must not close archive writers",
) )
} }
+16 -29
View File
@@ -351,23 +351,15 @@ func TestDeliverDatabase_ImmediateSuccess(
db := testWebhookDB(t) db := testWebhookDB(t)
// The database target archives for real now, so the engine // The database target archives for real, so the engine needs
// needs a webhook DB manager to locate the data directory. // the target in the main database and a data directory.
e := delivery.NewTestEngineWithDB( env := setupArchiveTest(t)
nil, tgt := env.seedDatabaseTarget(t, "")
database.NewTestWebhookDBManager(t.TempDir()),
slog.New(slog.NewTextHandler(
os.Stderr,
&slog.HandlerOptions{Level: slog.LevelDebug},
)),
&http.Client{Timeout: 5 * time.Second},
1,
)
event := seedEvent(t, db, `{"db":"target"}`) 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 var updated database.Delivery
@@ -1332,32 +1324,27 @@ func TestProcessDelivery_RoutesToCorrectHandler(
db := testWebhookDB(t) db := testWebhookDB(t)
// The database target archives for real now, so the engine // The database target archives for real, so the engine needs
// needs a webhook DB manager to locate the data directory. // the target in the main database and a data directory.
e := delivery.NewTestEngineWithDB( env := setupArchiveTest(t)
nil, archive := env.seedDatabaseTarget(t, "")
database.NewTestWebhookDBManager(t.TempDir()),
slog.New(slog.NewTextHandler(
os.Stderr,
&slog.HandlerOptions{Level: slog.LevelDebug},
)),
&http.Client{Timeout: 5 * time.Second},
1,
)
tests := []struct { tests := []struct {
name string name string
targetType database.TargetType targetType database.TargetType
targetID string
wantStatus database.DeliveryStatus wantStatus database.DeliveryStatus
}{ }{
{ {
"database target", "database target",
database.TargetTypeDatabase, database.TargetTypeDatabase,
archive.ID,
database.DeliveryStatusDelivered, database.DeliveryStatusDelivered,
}, },
{ {
"log target", "log target",
database.TargetTypeLog, database.TargetTypeLog,
uuid.New().String(),
database.DeliveryStatusDelivered, database.DeliveryStatusDelivered,
}, },
} }
@@ -1367,7 +1354,7 @@ func TestProcessDelivery_RoutesToCorrectHandler(
t.Parallel() t.Parallel()
runRoutingSubtest( runRoutingSubtest(
t, db, e, tt.targetType, t, db, env.eng, tt.targetType, tt.targetID,
tt.wantStatus, tt.wantStatus,
) )
}) })
@@ -1379,6 +1366,7 @@ func runRoutingSubtest(
db *gorm.DB, db *gorm.DB,
e *delivery.Engine, e *delivery.Engine,
targetType database.TargetType, targetType database.TargetType,
targetID string,
wantStatus database.DeliveryStatus, wantStatus database.DeliveryStatus,
) { ) {
t.Helper() t.Helper()
@@ -1386,8 +1374,7 @@ func runRoutingSubtest(
event := seedEvent(t, db, `{"routing":"test"}`) event := seedEvent(t, db, `{"routing":"test"}`)
dlv := seedDelivery( dlv := seedDelivery(
t, db, event.ID, t, db, event.ID, targetID,
uuid.New().String(),
database.DeliveryStatusPending, database.DeliveryStatusPending,
) )
+26 -20
View File
@@ -474,7 +474,7 @@ func NewTestCircuitBreaker(
type ExportArchivedEvent = archivedEvent type ExportArchivedEvent = archivedEvent
// ExportArchiveWriter wraps an archiveWriter so black-box tests // 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 { type ExportArchiveWriter struct {
w *archiveWriter w *archiveWriter
} }
@@ -549,6 +549,12 @@ func (e *ExportArchiveWriter) Evict() {
e.w.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 // HandleOpen reports whether the writer currently holds an open
// archive handle. // archive handle.
func (e *ExportArchiveWriter) HandleOpen() bool { func (e *ExportArchiveWriter) HandleOpen() bool {
@@ -568,16 +574,16 @@ func (e *ExportArchiveWriter) Same(
} }
// ExportArchiveWriterFor returns the archive writer the registry // ExportArchiveWriterFor returns the archive writer the registry
// currently caches for a webhook, or nil when none is cached. It // currently caches for a database target, or nil when none is
// never creates one, so a test can hold a reference to the very // cached. It never creates one, so a test can hold a reference to
// writer an eviction is about to detach. // the very writer an eviction is about to detach.
func (e *Engine) ExportArchiveWriterFor( func (e *Engine) ExportArchiveWriterFor(
webhookID string, targetID string,
) *ExportArchiveWriter { ) *ExportArchiveWriter {
e.dbTarget.mu.Lock() e.dbTarget.mu.Lock()
defer e.dbTarget.mu.Unlock() defer e.dbTarget.mu.Unlock()
w, ok := e.dbTarget.writers[webhookID] w, ok := e.dbTarget.writers[targetID]
if !ok { if !ok {
return nil return nil
} }
@@ -586,26 +592,26 @@ func (e *Engine) ExportArchiveWriterFor(
} }
// ExportHasArchiveWriter reports whether the database target // 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( func (e *Engine) ExportHasArchiveWriter(
webhookID string, targetID string,
) bool { ) bool {
e.dbTarget.mu.Lock() e.dbTarget.mu.Lock()
defer e.dbTarget.mu.Unlock() defer e.dbTarget.mu.Unlock()
_, ok := e.dbTarget.writers[webhookID] _, ok := e.dbTarget.writers[targetID]
return ok return ok
} }
// ExportArchiveHandleOpen reports whether the cached archive // 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. // returns false when no writer is cached.
func (e *Engine) ExportArchiveHandleOpen( func (e *Engine) ExportArchiveHandleOpen(
webhookID string, targetID string,
) bool { ) bool {
e.dbTarget.mu.Lock() e.dbTarget.mu.Lock()
w, ok := e.dbTarget.writers[webhookID] w, ok := e.dbTarget.writers[targetID]
e.dbTarget.mu.Unlock() e.dbTarget.mu.Unlock()
if !ok { if !ok {
@@ -619,12 +625,12 @@ func (e *Engine) ExportArchiveHandleOpen(
} }
// ExportEnsureArchiveWriter creates (if needed) and returns the // 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. // test can prime the registry the way a delivery would.
func (e *Engine) ExportEnsureArchiveWriter( func (e *Engine) ExportEnsureArchiveWriter(
webhookID string, targetID string,
) (string, error) { ) (string, error) {
w, err := e.dbTarget.writerFor(webhookID) w, err := e.dbTarget.writerFor(targetID)
if err != nil { if err != nil {
return "", err return "", err
} }
@@ -632,14 +638,14 @@ func (e *Engine) ExportEnsureArchiveWriter(
return w.path, nil 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 // as the idle sweep does, reporting whether the sweep had to
// create the entry. It lets a test drive the registry through the // create the entry. It lets a test drive the registry through the
// sweep's own entry point instead of choreographing goroutines. // sweep's own entry point instead of choreographing goroutines.
func (e *Engine) ExportSweepWriterFor( func (e *Engine) ExportSweepWriterFor(
webhookID string, targetID string,
) (*ExportArchiveWriter, bool, error) { ) (*ExportArchiveWriter, bool, error) {
w, created, err := e.dbTarget.sweepWriterFor(webhookID) w, created, err := e.dbTarget.sweepWriterFor(targetID)
if err != nil { if err != nil {
return nil, false, err return nil, false, err
} }
@@ -650,9 +656,9 @@ func (e *Engine) ExportSweepWriterFor(
// ExportReleaseSweepWriter releases a sweep-created registry entry // ExportReleaseSweepWriter releases a sweep-created registry entry
// exactly as a finished sweep does. // exactly as a finished sweep does.
func (e *Engine) ExportReleaseSweepWriter( 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 // NewTestArchiveSweeper builds an ArchiveSweeper backed by the
+196 -89
View File
@@ -4,6 +4,7 @@ import (
"context" "context"
"fmt" "fmt"
"path/filepath" "path/filepath"
"strings"
"sync" "sync"
"time" "time"
@@ -11,22 +12,75 @@ import (
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
) )
// databaseTarget is a no-retry target that archives the // archiveNameMaxLen is how many characters of a webhook or target
// full inbound event into a per-webhook archive SQLite file, // name an archive file name keeps.
// separate from the per-webhook event database. The event is const archiveNameMaxLen = 40
// already persisted in the per-webhook event DB by the time
// delivery runs; the database target additionally writes a // databaseTarget is a no-retry target that archives the full
// durable long-term copy into archive-{webhookID}.db and then // inbound event into the target's own archive SQLite file, separate
// records a single attempt whose outcome reflects whether the // from the per-webhook event database. The event is already
// archive write succeeded. See archiveWriter for the // persisted in the per-webhook event DB by the time delivery runs;
// close/reopen, auto-recreate, and expiry semantics. // 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 { type databaseTarget struct {
eng *Engine eng *Engine
// writers holds one archive writer per database target, keyed
// by target ID.
mu sync.Mutex mu sync.Mutex
writers map[string]*archiveWriter 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 // Deliver implements Target. It archives the event, then
// records one successful attempt and marks the delivery // records one successful attempt and marks the delivery
// delivered. An archiving error fails the delivery: the // 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 // archive database, honouring the optional per-target expiry
// parsed from the target config JSON. // parsed from the target config JSON.
func (t *databaseTarget) archive(d *database.Delivery) error { func (t *databaseTarget) archive(d *database.Delivery) error {
@@ -106,7 +160,7 @@ func (t *databaseTarget) archive(d *database.Delivery) error {
return err return err
} }
w, err := t.writerFor(webhookID) w, err := t.writerFor(d.TargetID)
if err != nil { if err != nil {
return err return err
} }
@@ -124,30 +178,31 @@ func (t *databaseTarget) archive(d *database.Delivery) error {
return w.write(row, expiry) return w.write(row, expiry)
} }
// writerFor returns the archiveWriter for a webhook, creating // writerFor returns the archive writer for a database target,
// and caching it on first use. Each webhook has one writer so // creating and caching it on first use. Each target has one writer
// its close/reopen debounce state is shared across concurrent // so its close/reopen debounce state is shared across concurrent
// deliveries. The archive file lives beside the per-webhook // deliveries, and so a rename and the idle sweep take the same lock
// event database in the data directory. // as its writes.
func (t *databaseTarget) writerFor( func (t *databaseTarget) writerFor(
webhookID string, targetID string,
) (*archiveWriter, error) { ) (*archiveWriter, error) {
path, err := t.archivePath(webhookID) t.mu.Lock()
defer t.mu.Unlock()
w, ok := t.writers[targetID]
if !ok {
var err error
w, err = t.newWriter(targetID)
if err != nil { if err != nil {
return nil, err return nil, err
} }
t.mu.Lock()
defer t.mu.Unlock()
if t.writers == nil { if t.writers == nil {
t.writers = make(map[string]*archiveWriter) t.writers = make(map[string]*archiveWriter)
} }
w, ok := t.writers[webhookID] t.writers[targetID] = w
if !ok {
w = newArchiveWriter(path, t.eng.log)
t.writers[webhookID] = w
} }
// A delivery claims the entry: even if the idle sweep created // 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 // sweepWriterFor returns the archive writer the idle sweep should
// prune a webhook through, together with whether the sweep itself // prune a target's archive through, together with whether the sweep
// created the registry entry. // itself created the registry entry.
// //
// The sweep must route its prune through the registered writer so // The sweep must route its prune through the registered writer so
// the writer's mutex orders it against concurrent writes, but it // the writer's mutex orders it against concurrent writes, but it
// must never leave a registry entry behind: a sweep that ran // 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 // re-create an entry that nothing will ever evict again, which is
// exactly the leak eviction exists to prevent. An entry the sweep // exactly the leak eviction exists to prevent. An entry the sweep
// creates is therefore marked sweep-owned and handed back to // creates is therefore marked sweep-owned and handed back to
// releaseSweepWriter when the sweep is done. // releaseSweepWriter when the sweep is done.
func (t *databaseTarget) sweepWriterFor( func (t *databaseTarget) sweepWriterFor(
webhookID string, targetID string,
) (*archiveWriter, bool, error) { ) (*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 { if err != nil {
return nil, false, err return nil, false, err
} }
t.mu.Lock()
defer t.mu.Unlock()
if t.writers == nil { if t.writers == nil {
t.writers = make(map[string]*archiveWriter) 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 w.sweepOwned = true
t.writers[webhookID] = w t.writers[targetID] = w
return w, true, nil return w, true, nil
} }
@@ -209,57 +263,95 @@ func (t *databaseTarget) sweepWriterFor(
// delivery that adopted the writer keeps a registered, evictable // delivery that adopted the writer keeps a registered, evictable
// one. // one.
func (t *databaseTarget) releaseSweepWriter( func (t *databaseTarget) releaseSweepWriter(
webhookID string, w *archiveWriter, targetID string, w *archiveWriter,
) { ) {
t.mu.Lock() t.mu.Lock()
defer t.mu.Unlock() defer t.mu.Unlock()
cur, ok := t.writers[webhookID] cur, ok := t.writers[targetID]
if !ok || cur != w || !cur.sweepOwned { if !ok || cur != w || !cur.sweepOwned {
return return
} }
delete(t.writers, webhookID) delete(t.writers, targetID)
} }
// archivePath returns the archive file path for a webhook: it // newWriter builds the writer for a database target's archive. The
// lives beside the per-webhook event database in the data // file lives beside the webhook's event database in the data
// directory. It does not touch the filesystem. // directory and is named for the webhook and the target as the main
func (t *databaseTarget) archivePath( // database has them now; from then on only rename changes the name
webhookID string, // the writer uses. It does not touch the archive file.
) (string, error) { func (t *databaseTarget) newWriter(
targetID string,
) (*archiveWriter, error) {
if t.eng.dbManager == nil { 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( err := t.eng.database.DB().
dir, fmt.Sprintf("archive-%s.db", webhookID), Preload("Webhook").
), nil First(&target, "id = ?", targetID).Error
if err != nil {
return nil, fmt.Errorf(
"loading database target %s: %w", targetID, err,
)
} }
// evict drops a webhook's archive writer from the registry and dir := filepath.Dir(t.eng.dbManager.DBPath(target.WebhookID))
// closes its handle, so a deleted webhook does not leave a name := ArchiveFileName(
// writer (and an open archive handle within its debounce target.Webhook.Name, target.Name, target.ID,
// window) alive for the process lifetime. )
w := newArchiveWriter(filepath.Join(dir, name), t.eng.log)
w.webhookID = target.WebhookID
return w, nil
}
// 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 // The map entry is removed under the registry lock, which is
// then released before the handle is closed under the writer's // then released before the handle is closed under the writer's
// own lock: that ordering keeps the registry available to other // 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 // closing under the writer's lock means eviction can never race
// a write. // a write.
// //
// Eviction is idempotent and silent for a webhook with no // Eviction is idempotent and silent for a target with no writer,
// writer, which is the common case: a webhook with no database // which is the common case: only a database target that has
// target never creates one. It never deletes the archive file. // received an event or been renamed has one. It never deletes the
func (t *databaseTarget) evict(webhookID string) { // archive file.
func (t *databaseTarget) evict(targetID string) {
t.mu.Lock() t.mu.Lock()
w, ok := t.writers[webhookID] w, ok := t.writers[targetID]
if ok { if ok {
delete(t.writers, webhookID) delete(t.writers, targetID)
} }
t.mu.Unlock() t.mu.Unlock()
@@ -272,13 +364,41 @@ func (t *databaseTarget) evict(webhookID string) {
t.eng.log.Info( t.eng.log.Info(
"evicted archive writer", "evicted archive writer",
"webhook_id", webhookID, "target_id", targetID,
"path", w.path, "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 // 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 // workers have returned. Closing the last handle on an archive
// moves the contents of its -wal into the .db and removes the // moves the contents of its -wal into the .db and removes the
// -wal, so a clean stop leaves each archive as a single file. // -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 // sweepArchive prunes one database target's archive of rows older
// expiry, without requiring a write. It returns nil (nothing to // than expiry, without requiring a write. A missing archive file is
// do) when the archive file does not exist, so a sweep never // left missing (see sweepExpired), so a sweep never creates an
// creates an archive for a webhook that has a database target // archive for a target that has never received an event.
// but has never received an event.
// //
// It also never leaves a registry entry behind: an entry it had // It also never leaves a registry entry behind: an entry it had
// to create to reach the writer's mutex is released again once // 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. // resurrect the writer the eviction just dropped.
func (t *databaseTarget) sweepWebhook( func (t *databaseTarget) sweepArchive(
webhookID string, expiry time.Duration, targetID string, expiry time.Duration,
) error { ) error {
path, err := t.archivePath(webhookID) w, created, err := t.sweepWriterFor(targetID)
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)
if err != nil { if err != nil {
return err return err
} }
if created { if created {
defer t.releaseSweepWriter(webhookID, w) defer t.releaseSweepWriter(targetID, w)
} }
return w.sweepExpired(expiry) return w.sweepExpired(expiry)
+81 -17
View File
@@ -4,8 +4,10 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"io/fs"
"log/slog" "log/slog"
"os" "os"
"path/filepath"
"sync" "sync"
"time" "time"
@@ -41,7 +43,7 @@ const (
var ( var (
// errArchiveMissingWebhookID is returned when an event to // 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( errArchiveMissingWebhookID = errors.New(
"cannot archive event without a webhook id", "cannot archive event without a webhook id",
) )
@@ -61,13 +63,19 @@ var (
) )
// errArchiveWriterEvicted is returned when a writer that has // errArchiveWriterEvicted is returned when a writer that has
// been evicted (its webhook was deleted, or its last database // been evicted (its target or its webhook was deleted) is used
// target was removed) is used again. An evicted writer is no // again. An evicted writer is no longer in the registry, so
// longer in the registry, so reopening its file would leak a // reopening its file would leak a handle nothing owns.
// handle nothing owns.
errArchiveWriterEvicted = errors.New( errArchiveWriterEvicted = errors.New(
"archive writer has been evicted", "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 // 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 // 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 // self-contained copy — independent of the per-webhook event
// database, which may prune events under its own retention. // database, which may prune events under its own retention.
type archivedEvent struct { type archivedEvent struct {
@@ -170,8 +178,8 @@ func ValidateArchiveExpiry(expiry string) error {
return nil return nil
} }
// archiveWriter owns one per-webhook archive SQLite file. It // archiveWriter owns one database target's archive SQLite file.
// serialises writes, and after each write closes and reopens // It serialises writes, and after each write closes and reopens
// the file (debounced to at most once per debounce window) so // the file (debounced to at most once per debounce window) so
// an operator can move the file away for offline archiving. The // an operator can move the file away for offline archiving. The
// next write recreates a moved or removed file, because the // next write recreates a moved or removed file, because the
@@ -187,16 +195,21 @@ type archiveWriter struct {
reopens int reopens int
// evicted marks a writer that has been removed from the // evicted marks a writer that has been removed from the
// per-webhook registry. Its handle is closed and it must // registry. Its handle is closed and it must never open the
// never open the file again: nothing holds it any more, so a // file again: nothing holds it any more, so a reopen would
// reopen would leak the handle for the process lifetime. // leak the handle for the process lifetime.
evicted bool 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 // 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 removes such an entry again when it is done, so a
// sweep can never leave — or resurrect — a registry entry // 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 // the writer clears the flag, handing the entry to the
// registry proper. // registry proper.
// //
@@ -385,11 +398,62 @@ func (w *archiveWriter) sweepExpired(expiry time.Duration) error {
return nil 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 // evict closes the writer's handle and marks it unusable. It is
// called when the writer leaves the registry, either because the // called when the writer leaves the registry, because its target
// webhook was deleted or because its last database target was // or its webhook was deleted, or at shutdown. The archive FILE is
// removed. The archive FILE is deliberately left on disk: it is // deliberately left on disk: it is long-term storage an operator
// long-term storage an operator may still want. // may still want.
func (w *archiveWriter) evict() { func (w *archiveWriter) evict() {
w.mu.Lock() w.mu.Lock()
defer w.mu.Unlock() defer w.mu.Unlock()
+93 -83
View File
@@ -17,85 +17,109 @@ import (
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
) )
// evictTestEngine builds an engine backed by a temporary data // deliverTo archives one event to a database target, leaving the
// directory and returns it along with that directory. // target's writer cached with its handle open.
func evictTestEngine(t *testing.T) (*delivery.Engine, string) { func deliverTo(
t *testing.T, env *archiveEnv, tgt *database.Target,
) {
t.Helper() 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) webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`) event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
eng.ExportDeliverDatabase(webhookDB, d) env.eng.ExportDeliverDatabase(
webhookDB, seedDatabaseTargetDelivery(t, webhookDB, event, tgt),
webhookID := event.WebhookID
require.True(
t, eng.ExportHasArchiveWriter(webhookID),
"a delivery should have cached an archive writer",
) )
}
// 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()
env := setupArchiveTest(t)
first := env.seedDatabaseTarget(t, "")
second := env.addDatabaseTarget(t, first.WebhookID, "")
other := env.seedDatabaseTarget(t, "")
for _, tgt := range []*database.Target{first, second, other} {
deliverTo(t, env, tgt)
require.True( require.True(
t, eng.ExportArchiveHandleOpen(webhookID), t, env.eng.ExportArchiveHandleOpen(tgt.ID),
"the writer should hold an open handle after a write", "the writer should hold an open handle after a write",
) )
}
eng.EvictWebhook(webhookID) env.eng.EvictWebhook(first.WebhookID)
for _, tgt := range []*database.Target{first, second} {
assert.False( assert.False(
t, eng.ExportHasArchiveWriter(webhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"eviction should remove the registry entry", "eviction should remove the registry entry",
) )
assert.False( assert.False(
t, eng.ExportArchiveHandleOpen(webhookID), t, env.eng.ExportArchiveHandleOpen(tgt.ID),
"eviction should close the archive handle", "eviction should close the archive handle",
) )
archivePath := filepath.Join(
dataDir, fmt.Sprintf("archive-%s.db", webhookID),
)
assert.FileExists( assert.FileExists(
t, archivePath, t, env.archivePath(tgt),
"eviction must not delete the archive file", "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, 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 // TestEvictWebhook_UnknownWebhookIsNoOp proves eviction is safe
// for the common case of a webhook that never had a database // for the common case of a webhook or target that never had an
// target, and that repeating it does not panic. // archive writer, and that repeating it does not panic.
func TestEvictWebhook_UnknownWebhookIsNoOp(t *testing.T) { func TestEvictWebhook_UnknownWebhookIsNoOp(t *testing.T) {
t.Parallel() t.Parallel()
eng, _ := evictTestEngine(t) env := setupArchiveTest(t)
assert.NotPanics(t, func() { assert.NotPanics(t, func() {
eng.EvictWebhook("no-such-webhook") env.eng.EvictWebhook("no-such-webhook")
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( assert.False(
t, eng.ExportHasArchiveWriter("no-such-webhook"), t, env.eng.ExportHasArchiveWriter("no-such-target"),
"eviction must not create a writer", "eviction must not create a writer",
) )
} }
@@ -289,17 +313,14 @@ func TestEvictWebhook_RacingWriteDoesNotReopenHandle(
) { ) {
t.Parallel() 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, "")
// Prime the registry so the test can hold the very writer the // Prime the registry so the test can hold the very writer the
// eviction is about to detach. // 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.NotNil(t, w)
require.True(t, w.HandleOpen()) require.True(t, w.HandleOpen())
@@ -309,7 +330,7 @@ func TestEvictWebhook_RacingWriteDoesNotReopenHandle(
// eviction has to contend for the writer's mutex. // eviction has to contend for the writer's mutex.
race.awaitFirstWrite() race.awaitFirstWrite()
eng.EvictWebhook(event.WebhookID) env.eng.EvictWebhook(tgt.WebhookID)
sawEvicted, otherErr := race.wait() sawEvicted, otherErr := race.wait()
@@ -324,41 +345,33 @@ func TestEvictWebhook_RacingWriteDoesNotReopenHandle(
"been evicted", "been evicted",
) )
assert.False( assert.False(
t, eng.ExportHasArchiveWriter(event.WebhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the registry entry must stay gone", "the registry entry must stay gone",
) )
} }
// TestEvictWebhook_LaterDeliveryRecreatesWriter proves eviction // 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. // subsequent delivery gets a brand new writer from the registry.
// It says nothing about the evicted writer itself — that is what // It says nothing about the evicted writer itself — that is what
// TestEvictedWriter_WriteDoesNotReopenFile covers. // TestEvictedWriter_WriteDoesNotReopenFile covers.
func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) { func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) {
t.Parallel() t.Parallel()
eng, _ := evictTestEngine(t) env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, "")
webhookDB := testWebhookDB(t) deliverTo(t, env, tgt)
event := seedEvent(t, webhookDB, `{"archived":true}`) require.True(t, env.eng.ExportHasArchiveWriter(tgt.ID))
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
eng.ExportDeliverDatabase(webhookDB, d) env.eng.EvictWebhook(tgt.WebhookID)
require.True(
t, eng.ExportHasArchiveWriter(event.WebhookID),
)
eng.EvictWebhook(event.WebhookID) // A fresh delivery for the same target gets a brand new
// A fresh delivery for the same webhook gets a brand new
// writer from the registry, so archiving keeps working. // writer from the registry, so archiving keeps working.
second := seedDatabaseTargetDelivery( deliverTo(t, env, tgt)
t, webhookDB, event, "",
)
eng.ExportDeliverDatabase(webhookDB, second)
assert.True( assert.True(
t, eng.ExportHasArchiveWriter(event.WebhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"a later delivery should recreate the writer", "a later delivery should recreate the writer",
) )
} }
@@ -370,19 +383,16 @@ func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) {
func TestEngineStop_WriteAfterStopIsRefused(t *testing.T) { func TestEngineStop_WriteAfterStopIsRefused(t *testing.T) {
t.Parallel() t.Parallel()
eng, _ := evictTestEngine(t) env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, "")
webhookDB := testWebhookDB(t) deliverTo(t, env, tgt)
event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
eng.ExportDeliverDatabase(webhookDB, d) w := env.eng.ExportArchiveWriterFor(tgt.ID)
w := eng.ExportArchiveWriterFor(event.WebhookID)
require.NotNil(t, w) require.NotNil(t, w)
require.True(t, w.HandleOpen()) 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) 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", "a refused write must not reopen the archive",
) )
assert.False( assert.False(
t, eng.ExportHasArchiveWriter(event.WebhookID), t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the stop should empty the registry", "the stop should empty the registry",
) )
+275 -40
View File
@@ -4,13 +4,12 @@ import (
"database/sql" "database/sql"
"fmt" "fmt"
"log/slog" "log/slog"
"net/http"
"os" "os"
"path/filepath" "path/filepath"
"strings"
"testing" "testing"
"time" "time"
"github.com/google/uuid"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
"gorm.io/driver/sqlite" "gorm.io/driver/sqlite"
@@ -74,25 +73,18 @@ func removeArchiveFiles(t *testing.T, path string) {
// TestDeliverDatabase_ArchivesEvent verifies that delivering to // TestDeliverDatabase_ArchivesEvent verifies that delivering to
// a database target marks the delivery delivered and archives // 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) { func TestDeliverDatabase_ArchivesEvent(t *testing.T) {
t.Parallel() t.Parallel()
dataDir := t.TempDir() env := setupArchiveTest(t)
dbMgr := database.NewTestWebhookDBManager(dataDir) tgt := env.seedDatabaseTarget(t, "")
e := delivery.NewTestEngineWithDB(
nil, dbMgr,
archiveTestLogger(),
&http.Client{Timeout: 5 * time.Second},
1,
)
webhookDB := testWebhookDB(t) webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`) 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 var updated database.Delivery
@@ -105,8 +97,7 @@ func TestDeliverDatabase_ArchivesEvent(t *testing.T) {
) )
archivePath := filepath.Join( archivePath := filepath.Join(
dataDir, env.dataDir, "archive-sweep-test-archive-"+tgt.ID+".db",
fmt.Sprintf("archive-%s.db", event.WebhookID),
) )
assert.FileExists(t, archivePath) assert.FileExists(t, archivePath)
@@ -288,31 +279,31 @@ func TestParseArchiveExpiry(t *testing.T) {
} }
} }
// seedDatabaseTargetDelivery seeds a pending delivery for a // seedDatabaseTargetDelivery seeds a pending delivery of an event
// database target with the given config JSON and returns the // to a database target and returns the in-memory delivery the
// in-memory delivery the target handler is invoked with. // target handler is invoked with.
func seedDatabaseTargetDelivery( func seedDatabaseTargetDelivery(
t *testing.T, t *testing.T,
webhookDB *gorm.DB, webhookDB *gorm.DB,
event database.Event, event database.Event,
config string, tgt *database.Target,
) *database.Delivery { ) *database.Delivery {
t.Helper() t.Helper()
dlv := seedDelivery( dlv := seedDelivery(
t, webhookDB, event.ID, uuid.New().String(), t, webhookDB, event.ID, tgt.ID,
database.DeliveryStatusPending, database.DeliveryStatusPending,
) )
d := &database.Delivery{ d := &database.Delivery{
EventID: event.ID, EventID: event.ID,
TargetID: dlv.TargetID, TargetID: tgt.ID,
Status: database.DeliveryStatusPending, Status: database.DeliveryStatusPending,
Event: event, Event: event,
Target: database.Target{ Target: database.Target{
Name: "test-db", Name: tgt.Name,
Type: database.TargetTypeDatabase, Type: database.TargetTypeDatabase,
Config: config, Config: tgt.Config,
}, },
} }
d.ID = dlv.ID d.ID = dlv.ID
@@ -330,22 +321,14 @@ func TestDeliverDatabase_ArchiveFailureFailsDelivery(
) { ) {
t.Parallel() t.Parallel()
dataDir := t.TempDir() env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, `{"expiry":"nonsense"}`)
e := delivery.NewTestEngineWithDB(
nil, database.NewTestWebhookDBManager(dataDir),
archiveTestLogger(),
&http.Client{Timeout: 5 * time.Second},
1,
)
webhookDB := testWebhookDB(t) webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":false}`) event := seedEvent(t, webhookDB, `{"archived":false}`)
d := seedDatabaseTargetDelivery( d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt)
t, webhookDB, event, `{"expiry":"nonsense"}`,
)
e.ExportDeliverDatabase(webhookDB, d) env.eng.ExportDeliverDatabase(webhookDB, d)
var updated database.Delivery var updated database.Delivery
@@ -373,10 +356,7 @@ func TestDeliverDatabase_ArchiveFailureFailsDelivery(
) )
assert.NoFileExists(t, assert.NoFileExists(t,
filepath.Join( env.archivePath(tgt),
dataDir,
fmt.Sprintf("archive-%s.db", event.WebhookID),
),
"no archive file should exist for a failed config", "no archive file should exist for a failed config",
) )
} }
@@ -400,3 +380,258 @@ 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, once
// the .db alone, once a lone -wal and once a lone -shm, and proves
// each time that 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()
for _, suffix := range archiveFileSuffixes() {
t.Run("planted .db"+suffix, func(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",
)
plantedPath := newPath + suffix
require.NoError(
t, os.WriteFile(plantedPath, []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(plantedPath)
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())
}
+3 -3
View File
@@ -63,7 +63,7 @@ type HandlersParams struct {
Session *session.Session Session *session.Session
Middleware *middleware.Middleware Middleware *middleware.Middleware
Notifier delivery.Notifier Notifier delivery.Notifier
Evictor delivery.WebhookEvictor Archives delivery.Archives
SSRFGuard *delivery.Guard SSRFGuard *delivery.Guard
Metrics *metrics.Set Metrics *metrics.Set
Registry *prometheus.Registry Registry *prometheus.Registry
@@ -80,7 +80,7 @@ type Handlers struct {
session *session.Session session *session.Session
mw *middleware.Middleware mw *middleware.Middleware
notifier delivery.Notifier notifier delivery.Notifier
evictor delivery.WebhookEvictor archives delivery.Archives
mtr *metrics.Set mtr *metrics.Set
templates map[string]*template.Template templates map[string]*template.Template
@@ -131,7 +131,7 @@ func New(
s.session = params.Session s.session = params.Session
s.mw = params.Middleware s.mw = params.Middleware
s.notifier = params.Notifier s.notifier = params.Notifier
s.evictor = params.Evictor s.archives = params.Archives
s.mtr = params.Metrics s.mtr = params.Metrics
s.ssrf = params.SSRFGuard s.ssrf = params.SSRFGuard
+82 -9
View File
@@ -3,6 +3,7 @@ package handlers_test
import ( import (
"context" "context"
"errors" "errors"
"fmt"
"html/template" "html/template"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
@@ -52,23 +53,73 @@ func (n *recordingNotifier) Tasks() []delivery.Task {
return out return out
} }
// recordingEvictor is a delivery.WebhookEvictor that records // recordingArchives is a delivery.Archives that records what it
// the webhook ids it was asked to evict, so a test can prove // was asked to do, so a test can prove that a deletion or rename
// that a deletion path reached the delivery engine. // path reached the delivery engine. After FailRenames, every
type recordingEvictor struct { // rename fails with the given error.
type recordingArchives struct {
mu sync.Mutex mu sync.Mutex
evicted []string 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() r.mu.Lock()
defer r.mu.Unlock() defer r.mu.Unlock()
r.evicted = append(r.evicted, webhookID) 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. // Evicted returns a copy of the recorded webhook ids.
func (r *recordingEvictor) Evicted() []string { func (r *recordingArchives) Evicted() []string {
r.mu.Lock() r.mu.Lock()
defer r.mu.Unlock() defer r.mu.Unlock()
@@ -78,6 +129,28 @@ func (r *recordingEvictor) Evicted() []string {
return out 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( func newTestApp(
t *testing.T, t *testing.T,
targets ...any, targets ...any,
@@ -104,10 +177,10 @@ func newTestApp(
func(n *recordingNotifier) delivery.Notifier { func(n *recordingNotifier) delivery.Notifier {
return n return n
}, },
func() *recordingEvictor { func() *recordingArchives {
return &recordingEvictor{} return &recordingArchives{}
}, },
func(r *recordingEvictor) delivery.WebhookEvictor { func(r *recordingArchives) delivery.Archives {
return r return r
}, },
metrics.NewRegistry, metrics.NewRegistry,
+59 -80
View File
@@ -15,6 +15,7 @@ import (
"gorm.io/gorm" "gorm.io/gorm"
"gorm.io/gorm/clause" "gorm.io/gorm/clause"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
"sneak.berlin/go/webhooker/internal/handlers" "sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/session" "sneak.berlin/go/webhooker/internal/session"
) )
@@ -79,6 +80,10 @@ func seedTarget(
// from a delete statement. // from a delete statement.
var errInjectedDelete = errors.New("injected delete failure") var errInjectedDelete = errors.New("injected delete failure")
// errInjectedSave is the failure failSaveOnTable reports from a
// save of an existing row.
var errInjectedSave = errors.New("injected save failure")
// seedEntrypoint inserts an entrypoint for a webhook. // seedEntrypoint inserts an entrypoint for a webhook.
func seedEntrypoint( func seedEntrypoint(
t *testing.T, t *testing.T,
@@ -146,19 +151,42 @@ func failDeleteOnTable(
) )
} }
// failSaveOnTable is failDeleteOnTable for saves: every update of
// an existing row in the named table fails.
func failSaveOnTable(
t *testing.T,
db *database.Database,
table string,
) {
t.Helper()
require.NoError(t, db.DB().Callback().Update().
Before("gorm:update").
Register(
"test:fail_save_"+table,
func(tx *gorm.DB) {
if tx.Statement.Table == table {
_ = tx.AddError(errInjectedSave)
}
},
),
)
}
// archivePathFor returns the archive database path the // archivePathFor returns the archive database path the
// delivery engine would use for a webhook: beside the webhook's // delivery engine would use for a database target: beside the
// event database in the data directory. // webhook's event database in the data directory.
func archivePathFor( func archivePathFor(
t *testing.T, t *testing.T,
mgr *database.WebhookDBManager, mgr *database.WebhookDBManager,
webhookID string, wh *database.Webhook,
tgt *database.Target,
) string { ) string {
t.Helper() t.Helper()
return filepath.Join( return filepath.Join(
filepath.Dir(mgr.DBPath(webhookID)), filepath.Dir(mgr.DBPath(wh.ID)),
"archive-"+webhookID+".db", delivery.ArchiveFileName(wh.Name, tgt.Name, tgt.ID),
) )
} }
@@ -195,8 +223,8 @@ func postRequest(
// TestHandleSourceDelete_EvictsArchiveWriter proves that // TestHandleSourceDelete_EvictsArchiveWriter proves that
// deleting a webhook reaches the delivery engine and releases // deleting a webhook reaches the delivery engine and releases
// the webhook's archive writer, exercised through the real // the webhook's archive writers, exercised through the real
// deletion handler rather than by calling the evictor directly. // deletion handler rather than by calling the engine directly.
func TestHandleSourceDelete_EvictsArchiveWriter(t *testing.T) { func TestHandleSourceDelete_EvictsArchiveWriter(t *testing.T) {
t.Parallel() t.Parallel()
@@ -204,7 +232,7 @@ func TestHandleSourceDelete_EvictsArchiveWriter(t *testing.T) {
h *handlers.Handlers h *handlers.Handlers
sess *session.Session sess *session.Session
db *database.Database db *database.Database
ev *recordingEvictor ev *recordingArchives
) )
app := newTestApp(t, &h, &sess, &db, &ev) app := newTestApp(t, &h, &sess, &db, &ev)
@@ -254,9 +282,10 @@ func TestHandleSourceDelete_KeepsArchiveFile(t *testing.T) {
t.Cleanup(app.RequireStop) t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db) wh := seedWebhook(t, db)
tgt := seedTarget(t, db, wh.ID, database.TargetTypeDatabase)
// Place an archive file where the delivery engine would. // Place an archive file where the delivery engine would.
archivePath := archivePathFor(t, mgr, wh.ID) archivePath := archivePathFor(t, mgr, wh, tgt)
require.NoError( require.NoError(
t, t,
writeArchivePlaceholder(archivePath), writeArchivePlaceholder(archivePath),
@@ -435,68 +464,17 @@ func TestHandleSourceDelete_RemovesConfigAndEventDatabase(
) )
} }
// TestHandleTargetDelete_EvictsWhenLastDatabaseTargetGone // TestHandleTargetDelete_EvictsThatTarget proves that deleting a
// proves that removing the last database target releases the // database target releases that target's archive writer and no
// archive writer. // other: the webhook's other database target keeps its own.
func TestHandleTargetDelete_EvictsWhenLastDatabaseTargetGone( func TestHandleTargetDelete_EvictsThatTarget(t *testing.T) {
t *testing.T,
) {
t.Parallel() t.Parallel()
var ( var (
h *handlers.Handlers h *handlers.Handlers
sess *session.Session sess *session.Session
db *database.Database db *database.Database
ev *recordingEvictor ev *recordingArchives
)
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
) )
app := newTestApp(t, &h, &sess, &db, &ev) app := newTestApp(t, &h, &sess, &db, &ev)
@@ -527,17 +505,17 @@ func TestHandleTargetDelete_KeepsWriterWhenDatabaseTargetRemains(
h.HandleTargetDelete().ServeHTTP(w, req) h.HandleTargetDelete().ServeHTTP(w, req)
require.Equal(t, http.StatusSeeOther, w.Code) require.Equal(t, http.StatusSeeOther, w.Code)
assert.Empty( assert.Equal(
t, ev.Evicted(), t, []string{doomed.ID}, ev.EvictedTargets(),
"a second database target still needs the writer", "deleting a database target should evict its writer",
) )
assert.Empty(t, ev.Evicted(), "the webhook is not deleted")
} }
// TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted proves // TestHandleTargetDelete_IgnoresAnotherWebhooksTarget proves that
// that deleting a target of an unrelated type leaves a // a target id from the URL that is not a target of the webhook
// still-needed archive writer alone: the webhook's database // deletes nothing and so evicts nothing.
// target is untouched, so its writer must stay. func TestHandleTargetDelete_IgnoresAnotherWebhooksTarget(
func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted(
t *testing.T, t *testing.T,
) { ) {
t.Parallel() t.Parallel()
@@ -546,7 +524,7 @@ func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted(
h *handlers.Handlers h *handlers.Handlers
sess *session.Session sess *session.Session
db *database.Database db *database.Database
ev *recordingEvictor ev *recordingArchives
) )
app := newTestApp(t, &h, &sess, &db, &ev) app := newTestApp(t, &h, &sess, &db, &ev)
@@ -555,19 +533,20 @@ func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted(
t.Cleanup(app.RequireStop) t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db) wh := seedWebhook(t, db)
seedTarget(t, db, wh.ID, database.TargetTypeDatabase) elsewhere := seedTarget(
other := seedTarget(t, db, wh.ID, database.TargetTypeLog) t, db, seedWebhook(t, db).ID, database.TargetTypeDatabase,
)
cookies := authenticatedCookies( cookies := authenticatedCookies(
t, sess, deleteTestUserID, deleteTestUsername, t, sess, deleteTestUserID, deleteTestUsername,
) )
req := postRequest( req := postRequest(
"/hook/"+wh.ID+"/targets/"+other.ID+"/delete", "/hook/"+wh.ID+"/targets/"+elsewhere.ID+"/delete",
cookies, cookies,
map[string]string{ map[string]string{
paramSourceID: wh.ID, paramSourceID: wh.ID,
paramTargetID: other.ID, paramTargetID: elsewhere.ID,
}, },
) )
w := httptest.NewRecorder() w := httptest.NewRecorder()
@@ -576,7 +555,7 @@ func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted(
require.Equal(t, http.StatusSeeOther, w.Code) require.Equal(t, http.StatusSeeOther, w.Code)
assert.Empty( assert.Empty(
t, ev.Evicted(), t, ev.EvictedTargets(),
"a surviving database target must keep its writer", "another webhook's target must not be evicted",
) )
} }
+86 -44
View File
@@ -549,6 +549,7 @@ func (h *Handlers) applyWebhookEdit(
return return
} }
oldName := webhook.Name
webhook.Name = name webhook.Name = name
webhook.Description = r.PostFormValue("description") webhook.Description = r.PostFormValue("description")
@@ -571,8 +572,40 @@ func (h *Handlers) applyWebhookEdit(
webhook.RetentionDays = retentionDays 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 { 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) h.serverError(w, r, "failed to update webhook", err)
return return
@@ -707,11 +740,11 @@ func (h *Handlers) commitWebhookDeletion(
return tx.Commit().Error return tx.Commit().Error
} }
// evictArchiveWriter asks the delivery engine to drop its // evictArchiveWriter asks the delivery engine to drop the cached
// cached archive writer for a webhook, closing the archive file // archive writers of a webhook's database targets, closing their
// handle. // 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 // database — which is per-webhook working storage and is
// hard-deleted with the webhook — an archive is explicitly // hard-deleted with the webhook — an archive is explicitly
// long-term storage that an operator may want to keep or move // 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 // deleting a webhook would be a surprising and unrecoverable
// data loss, so the file is left for the operator to handle. // data loss, so the file is left for the operator to handle.
func (h *Handlers) evictArchiveWriter(webhookID string) { func (h *Handlers) evictArchiveWriter(webhookID string) {
if h.evictor == nil { if h.archives == nil {
return return
} }
h.evictor.EvictWebhook(webhookID) h.archives.EvictWebhook(webhookID)
} }
// evictArchiveWriterIfUnused releases a webhook's archive // evictTargetArchiveWriter is evictArchiveWriter for one deleted
// writer once the webhook has no database target left to feed // target, and leaves its archive file on disk for the same reason.
// it. // A target that is not a database target has no writer, and
// // evicting it does nothing.
// It is called after any child resource of a webhook is func (h *Handlers) evictTargetArchiveWriter(targetID string) {
// deleted, and is correct without knowing which kind was: it if h.archives == nil {
// evicts only when no database target remains, so deleting one return
// 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 h.archives.EvictTarget(targetID)
// 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) { // renameWebhookArchives renames the archive file of every database
var remaining int64 // 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(). err := h.db.DB().
Model(&database.Target{}).
Where( Where(
"webhook_id = ? AND type = ?", "webhook_id = ? AND type = ?",
webhookID, database.TargetTypeDatabase, webhookID, database.TargetTypeDatabase,
). ).
Count(&remaining).Error Find(&targets).Error
if err != nil { if err != nil {
h.log.Error( return err
"failed to count remaining database targets", }
"webhook_id", webhookID,
"error", err, for i := range targets {
err = h.archives.Rename(
targets[i].ID, newName, targets[i].Name,
) )
if err != nil {
return return err
}
} }
if remaining > 0 { return nil
return
}
h.evictArchiveWriter(webhookID)
} }
// ownedWebhook resolves the request's sourceID parameter to a // 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 // HandleTargetDelete handles deleting a target. A deleted
// last database target of a webhook leaves its archive writer // database target's archive writer is evicted and its handle
// with nothing to write, so the writer is evicted and its // closed; the archive file is left on disk.
// handle closed; the archive file is left on disk.
func (h *Handlers) HandleTargetDelete() http.HandlerFunc { func (h *Handlers) HandleTargetDelete() http.HandlerFunc {
return h.deleteChildResource( return h.deleteChildResource(
"targetID", &database.Target{}, "targetID", &database.Target{},
"failed to delete target", "failed to delete target",
h.evictArchiveWriterIfUnused, h.evictTargetArchiveWriter,
) )
} }
// deleteChildResource returns a handler that deletes a child // deleteChildResource returns a handler that deletes a child
// resource (entrypoint or target) belonging to a webhook. The // resource (entrypoint or target) belonging to a webhook. The
// optional afterDelete hook runs with the webhook's id once the // optional afterDelete hook runs with the child's id once the
// delete has succeeded, before the redirect. // delete has removed it, before the redirect.
func (h *Handlers) deleteChildResource( func (h *Handlers) deleteChildResource(
idParam string, idParam string,
model any, model any,
errMsg string, errMsg string,
afterDelete func(webhookID string), afterDelete func(childID string),
) http.HandlerFunc { ) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) { return func(w http.ResponseWriter, r *http.Request) {
userID, ok := h.getUserID(r) userID, ok := h.getUserID(r)
@@ -1702,8 +1742,10 @@ func (h *Handlers) deleteChildResource(
return return
} }
if afterDelete != nil { // Only for a row this webhook really had: the id came from
afterDelete(webhook.ID) // the URL and may name another webhook's child.
if afterDelete != nil && result.RowsAffected > 0 {
afterDelete(childID)
} }
http.Redirect( http.Redirect(
+141 -1
View File
@@ -187,6 +187,7 @@ func storedRetentionDays(
type sourceTestEnv struct { type sourceTestEnv struct {
handlers *handlers.Handlers handlers *handlers.Handlers
db *database.Database db *database.Database
archives *recordingArchives
cookies []*http.Cookie cookies []*http.Cookie
} }
@@ -199,7 +200,9 @@ func setupSourceTest(t *testing.T) *sourceTestEnv {
var db *database.Database var db *database.Database
app := newTestApp(t, &h, &sess, &db) var archives *recordingArchives
app := newTestApp(t, &h, &sess, &db, &archives)
app.RequireStart() app.RequireStart()
t.Cleanup(app.RequireStop) t.Cleanup(app.RequireStop)
@@ -207,6 +210,7 @@ func setupSourceTest(t *testing.T) *sourceTestEnv {
return &sourceTestEnv{ return &sourceTestEnv{
handlers: h, handlers: h,
db: db, db: db,
archives: archives,
cookies: authenticatedCookies( cookies: authenticatedCookies(
t, sess, sourceTestUserID, "sourceuser", t, sess, sourceTestUserID, "sourceuser",
), ),
@@ -498,6 +502,142 @@ func TestHandleSourceEditSubmit_EmptyRetentionLeavesValueUnchanged(
assert.Equal(t, 7, storedRetentionDays(t, env.db, wh.ID)) 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_FailedSaveRenamesBack proves that when
// the archive is renamed but the new name cannot be saved, the
// archive is renamed back to the stored name and the stored name
// stays.
func TestHandleSourceEditSubmit_FailedSaveRenamesBack(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
tgt := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
failSaveOnTable(t, env.db, "webhooks")
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 // TestSourceEditForm_ForeverWebhookRoundTrips walks the exact path that
// the removed max="365" cap used to break: render the edit form for a // 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 // retain-forever webhook, confirm the pre-filled sentinel is not capped
+47
View File
@@ -1,6 +1,7 @@
package handlers package handlers
import ( import (
"errors"
"net/http" "net/http"
"github.com/go-chi/chi" "github.com/go-chi/chi"
@@ -150,11 +151,42 @@ func (h *Handlers) applyTargetEdit(
target.MaxRetries = retries target.MaxRetries = retries
} }
oldName := target.Name
target.Name = name target.Name = name
target.Config = configJSON target.Config = configJSON
// 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 err = h.db.DB().Save(target).Error
}
if err != nil { 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) h.serverError(w, r, "failed to update target", err)
return 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 // renderTargetEdit renders the target edit page with an optional
// error message. // error message.
func (h *Handlers) renderTargetEdit( func (h *Handlers) renderTargetEdit(
+97
View File
@@ -635,3 +635,100 @@ func assertWebhookOfAnotherUser404s(
assert.Equal(t, http.StatusNotFound, w.Code) 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,
)
}
// TestHandleTargetEditSubmit_FailedSaveRenamesBack proves that when a
// database target's archive is renamed but the new name cannot be
// saved, the archive is renamed back to the stored name and the
// stored name stays.
func TestHandleTargetEditSubmit_FailedSaveRenamesBack(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
failSaveOnTable(t, env.db, "targets")
form := url.Values{}
form.Set("name", renamedTargetName)
w := submitTargetEdit(env, wh.ID, archive.ID, form)
require.Equal(t, http.StatusInternalServerError, w.Code)
assert.Equal(
t, archive.Name, storedTarget(t, env, archive.ID).Name,
)
assert.Equal(
t,
[]archiveRename{
{archive.ID, wh.Name, renamedTargetName},
{archive.ID, wh.Name, archive.Name},
},
env.archives.Renames(),
)
}
+9 -3
View File
@@ -132,9 +132,15 @@ type noopNotifier struct{}
func (n *noopNotifier) Notify([]delivery.Task) {} 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, // newServerApp starts the real login path against dir: the handlers,
// the middleware that bounds password verification, the session store // the middleware that bounds password verification, the session store
@@ -163,7 +169,7 @@ func newServerApp(
healthcheck.New, healthcheck.New,
session.New, session.New,
func() delivery.Notifier { return &noopNotifier{} }, func() delivery.Notifier { return &noopNotifier{} },
func() delivery.WebhookEvictor { return &noopEvictor{} }, func() delivery.Archives { return &noopArchives{} },
metrics.NewRegistry, metrics.NewRegistry,
metrics.New, metrics.New,
middleware.New, middleware.New,
+12 -6
View File
@@ -47,12 +47,18 @@ type noopNotifier struct{}
func (n *noopNotifier) Notify([]delivery.Task) {} func (n *noopNotifier) Notify([]delivery.Task) {}
// noopEvictor satisfies handlers.New's delivery.WebhookEvictor // noopArchives satisfies handlers.New's delivery.Archives
// dependency. No test here checks what gets evicted, so it records // dependency. No test here checks what gets evicted or renamed, so
// nothing. // it records nothing.
type noopEvictor struct{} 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 // testEnv is the real router from routes.go plus the collaborators
// tests need to seed users and forge sessions. // tests need to seed users and forge sessions.
@@ -113,7 +119,7 @@ func newTestEnvWithConfig(
healthcheck.New, healthcheck.New,
session.New, session.New,
func() delivery.Notifier { return &noopNotifier{} }, func() delivery.Notifier { return &noopNotifier{} },
func() delivery.WebhookEvictor { return &noopEvictor{} }, func() delivery.Archives { return &noopArchives{} },
metrics.NewRegistry, metrics.NewRegistry,
metrics.New, metrics.New,
middleware.New, middleware.New,