From 13f881007a81dbffd6b0bad7e06e214f4161c196 Mon Sep 17 00:00:00 2001 From: sneak Date: Fri, 2 Oct 2026 23:41:42 +0000 Subject: [PATCH] Rotate a database target's archive monthly, daily or hourly (closes #379) A database target's rotation (none, monthly, daily or hourly) puts the UTC period of each event's receive time in its archive file name, so each file holds exactly its period's events. It is on the new webhook page, the add target form and the target edit form, and shown in the target list. Renames move every one of a target's files, the sweep prunes every file and deletes a rotated file it leaves empty, Download exports every file oldest first with each row's period, and the target list names the current file and totals the size of all of them. Model: opus-5-5 --- README.md | 119 ++++-- internal/delivery/archive_sweeper_test.go | 15 +- internal/delivery/engine.go | 13 +- internal/delivery/export_test.go | 19 +- internal/delivery/target_config_edit.go | 27 +- internal/delivery/target_config_view.go | 13 +- internal/delivery/target_config_view_test.go | 42 +- internal/delivery/target_database.go | 37 +- internal/delivery/target_database_archive.go | 279 +++++++++---- internal/delivery/target_database_export.go | 178 +++++--- .../delivery/target_database_export_test.go | 50 +++ internal/delivery/target_database_rotation.go | 198 +++++++++ .../delivery/target_database_rotation_test.go | 389 ++++++++++++++++++ internal/delivery/target_database_test.go | 1 + internal/delivery/target_headers_test.go | 7 +- internal/handlers/archive_expiry.go | 18 +- internal/handlers/archive_expiry_test.go | 25 +- internal/handlers/archive_rotation.go | 42 ++ internal/handlers/archive_rotation_test.go | 229 +++++++++++ internal/handlers/export_test.go | 13 +- internal/handlers/handlers_test.go | 41 +- .../handlers/source_create_targets_test.go | 4 +- internal/handlers/source_management.go | 79 ++-- internal/handlers/target_download.go | 7 +- internal/handlers/target_edit.go | 4 + internal/handlers/target_edit_test.go | 9 +- internal/handlers/target_list.go | 87 ++-- internal/server/alpine_browser_test.go | 64 ++- static/css/tailwind.css | 2 +- static/js/app.js | 7 +- templates/source_detail.html | 38 +- templates/sources_new.html | 15 +- templates/target_edit.html | 10 + 33 files changed, 1753 insertions(+), 328 deletions(-) create mode 100644 internal/delivery/target_database_rotation.go create mode 100644 internal/delivery/target_database_rotation_test.go create mode 100644 internal/handlers/archive_rotation.go create mode 100644 internal/handlers/archive_rotation_test.go diff --git a/README.md b/README.md index 588a927..2c9491b 100644 --- a/README.md +++ b/README.md @@ -736,7 +736,8 @@ The app runs as a non-root user (`webhooker`, UID 1000), exposes port The `/var/lib/webhooker` volume holds all SQLite databases: the main application database (`webhooker.db`), the per-webhook event databases (`events-{uuid}.db`), and any archive databases written by `database` -targets (`archive-{webhook_name}-{target_name}-{target_uuid}.db`). Mount +targets (`archive-{webhook_name}-{target_name}-{target_uuid}.db`, with +`-{period}` before `.db` for a target that rotates). Mount this as a persistent volume to preserve data across container restarts. @@ -976,9 +977,11 @@ is both the simplest and the only complete rule: - `events-{webhook_uuid}.db` — **one per webhook**. Events, deliveries, delivery results. - `archive-{webhook_name}-{target_name}-{target_uuid}.db` — **one per - `database` target**. Archived events. The two names are made safe for - a file name, and the file is renamed when the webhook or the target is - (see [Database Architecture](#database-architecture)). + `database` target**, or one per month, day or hour for a target that + rotates, with `-{period}` before `.db`. Archived events. The two names + are made safe for a file name, and the files are renamed when the + webhook or the target is (see + [Database Architecture](#database-architecture)). `{webhook_uuid}` and `{target_uuid}` are UUID primary keys in their canonical 36-character hyphenated form, so a real filename looks like @@ -1610,8 +1613,9 @@ The new webhook form can also give the webhook its first targets: an optional HTTP target URL creates an `http` target named `HTTP`, and the archive checkbox creates a `database` target named `Archive` whose `expiry` is the pruning chosen beside it (never, 1h, 12h, 24h, 30d, 90d -or 365d). Both are validated as on the add target form, and the webhook -and its targets are created together or not at all. +or 365d) and whose `rotation` is the rotation chosen below that (none, +monthly, daily or hourly). Both are validated as on the add target form, +and the webhook and its targets are created together or not at all. | Field | Type | Description | | ---------------- | ------- | ----------- | @@ -1720,12 +1724,14 @@ events should be forwarded. own archive database (`archive-{webhook_name}-{target_name}-{target_uuid}.db`) for long-term retention, with an optional creation-validated expiry (default: keep - forever). The new webhook form, the add target form and the target edit - form all offer the same expiries: never, 1h, 12h, 24h, 30d, 90d or 365d. - The target list shows the expiry in plain units, such as "30 days". No - external delivery and no retries; an archive write failure fails the - delivery. See the database target section under - "Per-Webhook Event Databases" for the full semantics. + forever) and rotation (default: none, one file). The new webhook form, + the add target form and the target edit form all offer the same + expiries: never, 1h, 12h, 24h, 30d, 90d or 365d, and the same + rotations: none, monthly, daily or hourly. The target list shows the + expiry in plain units, such as "30 days", and the rotation. No external + delivery and no retries; an archive write failure fails the delivery. + See the database target section under "Per-Webhook Event Databases" + for the full semantics. - **`log`** — Write the event to the application log (stdout). Useful for debugging. @@ -2088,15 +2094,31 @@ 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. Webhook edits, -target edits and target creation run one at a time, so no edit can -rename the file between another's rename and save, and 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. + +An optional `rotation` in the target's config JSON (e.g. +`{"rotation":"daily"}`) is `none`, the default, which keeps the one +file, or `monthly`, `daily` or `hourly`. A target that rotates writes +each event to a file named for the period of the event's receive time, +in UTC, put before the `.db`: +`archive-orders-eu-long-term-archive-{target_uuid}-2026-10.db` monthly, +`…-2026-10-01.db` daily and `…-2026-10-01-19.db` hourly. Each file holds +exactly its period's events, and the first event of a new period starts +the next file, so a finished period's file can be moved away like any +archive. A changed rotation applies from the next event: the files +already written keep their names and stay, pruned, shown and downloaded +with the rest, since every file named for the target is its archive, +whichever rotation wrote it. + +Renaming the webhook or the target renames every one of the target's +files, each keeping its period, under the same lock the archive writes +and the archive sweeper take. Webhook edits, target edits and target +creation run one at a time, so no edit can rename the files between +another's rename and save, and the names on disk match the UI. A rename +never replaces a file: if one already has a new name, the edit is +refused with an error naming that file, nothing is moved, 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 @@ -2135,18 +2157,27 @@ interleave with a write, and it leaves the archive closed afterwards so the move-the-file-away workflow keeps working. Archives with no expiry, or the expiry `never`, are not touched by the sweep at all. +For a target with several files, the write path prunes only the file it +writes to, and the sweep prunes every one of them. A file named for a +period that the sweep leaves empty is deleted, with any `-wal` and +`-shm` beside it; the file without a period is kept even when empty, as +it always has been. + Because each `database` target has its own archive file, a target's `expiry` governs only its own archive. Two `database` targets on one webhook with different expiries keep two archives, each pruned on its own schedule. -The webhook page shows, for each `database` target, its archive file's -name, its size on disk and when it was last written. The size counts -the `.db` and its `-wal` together, and the last write is the later of -their two modification times, since a write lands in the `-wal` first. -Both are read from the files' metadata; the archive is never opened. -Before the first write, and after the file has been moved away, the page -shows `not created yet` beside the name. +The webhook page shows, for each `database` target, the name of the +archive file an event received now would go to, the size on disk of all +the target's archive files together, how many there are when there is +more than one, and when the latest of them was last written. The size +counts each `.db` and its `-wal` together, and the last write is the +latest of their modification times, since a write lands in the `-wal` +first. All are read from the files' metadata; the archive is never +opened. While the named file does not exist — before the first write, +before the first event of a new period, and after the file has been +moved away — the page shows `not created yet` beside the name. Each `database` target on the webhook page has a **Download** button, which returns its archive as one gzipped JSON file, @@ -2154,22 +2185,27 @@ which returns its archive as one gzipped JSON file, names made safe as above and the time in UTC. The file holds one object: `webhook` and `target`, each an `id` and a `name`; `exported_at`; and `archived_events`, one object per archived row with -every column, keyed by column name. A body that is not valid UTF-8 is -written in base64, with `"body_encoding": "base64"` beside it. An -archive that does not exist yet, or was moved away, downloads with an -empty `archived_events`; the download never creates the file. +every column, keyed by column name. The rows come from every one of the +target's files: the file without a period first, then the others in +the order of their periods, oldest first, and each row from a file named +for a period has that `period` beside its columns. A body that is not +valid UTF-8 is written in base64, with `"body_encoding": "base64"` +beside it. An archive that does not exist yet, or was moved away, +downloads with an empty `archived_events`; the download never creates a +file. The download streams: each row is read and written out compressed before the next is read, so neither the archive nor the JSON is held in -memory. It reads on a connection of its own, inside one read-only -transaction, so the file holds the archive as it stood when the -download started, and archive writes go on meanwhile, since under WAL a -reader never blocks a writer. While it runs, the `-wal` cannot be -checkpointed past what it reads, so a long download lets the `-wal` -grow. It finds the file by the stored names under the lock that webhook -edits, target edits and target creation hold, and lets go once the file -is open: a rename during the download moves the file without affecting -it. +memory. When it starts it opens every one of the target's files, each on +a connection of its own inside one read-only transaction, so the +download holds the archive as it stood then, and archive writes go on +meanwhile, since under WAL a reader never blocks a writer. Each file is +closed once its rows are written out; until then its `-wal` cannot be +checkpointed past what the download reads, so a long download lets the +`-wal` grow. It finds the files by the stored names under the lock that +webhook edits, target edits and target creation hold, and lets go once +the files are open: a rename during the download moves the files without +affecting it. Deleting a webhook releases its archives: the delivery engine's cached archive writers are dropped and their file handles closed, so nothing @@ -3121,6 +3157,7 @@ webhooker/ │ │ ├── target_slack.go # Slack/Mattermost incoming-webhook target │ │ ├── target_database.go # Database archive target │ │ ├── target_database_archive.go # Archive file lifecycle and pruning +│ │ ├── target_database_rotation.go # Archive rotation and file names │ │ ├── target_database_export.go # Archive download as gzipped JSON │ │ ├── target_log.go # Log target (stdout) │ │ ├── target_config_view.go # Masked target config for templates diff --git a/internal/delivery/archive_sweeper_test.go b/internal/delivery/archive_sweeper_test.go index c596692..5563349 100644 --- a/internal/delivery/archive_sweeper_test.go +++ b/internal/delivery/archive_sweeper_test.go @@ -163,6 +163,17 @@ func (env *archiveEnv) seedArchiveRows( t.Helper() path := env.archivePath(tgt) + seedArchiveFile(t, path, tgt.WebhookID, archivedAt...) + + return path +} + +// seedArchiveFile creates the archive file at path and inserts one row +// per supplied archived-at timestamp, as seedArchiveRows does. +func seedArchiveFile( + t *testing.T, path, webhookID string, archivedAt ...time.Time, +) { + t.Helper() sqlDB, err := sql.Open( "sqlite", fmt.Sprintf("file:%s?mode=rwc", path), @@ -182,7 +193,7 @@ func (env *archiveEnv) seedArchiveRows( for i, at := range archivedAt { row := delivery.ExportArchivedEvent{ EventID: fmt.Sprintf("ev-%d", i), - WebhookID: tgt.WebhookID, + WebhookID: webhookID, Method: http.MethodPost, Body: `{"seeded":true}`, ArchivedAt: at, @@ -191,8 +202,6 @@ func (env *archiveEnv) seedArchiveRows( } require.NoError(t, sqlDB.Close()) - - return path } // archivedEventIDs returns the event ids currently stored in an diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index 9fe588f..44c47e4 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -284,12 +284,13 @@ func (e *Engine) EvictTarget(targetID string) { 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. +// Rename implements Archives. It renames every one of a database +// target's archive files to ArchiveFileName(webhookName, targetName, +// targetID), each keeping the period in its name, under the lock the +// target's archive writes and the idle sweep take. It never replaces +// a file: if one already has a 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 { diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index 40585aa..044cf3d 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -493,23 +493,32 @@ func NewExportArchiveWriter( return &ExportArchiveWriter{w: w} } -// Write archives a row through the writer. +// Write archives a row through the writer, into the file named +// without a period. func (e *ExportArchiveWriter) Write( row ExportArchivedEvent, expiry time.Duration, ) error { - return e.w.write(row, expiry) + return e.w.write(row, expiry, "") +} + +// WritePeriod archives a row through the writer, into the file for +// period. +func (e *ExportArchiveWriter) WritePeriod( + row ExportArchivedEvent, expiry time.Duration, period string, +) error { + return e.w.write(row, expiry, period) } // Open opens the archive file, pruning when expiry is positive. func (e *ExportArchiveWriter) Open(expiry time.Duration) error { - return e.w.open(expiry) + return e.w.open(e.w.path, expiry) } // Reopen closes and reopens the archive file. func (e *ExportArchiveWriter) Reopen( expiry time.Duration, ) error { - return e.w.reopen(expiry) + return e.w.reopen(e.w.path, expiry) } // SetNow replaces the clock the writer measures its reopen @@ -539,7 +548,7 @@ func (e *ExportArchiveWriter) Path() string { func (e *ExportArchiveWriter) OpenExisting( expiry time.Duration, ) error { - return e.w.openMode(archiveModeExisting, expiry) + return e.w.openMode(e.w.path, archiveModeExisting, expiry) } // SweepExpired runs an idle sweep of the archive. diff --git a/internal/delivery/target_config_edit.go b/internal/delivery/target_config_edit.go index 3753e3c..a62429b 100644 --- a/internal/delivery/target_config_edit.go +++ b/internal/delivery/target_config_edit.go @@ -41,6 +41,8 @@ type TargetConfigForm struct { Timeout string // Expiry is the database (archive) target's row expiry. Expiry string + // Rotation is the database (archive) target's rotation. + Rotation string } // NewTargetConfigForm parses a target's stored configuration into @@ -85,11 +87,13 @@ func NewTargetConfigForm( } } -// databaseConfigForm parses an archive target's optional expiry. -// An absent, empty or never expiry yields an empty expiry, on which -// the edit form starts at never; saving it unchanged stores never, -// which means the same as an empty expiry. An expiry that is set -// but not a valid duration is an error, not a blank field. +// databaseConfigForm parses an archive target's optional expiry and +// rotation. An absent, empty or never expiry yields an empty expiry, +// on which the edit form starts at never; saving it unchanged stores +// never, which means the same as an empty expiry. An absent rotation +// is empty too, and the form starts at none. An expiry that is set +// but not a valid duration, or a rotation that is not one of the +// four, is an error, not a blank field. func databaseConfigForm( configJSON string, ) (TargetConfigForm, error) { @@ -106,8 +110,15 @@ func databaseConfigForm( ) } + err = ValidateArchiveRotation(cfg.Rotation) + if err != nil { + return TargetConfigForm{}, err + } + + form := TargetConfigForm{Rotation: cfg.Rotation} + if cfg.Expiry == "" || cfg.Expiry == archiveExpiryNever { - return TargetConfigForm{}, nil + return form, nil } err = ValidateArchiveExpiry(cfg.Expiry) @@ -115,5 +126,7 @@ func databaseConfigForm( return TargetConfigForm{}, err } - return TargetConfigForm{Expiry: cfg.Expiry}, nil + form.Expiry = cfg.Expiry + + return form, nil } diff --git a/internal/delivery/target_config_view.go b/internal/delivery/target_config_view.go index 295cfb1..575706e 100644 --- a/internal/delivery/target_config_view.go +++ b/internal/delivery/target_config_view.go @@ -192,8 +192,9 @@ func maxRetriesField(t *database.Target) ConfigField { // databaseConfigFields describes an archive target by its // expiry in plain units, such as "30 days", or "never" when -// the archive is kept forever. An expiry that is set but not -// a valid duration is reported as unavailable rather than +// the archive is kept forever, and by its rotation. An expiry +// that is set but not a valid duration, or a rotation that is +// not one of the four, is reported as unavailable rather than // echoed back. func databaseConfigFields(configJSON string) []ConfigField { expiry, err := parseArchiveExpiry(configJSON) @@ -201,6 +202,11 @@ func databaseConfigFields(configJSON string) []ConfigField { return unavailableConfigFields() } + rotation, err := parseArchiveRotation(configJSON) + if err != nil { + return unavailableConfigFields() + } + value := archiveExpiryNever if expiry > 0 { value = plainDuration(expiry) @@ -209,6 +215,9 @@ func databaseConfigFields(configJSON string) []ConfigField { return []ConfigField{{ Label: "Archive Expiry", Value: value, + }, { + Label: "Archive Rotation", + Value: rotation, }} } diff --git a/internal/delivery/target_config_view_test.go b/internal/delivery/target_config_view_test.go index 5f8d987..923bd78 100644 --- a/internal/delivery/target_config_view_test.go +++ b/internal/delivery/target_config_view_test.go @@ -328,13 +328,49 @@ func TestNewTargetViews_Database(t *testing.T) { assert.Equal( t, - map[string]string{"Archive Expiry": tc.want}, + map[string]string{ + "Archive Expiry": tc.want, + "Archive Rotation": rotationNone, + }, fieldMap(view.Config), ) }) } } +// TestNewTargetViews_DatabaseRotation proves the target list shows a +// database target's rotation, none when it has none stored. +func TestNewTargetViews_DatabaseRotation(t *testing.T) { + t.Parallel() + + // Each stored config, and the rotation the list shows for it. + tests := map[string]string{ + "": rotationNone, + `{"rotation":""}`: rotationNone, + } + + for _, rotation := range []string{ + rotationNone, rotationMonthly, rotationDaily, rotationHourly, + } { + tests[`{"rotation":"`+rotation+`"}`] = rotation + } + + for config, want := range tests { + t.Run(config, func(t *testing.T) { + t.Parallel() + + view := viewFor(t, database.Target{ + Type: database.TargetTypeDatabase, + Config: config, + }) + + assert.Equal( + t, want, fieldMap(view.Config)["Archive Rotation"], + ) + }) + } +} + func TestNewTargetViews_Log(t *testing.T) { t.Parallel() @@ -382,6 +418,10 @@ func TestNewTargetViews_Unpresentable(t *testing.T) { Type: database.TargetTypeDatabase, Config: `{"expiry":"a fortnight"}`, }, + "invalid archive rotation": { + Type: database.TargetTypeDatabase, + Config: weeklyConfig, + }, } for name, target := range tests { diff --git a/internal/delivery/target_database.go b/internal/delivery/target_database.go index 483e317..862b7e8 100644 --- a/internal/delivery/target_database.go +++ b/internal/delivery/target_database.go @@ -20,7 +20,8 @@ const archiveNameMaxLen = 40 // from the per-webhook event database. The event is already // persisted in the per-webhook event DB by the time delivery runs; // the database target additionally writes a durable long-term copy -// into the file ArchiveFileName names and then records a single +// into the file ArchiveFileName names, with a period added when the +// target rotates (see archivePeriodPath), 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. @@ -146,8 +147,9 @@ func (t *databaseTarget) Deliver( } // archive writes the full event as a row into the target's -// archive database, honouring the optional per-target expiry -// parsed from the target config JSON. +// archive database, honouring the optional per-target expiry and +// rotation parsed from the target config JSON. With rotation, the +// event goes to the file for the period of its receive time. func (t *databaseTarget) archive(d *database.Delivery) error { webhookID := d.Event.WebhookID if webhookID == "" { @@ -159,6 +161,19 @@ func (t *databaseTarget) archive(d *database.Delivery) error { return err } + rotation, err := parseArchiveRotation(d.Target.Config) + if err != nil { + return err + } + + // An event whose stored row was gone before its delivery ran has + // no receive time (see Engine.hydrateEvent), and goes to the + // file for now. + receivedAt := d.Event.CreatedAt + if receivedAt.IsZero() { + receivedAt = time.Now() + } + w, err := t.writerFor(d.TargetID) if err != nil { return err @@ -174,7 +189,7 @@ func (t *databaseTarget) archive(d *database.Delivery) error { ContentType: d.Event.ContentType, } - return w.write(row, expiry) + return w.write(row, expiry, archivePeriod(rotation, receivedAt)) } // writerFor returns the archive writer for a database target, @@ -275,10 +290,10 @@ func (t *databaseTarget) releaseSweepWriter( delete(t.writers, targetID) } -// newWriter builds the writer for a database target's archive. The -// file is the one ArchivePath gives for the webhook and the target as +// newWriter builds the writer for a database target's archive. Its +// path is the one ArchivePath gives for the webhook and the target as // the main database names them now; from then on only rename changes -// the name the writer uses. It does not touch the archive file. +// the name the writer uses. It does not touch the archive files. func (t *databaseTarget) newWriter( targetID string, ) (*archiveWriter, error) { @@ -306,10 +321,10 @@ func (t *databaseTarget) newWriter( 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. +// rename moves every one of a database target's archive files 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 diff --git a/internal/delivery/target_database_archive.go b/internal/delivery/target_database_archive.go index c58664c..2f250ad 100644 --- a/internal/delivery/target_database_archive.go +++ b/internal/delivery/target_database_archive.go @@ -85,6 +85,10 @@ type databaseTargetConfig struct { // archived rows are pruned, or "never" (the default) to // keep them forever. Expiry string `json:"expiry"` + + // Rotation is none (the default), monthly, daily or hourly: see + // archivePeriod. + Rotation string `json:"rotation"` } // archivedEvent is one fully captured webhook event stored in a @@ -178,7 +182,7 @@ func ValidateArchiveExpiry(expiry string) error { return nil } -// archiveWriter owns one database target's archive SQLite file. +// archiveWriter owns one database target's archive SQLite files. // It serialises writes, and after each write closes and reopens // the file (debounced to at most once per debounce window) so // an operator can move the file away for offline archiving. The @@ -186,8 +190,16 @@ func ValidateArchiveExpiry(expiry string) error { // file is opened create-if-missing and its schema is migrated // on every open. type archiveWriter struct { - mu sync.Mutex - path string + mu sync.Mutex + + // path is the target's archive file as ArchivePath names it. A + // target that rotates writes to the files archivePeriodPath names + // for path and a period instead. + path string + + // current is the file db is open on. + current string + log *slog.Logger debounce time.Duration db *gorm.DB @@ -236,12 +248,14 @@ func newArchiveWriter( } } -// write appends the event as a row, then applies the debounced -// close/reopen. It recreates the archive file if it was moved -// or removed since the last open. A positive expiry prunes rows -// older than it on each (re)open. +// write appends the event as a row to the archive file for period +// (see archivePeriodPath), then applies the debounced close/reopen. +// When period names a different file from the one open, the open one +// is closed first. It recreates the archive file if it was moved or +// removed since the last open. A positive expiry prunes rows older +// than it on each (re)open. func (w *archiveWriter) write( - row archivedEvent, expiry time.Duration, + row archivedEvent, expiry time.Duration, period string, ) error { w.mu.Lock() defer w.mu.Unlock() @@ -252,8 +266,10 @@ func (w *archiveWriter) write( ) } - if w.db == nil || !fileExists(w.path) { - err := w.reopen(expiry) + file := archivePeriodPath(w.path, period) + + if w.db == nil || w.current != file || !fileExists(file) { + err := w.reopen(file, expiry) if err != nil { return err } @@ -264,41 +280,41 @@ func (w *archiveWriter) write( err := w.db.Create(&row).Error if err != nil { return fmt.Errorf( - "archiving event to %s: %w", w.path, err, + "archiving event to %s: %w", file, err, ) } if w.now().Sub(w.lastReopen) >= w.debounce { - return w.reopen(expiry) + return w.reopen(file, expiry) } return nil } -// open opens (creating if missing) the archive file, migrates +// open opens (creating if missing) an archive file, migrates // its schema, records the reopen time, and prunes expired rows // when expiry is positive. -func (w *archiveWriter) open(expiry time.Duration) error { - return w.openMode(archiveModeCreate, expiry) +func (w *archiveWriter) open(file string, expiry time.Duration) error { + return w.openMode(file, archiveModeCreate, expiry) } -// openMode opens the archive file with the given SQLite URI +// openMode opens an archive file with the given SQLite URI // mode, migrates its schema, records the reopen time, and // prunes expired rows when expiry is positive. The write path // passes archiveModeCreate so a missing file is recreated; the // idle sweep passes archiveModeExisting so a missing file is an // error rather than a newly conjured empty archive. func (w *archiveWriter) openMode( - mode string, expiry time.Duration, + file, mode string, expiry time.Duration, ) error { // Opened through database.OpenSQLite so an archive file carries // the same WAL journaling, busy timeout, immediate-transaction // locking, and pool bounds as every other database file. See // internal/database/sqlite_open.go. - sqlDB, err := database.OpenSQLite(w.path, mode) + sqlDB, err := database.OpenSQLite(file, mode) if err != nil { return fmt.Errorf( - "opening archive database %s: %w", w.path, err, + "opening archive database %s: %w", file, err, ) } @@ -314,7 +330,7 @@ func (w *archiveWriter) openMode( return fmt.Errorf( "connecting to archive database %s: %w", - w.path, err, + file, err, ) } @@ -323,11 +339,12 @@ func (w *archiveWriter) openMode( _ = sqlDB.Close() return fmt.Errorf( - "migrating archive database %s: %w", w.path, err, + "migrating archive database %s: %w", file, err, ) } w.db = gdb + w.current = file w.lastReopen = w.now() w.reopens++ @@ -338,12 +355,12 @@ func (w *archiveWriter) openMode( return nil } -// reopen closes any open handle and opens the file afresh. The +// reopen closes any open handle and opens file afresh. The // fresh open recreates the file if it was moved away. -func (w *archiveWriter) reopen(expiry time.Duration) error { +func (w *archiveWriter) reopen(file string, expiry time.Duration) error { w.close() - return w.open(expiry) + return w.open(file, expiry) } // close closes the underlying handle, if any. @@ -360,16 +377,17 @@ func (w *archiveWriter) close() { w.db = nil } -// sweepExpired prunes an archive that may have gone idle, with -// no write to trigger the usual on-reopen prune. It takes the -// writer's own mutex for the whole operation, so a sweep is -// ordered against concurrent writes rather than reaching around -// them to the file. +// sweepExpired prunes the target's archive files, which may have +// gone idle, with no write to trigger the usual on-reopen prune. It +// takes the writer's own mutex for the whole operation, so a sweep +// is ordered against concurrent writes rather than reaching around +// them to the files. // -// It never creates the archive file: a missing file is skipped, -// and the reopen uses archiveModeExisting so SQLite itself -// refuses to create one if the file disappears between the -// check and the open. +// It never creates an archive file: it prunes only the files +// archiveFiles finds, and opens each with archiveModeExisting so +// SQLite itself refuses to create one if the file disappears +// between the listing and the open. A file named for a period that +// the prune leaves empty is deleted. // // The archive is left CLOSED afterwards. An idle archive holding // no handle is what keeps the operator's move-the-file-away @@ -385,34 +403,84 @@ func (w *archiveWriter) sweepExpired(expiry time.Duration) error { ) } - if !fileExists(w.path) { - return nil - } - - // Drop any live handle first so the prune runs against a - // freshly opened file, matching the write path's semantics. - w.close() - - err := w.openMode(archiveModeExisting, expiry) + files, err := archiveFiles(w.path) if err != nil { return err } + var errs []error + + for _, file := range files { + err = w.sweepFile(file, expiry) + if err != nil { + errs = append(errs, err) + } + } + + return errors.Join(errs...) +} + +// sweepFile prunes one of the target's archive files for +// sweepExpired, which holds w.mu, and deletes the file, with its -wal +// and -shm, when it is named for a period and the prune leaves it +// empty. +func (w *archiveWriter) sweepFile( + file archiveFile, expiry time.Duration, +) error { + // Drop any live handle first so the prune runs against a + // freshly opened file, matching the write path's semantics. w.close() + err := w.openMode(file.path, archiveModeExisting, expiry) + if err != nil { + return err + } + + if file.period == "" { + w.close() + + return nil + } + + var rows int64 + + err = w.db.Model(&archivedEvent{}).Count(&rows).Error + + w.close() + + if err != nil { + return fmt.Errorf( + "counting rows in archive %s: %w", file.path, err, + ) + } + + if rows > 0 { + return nil + } + + for _, suffix := range []string{"", "-wal", "-shm"} { + err = os.Remove(file.path + suffix) + if err != nil && !errors.Is(err, fs.ErrNotExist) { + return fmt.Errorf("deleting empty archive file: %w", err) + } + } + + w.log.Info("deleted empty archive file", "path", file.path) + 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. +// rename gives every one of the target's archive files the new +// name, keeping the period in the name of each (see +// archivePeriodPath), and the writer uses the files under that name +// from now on. The handle is closed first, which folds the -wal into +// the .db; any -wal or -shm still beside a file (left by a crash) is +// moved with it, because SQLite finds them by name. A target with no +// files is not an error: the operator may have moved them away, and +// the next write creates its file under the new name. // -// If a file already has the new name, nothing is moved and the -// error is ErrArchiveNameTaken. If one file fails to move, those +// If a file already has one of the new names, nothing is moved and +// the error is ErrArchiveNameTaken. If one file fails to move, those // already moved are moved back before the error is returned, so the // archive is never split across two names. func (w *archiveWriter) rename(name string) error { @@ -430,38 +498,53 @@ func (w *archiveWriter) rename(name string) error { return nil } - suffixes := []string{"", "-wal", "-shm"} + files, err := archiveFiles(w.path) + if err != nil { + return err + } - for _, suffix := range suffixes { - if fileExists(path + suffix) { + // from[i] moves to to[i]. + var from, to []string + + for _, file := range files { + renamed := archivePeriodPath(path, file.period) + + for _, suffix := range []string{"", "-wal", "-shm"} { + from = append(from, file.path+suffix) + to = append(to, renamed+suffix) + } + } + + for _, taken := range to { + if fileExists(taken) { return fmt.Errorf( - "%w: %s", ErrArchiveNameTaken, name+suffix, + "%w: %s", ErrArchiveNameTaken, filepath.Base(taken), ) } } w.close() - for i, suffix := range suffixes { - err := os.Rename(w.path+suffix, path+suffix) + for i := range from { + err = os.Rename(from[i], to[i]) if err == nil || errors.Is(err, fs.ErrNotExist) { continue } - for _, moved := range suffixes[:i] { - backErr := os.Rename(path+moved, w.path+moved) + for j := range i { + backErr := os.Rename(to[j], from[j]) if backErr != nil && !errors.Is(backErr, fs.ErrNotExist) { w.log.Error( "failed to move archive file back", - "from", path+moved, - "to", w.path+moved, + "from", to[j], + "to", from[j], "error", backErr, ) } } return fmt.Errorf( - "renaming archive %s to %s: %w", w.path+suffix, path+suffix, err, + "renaming archive %s to %s: %w", from[i], to[i], err, ) } @@ -499,7 +582,7 @@ func (w *archiveWriter) prune(expiry time.Duration) { if res.Error != nil { w.log.Error( "failed to prune expired archive rows", - "path", w.path, + "path", w.current, "error", res.Error, ) @@ -509,48 +592,76 @@ func (w *archiveWriter) prune(expiry time.Duration) { if res.RowsAffected > 0 { w.log.Info( "pruned expired archive rows", - "path", w.path, + "path", w.current, "rows_deleted", res.RowsAffected, ) } } // ArchiveFileInfo is what the metadata of a database target's archive -// file says about it. +// files says about them. type ArchiveFileInfo struct { - // Size is the bytes on disk of the file and its -wal together. + // Files counts the files. + Files int + + // Size is the bytes on disk of the files and their -wal together. Size int64 - // Written is when the file or its -wal was last modified, whichever - // is later: a write lands in the -wal first. + // Written is when a file or a -wal was last modified, whichever is + // latest: a write lands in the -wal first. Written time.Time } -// StatArchive reads the metadata of the archive file at path and of -// its -wal, without opening the archive. With no file at path, which is -// so before the first write and after the operator moved it away, the +// StatArchive reads the metadata of a database target's archive +// files, given the path ArchivePath gives it (see archiveFiles), and +// of their -wal, without opening them. With no files, which is so +// before the first write and after the operator moved them away, the // error wraps fs.ErrNotExist. func StatArchive(path string) (ArchiveFileInfo, error) { - file, err := os.Stat(path) + files, err := archiveFiles(path) if err != nil { return ArchiveFileInfo{}, err } - info := ArchiveFileInfo{Size: file.Size(), Written: file.ModTime()} + var info ArchiveFileInfo - wal, err := os.Stat(path + "-wal") - if errors.Is(err, fs.ErrNotExist) { - return info, nil + for _, file := range files { + db, err := os.Stat(file.path) + if errors.Is(err, fs.ErrNotExist) { + continue + } + + if err != nil { + return ArchiveFileInfo{}, err + } + + info.Files++ + info.Size += db.Size() + + if db.ModTime().After(info.Written) { + info.Written = db.ModTime() + } + + wal, err := os.Stat(file.path + "-wal") + if errors.Is(err, fs.ErrNotExist) { + continue + } + + if err != nil { + return ArchiveFileInfo{}, err + } + + info.Size += wal.Size() + + if wal.ModTime().After(info.Written) { + info.Written = wal.ModTime() + } } - if err != nil { - return ArchiveFileInfo{}, err - } - - info.Size += wal.Size() - - if wal.ModTime().After(info.Written) { - info.Written = wal.ModTime() + if info.Files == 0 { + return ArchiveFileInfo{}, fmt.Errorf( + "no archive file for %s: %w", path, fs.ErrNotExist, + ) } return info, nil diff --git a/internal/delivery/target_database_export.go b/internal/delivery/target_database_export.go index 231aff7..d4d8f23 100644 --- a/internal/delivery/target_database_export.go +++ b/internal/delivery/target_database_export.go @@ -6,6 +6,7 @@ import ( "database/sql" "encoding/base64" "encoding/json" + "errors" "fmt" "io" "log/slog" @@ -50,21 +51,29 @@ func ArchiveExportFileName( at.UTC().Format("20060102T150405Z") + ".json.gz" } -// ArchiveExport is a database target's archive opened for download. -// It reads the file on its own connection, inside one read-only -// transaction, so it writes out the archive as it stood when +// ArchiveExport is a database target's archive opened for download: +// every one of its files, each read on its own connection inside one +// read-only transaction, so it writes out the archive as it stood when // OpenArchiveExport returned. // // Archives are in WAL mode, where a reader works from a snapshot and // never blocks a writer: archive writes go on while an export is open, -// and the export does not see them. SQLite cannot checkpoint the -wal -// past an open snapshot, so the -wal grows until the export is closed. +// and the export does not see them. SQLite cannot checkpoint a -wal +// past an open snapshot, so a file's -wal grows until the export has +// written out that file. type ArchiveExport struct { + files []*exportFile +} + +// exportFile is one archive file opened for an export. +type exportFile struct { db *sql.DB tx *gorm.DB - // empty is true when there is nothing to read: no file, or a file - // without the archive's table yet. + // period is the period in the file's name, "" for none. + period string + + // empty is true for a file without the archive's table yet. empty bool } @@ -74,26 +83,49 @@ type exportedName struct { Name string `json:"name"` } -// OpenArchiveExport opens the archive file at path for export and -// takes the snapshot the export reads. It never creates the file: with -// no file at path, the export has no rows. +// OpenArchiveExport opens every one of a database target's archive +// files for export, given the path ArchivePath gives it (see +// archiveFiles), and takes the snapshots the export reads. It never +// creates a file: with no files, the export has no rows. // -// Once it has returned, the file is open, so a rename or a move of it -// does not affect the export, which reads the same file under its new -// name. +// Once it has returned, the files are open, so a rename or a move of +// them does not affect the export, which reads the same files under +// their new names. // -// The transaction lasts as long as ctx does, so ctx must last for the +// The transactions last as long as ctx does, so ctx must last for the // whole export. func OpenArchiveExport( ctx context.Context, path string, log *slog.Logger, ) (*ArchiveExport, error) { - if !fileExists(path) { - return &ArchiveExport{empty: true}, nil + files, err := archiveFiles(path) + if err != nil { + return nil, err } - db, err := database.OpenSQLite(path, archiveModeExisting) + x := &ArchiveExport{} + + for _, file := range files { + f, err := openExportFile(ctx, file, log) + if err != nil { + _ = x.Close() + + return nil, err + } + + x.files = append(x.files, f) + } + + return x, nil +} + +// openExportFile opens one archive file for an export and takes its +// snapshot. +func openExportFile( + ctx context.Context, file archiveFile, log *slog.Logger, +) (*exportFile, error) { + db, err := database.OpenSQLite(file.path, archiveModeExisting) if err != nil { - return nil, fmt.Errorf("opening archive %s: %w", path, err) + return nil, fmt.Errorf("opening archive %s: %w", file.path, err) } gdb, err := gorm.Open( @@ -106,7 +138,7 @@ func OpenArchiveExport( if err != nil { _ = db.Close() - return nil, fmt.Errorf("opening archive %s: %w", path, err) + return nil, fmt.Errorf("opening archive %s: %w", file.path, err) } // ReadOnly makes the driver begin a deferred transaction in place @@ -116,7 +148,9 @@ func OpenArchiveExport( if tx.Error != nil { _ = db.Close() - return nil, fmt.Errorf("reading archive %s: %w", path, tx.Error) + return nil, fmt.Errorf( + "reading archive %s: %w", file.path, tx.Error, + ) } // The transaction's first read is what takes the snapshot. @@ -127,22 +161,27 @@ func OpenArchiveExport( _ = tx.Rollback() _ = db.Close() - return nil, fmt.Errorf("reading archive %s: %w", path, err) + return nil, fmt.Errorf("reading archive %s: %w", file.path, err) } - return &ArchiveExport{db: db, tx: tx, empty: tables == 0}, nil + return &exportFile{ + db: db, tx: tx, period: file.period, empty: tables == 0, + }, nil } // WriteGzipJSON writes the export to w as one gzipped JSON object: // webhook and target, each an id and a name; exported_at; and -// archived_events, one object per archived row, keyed by column name. -// A body that is not valid UTF-8 cannot be a JSON string, so it is -// written in base64, with "body_encoding": "base64" beside it. +// archived_events, one object per archived row, keyed by column name, +// the files in the order archiveFiles lists them. A row from a file +// named for a period has "period" beside its columns. A body that is +// not valid UTF-8 cannot be a JSON string, so it is written in base64, +// with "body_encoding": "base64" beside it. // // Each row is written out before the next is read, so neither the -// archive nor its JSON is ever held in memory whole. After an error -// the gzip stream is left unfinished, so what was written does not -// decompress as a whole file. +// archive nor its JSON is ever held in memory whole, and each file is +// closed once its rows are written. After an error the gzip stream is +// left unfinished, so what was written does not decompress as a whole +// file. func (x *ArchiveExport) WriteGzipJSON( ctx context.Context, w io.Writer, @@ -169,15 +208,30 @@ func (x *ArchiveExport) WriteGzipJSON( return zw.Close() } -// Close ends the export's transaction and closes its connection. +// Close ends the transactions of the files still open and closes +// their connections. func (x *ArchiveExport) Close() error { - if x.db == nil { + errs := make([]error, 0, len(x.files)) + + for _, f := range x.files { + errs = append(errs, f.close()) + } + + return errors.Join(errs...) +} + +// close ends the file's transaction and closes its connection, once. +func (f *exportFile) close() error { + if f.db == nil { return nil } - _ = x.tx.Rollback() + _ = f.tx.Rollback() - return x.db.Close() + err := f.db.Close() + f.db = nil + + return err } // writeJSON writes head with archived_events added as its last key, @@ -207,46 +261,72 @@ func (x *ArchiveExport) writeJSON( return err } -// writeRows writes each archived row to w, oldest first, one per line, -// separated by commas. +// writeRows writes the archived rows of each file to w, one per line, +// separated by commas, and closes each file once its rows are written. func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error { - if x.empty { - return nil + sep := "\n" + + for _, f := range x.files { + var err error + + sep, err = f.writeRows(ctx, w, sep) + if err != nil { + return err + } + + err = f.close() + if err != nil { + return err + } } - rows, err := x.tx.WithContext(ctx). + return nil +} + +// writeRows writes the file's archived rows to w, oldest first, the +// first after sep and each other after ",\n". It returns what goes +// before the next row: sep again when the file had no rows. +func (f *exportFile) writeRows( + ctx context.Context, w io.Writer, sep string, +) (string, error) { + if f.empty { + return sep, nil + } + + rows, err := f.tx.WithContext(ctx). Model(&archivedEvent{}).Order("id").Rows() if err != nil { - return err + return "", err } defer func() { _ = rows.Close() }() - for sep := "\n"; rows.Next(); sep = ",\n" { + for ; rows.Next(); sep = ",\n" { var ev archivedEvent - err = x.tx.ScanRows(rows, &ev) + err = f.tx.ScanRows(rows, &ev) if err != nil { - return err + return "", err } _, err = io.WriteString(w, sep) if err != nil { - return err + return "", err } - err = writeRow(w, &ev) + err = writeRow(w, &ev, f.period) if err != nil { - return err + return "", err } } - return rows.Err() + return sep, rows.Err() } // writeRow writes an archived row to w as a JSON object keyed by -// column name, its body in base64 when it is not valid UTF-8. -func writeRow(w io.Writer, ev *archivedEvent) error { +// column name, its body in base64 when it is not valid UTF-8, with +// the period of its file beside them unless that is "". +func writeRow(w io.Writer, ev *archivedEvent, period string) error { row := map[string]any{ "id": ev.ID, "event_id": ev.EventID, @@ -264,6 +344,10 @@ func writeRow(w io.Writer, ev *archivedEvent) error { row["body_encoding"] = "base64" } + if period != "" { + row["period"] = period + } + line, err := json.Marshal(row) if err != nil { return err diff --git a/internal/delivery/target_database_export_test.go b/internal/delivery/target_database_export_test.go index e3d63fe..8d30e00 100644 --- a/internal/delivery/target_database_export_test.go +++ b/internal/delivery/target_database_export_test.go @@ -301,6 +301,56 @@ func TestArchiveExport_SurvivesRename(t *testing.T) { ) } +// TestArchiveExport_EveryFileOldestFirst writes a row to a target's +// file without a period and to its files for a month, an hour and a +// day, and proves the export holds every row: the file without a +// period first, then the others oldest period first, each row from a +// file named for a period carrying that period. A file made after the +// export opened is not in it. +func TestArchiveExport_EveryFileOldestFirst(t *testing.T) { + t.Parallel() + + path := filepath.Join(t.TempDir(), "archive-wh.db") + w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0) + + // Written in an order that is not the export's. + for _, period := range []string{nextDayPeriod, "", hourPeriod, "2026-03"} { + require.NoError(t, w.WritePeriod( + delivery.ExportArchivedEvent{EventID: "in-" + period}, 0, period, + )) + } + + export, err := delivery.OpenArchiveExport( + t.Context(), path, archiveTestLogger(), + ) + require.NoError(t, err) + + defer func() { require.NoError(t, export.Close()) }() + + require.NoError(t, w.WritePeriod( + delivery.ExportArchivedEvent{EventID: "later"}, 0, "2026-03-06", + )) + + events := exportedEvents(t, writeExport(t, export)) + ids := make([]string, 0, len(events)) + periods := make([]any, 0, len(events)) + + for _, ev := range events { + ids = append(ids, fmt.Sprint(ev["event_id"])) + periods = append(periods, ev["period"]) + } + + assert.Equal(t, + []string{"in-", "in-2026-03", "in-" + hourPeriod, "in-" + nextDayPeriod}, + ids, + ) + assert.Equal(t, + []any{nil, "2026-03", hourPeriod, nextDayPeriod}, periods, + ) + assert.NotContains(t, events[0], "period", + "a row from the file without a period has no period") +} + // heapPeak is an io.Writer that discards what it is given and records // the largest heap it saw at a write. It collects garbage before each // reading, so the heap it reads is what is still held. diff --git a/internal/delivery/target_database_rotation.go b/internal/delivery/target_database_rotation.go new file mode 100644 index 0000000..e16180e --- /dev/null +++ b/internal/delivery/target_database_rotation.go @@ -0,0 +1,198 @@ +package delivery + +import ( + "encoding/json" + "errors" + "fmt" + "os" + "path/filepath" + "slices" + "strings" + "time" + + "sneak.berlin/go/webhooker/internal/database" +) + +// The archive rotations: how often a database target starts a new +// archive file. Every rotation but none puts the period of an event's +// receive time, in UTC, in the name of the file the event goes to, +// written in the layout beside it. +const ( + archiveRotationNone = "none" + archiveRotationMonthly = "monthly" + archiveRotationDaily = "daily" + archiveRotationHourly = "hourly" + + archiveMonthLayout = "2006-01" + archiveDayLayout = "2006-01-02" + archiveHourLayout = "2006-01-02-15" +) + +// errArchiveRotationUnknown is returned for a rotation that is not one +// of the four. +var errArchiveRotationUnknown = errors.New( + "rotation must be none, monthly, daily or hourly", +) + +// archiveFile is one of a database target's archive files, and the +// period in its name: "" for the file named without one. +type archiveFile struct { + path string + period string +} + +// ValidateArchiveRotation checks a user-supplied archive rotation for +// a database target: empty or none (both meaning one file), monthly, +// daily or hourly. +func ValidateArchiveRotation(rotation string) error { + switch rotation { + case "", archiveRotationNone, archiveRotationMonthly, + archiveRotationDaily, archiveRotationHourly: + return nil + default: + return fmt.Errorf("%w: %q", errArchiveRotationUnknown, rotation) + } +} + +// parseArchiveRotation reads the rotation from a database target's +// config JSON. An empty config or an empty rotation is none. +func parseArchiveRotation(configJSON string) (string, error) { + if configJSON == "" { + return archiveRotationNone, nil + } + + var cfg databaseTargetConfig + + err := json.Unmarshal([]byte(configJSON), &cfg) + if err != nil { + return "", fmt.Errorf("parsing database target config: %w", err) + } + + err = ValidateArchiveRotation(cfg.Rotation) + if err != nil { + return "", err + } + + if cfg.Rotation == "" { + return archiveRotationNone, nil + } + + return cfg.Rotation, nil +} + +// archivePeriod returns the period, in UTC, that a rotation puts an +// event received at receivedAt in: "2026-10" for monthly, +// "2026-10-01" for daily, "2026-10-01-19" for hourly, and "" for none. +func archivePeriod(rotation string, receivedAt time.Time) string { + switch rotation { + case archiveRotationMonthly: + return receivedAt.UTC().Format(archiveMonthLayout) + case archiveRotationDaily: + return receivedAt.UTC().Format(archiveDayLayout) + case archiveRotationHourly: + return receivedAt.UTC().Format(archiveHourLayout) + default: + return "" + } +} + +// archivePeriodPath returns the path of a database target's archive +// file for a period: path, as ArchivePath gives it, with "-" and the +// period put before its ".db". The period "" gives path itself. +func archivePeriodPath(path, period string) string { + if period == "" { + return path + } + + return strings.TrimSuffix(path, ".db") + "-" + period + ".db" +} + +// ArchivePathAt returns the archive file a database target writes an +// event received at receivedAt to: ArchivePath's file, with the +// period in its name when the target rotates. +func ArchivePathAt( + dbMgr *database.WebhookDBManager, + webhook *database.Webhook, + target *database.Target, + receivedAt time.Time, +) (string, error) { + rotation, err := parseArchiveRotation(target.Config) + if err != nil { + return "", err + } + + return archivePeriodPath( + ArchivePath(dbMgr, webhook, target), + archivePeriod(rotation, receivedAt), + ), nil +} + +// archiveFiles lists the database target's archive files that exist, +// given the path ArchivePath gives it: the file at path, then each +// file archivePeriodPath names for path and a period, oldest period +// first. Which rotation wrote a file does not matter, so the files of +// an earlier rotation setting are listed too. +func archiveFiles(path string) ([]archiveFile, error) { + dir := filepath.Dir(path) + + entries, err := os.ReadDir(dir) + if err != nil { + return nil, fmt.Errorf("listing archive files: %w", err) + } + + stem := strings.TrimSuffix(filepath.Base(path), ".db") + + var files []archiveFile + + for _, entry := range entries { + period, ok := archiveFilePeriod(stem, entry.Name()) + if ok { + files = append(files, archiveFile{ + path: filepath.Join(dir, entry.Name()), + period: period, + }) + } + } + + // A month sorts before the days and hours in it. + slices.SortFunc(files, func(a, b archiveFile) int { + return strings.Compare(a.period, b.period) + }) + + return files, nil +} + +// archiveFilePeriod reports whether name is the name of an archive +// file of the target whose file name without a period is stem+".db", +// and the period in it. +func archiveFilePeriod(stem, name string) (string, bool) { + rest, ok := strings.CutPrefix(name, stem) + if !ok { + return "", false + } + + rest, ok = strings.CutSuffix(rest, ".db") + if !ok { + return "", false + } + + if rest == "" { + return "", true + } + + period, ok := strings.CutPrefix(rest, "-") + if !ok { + return "", false + } + + for _, layout := range []string{ + archiveMonthLayout, archiveDayLayout, archiveHourLayout, + } { + _, err := time.Parse(layout, period) + if err == nil { + return period, true + } + } + + return "", false +} diff --git a/internal/delivery/target_database_rotation_test.go b/internal/delivery/target_database_rotation_test.go new file mode 100644 index 0000000..8bfb7c3 --- /dev/null +++ b/internal/delivery/target_database_rotation_test.go @@ -0,0 +1,389 @@ +package delivery_test + +import ( + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "sneak.berlin/go/webhooker/internal/database" + "sneak.berlin/go/webhooker/internal/delivery" +) + +// The archive rotations, and the configs of a daily target and of one +// whose rotation is not one of the four. +const ( + rotationNone = "none" + rotationMonthly = "monthly" + rotationDaily = "daily" + rotationHourly = "hourly" + + dailyConfig = `{"rotation":"daily"}` + weeklyConfig = `{"rotation":"weekly"}` +) + +// The periods the tests archive into most: two days, and an hour of +// the first. +const ( + dayPeriod = "2026-03-04" + nextDayPeriod = "2026-03-05" + hourPeriod = "2026-03-04-05" +) + +// periodPath returns the archive file for a period of the target whose +// file without a period is path. +func periodPath(path, period string) string { + return strings.TrimSuffix(path, ".db") + "-" + period + ".db" +} + +// deliverReceivedAt delivers to a database target an event whose +// receive time is receivedAt, and returns the event's id. The receive +// time is what decides a rotated archive's file, so setting it is how +// these tests move the clock across a period boundary. +func (env *archiveEnv) deliverReceivedAt( + t *testing.T, tgt *database.Target, receivedAt time.Time, +) string { + t.Helper() + + webhookDB := testWebhookDB(t) + event := seedEvent(t, webhookDB, `{"n":1}`) + event.CreatedAt = receivedAt + + env.eng.ExportDeliverDatabase( + webhookDB, seedDatabaseTargetDelivery(t, webhookDB, event, tgt), + ) + + return event.ID +} + +// TestDeliverDatabase_RotatesAtEachPeriodBoundary delivers, for each +// rotation, an event received in the last second of a period and one +// received in the first second of the next, and checks each lands in +// the file named for its own period, in UTC. Rotation none keeps both +// in the one file. +func TestDeliverDatabase_RotatesAtEachPeriodBoundary(t *testing.T) { + t.Parallel() + + berlin := time.FixedZone("CEST", 2*60*60) + + cases := []struct { + name string + rotation string + before, after time.Time + // periods are the periods of before and after. + periods [2]string + }{ + { + rotationMonthly, rotationMonthly, + time.Date(2026, 1, 31, 23, 59, 59, 0, time.UTC), + time.Date(2026, 2, 1, 0, 0, 0, 0, time.UTC), + [2]string{"2026-01", "2026-02"}, + }, + { + rotationDaily, rotationDaily, + time.Date(2026, 3, 4, 23, 59, 59, 0, time.UTC), + time.Date(2026, 3, 5, 0, 0, 0, 0, time.UTC), + [2]string{dayPeriod, nextDayPeriod}, + }, + { + // The same instants, received in a zone two hours ahead + // of UTC, where they fall on 5 March: the period is UTC's. + "daily in another zone", rotationDaily, + time.Date(2026, 3, 5, 1, 59, 59, 0, berlin), + time.Date(2026, 3, 5, 2, 0, 0, 0, berlin), + [2]string{dayPeriod, nextDayPeriod}, + }, + { + rotationHourly, rotationHourly, + time.Date(2026, 3, 4, 5, 59, 59, 0, time.UTC), + time.Date(2026, 3, 4, 6, 0, 0, 0, time.UTC), + [2]string{hourPeriod, "2026-03-04-06"}, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, `{"rotation":"`+tc.rotation+`"}`) + path := env.archivePath(tgt) + + first := env.deliverReceivedAt(t, tgt, tc.before) + second := env.deliverReceivedAt(t, tgt, tc.after) + + assert.Equal(t, []string{first}, + archivedEventIDs(t, periodPath(path, tc.periods[0]))) + assert.Equal(t, []string{second}, + archivedEventIDs(t, periodPath(path, tc.periods[1]))) + assert.NoFileExists(t, path, + "a rotated target never writes the file without a period") + }) + } + + t.Run(rotationNone, func(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, `{"rotation":"`+rotationNone+`"}`) + + first := env.deliverReceivedAt(t, tgt, cases[0].before) + second := env.deliverReceivedAt(t, tgt, cases[0].after) + + assert.ElementsMatch(t, []string{first, second}, + archivedEventIDs(t, env.archivePath(tgt))) + }) +} + +// TestDeliverDatabase_EventWithoutReceiveTime proves an event whose +// receive time is not known goes to the file for the time it is +// archived, rather than to one for the year 1. +func TestDeliverDatabase_EventWithoutReceiveTime(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, `{"rotation":"`+rotationMonthly+`"}`) + path := env.archivePath(tgt) + + before := time.Now().UTC().Format("2006-01") + id := env.deliverReceivedAt(t, tgt, time.Time{}) + + file := periodPath(path, time.Now().UTC().Format("2006-01")) + + _, err := os.Stat(file) + if err != nil { + // The month turned during the delivery. + file = periodPath(path, before) + } + + assert.Equal(t, []string{id}, archivedEventIDs(t, file)) + assert.NoFileExists(t, periodPath(path, "0001-01")) +} + +// TestDeliverDatabase_RotationChangeKeepsOldFiles changes a target's +// rotation from none to daily between two events, and checks the +// second goes to the daily file while the first stays where it was. +func TestDeliverDatabase_RotationChangeKeepsOldFiles(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, "") + path := env.archivePath(tgt) + at := time.Date(2026, 3, 4, 12, 0, 0, 0, time.UTC) + + first := env.deliverReceivedAt(t, tgt, at) + + tgt.Config = dailyConfig + second := env.deliverReceivedAt(t, tgt, at) + + assert.Equal(t, []string{first}, archivedEventIDs(t, path)) + assert.Equal(t, []string{second}, + archivedEventIDs(t, periodPath(path, dayPeriod))) +} + +// TestArchiveSweep_PrunesEveryFile gives a daily target three files: +// the file without a period, left from before it rotated, and two +// daily files. Each holds a row older than the expiry, and one daily +// file also a newer row. The sweep prunes the old row from every file, +// deletes the daily file it leaves empty, and keeps the file without a +// period although it is empty too. +func TestArchiveSweep_PrunesEveryFile(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget( + t, `{"expiry":"1h","rotation":"`+rotationDaily+`"}`, + ) + path := env.archivePath(tgt) + emptied := periodPath(path, dayPeriod) + kept := periodPath(path, nextDayPeriod) + + now := time.Now() + old := now.Add(-48 * time.Hour) + + seedArchiveFile(t, path, tgt.WebhookID, old) + seedArchiveFile(t, emptied, tgt.WebhookID, old) + seedArchiveFile(t, kept, tgt.WebhookID, old, now.Add(-time.Minute)) + + env.sweeper.ExportSweep(t.Context()) + + assert.Empty(t, archivedEventIDs(t, path)) + assert.Equal(t, []string{sweepRowNew}, archivedEventIDs(t, kept)) + + for _, suffix := range archiveFileSuffixes() { + assert.NoFileExists(t, emptied+suffix) + } +} + +// TestRename_MovesEveryFile renames a daily target that also has a +// file without a period, and checks every file moves to the new name +// with its period, rows and all, and that a later write uses the new +// name. +func TestRename_MovesEveryFile(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, "") + oldPath := env.archivePath(tgt) + day := time.Date(2026, 3, 4, 12, 0, 0, 0, time.UTC) + + unrotated := env.deliverReceivedAt(t, tgt, day) + + tgt.Config = dailyConfig + first := env.deliverReceivedAt(t, tgt, day) + second := env.deliverReceivedAt(t, tgt, day.Add(24*time.Hour)) + + require.NoError(t, env.eng.Rename(tgt.ID, "Orders", "Long Term")) + + newPath := filepath.Join( + env.dataDir, "archive-orders-long-term-"+tgt.ID+".db", + ) + + for _, old := range []string{ + oldPath, + periodPath(oldPath, dayPeriod), + periodPath(oldPath, nextDayPeriod), + } { + assert.NoFileExists(t, old) + } + + assert.Equal(t, []string{unrotated}, archivedEventIDs(t, newPath)) + assert.Equal(t, []string{first}, + archivedEventIDs(t, periodPath(newPath, dayPeriod))) + assert.Equal(t, []string{second}, + archivedEventIDs(t, periodPath(newPath, nextDayPeriod))) + + third := env.deliverReceivedAt(t, tgt, day.Add(48*time.Hour)) + assert.Equal(t, []string{third}, + archivedEventIDs(t, periodPath(newPath, "2026-03-06"))) +} + +// TestRename_NeverReplacesARotatedFile plants a file at the new name +// of a target's daily file, and proves the rename is refused and moves +// none of the target's files. +func TestRename_NeverReplacesARotatedFile(t *testing.T) { + t.Parallel() + + env := setupArchiveTest(t) + tgt := env.seedDatabaseTarget(t, dailyConfig) + oldPath := env.archivePath(tgt) + day := time.Date(2026, 3, 4, 12, 0, 0, 0, time.UTC) + + env.deliverReceivedAt(t, tgt, day) + env.deliverReceivedAt(t, tgt, day.Add(24*time.Hour)) + + newPath := filepath.Join( + env.dataDir, "archive-orders-long-term-"+tgt.ID+".db", + ) + planted := periodPath(newPath, nextDayPeriod) + require.NoError(t, os.WriteFile(planted, []byte("planted"), 0o600)) + + require.ErrorIs( + t, env.eng.Rename(tgt.ID, "Orders", "Long Term"), + delivery.ErrArchiveNameTaken, + ) + + assert.FileExists(t, periodPath(oldPath, dayPeriod)) + assert.FileExists(t, periodPath(oldPath, nextDayPeriod)) + assert.NoFileExists(t, periodPath(newPath, dayPeriod)) +} + +// TestStatArchive_EveryFile proves StatArchive counts and adds up every +// one of a target's files, takes the latest write of any of them, and +// leaves out files whose names only look like the target's. +func TestStatArchive_EveryFile(t *testing.T) { + t.Parallel() + + dir := t.TempDir() + path := filepath.Join(dir, "archive-wh.db") + files := []string{ + path, periodPath(path, "2026-03"), periodPath(path, hourPeriod), + } + written := time.Date(2026, 3, 4, 5, 6, 7, 0, time.UTC) + + var size int64 + + for i, file := range files { + require.NoError(t, os.WriteFile(file, make([]byte, 100*(i+1)), 0o600)) + + size += int64(100 * (i + 1)) + at := written.Add(-time.Duration(i) * time.Hour) + require.NoError(t, os.Chtimes(file, at, at)) + } + + for _, other := range []string{ + "archive-wh-2026-13.db", "archive-wh-2026-3.db", + "archive-wh-other.db", "archive-wh-2026-03.json", + "archive-whx.db", + } { + require.NoError(t, + os.WriteFile(filepath.Join(dir, other), []byte("x"), 0o600)) + } + + got, err := delivery.StatArchive(path) + require.NoError(t, err) + assert.Equal(t, len(files), got.Files) + assert.Equal(t, size, got.Size) + assert.True(t, written.Equal(got.Written), got.Written) +} + +// TestArchivePathAt names the file each rotation writes an event to. +func TestArchivePathAt(t *testing.T) { + t.Parallel() + + dataDir := t.TempDir() + dbMgr := database.NewTestWebhookDBManager(dataDir) + webhook := &database.Webhook{ + BaseModel: database.BaseModel{ID: "wh-id"}, Name: "Orders", + } + at := time.Date(2026, 10, 1, 19, 30, 0, 0, time.UTC) + + cases := map[string]string{ + "": "", + `{"rotation":"` + rotationNone + `"}`: "", + `{"rotation":"` + rotationMonthly + `"}`: "-2026-10", + dailyConfig: "-2026-10-01", + `{"rotation":"` + rotationHourly + `"}`: "-2026-10-01-19", + } + + for config, period := range cases { + target := &database.Target{ + BaseModel: database.BaseModel{ID: "tgt-id"}, + Name: "Archive", + Config: config, + } + + got, err := delivery.ArchivePathAt(dbMgr, webhook, target, at) + require.NoError(t, err, config) + assert.Equal(t, + filepath.Join( + dataDir, "archive-orders-archive-tgt-id"+period+".db", + ), + got, config, + ) + } + + _, err := delivery.ArchivePathAt(dbMgr, webhook, &database.Target{ + Config: weeklyConfig, + }, at) + require.Error(t, err) +} + +// TestValidateArchiveRotation accepts the four rotations, and empty, +// and refuses anything else. +func TestValidateArchiveRotation(t *testing.T) { + t.Parallel() + + for _, ok := range []string{ + "", rotationNone, rotationMonthly, rotationDaily, rotationHourly, + } { + require.NoError(t, delivery.ValidateArchiveRotation(ok), ok) + } + + for _, bad := range []string{"weekly", "Daily", "hourly "} { + require.Error(t, delivery.ValidateArchiveRotation(bad), bad) + } +} diff --git a/internal/delivery/target_database_test.go b/internal/delivery/target_database_test.go index 516e579..8b65fe4 100644 --- a/internal/delivery/target_database_test.go +++ b/internal/delivery/target_database_test.go @@ -218,6 +218,7 @@ func TestStatArchive(t *testing.T) { got, err := delivery.StatArchive(path) require.NoError(t, err) + assert.Equal(t, 1, got.Files) assert.Equal(t, file.Size()+wal.Size(), got.Size) assert.True(t, written.Equal(got.Written), got.Written) diff --git a/internal/delivery/target_headers_test.go b/internal/delivery/target_headers_test.go index ac72b86..4250dee 100644 --- a/internal/delivery/target_headers_test.go +++ b/internal/delivery/target_headers_test.go @@ -204,10 +204,11 @@ func TestNewTargetConfigForm(t *testing.T) { form, err = delivery.NewTargetConfigForm(&database.Target{ Type: database.TargetTypeDatabase, - Config: `{"expiry":"720h"}`, + Config: `{"expiry":"720h","rotation":"daily"}`, }) require.NoError(t, err) assert.Equal(t, "720h", form.Expiry) + assert.Equal(t, rotationDaily, form.Rotation) form, err = delivery.NewTargetConfigForm(&database.Target{ Type: database.TargetTypeLog, @@ -248,6 +249,10 @@ func TestNewTargetConfigForm_UnreadableConfigErrors(t *testing.T) { Type: database.TargetTypeDatabase, Config: `{"expiry":"soon"}`, }, + { + Type: database.TargetTypeDatabase, + Config: weeklyConfig, + }, {Type: database.TargetType("nope")}, } diff --git a/internal/handlers/archive_expiry.go b/internal/handlers/archive_expiry.go index 060737d..38a328b 100644 --- a/internal/handlers/archive_expiry.go +++ b/internal/handlers/archive_expiry.go @@ -10,10 +10,10 @@ const ( tmplKeyArchiveExpiryChoices = "ArchiveExpiryChoices" ) -// archiveExpiryChoice is one entry of a database target's archive -// expiry select: the expiry stored, the label shown, and whether the -// select starts on it. -type archiveExpiryChoice struct { +// archiveChoice is one entry of a database target's archive expiry +// or archive rotation select: the value stored, the label shown, and +// whether the select starts on it. +type archiveChoice struct { Value string Label string Selected bool @@ -21,8 +21,8 @@ type archiveExpiryChoice struct { // archiveExpiryChoices lists the archive expiries offered by the new // webhook page, the add target form and the target edit form. -func archiveExpiryChoices() []archiveExpiryChoice { - return []archiveExpiryChoice{ +func archiveExpiryChoices() []archiveChoice { + return []archiveChoice{ {Value: archiveExpiryNever, Label: archiveExpiryNever}, {Value: "1h", Label: "1h"}, {Value: "12h", Label: "12h"}, @@ -37,7 +37,7 @@ func archiveExpiryChoices() []archiveExpiryChoice { // empty expiry selects never. An expiry that is not one of the // choices comes first as its own selected entry, so saving the form // unchanged keeps it. -func archiveExpiryOptions(expiry string) []archiveExpiryChoice { +func archiveExpiryOptions(expiry string) []archiveChoice { if expiry == "" { expiry = archiveExpiryNever } @@ -52,7 +52,7 @@ func archiveExpiryOptions(expiry string) []archiveExpiryChoice { } } - own := archiveExpiryChoice{Value: expiry, Label: expiry, Selected: true} + own := archiveChoice{Value: expiry, Label: expiry, Selected: true} - return append([]archiveExpiryChoice{own}, options...) + return append([]archiveChoice{own}, options...) } diff --git a/internal/handlers/archive_expiry_test.go b/internal/handlers/archive_expiry_test.go index 70bf476..9484543 100644 --- a/internal/handlers/archive_expiry_test.go +++ b/internal/handlers/archive_expiry_test.go @@ -5,6 +5,7 @@ import ( "net/http/httptest" "net/url" "regexp" + "strings" "testing" "github.com/stretchr/testify/assert" @@ -49,21 +50,39 @@ func expiryShown( } // expirySelected returns the target edit page and the expiries its -// select starts on. +// expiry select starts on. func expirySelected( t *testing.T, env *sourceTestEnv, webhookID, targetID string, ) (string, []string) { t.Helper() + page := targetEditPage(t, env, webhookID, targetID) + + return page, selectedIn(page, "expiry") +} + +// targetEditPage returns a target's edit page. +func targetEditPage( + t *testing.T, env *sourceTestEnv, webhookID, targetID string, +) string { + t.Helper() + w := serveTarget( env, http.MethodGet, "/hook/"+webhookID+"/targets/"+targetID+"/edit", nil, ) require.Equal(t, http.StatusOK, w.Code) - page := w.Body.String() + return w.Body.String() +} - return page, matched(`