1 Commits
Author SHA1 Message Date
sneak 13f881007a Rotate a database target's archive monthly, daily or hourly (closes #379)
check / check (push) Successful in 3m42s
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
2026-10-02 23:51:42 +00:00
33 changed files with 1753 additions and 328 deletions
+78 -41
View File
@@ -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 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-{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 this as a persistent volume to
preserve data across container restarts. 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, - `events-{webhook_uuid}.db` — **one per webhook**. Events, deliveries,
delivery results. delivery results.
- `archive-{webhook_name}-{target_name}-{target_uuid}.db` — **one per - `archive-{webhook_name}-{target_name}-{target_uuid}.db` — **one per
`database` target**. Archived events. The two names are made safe for `database` target**, or one per month, day or hour for a target that
a file name, and the file is renamed when the webhook or the target is rotates, with `-{period}` before `.db`. Archived events. The two names
(see [Database Architecture](#database-architecture)). 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 `{webhook_uuid}` and `{target_uuid}` are UUID primary keys in their
canonical 36-character hyphenated form, so a real filename looks like 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 optional HTTP target URL creates an `http` target named `HTTP`, and the
archive checkbox creates a `database` target named `Archive` whose archive checkbox creates a `database` target named `Archive` whose
`expiry` is the pruning chosen beside it (never, 1h, 12h, 24h, 30d, 90d `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 or 365d) and whose `rotation` is the rotation chosen below that (none,
and its targets are created together or not at all. 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 | | Field | Type | Description |
| ---------------- | ------- | ----------- | | ---------------- | ------- | ----------- |
@@ -1720,12 +1724,14 @@ events should be forwarded.
own archive database own archive database
(`archive-{webhook_name}-{target_name}-{target_uuid}.db`) for long-term (`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). The new webhook form, the add target form and the target edit forever) and rotation (default: none, one file). The new webhook form,
form all offer the same expiries: never, 1h, 12h, 24h, 30d, 90d or 365d. the add target form and the target edit form all offer the same
The target list shows the expiry in plain units, such as "30 days". No expiries: never, 1h, 12h, 24h, 30d, 90d or 365d, and the same
external delivery and no retries; an archive write failure fails the rotations: none, monthly, daily or hourly. The target list shows the
delivery. See the database target section under expiry in plain units, such as "30 days", and the rotation. No external
"Per-Webhook Event Databases" for the full semantics. 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 - **`log`** — Write the event to the application log (stdout). Useful
for debugging. 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 when nothing is left. The target UUID keeps the file name unique. A
webhook named `Orders (EU)` with a target named `Long-term archive` webhook named `Orders (EU)` with a target named `Long-term archive`
archives into `archive-orders-eu-long-term-archive-{target_uuid}.db`. 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, An optional `rotation` in the target's config JSON (e.g.
target edits and target creation run one at a time, so no edit can `{"rotation":"daily"}`) is `none`, the default, which keeps the one
rename the file between another's rename and save, and the name on disk file, or `monthly`, `daily` or `hourly`. A target that rotates writes
matches the UI. A rename never replaces a file: if one already has each event to a file named for the period of the event's receive time,
the new name, the edit is refused with an error naming that file, and in UTC, put before the `.db`:
the stored name stays. If the archive is not there (the operator moved `archive-orders-eu-long-term-archive-{target_uuid}-2026-10.db` monthly,
it away), the rename is not an error, and the next write creates the `…-2026-10-01.db` daily and `…-2026-10-01-19.db` hourly. Each file holds
file under the new name. 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 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 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, the move-the-file-away workflow keeps working. Archives with no expiry,
or the expiry `never`, are not touched by the sweep at all. 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 Because each `database` target has its own archive file, a target's
`expiry` governs only its own archive. Two `database` targets on one `expiry` governs only its own archive. Two `database` targets on one
webhook with different expiries keep two archives, each pruned on its webhook with different expiries keep two archives, each pruned on its
own schedule. own schedule.
The webhook page shows, for each `database` target, its archive file's The webhook page shows, for each `database` target, the name of the
name, its size on disk and when it was last written. The size counts archive file an event received now would go to, the size on disk of all
the `.db` and its `-wal` together, and the last write is the later of the target's archive files together, how many there are when there is
their two modification times, since a write lands in the `-wal` first. more than one, and when the latest of them was last written. The size
Both are read from the files' metadata; the archive is never opened. counts each `.db` and its `-wal` together, and the last write is the
Before the first write, and after the file has been moved away, the page latest of their modification times, since a write lands in the `-wal`
shows `not created yet` beside the name. 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, Each `database` target on the webhook page has a **Download** button,
which returns its archive as one gzipped JSON file, 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 names made safe as above and the time in UTC. The file holds one
object: `webhook` and `target`, each an `id` and a `name`; object: `webhook` and `target`, each an `id` and a `name`;
`exported_at`; and `archived_events`, one object per archived row with `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 every column, keyed by column name. The rows come from every one of the
written in base64, with `"body_encoding": "base64"` beside it. An target's files: the file without a period first, then the others in
archive that does not exist yet, or was moved away, downloads with an the order of their periods, oldest first, and each row from a file named
empty `archived_events`; the download never creates the file. 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 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 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 memory. When it starts it opens every one of the target's files, each on
transaction, so the file holds the archive as it stood when the a connection of its own inside one read-only transaction, so the
download started, and archive writes go on meanwhile, since under WAL a download holds the archive as it stood then, and archive writes go on
reader never blocks a writer. While it runs, the `-wal` cannot be meanwhile, since under WAL a reader never blocks a writer. Each file is
checkpointed past what it reads, so a long download lets the `-wal` closed once its rows are written out; until then its `-wal` cannot be
grow. It finds the file by the stored names under the lock that webhook checkpointed past what the download reads, so a long download lets the
edits, target edits and target creation hold, and lets go once the file `-wal` grow. It finds the files by the stored names under the lock that
is open: a rename during the download moves the file without affecting webhook edits, target edits and target creation hold, and lets go once
it. 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 Deleting a webhook releases its archives: the delivery engine's cached
archive writers are dropped and their file handles closed, so nothing 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_slack.go # Slack/Mattermost incoming-webhook target
│ │ ├── target_database.go # Database archive target │ │ ├── target_database.go # Database archive target
│ │ ├── target_database_archive.go # Archive file lifecycle and pruning │ │ ├── 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_database_export.go # Archive download as gzipped JSON
│ │ ├── target_log.go # Log target (stdout) │ │ ├── target_log.go # Log target (stdout)
│ │ ├── target_config_view.go # Masked target config for templates │ │ ├── target_config_view.go # Masked target config for templates
+12 -3
View File
@@ -163,6 +163,17 @@ func (env *archiveEnv) seedArchiveRows(
t.Helper() t.Helper()
path := env.archivePath(tgt) 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( sqlDB, err := sql.Open(
"sqlite", fmt.Sprintf("file:%s?mode=rwc", path), "sqlite", fmt.Sprintf("file:%s?mode=rwc", path),
@@ -182,7 +193,7 @@ func (env *archiveEnv) 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: tgt.WebhookID, WebhookID: webhookID,
Method: http.MethodPost, Method: http.MethodPost,
Body: `{"seeded":true}`, Body: `{"seeded":true}`,
ArchivedAt: at, ArchivedAt: at,
@@ -191,8 +202,6 @@ func (env *archiveEnv) seedArchiveRows(
} }
require.NoError(t, sqlDB.Close()) require.NoError(t, sqlDB.Close())
return path
} }
// archivedEventIDs returns the event ids currently stored in an // archivedEventIDs returns the event ids currently stored in an
+7 -6
View File
@@ -284,12 +284,13 @@ func (e *Engine) EvictTarget(targetID string) {
e.dbTarget.evict(targetID) e.dbTarget.evict(targetID)
} }
// Rename implements Archives. It renames a database target's // Rename implements Archives. It renames every one of a database
// archive file to ArchiveFileName(webhookName, targetName, // target's archive files to ArchiveFileName(webhookName, targetName,
// targetID), under the lock the target's archive writes and the // targetID), each keeping the period in its name, under the lock the
// idle sweep take. It never replaces a file: if one already has the // target's archive writes and the idle sweep take. It never replaces
// new name, the error is ErrArchiveNameTaken. The caller renames // a file: if one already has a new name, the error is
// before it saves the new name: see databaseTarget.rename. // ErrArchiveNameTaken. The caller renames before it saves the new
// name: see databaseTarget.rename.
func (e *Engine) Rename( func (e *Engine) Rename(
targetID, webhookName, targetName string, targetID, webhookName, targetName string,
) error { ) error {
+14 -5
View File
@@ -493,23 +493,32 @@ func NewExportArchiveWriter(
return &ExportArchiveWriter{w: w} 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( func (e *ExportArchiveWriter) Write(
row ExportArchivedEvent, expiry time.Duration, row ExportArchivedEvent, expiry time.Duration,
) error { ) 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. // Open opens the archive file, pruning when expiry is positive.
func (e *ExportArchiveWriter) Open(expiry time.Duration) error { 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. // Reopen closes and reopens the archive file.
func (e *ExportArchiveWriter) Reopen( func (e *ExportArchiveWriter) Reopen(
expiry time.Duration, expiry time.Duration,
) error { ) error {
return e.w.reopen(expiry) return e.w.reopen(e.w.path, expiry)
} }
// SetNow replaces the clock the writer measures its reopen // SetNow replaces the clock the writer measures its reopen
@@ -539,7 +548,7 @@ func (e *ExportArchiveWriter) Path() string {
func (e *ExportArchiveWriter) OpenExisting( func (e *ExportArchiveWriter) OpenExisting(
expiry time.Duration, expiry time.Duration,
) error { ) error {
return e.w.openMode(archiveModeExisting, expiry) return e.w.openMode(e.w.path, archiveModeExisting, expiry)
} }
// SweepExpired runs an idle sweep of the archive. // SweepExpired runs an idle sweep of the archive.
+20 -7
View File
@@ -41,6 +41,8 @@ type TargetConfigForm struct {
Timeout string Timeout string
// Expiry is the database (archive) target's row expiry. // Expiry is the database (archive) target's row expiry.
Expiry string Expiry string
// Rotation is the database (archive) target's rotation.
Rotation string
} }
// NewTargetConfigForm parses a target's stored configuration into // NewTargetConfigForm parses a target's stored configuration into
@@ -85,11 +87,13 @@ func NewTargetConfigForm(
} }
} }
// databaseConfigForm parses an archive target's optional expiry. // databaseConfigForm parses an archive target's optional expiry and
// An absent, empty or never expiry yields an empty expiry, on which // rotation. An absent, empty or never expiry yields an empty expiry,
// the edit form starts at never; saving it unchanged stores never, // on which the edit form starts at never; saving it unchanged stores
// which means the same as an empty expiry. An expiry that is set // never, which means the same as an empty expiry. An absent rotation
// but not a valid duration is an error, not a blank field. // 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( func databaseConfigForm(
configJSON string, configJSON string,
) (TargetConfigForm, error) { ) (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 { if cfg.Expiry == "" || cfg.Expiry == archiveExpiryNever {
return TargetConfigForm{}, nil return form, nil
} }
err = ValidateArchiveExpiry(cfg.Expiry) err = ValidateArchiveExpiry(cfg.Expiry)
@@ -115,5 +126,7 @@ func databaseConfigForm(
return TargetConfigForm{}, err return TargetConfigForm{}, err
} }
return TargetConfigForm{Expiry: cfg.Expiry}, nil form.Expiry = cfg.Expiry
return form, nil
} }
+11 -2
View File
@@ -192,8 +192,9 @@ func maxRetriesField(t *database.Target) ConfigField {
// databaseConfigFields describes an archive target by its // databaseConfigFields describes an archive target by its
// expiry in plain units, such as "30 days", or "never" when // expiry in plain units, such as "30 days", or "never" when
// the archive is kept forever. An expiry that is set but not // the archive is kept forever, and by its rotation. An expiry
// a valid duration is reported as unavailable rather than // 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. // echoed back.
func databaseConfigFields(configJSON string) []ConfigField { func databaseConfigFields(configJSON string) []ConfigField {
expiry, err := parseArchiveExpiry(configJSON) expiry, err := parseArchiveExpiry(configJSON)
@@ -201,6 +202,11 @@ func databaseConfigFields(configJSON string) []ConfigField {
return unavailableConfigFields() return unavailableConfigFields()
} }
rotation, err := parseArchiveRotation(configJSON)
if err != nil {
return unavailableConfigFields()
}
value := archiveExpiryNever value := archiveExpiryNever
if expiry > 0 { if expiry > 0 {
value = plainDuration(expiry) value = plainDuration(expiry)
@@ -209,6 +215,9 @@ func databaseConfigFields(configJSON string) []ConfigField {
return []ConfigField{{ return []ConfigField{{
Label: "Archive Expiry", Label: "Archive Expiry",
Value: value, Value: value,
}, {
Label: "Archive Rotation",
Value: rotation,
}} }}
} }
+41 -1
View File
@@ -328,13 +328,49 @@ func TestNewTargetViews_Database(t *testing.T) {
assert.Equal( assert.Equal(
t, t,
map[string]string{"Archive Expiry": tc.want}, map[string]string{
"Archive Expiry": tc.want,
"Archive Rotation": rotationNone,
},
fieldMap(view.Config), 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) { func TestNewTargetViews_Log(t *testing.T) {
t.Parallel() t.Parallel()
@@ -382,6 +418,10 @@ func TestNewTargetViews_Unpresentable(t *testing.T) {
Type: database.TargetTypeDatabase, Type: database.TargetTypeDatabase,
Config: `{"expiry":"a fortnight"}`, Config: `{"expiry":"a fortnight"}`,
}, },
"invalid archive rotation": {
Type: database.TargetTypeDatabase,
Config: weeklyConfig,
},
} }
for name, target := range tests { for name, target := range tests {
+26 -11
View File
@@ -20,7 +20,8 @@ const archiveNameMaxLen = 40
// from the per-webhook event database. The event is already // from the per-webhook event database. The event is already
// persisted in the per-webhook event DB by the time delivery runs; // persisted in the per-webhook event DB by the time delivery runs;
// the database target additionally writes a durable long-term copy // 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 // attempt whose outcome reflects whether the archive write
// succeeded. See archiveWriter for the close/reopen, auto-recreate, // succeeded. See archiveWriter for the close/reopen, auto-recreate,
// and expiry semantics. // and expiry semantics.
@@ -146,8 +147,9 @@ func (t *databaseTarget) Deliver(
} }
// archive writes the full event as a row into the target'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 and
// parsed from the target config JSON. // 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 { func (t *databaseTarget) archive(d *database.Delivery) error {
webhookID := d.Event.WebhookID webhookID := d.Event.WebhookID
if webhookID == "" { if webhookID == "" {
@@ -159,6 +161,19 @@ func (t *databaseTarget) archive(d *database.Delivery) error {
return err 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) w, err := t.writerFor(d.TargetID)
if err != nil { if err != nil {
return err return err
@@ -174,7 +189,7 @@ func (t *databaseTarget) archive(d *database.Delivery) error {
ContentType: d.Event.ContentType, 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, // writerFor returns the archive writer for a database target,
@@ -275,10 +290,10 @@ func (t *databaseTarget) releaseSweepWriter(
delete(t.writers, targetID) delete(t.writers, targetID)
} }
// newWriter builds the writer for a database target's archive. The // newWriter builds the writer for a database target's archive. Its
// file is the one ArchivePath gives for the webhook and the target as // 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 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( func (t *databaseTarget) newWriter(
targetID string, targetID string,
) (*archiveWriter, error) { ) (*archiveWriter, error) {
@@ -306,10 +321,10 @@ func (t *databaseTarget) newWriter(
return w, nil return w, nil
} }
// rename moves a database target's archive file to the name for // rename moves every one of a database target's archive files to the
// webhookName and targetName. It goes through the target's writer, // name for webhookName and targetName. It goes through the target's
// so the move holds the lock that writes and the idle sweep take, // writer, so the move holds the lock that writes and the idle sweep
// and later writes use the new name. // take, and later writes use the new name.
// //
// The writer is created if there is none, and it stays cached. The // 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 // handlers rename before they save the new name, so until the save
+184 -73
View File
@@ -85,6 +85,10 @@ type databaseTargetConfig struct {
// archived rows are pruned, or "never" (the default) to // archived rows are pruned, or "never" (the default) to
// keep them forever. // keep them forever.
Expiry string `json:"expiry"` 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 // archivedEvent is one fully captured webhook event stored in a
@@ -178,7 +182,7 @@ func ValidateArchiveExpiry(expiry string) error {
return nil 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 // 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
@@ -187,7 +191,15 @@ func ValidateArchiveExpiry(expiry string) error {
// on every open. // on every open.
type archiveWriter struct { type archiveWriter struct {
mu sync.Mutex 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 path string
// current is the file db is open on.
current string
log *slog.Logger log *slog.Logger
debounce time.Duration debounce time.Duration
db *gorm.DB db *gorm.DB
@@ -236,12 +248,14 @@ func newArchiveWriter(
} }
} }
// write appends the event as a row, then applies the debounced // write appends the event as a row to the archive file for period
// close/reopen. It recreates the archive file if it was moved // (see archivePeriodPath), then applies the debounced close/reopen.
// or removed since the last open. A positive expiry prunes rows // When period names a different file from the one open, the open one
// older than it on each (re)open. // 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( func (w *archiveWriter) write(
row archivedEvent, expiry time.Duration, row archivedEvent, expiry time.Duration, period string,
) error { ) error {
w.mu.Lock() w.mu.Lock()
defer w.mu.Unlock() defer w.mu.Unlock()
@@ -252,8 +266,10 @@ func (w *archiveWriter) write(
) )
} }
if w.db == nil || !fileExists(w.path) { file := archivePeriodPath(w.path, period)
err := w.reopen(expiry)
if w.db == nil || w.current != file || !fileExists(file) {
err := w.reopen(file, expiry)
if err != nil { if err != nil {
return err return err
} }
@@ -264,41 +280,41 @@ func (w *archiveWriter) write(
err := w.db.Create(&row).Error err := w.db.Create(&row).Error
if err != nil { if err != nil {
return fmt.Errorf( 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 { if w.now().Sub(w.lastReopen) >= w.debounce {
return w.reopen(expiry) return w.reopen(file, expiry)
} }
return nil 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 // its schema, records the reopen time, and prunes expired rows
// when expiry is positive. // when expiry is positive.
func (w *archiveWriter) open(expiry time.Duration) error { func (w *archiveWriter) open(file string, expiry time.Duration) error {
return w.openMode(archiveModeCreate, expiry) 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 // mode, migrates its schema, records the reopen time, and
// prunes expired rows when expiry is positive. The write path // prunes expired rows when expiry is positive. The write path
// passes archiveModeCreate so a missing file is recreated; the // passes archiveModeCreate so a missing file is recreated; the
// idle sweep passes archiveModeExisting so a missing file is an // idle sweep passes archiveModeExisting so a missing file is an
// error rather than a newly conjured empty archive. // error rather than a newly conjured empty archive.
func (w *archiveWriter) openMode( func (w *archiveWriter) openMode(
mode string, expiry time.Duration, file, mode string, expiry time.Duration,
) error { ) error {
// Opened through database.OpenSQLite so an archive file carries // Opened through database.OpenSQLite so an archive file carries
// the same WAL journaling, busy timeout, immediate-transaction // the same WAL journaling, busy timeout, immediate-transaction
// locking, and pool bounds as every other database file. See // locking, and pool bounds as every other database file. See
// internal/database/sqlite_open.go. // internal/database/sqlite_open.go.
sqlDB, err := database.OpenSQLite(w.path, mode) sqlDB, err := database.OpenSQLite(file, mode)
if err != nil { if err != nil {
return fmt.Errorf( 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( return fmt.Errorf(
"connecting to archive database %s: %w", "connecting to archive database %s: %w",
w.path, err, file, err,
) )
} }
@@ -323,11 +339,12 @@ func (w *archiveWriter) openMode(
_ = sqlDB.Close() _ = sqlDB.Close()
return fmt.Errorf( return fmt.Errorf(
"migrating archive database %s: %w", w.path, err, "migrating archive database %s: %w", file, err,
) )
} }
w.db = gdb w.db = gdb
w.current = file
w.lastReopen = w.now() w.lastReopen = w.now()
w.reopens++ w.reopens++
@@ -338,12 +355,12 @@ func (w *archiveWriter) openMode(
return nil 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. // 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() w.close()
return w.open(expiry) return w.open(file, expiry)
} }
// close closes the underlying handle, if any. // close closes the underlying handle, if any.
@@ -360,16 +377,17 @@ func (w *archiveWriter) close() {
w.db = nil w.db = nil
} }
// sweepExpired prunes an archive that may have gone idle, with // sweepExpired prunes the target's archive files, which may have
// no write to trigger the usual on-reopen prune. It takes the // gone idle, with no write to trigger the usual on-reopen prune. It
// writer's own mutex for the whole operation, so a sweep is // takes the writer's own mutex for the whole operation, so a sweep
// ordered against concurrent writes rather than reaching around // is ordered against concurrent writes rather than reaching around
// them to the file. // them to the files.
// //
// It never creates the archive file: a missing file is skipped, // It never creates an archive file: it prunes only the files
// and the reopen uses archiveModeExisting so SQLite itself // archiveFiles finds, and opens each with archiveModeExisting so
// refuses to create one if the file disappears between the // SQLite itself refuses to create one if the file disappears
// check and the open. // 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 // The archive is left CLOSED afterwards. An idle archive holding
// no handle is what keeps the operator's move-the-file-away // 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) { files, err := archiveFiles(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)
if err != nil { if err != nil {
return err 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() 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 return nil
} }
// rename gives the archive file a new name in the same directory, // rename gives every one of the target's archive files the new
// and the writer uses the file under that name from now on. The // name, keeping the period in the name of each (see
// handle is closed first, which folds the -wal into the .db; any // archivePeriodPath), and the writer uses the files under that name
// -wal or -shm still beside the file (left by a crash) is moved with // from now on. The handle is closed first, which folds the -wal into
// it, because SQLite finds them by name. A missing file is not an // the .db; any -wal or -shm still beside a file (left by a crash) is
// error: the operator may have moved it away, and the next write // moved with it, because SQLite finds them by name. A target with no
// creates it under the new name. // 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 // If a file already has one of the new names, nothing is moved and
// error is ErrArchiveNameTaken. If one file fails to move, those // the error is ErrArchiveNameTaken. If one file fails to move, those
// already moved are moved back before the error is returned, so the // already moved are moved back before the error is returned, so the
// archive is never split across two names. // archive is never split across two names.
func (w *archiveWriter) rename(name string) error { func (w *archiveWriter) rename(name string) error {
@@ -430,38 +498,53 @@ func (w *archiveWriter) rename(name string) error {
return nil return nil
} }
suffixes := []string{"", "-wal", "-shm"} files, err := archiveFiles(w.path)
if err != nil {
return err
}
for _, suffix := range suffixes { // from[i] moves to to[i].
if fileExists(path + suffix) { 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( return fmt.Errorf(
"%w: %s", ErrArchiveNameTaken, name+suffix, "%w: %s", ErrArchiveNameTaken, filepath.Base(taken),
) )
} }
} }
w.close() w.close()
for i, suffix := range suffixes { for i := range from {
err := os.Rename(w.path+suffix, path+suffix) err = os.Rename(from[i], to[i])
if err == nil || errors.Is(err, fs.ErrNotExist) { if err == nil || errors.Is(err, fs.ErrNotExist) {
continue continue
} }
for _, moved := range suffixes[:i] { for j := range i {
backErr := os.Rename(path+moved, w.path+moved) backErr := os.Rename(to[j], from[j])
if backErr != nil && !errors.Is(backErr, fs.ErrNotExist) { if backErr != nil && !errors.Is(backErr, fs.ErrNotExist) {
w.log.Error( w.log.Error(
"failed to move archive file back", "failed to move archive file back",
"from", path+moved, "from", to[j],
"to", w.path+moved, "to", from[j],
"error", backErr, "error", backErr,
) )
} }
} }
return fmt.Errorf( 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 { if res.Error != nil {
w.log.Error( w.log.Error(
"failed to prune expired archive rows", "failed to prune expired archive rows",
"path", w.path, "path", w.current,
"error", res.Error, "error", res.Error,
) )
@@ -509,38 +592,59 @@ func (w *archiveWriter) prune(expiry time.Duration) {
if res.RowsAffected > 0 { if res.RowsAffected > 0 {
w.log.Info( w.log.Info(
"pruned expired archive rows", "pruned expired archive rows",
"path", w.path, "path", w.current,
"rows_deleted", res.RowsAffected, "rows_deleted", res.RowsAffected,
) )
} }
} }
// ArchiveFileInfo is what the metadata of a database target's archive // ArchiveFileInfo is what the metadata of a database target's archive
// file says about it. // files says about them.
type ArchiveFileInfo struct { 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 Size int64
// Written is when the file or its -wal was last modified, whichever // Written is when a file or a -wal was last modified, whichever is
// is later: a write lands in the -wal first. // latest: a write lands in the -wal first.
Written time.Time Written time.Time
} }
// StatArchive reads the metadata of the archive file at path and of // StatArchive reads the metadata of a database target's archive
// its -wal, without opening the archive. With no file at path, which is // files, given the path ArchivePath gives it (see archiveFiles), and
// so before the first write and after the operator moved it away, the // 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. // error wraps fs.ErrNotExist.
func StatArchive(path string) (ArchiveFileInfo, error) { func StatArchive(path string) (ArchiveFileInfo, error) {
file, err := os.Stat(path) files, err := archiveFiles(path)
if err != nil { if err != nil {
return ArchiveFileInfo{}, err return ArchiveFileInfo{}, err
} }
info := ArchiveFileInfo{Size: file.Size(), Written: file.ModTime()} var info ArchiveFileInfo
wal, err := os.Stat(path + "-wal") for _, file := range files {
db, err := os.Stat(file.path)
if errors.Is(err, fs.ErrNotExist) { if errors.Is(err, fs.ErrNotExist) {
return info, nil 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 { if err != nil {
@@ -552,6 +656,13 @@ func StatArchive(path string) (ArchiveFileInfo, error) {
if wal.ModTime().After(info.Written) { if wal.ModTime().After(info.Written) {
info.Written = wal.ModTime() info.Written = wal.ModTime()
} }
}
if info.Files == 0 {
return ArchiveFileInfo{}, fmt.Errorf(
"no archive file for %s: %w", path, fs.ErrNotExist,
)
}
return info, nil return info, nil
} }
+132 -48
View File
@@ -6,6 +6,7 @@ import (
"database/sql" "database/sql"
"encoding/base64" "encoding/base64"
"encoding/json" "encoding/json"
"errors"
"fmt" "fmt"
"io" "io"
"log/slog" "log/slog"
@@ -50,21 +51,29 @@ func ArchiveExportFileName(
at.UTC().Format("20060102T150405Z") + ".json.gz" at.UTC().Format("20060102T150405Z") + ".json.gz"
} }
// ArchiveExport is a database target's archive opened for download. // ArchiveExport is a database target's archive opened for download:
// It reads the file on its own connection, inside one read-only // every one of its files, each read on its own connection inside one
// transaction, so it writes out the archive as it stood when // read-only transaction, so it writes out the archive as it stood when
// OpenArchiveExport returned. // OpenArchiveExport returned.
// //
// Archives are in WAL mode, where a reader works from a snapshot and // 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, // 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 // and the export does not see them. SQLite cannot checkpoint a -wal
// past an open snapshot, so the -wal grows until the export is closed. // past an open snapshot, so a file's -wal grows until the export has
// written out that file.
type ArchiveExport struct { type ArchiveExport struct {
files []*exportFile
}
// exportFile is one archive file opened for an export.
type exportFile struct {
db *sql.DB db *sql.DB
tx *gorm.DB tx *gorm.DB
// empty is true when there is nothing to read: no file, or a file // period is the period in the file's name, "" for none.
// without the archive's table yet. period string
// empty is true for a file without the archive's table yet.
empty bool empty bool
} }
@@ -74,26 +83,49 @@ type exportedName struct {
Name string `json:"name"` Name string `json:"name"`
} }
// OpenArchiveExport opens the archive file at path for export and // OpenArchiveExport opens every one of a database target's archive
// takes the snapshot the export reads. It never creates the file: with // files for export, given the path ArchivePath gives it (see
// no file at path, the export has no rows. // 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 // Once it has returned, the files are open, so a rename or a move of
// does not affect the export, which reads the same file under its new // them does not affect the export, which reads the same files under
// name. // 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. // whole export.
func OpenArchiveExport( func OpenArchiveExport(
ctx context.Context, path string, log *slog.Logger, ctx context.Context, path string, log *slog.Logger,
) (*ArchiveExport, error) { ) (*ArchiveExport, error) {
if !fileExists(path) { files, err := archiveFiles(path)
return &ArchiveExport{empty: true}, nil 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 { if err != nil {
return nil, fmt.Errorf("opening archive %s: %w", path, err) _ = 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", file.path, err)
} }
gdb, err := gorm.Open( gdb, err := gorm.Open(
@@ -106,7 +138,7 @@ func OpenArchiveExport(
if err != nil { if err != nil {
_ = db.Close() _ = 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 // ReadOnly makes the driver begin a deferred transaction in place
@@ -116,7 +148,9 @@ func OpenArchiveExport(
if tx.Error != nil { if tx.Error != nil {
_ = db.Close() _ = 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. // The transaction's first read is what takes the snapshot.
@@ -127,22 +161,27 @@ func OpenArchiveExport(
_ = tx.Rollback() _ = tx.Rollback()
_ = db.Close() _ = 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: // WriteGzipJSON writes the export to w as one gzipped JSON object:
// webhook and target, each an id and a name; exported_at; and // webhook and target, each an id and a name; exported_at; and
// archived_events, one object per archived row, keyed by column name. // 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 // the files in the order archiveFiles lists them. A row from a file
// written in base64, with "body_encoding": "base64" beside it. // 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 // 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 // archive nor its JSON is ever held in memory whole, and each file is
// the gzip stream is left unfinished, so what was written does not // closed once its rows are written. After an error the gzip stream is
// decompress as a whole file. // left unfinished, so what was written does not decompress as a whole
// file.
func (x *ArchiveExport) WriteGzipJSON( func (x *ArchiveExport) WriteGzipJSON(
ctx context.Context, ctx context.Context,
w io.Writer, w io.Writer,
@@ -169,15 +208,30 @@ func (x *ArchiveExport) WriteGzipJSON(
return zw.Close() 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 { 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 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, // writeJSON writes head with archived_events added as its last key,
@@ -207,46 +261,72 @@ func (x *ArchiveExport) writeJSON(
return err return err
} }
// writeRows writes each archived row to w, oldest first, one per line, // writeRows writes the archived rows of each file to w, one per line,
// separated by commas. // separated by commas, and closes each file once its rows are written.
func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error { func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error {
if x.empty { sep := "\n"
return nil
}
rows, err := x.tx.WithContext(ctx). for _, f := range x.files {
Model(&archivedEvent{}).Order("id").Rows() var err error
sep, err = f.writeRows(ctx, w, sep)
if err != nil { if err != nil {
return err return err
} }
err = f.close()
if err != nil {
return err
}
}
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
}
defer func() { _ = rows.Close() }() defer func() { _ = rows.Close() }()
for sep := "\n"; rows.Next(); sep = ",\n" { for ; rows.Next(); sep = ",\n" {
var ev archivedEvent var ev archivedEvent
err = x.tx.ScanRows(rows, &ev) err = f.tx.ScanRows(rows, &ev)
if err != nil { if err != nil {
return err return "", err
} }
_, err = io.WriteString(w, sep) _, err = io.WriteString(w, sep)
if err != nil { if err != nil {
return err return "", err
} }
err = writeRow(w, &ev) err = writeRow(w, &ev, f.period)
if err != nil { 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 // 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. // column name, its body in base64 when it is not valid UTF-8, with
func writeRow(w io.Writer, ev *archivedEvent) error { // the period of its file beside them unless that is "".
func writeRow(w io.Writer, ev *archivedEvent, period string) error {
row := map[string]any{ row := map[string]any{
"id": ev.ID, "id": ev.ID,
"event_id": ev.EventID, "event_id": ev.EventID,
@@ -264,6 +344,10 @@ func writeRow(w io.Writer, ev *archivedEvent) error {
row["body_encoding"] = "base64" row["body_encoding"] = "base64"
} }
if period != "" {
row["period"] = period
}
line, err := json.Marshal(row) line, err := json.Marshal(row)
if err != nil { if err != nil {
return err return err
@@ -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 // 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 // the largest heap it saw at a write. It collects garbage before each
// reading, so the heap it reads is what is still held. // reading, so the heap it reads is what is still held.
@@ -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
}
@@ -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)
}
}
@@ -218,6 +218,7 @@ func TestStatArchive(t *testing.T) {
got, err := delivery.StatArchive(path) got, err := delivery.StatArchive(path)
require.NoError(t, err) require.NoError(t, err)
assert.Equal(t, 1, got.Files)
assert.Equal(t, file.Size()+wal.Size(), got.Size) assert.Equal(t, file.Size()+wal.Size(), got.Size)
assert.True(t, written.Equal(got.Written), got.Written) assert.True(t, written.Equal(got.Written), got.Written)
+6 -1
View File
@@ -204,10 +204,11 @@ func TestNewTargetConfigForm(t *testing.T) {
form, err = delivery.NewTargetConfigForm(&database.Target{ form, err = delivery.NewTargetConfigForm(&database.Target{
Type: database.TargetTypeDatabase, Type: database.TargetTypeDatabase,
Config: `{"expiry":"720h"}`, Config: `{"expiry":"720h","rotation":"daily"}`,
}) })
require.NoError(t, err) require.NoError(t, err)
assert.Equal(t, "720h", form.Expiry) assert.Equal(t, "720h", form.Expiry)
assert.Equal(t, rotationDaily, form.Rotation)
form, err = delivery.NewTargetConfigForm(&database.Target{ form, err = delivery.NewTargetConfigForm(&database.Target{
Type: database.TargetTypeLog, Type: database.TargetTypeLog,
@@ -248,6 +249,10 @@ func TestNewTargetConfigForm_UnreadableConfigErrors(t *testing.T) {
Type: database.TargetTypeDatabase, Type: database.TargetTypeDatabase,
Config: `{"expiry":"soon"}`, Config: `{"expiry":"soon"}`,
}, },
{
Type: database.TargetTypeDatabase,
Config: weeklyConfig,
},
{Type: database.TargetType("nope")}, {Type: database.TargetType("nope")},
} }
+9 -9
View File
@@ -10,10 +10,10 @@ const (
tmplKeyArchiveExpiryChoices = "ArchiveExpiryChoices" tmplKeyArchiveExpiryChoices = "ArchiveExpiryChoices"
) )
// archiveExpiryChoice is one entry of a database target's archive // archiveChoice is one entry of a database target's archive expiry
// expiry select: the expiry stored, the label shown, and whether the // or archive rotation select: the value stored, the label shown, and
// select starts on it. // whether the select starts on it.
type archiveExpiryChoice struct { type archiveChoice struct {
Value string Value string
Label string Label string
Selected bool Selected bool
@@ -21,8 +21,8 @@ type archiveExpiryChoice struct {
// archiveExpiryChoices lists the archive expiries offered by the new // archiveExpiryChoices lists the archive expiries offered by the new
// webhook page, the add target form and the target edit form. // webhook page, the add target form and the target edit form.
func archiveExpiryChoices() []archiveExpiryChoice { func archiveExpiryChoices() []archiveChoice {
return []archiveExpiryChoice{ return []archiveChoice{
{Value: archiveExpiryNever, Label: archiveExpiryNever}, {Value: archiveExpiryNever, Label: archiveExpiryNever},
{Value: "1h", Label: "1h"}, {Value: "1h", Label: "1h"},
{Value: "12h", Label: "12h"}, {Value: "12h", Label: "12h"},
@@ -37,7 +37,7 @@ func archiveExpiryChoices() []archiveExpiryChoice {
// empty expiry selects never. An expiry that is not one of the // empty expiry selects never. An expiry that is not one of the
// choices comes first as its own selected entry, so saving the form // choices comes first as its own selected entry, so saving the form
// unchanged keeps it. // unchanged keeps it.
func archiveExpiryOptions(expiry string) []archiveExpiryChoice { func archiveExpiryOptions(expiry string) []archiveChoice {
if expiry == "" { if expiry == "" {
expiry = archiveExpiryNever 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...)
} }
+22 -3
View File
@@ -5,6 +5,7 @@ import (
"net/http/httptest" "net/http/httptest"
"net/url" "net/url"
"regexp" "regexp"
"strings"
"testing" "testing"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
@@ -49,21 +50,39 @@ func expiryShown(
} }
// expirySelected returns the target edit page and the expiries its // expirySelected returns the target edit page and the expiries its
// select starts on. // expiry select starts on.
func expirySelected( func expirySelected(
t *testing.T, env *sourceTestEnv, webhookID, targetID string, t *testing.T, env *sourceTestEnv, webhookID, targetID string,
) (string, []string) { ) (string, []string) {
t.Helper() 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( w := serveTarget(
env, http.MethodGet, env, http.MethodGet,
"/hook/"+webhookID+"/targets/"+targetID+"/edit", nil, "/hook/"+webhookID+"/targets/"+targetID+"/edit", nil,
) )
require.Equal(t, http.StatusOK, w.Code) require.Equal(t, http.StatusOK, w.Code)
page := w.Body.String() return w.Body.String()
}
return page, matched(`<option value="([^"]*)" selected>`, page) // selectedIn returns the values the select named name on page starts
// on.
func selectedIn(page, name string) []string {
_, rest, _ := strings.Cut(page, `<select id="`+name+`" name="`+name+`"`)
options, _, _ := strings.Cut(rest, "</select>")
return matched(`<option value="([^"]*)" selected>`, options)
} }
// TestArchiveExpiryChoices adds a database target with each archive // TestArchiveExpiryChoices adds a database target with each archive
+42
View File
@@ -0,0 +1,42 @@
package handlers
const (
// archiveRotationNone is the archive rotation that keeps a
// database target's archive in one file. A stored empty rotation
// means the same.
archiveRotationNone = "none"
// tmplKeyArchiveRotationChoices is the template data key for the
// entries of a page's archive rotation select.
tmplKeyArchiveRotationChoices = "ArchiveRotationChoices"
)
// archiveRotationChoices lists the archive rotations offered by the
// new webhook page, the add target form and the target edit form.
func archiveRotationChoices() []archiveChoice {
return []archiveChoice{
{Value: archiveRotationNone, Label: archiveRotationNone},
{Value: "monthly", Label: "monthly"},
{Value: "daily", Label: "daily"},
{Value: "hourly", Label: "hourly"},
}
}
// archiveRotationOptions returns the choices with rotation selected;
// an empty rotation, or one that is not a choice, selects none. A
// stored rotation is always a choice: the forms refuse any other.
func archiveRotationOptions(rotation string) []archiveChoice {
options := archiveRotationChoices()
for i := range options {
if options[i].Value == rotation {
options[i].Selected = true
return options
}
}
options[0].Selected = true
return options
}
+229
View File
@@ -0,0 +1,229 @@
package handlers_test
import (
"net/http"
"net/http/httptest"
"net/url"
"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"
)
// rotationNone is the archive rotation that keeps one file.
const rotationNone = "none"
// rotationShown returns the archive rotations the webhook page's
// target list shows.
func rotationShown(
t *testing.T, env *sourceTestEnv, webhookID string,
) []string {
t.Helper()
return matched(
`Archive Rotation:</span>\s*<span>([^<]*)</span>`,
renderedPage(t, env, webhookID),
)
}
// TestArchiveRotationChoices adds a database target with each archive
// rotation the forms offer, and checks that it is stored as chosen,
// shown in the target list, and that the target edit form starts on
// it. It then edits the target to hourly, keeping its expiry.
func TestArchiveRotationChoices(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
for _, rotation := range []string{rotationNone, "monthly", "daily", "hourly"} {
t.Run(rotation, func(t *testing.T) {
t.Parallel()
webhook := seedWebhookWithRetention(t, env.db, 30)
form := url.Values{}
form.Set("name", "archive")
form.Set("type", string(database.TargetTypeDatabase))
form.Set("expiry", "720h")
form.Set("rotation", rotation)
w := serveTarget(
env, http.MethodPost, "/hook/"+webhook.ID+"/targets", form,
)
require.Equal(t, http.StatusSeeOther, w.Code, w.Body.String())
targets := targetsForWebhook(t, env.db, webhook.ID)
require.Len(t, targets, 1)
assert.JSONEq(t,
`{"expiry":"720h","rotation":"`+rotation+`"}`,
targets[0].Config,
)
assert.Equal(t,
[]string{rotation}, rotationShown(t, env, webhook.ID))
page := targetEditPage(t, env, webhook.ID, targets[0].ID)
assert.Equal(t, []string{rotation}, selectedIn(page, "rotation"))
form = url.Values{}
form.Set("name", "archive")
form.Set("expiry", "720h")
form.Set("rotation", "hourly")
w = submitTargetEdit(env, webhook.ID, targets[0].ID, form)
require.Equal(t, http.StatusSeeOther, w.Code, w.Body.String())
assert.JSONEq(t,
`{"expiry":"720h","rotation":"hourly"}`,
storedTarget(t, env, targets[0].ID).Config,
)
})
}
}
// TestArchiveRotationEditStartsOnNone checks the edit form of a
// database target with no rotation stored starts on none.
func TestArchiveRotationEditStartsOnNone(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
webhook := seedWebhookWithRetention(t, env.db, 30)
target := seedConfiguredTarget(
t, env.db, webhook.ID, database.TargetTypeDatabase, "",
)
page := targetEditPage(t, env, webhook.ID, target.ID)
assert.Equal(t, []string{rotationNone}, selectedIn(page, "rotation"))
assert.Equal(t, []string{rotationNone}, rotationShown(t, env, webhook.ID))
}
// TestArchiveRotationRefused proves a rotation that is not one of the
// four is refused on the add target form and the target edit form, and
// that nothing is stored.
func TestArchiveRotationRefused(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
webhook := seedWebhookWithRetention(t, env.db, 30)
form := url.Values{}
form.Set("name", "archive")
form.Set("type", string(database.TargetTypeDatabase))
form.Set("rotation", "weekly")
w := serveTarget(
env, http.MethodPost, "/hook/"+webhook.ID+"/targets", form,
)
assert.Equal(t, http.StatusBadRequest, w.Code)
assert.Contains(t, w.Body.String(), "Invalid archive rotation")
assert.Empty(t, targetsForWebhook(t, env.db, webhook.ID))
target := seedConfiguredTarget(
t, env.db, webhook.ID, database.TargetTypeDatabase,
`{"rotation":"daily"}`,
)
form.Del("type")
w = submitTargetEdit(env, webhook.ID, target.ID, form)
assert.Equal(t, http.StatusBadRequest, w.Code)
assert.Contains(t, w.Body.String(), "Invalid archive rotation")
assert.JSONEq(t,
`{"rotation":"daily"}`, storedTarget(t, env, target.ID).Config,
)
}
// TestHandleSourceCreateSubmit_ArchiveRotation proves the new webhook
// page's archive rotation is stored on the archive target it creates.
func TestHandleSourceCreateSubmit_ArchiveRotation(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
form := url.Values{}
form.Set("name", "rotated")
form.Set("archive", "on")
form.Set("archive_expiry", "720h")
form.Set("archive_rotation", "daily")
w := submitCreateForm(env, form)
require.Equal(t, http.StatusSeeOther, w.Code, w.Body.String())
var webhook database.Webhook
require.NoError(t, env.db.DB().
Where("name = ?", "rotated").First(&webhook).Error)
targets := targetsForWebhook(t, env.db, webhook.ID)
require.Len(t, targets, 1)
assert.JSONEq(t,
`{"expiry":"720h","rotation":"daily"}`, targets[0].Config,
)
}
// TestArchiveFileView_Rotated describes a daily target's archive files
// at two times. On a day that has a file, the view names that file;
// on the next, before any event, it names the file the next event
// will go to, not created yet. Both times the size is of every file
// together and the last write the latest of them.
func TestArchiveFileView_Rotated(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
webhook := seedWebhookWithRetention(t, env.db, 30)
target := seedConfiguredTarget(
t, env.db, webhook.ID, database.TargetTypeDatabase,
`{"rotation":"daily"}`,
)
path := delivery.ArchivePath(env.dbMgr, &webhook, target)
stem := strings.TrimSuffix(path, ".db")
written := time.Date(2026, 10, 2, 9, 0, 0, 0, time.UTC)
for i, day := range []string{"2026-10-01", "2026-10-02"} {
file := stem + "-" + day + ".db"
require.NoError(t, os.WriteFile(file, make([]byte, 1000), 0o600))
at := written.Add(time.Duration(i-1) * 24 * time.Hour)
require.NoError(t, os.Chtimes(file, at, at))
}
view := env.handlers.ArchiveFileViewForTest(
&webhook, target, time.Date(2026, 10, 2, 23, 0, 0, 0, time.UTC),
)
assert.Equal(t, filepath.Base(stem)+"-2026-10-02.db", view.Name)
assert.Empty(t, view.Note)
assert.Equal(t, 2, view.Files)
assert.Equal(t, "2.0 kB", view.Size)
assert.Equal(t, "2026-10-02 09:00:00 UTC", view.WrittenUTC)
view = env.handlers.ArchiveFileViewForTest(
&webhook, target, time.Date(2026, 10, 3, 0, 0, 0, 0, time.UTC),
)
assert.Equal(t, filepath.Base(stem)+"-2026-10-03.db", view.Name)
assert.Equal(t, "not created yet", view.Note)
assert.Equal(t, 2, view.Files)
assert.Equal(t, "2.0 kB", view.Size)
page := targetList(t, renderedPage(t, env, webhook.ID))
assert.Contains(t, page, "Archive Size: 2.0 kB in 2 files")
}
// renderedPage returns the webhook page.
func renderedPage(t *testing.T, env *sourceTestEnv, webhookID string) string {
t.Helper()
w := httptest.NewRecorder()
env.handlers.HandleSourceDetail().ServeHTTP(w, getRequest(
t, "/hook/"+webhookID, env.cookies,
map[string]string{sourceIDParam: webhookID},
))
require.Equal(t, http.StatusOK, w.Code)
return w.Body.String()
}
+11 -2
View File
@@ -167,7 +167,16 @@ func (s *Handlers) BuildHTTPTargetConfigForTest(
// buildDatabaseTargetConfig for use in the handlers_test // buildDatabaseTargetConfig for use in the handlers_test
// package. // package.
func BuildDatabaseTargetConfigForTest( func BuildDatabaseTargetConfigForTest(
expiry string, expiry, rotation string,
) (string, string, error) { ) (string, string, error) {
return buildDatabaseTargetConfig(expiry) return buildDatabaseTargetConfig(expiry, rotation)
}
// ArchiveFileViewForTest exposes archiveFileView, which describes a
// database target's archive files as the target list shows them at
// now.
func (s *Handlers) ArchiveFileViewForTest(
webhook *database.Webhook, target *database.Target, now time.Time,
) *ArchiveFileView {
return s.archiveFileView(webhook, target, now)
} }
+36 -5
View File
@@ -436,23 +436,37 @@ func TestRenderTemplateMidRenderErrorSendsNoPartialBody(t *testing.T) {
func TestBuildDatabaseTargetConfig_Valid(t *testing.T) { func TestBuildDatabaseTargetConfig_Valid(t *testing.T) {
t.Parallel() t.Parallel()
// Empty expiry: the keep-forever default, empty config. // Empty expiry and rotation: the keep-forever, one-file default,
cfg, errMsg, err := handlers.BuildDatabaseTargetConfigForTest("") // empty config.
cfg, errMsg, err := handlers.BuildDatabaseTargetConfigForTest("", "")
require.NoError(t, err) require.NoError(t, err)
assert.Empty(t, errMsg) assert.Empty(t, errMsg)
assert.Empty(t, cfg) assert.Empty(t, cfg)
// Explicit never is stored as config. // Explicit never is stored as config.
cfg, errMsg, err = handlers.BuildDatabaseTargetConfigForTest("never") cfg, errMsg, err = handlers.BuildDatabaseTargetConfigForTest("never", "")
require.NoError(t, err) require.NoError(t, err)
assert.Empty(t, errMsg) assert.Empty(t, errMsg)
assert.JSONEq(t, `{"expiry":"never"}`, cfg) assert.JSONEq(t, `{"expiry":"never"}`, cfg)
// A positive duration is stored as config. // A positive duration is stored as config.
cfg, errMsg, err = handlers.BuildDatabaseTargetConfigForTest("720h") cfg, errMsg, err = handlers.BuildDatabaseTargetConfigForTest("720h", "")
require.NoError(t, err) require.NoError(t, err)
assert.Empty(t, errMsg) assert.Empty(t, errMsg)
assert.JSONEq(t, `{"expiry":"720h"}`, cfg) assert.JSONEq(t, `{"expiry":"720h"}`, cfg)
// A rotation is stored as config, with or without an expiry.
cfg, errMsg, err = handlers.BuildDatabaseTargetConfigForTest("", "daily")
require.NoError(t, err)
assert.Empty(t, errMsg)
assert.JSONEq(t, `{"rotation":"daily"}`, cfg)
cfg, errMsg, err = handlers.BuildDatabaseTargetConfigForTest(
"720h", "hourly",
)
require.NoError(t, err)
assert.Empty(t, errMsg)
assert.JSONEq(t, `{"expiry":"720h","rotation":"hourly"}`, cfg)
} }
func TestBuildDatabaseTargetConfig_RejectsBadExpiry( func TestBuildDatabaseTargetConfig_RejectsBadExpiry(
@@ -461,7 +475,7 @@ func TestBuildDatabaseTargetConfig_RejectsBadExpiry(
t.Parallel() t.Parallel()
for _, bad := range []string{"nonsense", "7d", "-5h"} { for _, bad := range []string{"nonsense", "7d", "-5h"} {
cfg, errMsg, err := handlers.BuildDatabaseTargetConfigForTest(bad) cfg, errMsg, err := handlers.BuildDatabaseTargetConfigForTest(bad, "")
require.NoError(t, err) require.NoError(t, err)
assert.Contains( assert.Contains(
@@ -471,3 +485,20 @@ func TestBuildDatabaseTargetConfig_RejectsBadExpiry(
assert.Empty(t, cfg) assert.Empty(t, cfg)
} }
} }
func TestBuildDatabaseTargetConfig_RejectsBadRotation(
t *testing.T,
) {
t.Parallel()
for _, bad := range []string{"weekly", "Daily", " none"} {
cfg, errMsg, err := handlers.BuildDatabaseTargetConfigForTest("", bad)
require.NoError(t, err)
assert.Contains(
t, errMsg, "Invalid archive rotation",
"rotation %q should be refused", bad,
)
assert.Empty(t, cfg)
}
}
@@ -120,7 +120,7 @@ func TestHandleSourceCreateSubmit_CreatesRequestedTargets(t *testing.T) {
// retention, each with archive on. Nothing is created, and the form // retention, each with archive on. Nothing is created, and the form
// comes back with the reason and every value entered: name, // comes back with the reason and every value entered: name,
// description, retention, URL, the checked archive box and the pruning // description, retention, URL, the checked archive box and the pruning
// choice. // and rotation choices.
func TestHandleSourceCreateSubmit_RefusedFormKeepsEveryValue( func TestHandleSourceCreateSubmit_RefusedFormKeepsEveryValue(
t *testing.T, t *testing.T,
) { ) {
@@ -152,6 +152,7 @@ func TestHandleSourceCreateSubmit_RefusedFormKeepsEveryValue(
form.Set("http_url", tc.httpURL) form.Set("http_url", tc.httpURL)
form.Set("archive", "on") form.Set("archive", "on")
form.Set("archive_expiry", "2160h") form.Set("archive_expiry", "2160h")
form.Set("archive_rotation", "hourly")
w := submitCreateForm(env, form) w := submitCreateForm(env, form)
require.Equal(t, http.StatusBadRequest, w.Code) require.Equal(t, http.StatusBadRequest, w.Code)
@@ -166,6 +167,7 @@ func TestHandleSourceCreateSubmit_RefusedFormKeepsEveryValue(
assert.Contains(t, page, `name="archive" value="on" checked`) assert.Contains(t, page, `name="archive" value="on" checked`)
assert.Contains(t, page, `x-data="collapsible" data-open`) assert.Contains(t, page, `x-data="collapsible" data-open`)
assert.Contains(t, page, `<option value="2160h" selected>`) assert.Contains(t, page, `<option value="2160h" selected>`)
assert.Contains(t, page, `<option value="hourly" selected>`)
assertNothingCreated(t, env.db) assertNothingCreated(t, env.db)
}) })
+42 -13
View File
@@ -296,9 +296,10 @@ type sourceFormInput struct {
// destination. // destination.
HTTPURL string HTTPURL string
// Archive asks for a database (archive) target, whose rows expire // Archive asks for a database (archive) target, whose rows expire
// after ArchiveExpiry. // after ArchiveExpiry and whose files rotate by ArchiveRotation.
Archive bool Archive bool
ArchiveExpiry string ArchiveExpiry string
ArchiveRotation string
} }
// newSourceFormData builds the template data for the webhook creation // newSourceFormData builds the template data for the webhook creation
@@ -315,6 +316,9 @@ func newSourceFormData(
tmplKeyArchiveExpiryChoices: archiveExpiryOptions( tmplKeyArchiveExpiryChoices: archiveExpiryOptions(
in.ArchiveExpiry, in.ArchiveExpiry,
), ),
tmplKeyArchiveRotationChoices: archiveRotationOptions(
in.ArchiveRotation,
),
} }
} }
@@ -347,6 +351,7 @@ func (h *Handlers) HandleSourceCreateSubmit() http.HandlerFunc {
HTTPURL: r.PostFormValue("http_url"), HTTPURL: r.PostFormValue("http_url"),
Archive: r.PostFormValue("archive") != "", Archive: r.PostFormValue("archive") != "",
ArchiveExpiry: r.PostFormValue("archive_expiry"), ArchiveExpiry: r.PostFormValue("archive_expiry"),
ArchiveRotation: r.PostFormValue("archive_rotation"),
} }
refuse := func(errMsg string) { refuse := func(errMsg string) {
@@ -420,6 +425,7 @@ func (h *Handlers) newWebhookTargets(
Name: "Archive", Name: "Archive",
Type: database.TargetTypeDatabase, Type: database.TargetTypeDatabase,
Expiry: in.ArchiveExpiry, Expiry: in.ArchiveExpiry,
Rotation: in.ArchiveRotation,
}) })
} }
@@ -626,9 +632,10 @@ func (h *Handlers) renderSourceDetail(
"Stats": h.loadWebhookStats(webhook.ID, entrypoints, targets), "Stats": h.loadWebhookStats(webhook.ID, entrypoints, targets),
tmplKeyTargetForm: targetForm, tmplKeyTargetForm: targetForm,
"TargetError": targetErr, "TargetError": targetErr,
// The add target form's select starts on its expiry // The add target form's selects start on its expiry and
// through Alpine, so no choice is selected here. // rotation through Alpine, so no choice is selected here.
tmplKeyArchiveExpiryChoices: archiveExpiryChoices(), tmplKeyArchiveExpiryChoices: archiveExpiryChoices(),
tmplKeyArchiveRotationChoices: archiveRotationChoices(),
} }
status := http.StatusOK status := http.StatusOK
@@ -1789,6 +1796,8 @@ type targetFormInput struct {
MaxRetries string MaxRetries string
// Expiry is a database (archive) target's row expiry. // Expiry is a database (archive) target's row expiry.
Expiry string Expiry string
// Rotation is a database (archive) target's rotation.
Rotation string
} }
// targetFormInputFrom reads a target form from a request body. The // targetFormInputFrom reads a target form from a request body. The
@@ -1812,6 +1821,7 @@ func targetFormInputFrom(r *http.Request) targetFormInput {
Timeout: r.PostFormValue("timeout"), Timeout: r.PostFormValue("timeout"),
MaxRetries: r.PostFormValue("max_retries"), MaxRetries: r.PostFormValue("max_retries"),
Expiry: r.PostFormValue("expiry"), Expiry: r.PostFormValue("expiry"),
Rotation: r.PostFormValue("rotation"),
} }
} }
@@ -1832,7 +1842,7 @@ func (h *Handlers) buildTargetConfig(
case database.TargetTypeSlack: case database.TargetTypeSlack:
return h.buildSlackTargetConfig(ctx, in.URL) return h.buildSlackTargetConfig(ctx, in.URL)
case database.TargetTypeDatabase: case database.TargetTypeDatabase:
return buildDatabaseTargetConfig(in.Expiry) return buildDatabaseTargetConfig(in.Expiry, in.Rotation)
case database.TargetTypeLog: case database.TargetTypeLog:
return "", "", nil return "", "", nil
default: default:
@@ -1953,22 +1963,41 @@ func marshalTargetConfig(cfg any) (string, error) {
} }
// buildDatabaseTargetConfig builds config JSON for a database // buildDatabaseTargetConfig builds config JSON for a database
// (archive) target. The optional expiry is validated here, at // (archive) target. The optional expiry and rotation are validated
// creation time, so an unparseable value is refused instead of // here, at creation time, so a bad value is refused instead of
// failing every subsequent delivery. An empty expiry yields an // failing every subsequent delivery. Each is stored only when set,
// empty config (the keep-forever default). // and with neither the config is empty (the keep-forever, one-file
func buildDatabaseTargetConfig(expiry string) (string, string, error) { // default).
func buildDatabaseTargetConfig(
expiry, rotation string,
) (string, string, error) {
expiry = strings.TrimSpace(expiry) expiry = strings.TrimSpace(expiry)
if expiry == "" {
return "", "", nil
}
err := delivery.ValidateArchiveExpiry(expiry) err := delivery.ValidateArchiveExpiry(expiry)
if err != nil { if err != nil {
return "", fmt.Sprintf("Invalid archive expiry: %v", err), nil return "", fmt.Sprintf("Invalid archive expiry: %v", err), nil
} }
configJSON, err := marshalTargetConfig(map[string]any{"expiry": expiry}) err = delivery.ValidateArchiveRotation(rotation)
if err != nil {
return "", fmt.Sprintf("Invalid archive rotation: %v", err), nil
}
cfg := map[string]any{}
if expiry != "" {
cfg["expiry"] = expiry
}
if rotation != "" {
cfg["rotation"] = rotation
}
if len(cfg) == 0 {
return "", "", nil
}
configJSON, err := marshalTargetConfig(cfg)
return configJSON, "", err return configJSON, "", err
} }
+4 -3
View File
@@ -86,9 +86,10 @@ func (d downloadWriter) Write(b []byte) (int, error) {
// written the response. // written the response.
// //
// It holds renameMu, which every archive rename runs under, while it // It holds renameMu, which every archive rename runs under, while it
// reads the stored names and opens the file, so the file it opens is // reads the stored names and opens the files, so the files it opens
// the one those names give. It lets go before the export is streamed: // are the ones those names give. It lets go before the export is
// once the file is open, a rename does not affect the export. // streamed: once the files are open, a rename does not affect the
// export.
func (h *Handlers) openTargetArchive( func (h *Handlers) openTargetArchive(
ctx context.Context, ctx context.Context,
w http.ResponseWriter, w http.ResponseWriter,
+4
View File
@@ -83,6 +83,7 @@ func (h *Handlers) HandleTargetEdit() http.HandlerFunc {
Timeout: cfg.Timeout, Timeout: cfg.Timeout,
MaxRetries: strconv.Itoa(target.MaxRetries), MaxRetries: strconv.Itoa(target.MaxRetries),
Expiry: cfg.Expiry, Expiry: cfg.Expiry,
Rotation: cfg.Rotation,
} }
h.renderTargetEdit( h.renderTargetEdit(
@@ -242,6 +243,9 @@ func (h *Handlers) renderTargetEdit(
tmplKeyMaxTimeout: delivery.MaxTargetTimeoutSeconds, tmplKeyMaxTimeout: delivery.MaxTargetTimeoutSeconds,
tmplKeyError: errMsg, tmplKeyError: errMsg,
tmplKeyArchiveExpiryChoices: archiveExpiryOptions(form.Expiry), tmplKeyArchiveExpiryChoices: archiveExpiryOptions(form.Expiry),
tmplKeyArchiveRotationChoices: archiveRotationOptions(
form.Rotation,
),
} }
h.renderTemplateStatus(w, r, targetEditTemplate, data, status) h.renderTemplateStatus(w, r, targetEditTemplate, data, status)
+5 -4
View File
@@ -591,7 +591,8 @@ func TestHandleTargetEditSubmit_RefusedFormComesBack(t *testing.T) {
"Invalid max retries", "Invalid max retries",
}, },
{ {
database.TargetTypeDatabase, "name=edited&expiry=7d", database.TargetTypeDatabase,
"name=edited&expiry=7d&rotation=daily",
"Invalid archive expiry", "Invalid archive expiry",
}, },
{database.TargetTypeLog, "name=", "Name is required"}, {database.TargetTypeLog, "name=", "Name is required"},
@@ -613,15 +614,15 @@ func TestHandleTargetEditSubmit_RefusedFormComesBack(t *testing.T) {
page := w.Body.String() page := w.Body.String()
assert.Contains(t, page, `class="alert-error">`+tc.reason) assert.Contains(t, page, `class="alert-error">`+tc.reason)
// headers is the form's one textarea and expiry its one // headers is the form's one textarea, and expiry and
// select; every other field is an input. // rotation its selects; every other field is an input.
for field := range form { for field := range form {
shown := `name="` + field + `" value="` + form.Get(field) + `"` shown := `name="` + field + `" value="` + form.Get(field) + `"`
switch field { switch field {
case "headers": case "headers":
shown = ">" + form.Get(field) + "</textarea>" shown = ">" + form.Get(field) + "</textarea>"
case "expiry": case "expiry", "rotation":
shown = `<option value="` + form.Get(field) + `" selected>` shown = `<option value="` + form.Get(field) + `" selected>`
} }
+54 -23
View File
@@ -4,6 +4,7 @@ import (
"errors" "errors"
"fmt" "fmt"
"io/fs" "io/fs"
"os"
"path/filepath" "path/filepath"
"time" "time"
@@ -21,8 +22,8 @@ type TargetRowView struct {
// and is nil when the webhook's event database could not be read. // and is nil when the webhook's event database could not be read.
Deliveries *TargetDeliveries Deliveries *TargetDeliveries
// Archive is a database target's archive file, and nil for a target // Archive is a database target's archive files, and nil for a
// of any other type. // target of any other type.
Archive *ArchiveFileView Archive *ArchiveFileView
} }
@@ -39,17 +40,23 @@ type TargetDeliveries struct {
} }
// ArchiveFileView is what a database target's row shows about its // ArchiveFileView is what a database target's row shows about its
// archive file. // archive files.
type ArchiveFileView struct { type ArchiveFileView struct {
// Name is the name of the file an event received now goes to.
Name string Name string
// Note stands in for the size and the last write when there are // Note says that file does not exist yet, or that the files could
// none to show, and is empty when there are. // not be read, and is empty otherwise.
Note string Note string
// Size is the size on disk. Written is how long ago the file was // Files counts the target's archive files, and is 0 when there are
// last written, and WrittenUTC the full time the page shows on // none, or they could not be read, and so no size or last write to
// hover. // show.
Files int
// Size is the size on disk of all the files together. Written is
// how long ago the latest of them was last written, and WrittenUTC
// the full time the page shows on hover.
Size string Size string
Written string Written string
WrittenUTC string WrittenUTC string
@@ -62,6 +69,7 @@ func (h *Handlers) targetRows(
) []TargetRowView { ) []TargetRowView {
views := delivery.NewTargetViews(targets) views := delivery.NewTargetViews(targets)
rows := make([]TargetRowView, len(views)) rows := make([]TargetRowView, len(views))
now := time.Now()
deliveries, err := h.loadTargetDeliveries(webhook.ID) deliveries, err := h.loadTargetDeliveries(webhook.ID)
if err != nil { if err != nil {
@@ -82,7 +90,7 @@ func (h *Handlers) targetRows(
} }
if targets[i].Type == database.TargetTypeDatabase { if targets[i].Type == database.TargetTypeDatabase {
rows[i].Archive = h.archiveFileView(webhook, &targets[i]) rows[i].Archive = h.archiveFileView(webhook, &targets[i], now)
} }
} }
@@ -147,22 +155,50 @@ func readTargetDeliveries(
return byTarget, nil return byTarget, nil
} }
// archiveFileView describes a database target's archive file from the // archiveFileView describes a database target's archive files from
// file's metadata alone; the archive is never opened. The file is found // their metadata alone; the archive is never opened. It names the file
// by the name the archive writer uses, so it follows a rename of the // an event received at now goes to. The files are found by the name
// webhook or the target. // the archive writer uses, so they follow a rename of the webhook or
// the target.
func (h *Handlers) archiveFileView( func (h *Handlers) archiveFileView(
webhook *database.Webhook, target *database.Target, webhook *database.Webhook, target *database.Target, now time.Time,
) *ArchiveFileView { ) *ArchiveFileView {
path := delivery.ArchivePath(h.dbMgr, webhook, target) current, err := delivery.ArchivePathAt(h.dbMgr, webhook, target, now)
view := &ArchiveFileView{Name: filepath.Base(path)} if err != nil {
return h.archiveUnreadable(&ArchiveFileView{}, target, err)
}
file, err := delivery.StatArchive(path) view := &ArchiveFileView{Name: filepath.Base(current)}
_, statErr := os.Stat(current)
if errors.Is(statErr, fs.ErrNotExist) {
view.Note = "not created yet"
}
files, err := delivery.StatArchive(
delivery.ArchivePath(h.dbMgr, webhook, target),
)
switch { switch {
case errors.Is(err, fs.ErrNotExist): case errors.Is(err, fs.ErrNotExist):
view.Note = "not created yet" // No files: no size or last write to show.
case err != nil: case err != nil:
return h.archiveUnreadable(view, target, err)
default:
view.Files = files.Files
view.Size = humanize.Bytes(uint64(files.Size)) //nolint:gosec // never negative
view.Written = humanize.Time(files.Written)
view.WrittenUTC = files.Written.UTC().Format(time.DateTime) + " UTC"
}
return view
}
// archiveUnreadable logs why a database target's archive files could
// not be described, and returns view saying so.
func (h *Handlers) archiveUnreadable(
view *ArchiveFileView, target *database.Target, err error,
) *ArchiveFileView {
h.log.Error( h.log.Error(
"failed to read archive file metadata", "failed to read archive file metadata",
"target_id", target.ID, "target_id", target.ID,
@@ -170,11 +206,6 @@ func (h *Handlers) archiveFileView(
) )
view.Note = "could not be read" view.Note = "could not be read"
default:
view.Size = humanize.Bytes(uint64(file.Size)) //nolint:gosec // never negative
view.Written = humanize.Time(file.Written)
view.WrittenUTC = file.Written.UTC().Format(time.DateTime) + " UTC"
}
return view return view
} }
+41 -21
View File
@@ -71,7 +71,7 @@ func TestAlpineRunsUnderTheSecurityPolicy(t *testing.T) {
// this order. A new check is one more line here. // this order. A new check is one more line here.
checkAddEntrypoint(ctx, t, page) checkAddEntrypoint(ctx, t, page)
checkAddEachTargetType(ctx, t, page) checkAddEachTargetType(ctx, t, page)
checkArchiveExpiry(ctx, t, page) checkArchiveChoices(ctx, t, page)
checkRefusedTarget(ctx, t, page) checkRefusedTarget(ctx, t, page)
checkTargetDeliveries(ctx, t, page, target.Name, checkTargetDeliveries(ctx, t, page, target.Name,
"0 in total, 0 in the last 24 hours", "0 in total, 0 in the last 24 hours",
@@ -321,8 +321,8 @@ func checkAddEachTargetType(ctx context.Context, t *testing.T, url string) {
map[string]string{"url": publicTargetURL}, map[string]string{"url": publicTargetURL},
}, },
{ {
"database", "csrf_token name type expiry", "database", "csrf_token name type expiry rotation",
map[string]string{"expiry": "720h"}, map[string]string{"expiry": "720h", "rotation": "daily"},
}, },
{"log", "csrf_token name type", nil}, {"log", "csrf_token name type", nil},
} }
@@ -422,14 +422,18 @@ func chooseTargetType(ctx context.Context, t *testing.T, targetType string) {
"%s: Add still shows while the form is open", targetType) "%s: Add still shows while the form is open", targetType)
} }
// checkArchiveExpiry loads a webhook page and checks that the add // checkArchiveChoices loads a webhook page and checks that the add
// target form's archive expiry starts on never, that the database // target form's archive expiry starts on never and its archive
// target checkAddTarget added with 720h is listed as 30 days, and that // rotation on none, that the database target checkAddTarget added with
// its edit form starts on 720h. // 720h and daily is listed as 30 days and daily, and that its edit
func checkArchiveExpiry(ctx context.Context, t *testing.T, url string) { // form starts on 720h and daily.
func checkArchiveChoices(ctx context.Context, t *testing.T, url string) {
t.Helper() t.Helper()
const expiry = `form[action$="/targets"] select[name="expiry"]` const (
expiry = `form[action$="/targets"] select[name="expiry"]`
rotation = `form[action$="/targets"] select[name="rotation"]`
)
row := `//span[text()="added-database"]/ancestor::div[@class="p-4"][1]` row := `//span[text()="added-database"]/ancestor::div[@class="p-4"][1]`
@@ -437,26 +441,36 @@ func checkArchiveExpiry(ctx context.Context, t *testing.T, url string) {
chooseTargetType(ctx, t, "database") chooseTargetType(ctx, t, "database")
var start, edited string var startExpiry, startRotation, editedExpiry, editedRotation string
require.NoError(t, chromedp.Run( require.NoError(t, chromedp.Run(
ctx, chromedp.Value(expiry, &start, chromedp.ByQuery), ctx,
chromedp.Value(expiry, &startExpiry, chromedp.ByQuery),
chromedp.Value(rotation, &startRotation, chromedp.ByQuery),
)) ))
assert.Equal(t, "never", start, assert.Equal(t, "never", startExpiry,
"the add target form's archive expiry does not start on never") "the add target form's archive expiry does not start on never")
assert.Equal(t, "none", startRotation,
"the add target form's archive rotation does not start on none")
assert.True(t, shown(ctx, row+`//span[text()="Archive Expiry:"]`+ assert.True(t, shown(ctx, row+`//span[text()="Archive Expiry:"]`+
`/following-sibling::span[text()="30 days"]`), `/following-sibling::span[text()="30 days"]`),
"a database target added with 720h is not listed as 30 days") "a database target added with 720h is not listed as 30 days")
assert.True(t, shown(ctx, row+`//span[text()="Archive Rotation:"]`+
`/following-sibling::span[text()="daily"]`),
"a database target added with daily is not listed as daily")
click(ctx, t, row+`//a[text()="Edit"]`) click(ctx, t, row+`//a[text()="Edit"]`)
require.NoError(t, chromedp.Run( require.NoError(t, chromedp.Run(
ctx, ctx,
chromedp.WaitReady("#expiry", chromedp.ByQuery), chromedp.WaitReady("#expiry", chromedp.ByQuery),
chromedp.Value("#expiry", &edited, chromedp.ByQuery), chromedp.Value("#expiry", &editedExpiry, chromedp.ByQuery),
chromedp.Value("#rotation", &editedRotation, chromedp.ByQuery),
)) ))
assert.Equal(t, "720h", edited, assert.Equal(t, "720h", editedExpiry,
"the edit form does not start on the stored archive expiry") "the edit form does not start on the stored archive expiry")
assert.Equal(t, "daily", editedRotation,
"the edit form does not start on the stored archive rotation")
} }
// checkRefusedTarget submits an http target the server refuses, a // checkRefusedTarget submits an http target the server refuses, a
@@ -892,7 +906,9 @@ func checkNewWebhookTargets(
pruningChoice, expiry, chromedp.BySearch, pruningChoice, expiry, chromedp.BySearch,
))) )))
want[database.TargetTypeDatabase] = `{"expiry":"` + expiry + `"}` // The rotation choice is left on none.
want[database.TargetTypeDatabase] = `{"expiry":"` + expiry +
`","rotation":"none"}`
} }
click(ctx, t, createButton) click(ctx, t, createButton)
@@ -933,8 +949,8 @@ func targetConfigs(
// checkRefusedNewWebhook submits the new webhook page with archive // checkRefusedNewWebhook submits the new webhook page with archive
// checked and an HTTP target URL the server refuses, a loopback // checked and an HTTP target URL the server refuses, a loopback
// destination, and checks that the page comes back with the reason and // destination, and checks that the page comes back with the reason and
// every value entered, archive still checked and its pruning choice // every value entered, archive still checked and its pruning and
// showing. // rotation choices showing.
func checkRefusedNewWebhook(ctx context.Context, t *testing.T, url string) { func checkRefusedNewWebhook(ctx context.Context, t *testing.T, url string) {
t.Helper() t.Helper()
@@ -951,16 +967,18 @@ func checkRefusedNewWebhook(ctx context.Context, t *testing.T, url string) {
click(ctx, t, archiveBox) click(ctx, t, archiveBox)
require.True(t, shown(ctx, pruningChoice), require.True(t, shown(ctx, pruningChoice),
"checking archive does not show the pruning choice") "checking archive does not show the pruning choice")
require.NoError(t, chromedp.Run(ctx, chromedp.SetValue( require.NoError(t, chromedp.Run(
pruningChoice, "2160h", chromedp.BySearch, ctx,
))) chromedp.SetValue(pruningChoice, "2160h", chromedp.BySearch),
chromedp.SetValue("#archive_rotation", "monthly", chromedp.ByQuery),
))
click(ctx, t, createButton) click(ctx, t, createButton)
assert.True(t, shown(ctx, `//div[@class="alert-error"]`), assert.True(t, shown(ctx, `//div[@class="alert-error"]`),
"a refused webhook does not show the reason") "a refused webhook does not show the reason")
var ( var (
name, description, retention, typed, expiry string name, description, retention, typed, expiry, rotation string
checked bool checked bool
) )
@@ -971,6 +989,7 @@ func checkRefusedNewWebhook(ctx context.Context, t *testing.T, url string) {
chromedp.Value("#retention_days", &retention, chromedp.ByQuery), chromedp.Value("#retention_days", &retention, chromedp.ByQuery),
chromedp.Value("#http_url", &typed, chromedp.ByQuery), chromedp.Value("#http_url", &typed, chromedp.ByQuery),
chromedp.Value("#archive_expiry", &expiry, chromedp.ByQuery), chromedp.Value("#archive_expiry", &expiry, chromedp.ByQuery),
chromedp.Value("#archive_rotation", &rotation, chromedp.ByQuery),
chromedp.Evaluate(archiveIsOn, &checked), chromedp.Evaluate(archiveIsOn, &checked),
)) ))
@@ -981,6 +1000,7 @@ func checkRefusedNewWebhook(ctx context.Context, t *testing.T, url string) {
assert.True(t, checked, "archive is no longer checked") assert.True(t, checked, "archive is no longer checked")
assert.True(t, shown(ctx, pruningChoice), "the pruning choice is hidden") assert.True(t, shown(ctx, pruningChoice), "the pruning choice is hidden")
assert.Equal(t, "2160h", expiry, "the pruning chosen is lost") assert.Equal(t, "2160h", expiry, "the pruning chosen is lost")
assert.Equal(t, "monthly", rotation, "the rotation chosen is lost")
} }
// checkMobileMenu loads a page in a phone-sized window and checks that // checkMobileMenu loads a page in a phone-sized window and checks that
File diff suppressed because one or more lines are too long
+5 -2
View File
@@ -72,8 +72,8 @@ document.addEventListener("alpine:init", function () {
// Something a click shows and hides: the mobile menu, an add form, // Something a click shows and hides: the mobile menu, an add form,
// an entrypoint's edit form, an event in the event log or in the // an entrypoint's edit form, an event in the event log or in the
// recent events, a delivery's attempts, the new webhook page's // recent events, a delivery's attempts, the new webhook page's
// archive pruning choice. It starts hidden, or shown when its // archive pruning and rotation choices. It starts hidden, or shown
// element has the data-open attribute. // when its element has the data-open attribute.
window.Alpine.data("collapsible", function () { window.Alpine.data("collapsible", function () {
return { return {
open: false, open: false,
@@ -116,6 +116,7 @@ document.addEventListener("alpine:init", function () {
timeout: "", timeout: "",
maxRetries: "", maxRetries: "",
expiry: "", expiry: "",
rotation: "",
init() { init() {
const refused = this.$root.dataset; const refused = this.$root.dataset;
@@ -127,6 +128,7 @@ document.addEventListener("alpine:init", function () {
this.timeout = refused.timeout; this.timeout = refused.timeout;
this.maxRetries = refused.maxRetries; this.maxRetries = refused.maxRetries;
this.expiry = refused.expiry; this.expiry = refused.expiry;
this.rotation = refused.rotation;
}, },
add() { add() {
this.choosing = true; this.choosing = true;
@@ -145,6 +147,7 @@ document.addEventListener("alpine:init", function () {
this.timeout = ""; this.timeout = "";
this.maxRetries = ""; this.maxRetries = "";
this.expiry = ""; this.expiry = "";
this.rotation = "";
this.$refs.form.reset(); this.$refs.form.reset();
}, },
get filling() { get filling() {
+18 -4
View File
@@ -129,7 +129,8 @@
data-headers="{{.TargetForm.Headers}}" data-headers="{{.TargetForm.Headers}}"
data-timeout="{{.TargetForm.Timeout}}" data-timeout="{{.TargetForm.Timeout}}"
data-max-retries="{{.TargetForm.MaxRetries}}" data-max-retries="{{.TargetForm.MaxRetries}}"
data-expiry="{{.TargetForm.Expiry}}"> data-expiry="{{.TargetForm.Expiry}}"
data-rotation="{{.TargetForm.Rotation}}">
<div class="p-4 border-b border-gray-200 flex justify-between items-center"> <div class="p-4 border-b border-gray-200 flex justify-between items-center">
<h2 class="text-lg font-medium text-gray-900">Targets</h2> <h2 class="text-lg font-medium text-gray-900">Targets</h2>
<button type="button" @click="add" x-show="closed" class="btn-small"> <button type="button" @click="add" x-show="closed" class="btn-small">
@@ -201,8 +202,9 @@
</div> </div>
</template> </template>
<template x-if="isDatabase"> <template x-if="isDatabase">
<div> <div class="space-y-3">
<input type="hidden" name="type" value="database"> <input type="hidden" name="type" value="database">
<div>
<div class="flex gap-2 items-center"> <div class="flex gap-2 items-center">
<label class="text-sm text-gray-700">Archive expiry:</label> <label class="text-sm text-gray-700">Archive expiry:</label>
<select name="expiry" :value="expiry" class="input text-sm w-24"> <select name="expiry" :value="expiry" class="input text-sm w-24">
@@ -213,6 +215,18 @@
</div> </div>
<p class="text-xs text-gray-500 mt-1">Archived events older than this are deleted from the archive; never keeps them all.</p> <p class="text-xs text-gray-500 mt-1">Archived events older than this are deleted from the archive; never keeps them all.</p>
</div> </div>
<div>
<div class="flex gap-2 items-center">
<label class="text-sm text-gray-700">Archive rotation:</label>
<select name="rotation" :value="rotation" class="input text-sm w-28">
{{range .ArchiveRotationChoices}}
<option value="{{.Value}}">{{.Label}}</option>
{{end}}
</select>
</div>
<p class="text-xs text-gray-500 mt-1">Starts a new archive file each month, day or hour (UTC), named for it; none keeps one file.</p>
</div>
</div>
</template> </template>
<template x-if="isLog"> <template x-if="isLog">
<div> <div>
@@ -267,10 +281,10 @@
<span class="break-all">{{.Name}}</span> <span class="break-all">{{.Name}}</span>
{{with .Note}}<span>({{.}})</span>{{end}} {{with .Note}}<span>({{.}})</span>{{end}}
</div> </div>
{{if .Size}} {{if .Files}}
<div class="text-xs text-gray-500 mt-1"> <div class="text-xs text-gray-500 mt-1">
<span class="font-medium text-gray-700">Archive Size:</span> <span class="font-medium text-gray-700">Archive Size:</span>
<span>{{.Size}}</span> <span>{{.Size}}{{if gt .Files 1}} in {{.Files}} files{{end}}</span>
</div> </div>
<div class="text-xs text-gray-500 mt-1"> <div class="text-xs text-gray-500 mt-1">
<span class="font-medium text-gray-700">Last Written:</span> <span class="font-medium text-gray-700">Last Written:</span>
+12 -3
View File
@@ -38,9 +38,9 @@
<p class="text-xs text-gray-500 mt-1">Optional. When filled in, the webhook is created with an HTTP target that delivers each event to this URL.</p> <p class="text-xs text-gray-500 mt-1">Optional. When filled in, the webhook is created with an HTTP target that delivers each event to this URL.</p>
</div> </div>
<!-- The checkbox shows the pruning choice while checked. With <!-- The checkbox shows the pruning and rotation choices while
autocomplete="off", going back to the page does not checked. With autocomplete="off", going back to the page
check the box again with the choice hidden. --> does not check the box again with the choices hidden. -->
<div class="form-group" x-data="collapsible"{{if .Form.Archive}} data-open{{end}}> <div class="form-group" x-data="collapsible"{{if .Form.Archive}} data-open{{end}}>
<label class="flex items-center gap-2 text-sm font-medium text-gray-700"> <label class="flex items-center gap-2 text-sm font-medium text-gray-700">
<input type="checkbox" name="archive" value="on"{{if .Form.Archive}} checked{{end}} autocomplete="off" @change="toggle" class="h-4 w-4"> <input type="checkbox" name="archive" value="on"{{if .Form.Archive}} checked{{end}} autocomplete="off" @change="toggle" class="h-4 w-4">
@@ -55,6 +55,15 @@
{{end}} {{end}}
</select> </select>
<p class="text-xs text-gray-500 mt-1">Archived events older than this are deleted from the archive; never keeps them all.</p> <p class="text-xs text-gray-500 mt-1">Archived events older than this are deleted from the archive; never keeps them all.</p>
<div class="mt-3">
<label for="archive_rotation" class="label">Archive rotation</label>
<select id="archive_rotation" name="archive_rotation" class="input">
{{range .ArchiveRotationChoices}}
<option value="{{.Value}}"{{if .Selected}} selected{{end}}>{{.Label}}</option>
{{end}}
</select>
<p class="text-xs text-gray-500 mt-1">Starts a new archive file each month, day or hour (UTC), named for it; none keeps one file.</p>
</div>
</div> </div>
</div> </div>
+10
View File
@@ -67,6 +67,16 @@
</select> </select>
<p class="text-xs text-gray-500 mt-1">Archived events older than this are deleted from the archive; never keeps them all.</p> <p class="text-xs text-gray-500 mt-1">Archived events older than this are deleted from the archive; never keeps them all.</p>
</div> </div>
<div class="form-group">
<label for="rotation" class="label">Archive Rotation</label>
<select id="rotation" name="rotation" class="input">
{{range .ArchiveRotationChoices}}
<option value="{{.Value}}"{{if .Selected}} selected{{end}}>{{.Label}}</option>
{{end}}
</select>
<p class="text-xs text-gray-500 mt-1">Starts a new archive file each month, day or hour (UTC), named for it; none keeps one file. A change applies from the next event; existing files keep their names.</p>
</div>
{{end}} {{end}}
{{if or (eq .Target.Type "http") (eq .Target.Type "slack")}} {{if or (eq .Target.Type "http") (eq .Target.Type "slack")}}