From ee7c62607104a7dacba5082cc46098f57049271c Mon Sep 17 00:00:00 2001 From: clawbot Date: Fri, 7 Aug 2026 22:50:08 +0200 Subject: [PATCH] Implement the database archiving target (closes #43) (#84) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Implements the `databaseTarget` as a real archiving target, replacing the always-successful stub. Delivering to a `database` target now writes the full event into a per-webhook archive SQLite file for long-term storage. ## Archive-writer semantics - **Separate file:** each webhook's full events are written as rows into `archive-{webhookID}.db` under the data dir, distinct from the per-webhook event DB (`events-{webhookID}.db`). The file and its schema are created on first write if missing. Each row carries the full event: body, headers, method, content type, webhook id, entrypoint id, event id, and an archived-at timestamp. - **Close/reopen with debounce:** after each write the archive handle is closed and reopened, unless the last (re)open was less than one second ago. This lets an operator move the archive file away for offline archiving while bounding file churn under load. A per-webhook `archiveWriter` owns this debounce state and serialises writes. - **Auto-recreate:** the file is opened create-if-missing (`mode=rwc`) and its schema re-migrated on every open, so if the archive was moved or removed since the last open, the next write recreates it. The writer also detects a missing file before writing and reopens first, so a moved-away file is recreated rather than lost. - **Optional expiry, validated at creation:** an optional `expiry` in the target's config JSON (e.g. `{"expiry":"720h"}`) is validated when the target is created (`ValidateArchiveExpiry`; bad values are rejected with a 400 at the add-target form, the Slack URL precedent). The default (missing, empty, or `"never"`) keeps rows forever with no pruning. When a positive duration is set, rows older than it (measured from each row's archived-at time) are pruned on every (re)open; because the file is reopened after writes, prune-on-open keeps the archive swept without a separate background sweeper. A set-but-invalid expiry in a stored config (unparseable, zero, or negative) is an error at delivery time too — never a silent default. - **No-retry, fail-loud:** the target performs a single attempt with no retries. On success it records one successful attempt and marks the delivery delivered. If the archive write fails, the attempt is recorded as failed with the error and the delivery is marked failed — archiving errors never report success. ## Scope - `internal/delivery/target_database.go` — the `databaseTarget` (no-retry) archives via a per-webhook writer registry; an archive error records a failed attempt and marks the delivery failed. - `internal/delivery/target_database_archive.go` (new) — the `archiveWriter`, the archived-row model, config/expiry parsing (fail-loud on set-but-invalid values), `ValidateArchiveExpiry`, and prune-on-open. - `internal/handlers/source_management.go` — database targets get a creation-validated `expiry` config (`buildDatabaseTargetConfig`); the expiry form value is read where the request body is bounded and bad values are rejected with a 400 at target creation. - `templates/source_detail.html` — the add-target form shows an expiry field for database targets. - `README.md` — the database-target documentation describes the archiving semantics. - `internal/delivery/export_test.go`, `internal/delivery/target_database_test.go`, `internal/handlers` tests — tests and their exported shims. No changes to the `Target` interface or other targets. ## Tests - a row is archived (both at the writer level and end-to-end through `Deliver`) - a forced archive failure (bad stored expiry config) yields a `Failed` delivery with a non-success `DeliveryResult` carrying the error and no archive file created - the file is recreated after removal, with only the post-removal row - the one-second reopen debounce (rapid writes reopen once; a write after the window reopens again) - expiry pruning removes rows older than the configured expiry - expiry config parsing (empty / `never` / duration accepted; unparseable, zero, and negative values error) - expiry validation at target creation (`TestValidateArchiveExpiry`; valid values build the config, bad values get a 400) ## Validation `docker build .` exits 0 (fmt-check, lint, test, build all pass). Closes #43 Co-authored-by: sneak Reviewed-on: https://git.eeqj.de/sneak/webhooker/pulls/84 Co-authored-by: clawbot Co-committed-by: clawbot --- README.md | 31 +- internal/delivery/engine_test.go | 48 ++- internal/delivery/export_test.go | 61 +++ internal/delivery/target_database.go | 115 +++++- internal/delivery/target_database_archive.go | 312 +++++++++++++++ internal/delivery/target_database_test.go | 395 +++++++++++++++++++ internal/handlers/export_test.go | 10 + internal/handlers/handlers_test.go | 54 +++ internal/handlers/source_management.go | 53 ++- templates/source_detail.html | 4 + 10 files changed, 1045 insertions(+), 38 deletions(-) create mode 100644 internal/delivery/target_database_archive.go create mode 100644 internal/delivery/target_database_test.go diff --git a/README.md b/README.md index 4621dd8..1514b61 100644 --- a/README.md +++ b/README.md @@ -363,10 +363,12 @@ events should be forwarded. greater than 0, failed deliveries are retried with exponential backoff up to `max_retries` attempts, protected by a per-target circuit breaker. -- **`database`** — Confirm the event is stored in the webhook's - per-webhook database (no external delivery). Since events are always - written to the per-webhook DB on ingestion, this target marks delivery - as immediately successful. Useful for ensuring durable event archival. +- **`database`** — Archive the full event as a row into a separate + per-webhook archive database (`archive-{webhookID}.db`) for long-term + retention, with an optional creation-validated expiry (default: keep + forever). No external delivery and no retries; an archive write + failure fails the delivery. See the database target section under + "Per-Webhook Event Databases" for the full semantics. - **`log`** — Write the event to the application log (stdout). Useful for debugging. @@ -512,11 +514,22 @@ This separation provides: page cache, and its own lock, so concurrent event ingestion across webhooks won't contend. -The **database target type** leverages this architecture: since events -are already stored in the per-webhook database by design, the database -target simply marks the delivery as immediately successful. The -per-webhook DB IS the dedicated event database — that's the whole point -of the database target type. +The **database target type** builds on this architecture to provide +long-term archiving, separate from the per-webhook event database (which +may prune events under its own retention). Delivering to a database +target writes the full event — body, headers, method, content type, and +webhook/entrypoint/event identifiers — as a row into a dedicated archive +database, `archive-{webhookID}.db`, stored under the data directory +beside the event database. After each write the archive handle is closed +and reopened, debounced to at most once per second, so an operator can +move the archive file away for offline archiving without stopping the +service; a moved or removed archive file is recreated automatically on +the next write. An optional `expiry` in the target's config JSON (e.g. +`{"expiry":"720h"}`) is validated when the target is created — the +default (unset or the literal `never`) keeps rows forever — and rows +older than the expiry are pruned each time the archive is (re)opened. An +archive write failure is never silent success: the delivery records a +failed attempt with the error and is marked failed. The **Slack target type** sends webhook events as formatted messages to any Slack-compatible incoming webhook URL (works with Slack, Mattermost, diff --git a/internal/delivery/engine_test.go b/internal/delivery/engine_test.go index 25b2c8b..9c58ae0 100644 --- a/internal/delivery/engine_test.go +++ b/internal/delivery/engine_test.go @@ -342,33 +342,29 @@ func TestDeliverDatabase_ImmediateSuccess( t.Parallel() db := testWebhookDB(t) - e := testEngine(t, 1) - event := seedEvent(t, db, `{"db":"target"}`) - - dlv := seedDelivery( - t, db, event.ID, uuid.New().String(), - database.DeliveryStatusPending, + // The database target archives for real now, so the engine + // needs a webhook DB manager to locate the data directory. + e := delivery.NewTestEngineWithDB( + nil, + database.NewTestWebhookDBManager(t.TempDir()), + slog.New(slog.NewTextHandler( + os.Stderr, + &slog.HandlerOptions{Level: slog.LevelDebug}, + )), + &http.Client{Timeout: 5 * time.Second}, + 1, ) - d := &database.Delivery{ - EventID: event.ID, - TargetID: dlv.TargetID, - Status: database.DeliveryStatusPending, - Event: event, - Target: database.Target{ - Name: "test-db", - Type: database.TargetTypeDatabase, - }, - } - d.ID = dlv.ID + event := seedEvent(t, db, `{"db":"target"}`) + d := seedDatabaseTargetDelivery(t, db, event, "") e.ExportDeliverDatabase(db, d) var updated database.Delivery require.NoError(t, db.First( - &updated, "id = ?", dlv.ID, + &updated, "id = ?", d.ID, ).Error) assert.Equal(t, @@ -379,7 +375,7 @@ func TestDeliverDatabase_ImmediateSuccess( var result database.DeliveryResult require.NoError(t, db.Where( - "delivery_id = ?", dlv.ID, + "delivery_id = ?", d.ID, ).First(&result).Error) assert.True(t, result.Success) @@ -1158,7 +1154,19 @@ func TestProcessDelivery_RoutesToCorrectHandler( t.Parallel() db := testWebhookDB(t) - e := testEngine(t, 1) + + // The database target archives for real now, so the engine + // needs a webhook DB manager to locate the data directory. + e := delivery.NewTestEngineWithDB( + nil, + database.NewTestWebhookDBManager(t.TempDir()), + slog.New(slog.NewTextHandler( + os.Stderr, + &slog.HandlerOptions{Level: slog.LevelDebug}, + )), + &http.Client{Timeout: 5 * time.Second}, + 1, + ) tests := []struct { name string diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index ec5de60..739eb03 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -273,3 +273,64 @@ func NewTestCircuitBreaker( cooldown: cooldown, } } + +// ExportArchivedEvent aliases the archive row type so black-box +// tests can construct and read archive rows. +type ExportArchivedEvent = archivedEvent + +// ExportArchiveWriter wraps an archiveWriter so black-box tests +// can exercise the per-webhook archive file mechanics. +type ExportArchiveWriter struct { + w *archiveWriter +} + +// NewExportArchiveWriter builds an archive writer for tests, +// optionally overriding the reopen debounce (a non-positive +// debounce keeps the production default). +func NewExportArchiveWriter( + path string, log *slog.Logger, debounce time.Duration, +) *ExportArchiveWriter { + w := newArchiveWriter(path, log) + if debounce > 0 { + w.debounce = debounce + } + + return &ExportArchiveWriter{w: w} +} + +// Write archives a row through the writer. +func (e *ExportArchiveWriter) Write( + row ExportArchivedEvent, expiry time.Duration, +) error { + return e.w.write(row, expiry) +} + +// Open opens the archive file, pruning when expiry is positive. +func (e *ExportArchiveWriter) Open(expiry time.Duration) error { + return e.w.open(expiry) +} + +// Reopen closes and reopens the archive file. +func (e *ExportArchiveWriter) Reopen( + expiry time.Duration, +) error { + return e.w.reopen(expiry) +} + +// Reopens reports how many times the file has been opened. +func (e *ExportArchiveWriter) Reopens() int { + return e.w.reopens +} + +// DB returns the writer's current open handle for row +// inspection in tests. +func (e *ExportArchiveWriter) DB() *gorm.DB { + return e.w.db +} + +// ExportParseArchiveExpiry exposes parseArchiveExpiry. +func ExportParseArchiveExpiry( + configJSON string, +) (time.Duration, error) { + return parseArchiveExpiry(configJSON) +} diff --git a/internal/delivery/target_database.go b/internal/delivery/target_database.go index a8b1d6a..0443aaa 100644 --- a/internal/delivery/target_database.go +++ b/internal/delivery/target_database.go @@ -2,21 +2,38 @@ package delivery import ( "context" + "fmt" + "path/filepath" + "sync" "gorm.io/gorm" "sneak.berlin/go/webhooker/internal/database" ) -// databaseTarget is a fire-and-forget target: the event is -// already persisted in the per-webhook database by the time -// delivery runs, so the target records a single successful -// attempt. (Durable archiving to a separate store is tracked -// as its own work.) +// databaseTarget is a no-retry target that archives the +// full inbound event into a per-webhook archive SQLite file, +// separate from the per-webhook event database. The event is +// already persisted in the per-webhook event DB by the time +// delivery runs; the database target additionally writes a +// durable long-term copy into archive-{webhookID}.db and then +// records a single attempt whose outcome reflects whether the +// archive write succeeded. See archiveWriter for the +// close/reopen, auto-recreate, and expiry semantics. type databaseTarget struct { eng *Engine + + mu sync.Mutex + writers map[string]*archiveWriter } -// Deliver implements Target. +// Deliver implements Target. It archives the event, then +// records one successful attempt and marks the delivery +// delivered. An archiving error fails the delivery: the +// attempt is recorded as failed with the error and the +// delivery is marked failed, so a target that could not do +// its one job (archiving) never reports success. The target +// does not retry; the event remains durably stored in the +// per-webhook event database. func (t *databaseTarget) Deliver( _ context.Context, webhookDB *gorm.DB, @@ -24,6 +41,27 @@ func (t *databaseTarget) Deliver( _ *Task, _ Scheduler, ) { + err := t.archive(d) + if err != nil { + t.eng.log.Error( + "failed to archive event to database target", + "delivery_id", d.ID, + "event_id", d.EventID, + "error", err, + ) + + t.eng.recordResult( + webhookDB, d, 1, false, 0, "", + err.Error(), 0, + ) + + t.eng.updateDeliveryStatus( + webhookDB, d, database.DeliveryStatusFailed, + ) + + return + } + t.eng.recordResult( webhookDB, d, 1, true, 0, "", "", 0, ) @@ -32,3 +70,68 @@ func (t *databaseTarget) Deliver( webhookDB, d, database.DeliveryStatusDelivered, ) } + +// archive writes the full event as a row into the webhook's +// archive database, honouring the optional per-target expiry +// parsed from the target config JSON. +func (t *databaseTarget) archive(d *database.Delivery) error { + webhookID := d.Event.WebhookID + if webhookID == "" { + return errArchiveMissingWebhookID + } + + expiry, err := parseArchiveExpiry(d.Target.Config) + if err != nil { + return err + } + + w, err := t.writerFor(webhookID) + if err != nil { + return err + } + + row := archivedEvent{ + EventID: d.Event.ID, + WebhookID: webhookID, + EntrypointID: d.Event.EntrypointID, + Method: d.Event.Method, + Headers: d.Event.Headers, + Body: d.Event.Body, + ContentType: d.Event.ContentType, + } + + return w.write(row, expiry) +} + +// writerFor returns the archiveWriter for a webhook, creating +// and caching it on first use. Each webhook has one writer so +// its close/reopen debounce state is shared across concurrent +// deliveries. The archive file lives beside the per-webhook +// event database in the data directory. +func (t *databaseTarget) writerFor( + webhookID string, +) (*archiveWriter, error) { + if t.eng.dbManager == nil { + return nil, errArchiveNoDataDir + } + + dir := filepath.Dir(t.eng.dbManager.DBPath(webhookID)) + path := filepath.Join( + dir, fmt.Sprintf("archive-%s.db", webhookID), + ) + + t.mu.Lock() + defer t.mu.Unlock() + + if t.writers == nil { + t.writers = make(map[string]*archiveWriter) + } + + w, ok := t.writers[webhookID] + if !ok { + w = newArchiveWriter(path, t.eng.log) + t.writers[webhookID] = w + } + + return w, nil +} diff --git a/internal/delivery/target_database_archive.go b/internal/delivery/target_database_archive.go new file mode 100644 index 0000000..547a0af --- /dev/null +++ b/internal/delivery/target_database_archive.go @@ -0,0 +1,312 @@ +package delivery + +import ( + "database/sql" + "encoding/json" + "errors" + "fmt" + "log/slog" + "os" + "sync" + "time" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +// archiveExpiryNever is the expiry sentinel (and default) that +// disables pruning so archived rows are kept forever. +const archiveExpiryNever = "never" + +// archiveReopenDebounce bounds how often an archive file is +// closed and reopened. After each write the handle is closed +// and reopened so an operator can move the file away for +// offline archiving, but never more than once per this window. +const archiveReopenDebounce = time.Second + +var ( + // errArchiveMissingWebhookID is returned when an event to + // archive has no webhook id to key its archive file on. + errArchiveMissingWebhookID = errors.New( + "cannot archive event without a webhook id", + ) + + // errArchiveNoDataDir is returned when the database target + // has no webhook database manager and so cannot locate the + // data directory for archive files. + errArchiveNoDataDir = errors.New( + "database target has no data directory", + ) + + // errArchiveExpiryNotPositive is returned when a + // user-supplied archive expiry parses as a duration but is + // zero or negative; "never" is the way to disable pruning. + errArchiveExpiryNotPositive = errors.New( + "expiry must be a positive duration or \"never\"", + ) +) + +// databaseTargetConfig is the optional per-target JSON config +// for a database (archive) target. +type databaseTargetConfig struct { + // Expiry is a Go duration (e.g. "720h") after which + // archived rows are pruned, or "never" (the default) to + // keep them forever. + Expiry string `json:"expiry"` +} + +// archivedEvent is one fully captured webhook event stored in a +// per-webhook archive database for long-term retention. It is a +// self-contained copy — independent of the per-webhook event +// database, which may prune events under its own retention. +type archivedEvent struct { + ID uint `gorm:"primaryKey;autoIncrement"` + EventID string `gorm:"index"` + WebhookID string + EntrypointID string + Method string + Headers string + Body string + ContentType string + + // ArchivedAt is when the row was archived and is the age + // basis for expiry pruning. + ArchivedAt time.Time `gorm:"index"` +} + +// parseArchiveExpiry reads the optional expiry from a database +// target's config JSON. An empty config, an empty expiry, or +// the literal "never" all mean keep forever, returned as a zero +// duration. Any other value must parse as a positive Go +// duration; a set-but-invalid value (unparseable, zero, or +// negative) is an error rather than a silent default, matching +// ValidateArchiveExpiry at target creation. +func parseArchiveExpiry( + configJSON string, +) (time.Duration, error) { + if configJSON == "" { + return 0, nil + } + + var cfg databaseTargetConfig + + err := json.Unmarshal([]byte(configJSON), &cfg) + if err != nil { + return 0, fmt.Errorf( + "parsing database target config: %w", err, + ) + } + + if cfg.Expiry == "" || cfg.Expiry == archiveExpiryNever { + return 0, nil + } + + dur, err := time.ParseDuration(cfg.Expiry) + if err != nil { + return 0, fmt.Errorf( + "parsing archive expiry %q: %w", cfg.Expiry, err, + ) + } + + if dur <= 0 { + return 0, fmt.Errorf( + "%w: %q", errArchiveExpiryNotPositive, cfg.Expiry, + ) + } + + return dur, nil +} + +// ValidateArchiveExpiry checks a user-supplied archive expiry +// for a database target at configuration time. Valid values are +// empty, "never" (both meaning keep forever), or a positive Go +// duration such as "720h". Anything else is an error, so a bad +// expiry is rejected when the target is created rather than +// failing every subsequent delivery. +func ValidateArchiveExpiry(expiry string) error { + if expiry == "" || expiry == archiveExpiryNever { + return nil + } + + dur, err := time.ParseDuration(expiry) + if err != nil { + return fmt.Errorf( + "expiry must be %q or a Go duration "+ + "such as \"720h\": %w", + archiveExpiryNever, err, + ) + } + + if dur <= 0 { + return fmt.Errorf( + "%w: %q", errArchiveExpiryNotPositive, expiry, + ) + } + + return nil +} + +// archiveWriter owns one per-webhook archive SQLite file. It +// serialises writes, and after each write closes and reopens +// the file (debounced to at most once per debounce window) so +// an operator can move the file away for offline archiving. The +// next write recreates a moved or removed file, because the +// file is opened create-if-missing and its schema is migrated +// on every open. +type archiveWriter struct { + mu sync.Mutex + path string + log *slog.Logger + debounce time.Duration + db *gorm.DB + lastReopen time.Time + reopens int +} + +// newArchiveWriter builds an archiveWriter for a file path with +// the default reopen debounce. +func newArchiveWriter( + path string, log *slog.Logger, +) *archiveWriter { + return &archiveWriter{ + path: path, + log: log, + debounce: archiveReopenDebounce, + } +} + +// write appends the event as a row, then applies the debounced +// close/reopen. It recreates the archive file if it was moved +// or removed since the last open. A positive expiry prunes rows +// older than it on each (re)open. +func (w *archiveWriter) write( + row archivedEvent, expiry time.Duration, +) error { + w.mu.Lock() + defer w.mu.Unlock() + + if w.db == nil || !fileExists(w.path) { + err := w.reopen(expiry) + if err != nil { + return err + } + } + + row.ArchivedAt = time.Now() + + err := w.db.Create(&row).Error + if err != nil { + return fmt.Errorf( + "archiving event to %s: %w", w.path, err, + ) + } + + if time.Since(w.lastReopen) >= w.debounce { + return w.reopen(expiry) + } + + return nil +} + +// open opens (creating if missing) the archive file, migrates +// its schema, records the reopen time, and prunes expired rows +// when expiry is positive. +func (w *archiveWriter) open(expiry time.Duration) error { + dbURL := fmt.Sprintf("file:%s?mode=rwc", w.path) + + sqlDB, err := sql.Open("sqlite", dbURL) + if err != nil { + return fmt.Errorf( + "opening archive database %s: %w", w.path, err, + ) + } + + gdb, err := gorm.Open( + sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, + ) + if err != nil { + _ = sqlDB.Close() + + return fmt.Errorf( + "connecting to archive database %s: %w", + w.path, err, + ) + } + + err = gdb.AutoMigrate(&archivedEvent{}) + if err != nil { + _ = sqlDB.Close() + + return fmt.Errorf( + "migrating archive database %s: %w", w.path, err, + ) + } + + w.db = gdb + w.lastReopen = time.Now() + w.reopens++ + + if expiry > 0 { + w.prune(expiry) + } + + return nil +} + +// reopen closes any open handle and opens the file afresh. The +// fresh open recreates the file if it was moved away. +func (w *archiveWriter) reopen(expiry time.Duration) error { + w.close() + + return w.open(expiry) +} + +// close closes the underlying handle, if any. +func (w *archiveWriter) close() { + if w.db == nil { + return + } + + sqlDB, err := w.db.DB() + if err == nil { + _ = sqlDB.Close() + } + + w.db = nil +} + +// prune deletes archived rows older than expiry, measured from +// each row's archived time. It runs on every (re)open, and +// because the file is reopened after writes this keeps the +// archive swept without a separate background sweeper. Failures +// are logged, not fatal: a prune error must not stop archiving. +func (w *archiveWriter) prune(expiry time.Duration) { + cutoff := time.Now().Add(-expiry) + + res := w.db.Where("archived_at < ?", cutoff). + Delete(&archivedEvent{}) + if res.Error != nil { + w.log.Error( + "failed to prune expired archive rows", + "path", w.path, + "error", res.Error, + ) + + return + } + + if res.RowsAffected > 0 { + w.log.Info( + "pruned expired archive rows", + "path", w.path, + "rows_deleted", res.RowsAffected, + ) + } +} + +// fileExists reports whether a path currently exists. +func fileExists(path string) bool { + _, err := os.Stat(path) + + return err == nil +} diff --git a/internal/delivery/target_database_test.go b/internal/delivery/target_database_test.go new file mode 100644 index 0000000..3ee38ac --- /dev/null +++ b/internal/delivery/target_database_test.go @@ -0,0 +1,395 @@ +package delivery_test + +import ( + "database/sql" + "fmt" + "log/slog" + "net/http" + "os" + "path/filepath" + "testing" + "time" + + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/driver/sqlite" + "gorm.io/gorm" + _ "modernc.org/sqlite" // Pure Go SQLite driver. + "sneak.berlin/go/webhooker/internal/database" + "sneak.berlin/go/webhooker/internal/delivery" +) + +func archiveTestLogger() *slog.Logger { + return slog.New(slog.NewTextHandler( + os.Stderr, + &slog.HandlerOptions{Level: slog.LevelDebug}, + )) +} + +// openArchiveDBForRead opens an archive file read-only so a +// test can inspect the rows the writer persisted. +func openArchiveDBForRead( + t *testing.T, path string, +) *gorm.DB { + t.Helper() + + sqlDB, err := sql.Open( + "sqlite", + fmt.Sprintf("file:%s?mode=ro", path), + ) + require.NoError(t, err) + + t.Cleanup(func() { _ = sqlDB.Close() }) + + gdb, err := gorm.Open( + sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, + ) + require.NoError(t, err) + + return gdb +} + +// removeArchiveFiles simulates an operator moving the archive +// away by deleting the SQLite file and its sidecar files. +func removeArchiveFiles(t *testing.T, path string) { + t.Helper() + + for _, suffix := range []string{ + "", "-wal", "-shm", "-journal", + } { + err := os.Remove(path + suffix) + if err != nil && !os.IsNotExist(err) { + t.Fatalf("removing %s%s: %v", path, suffix, err) + } + } +} + +// TestDeliverDatabase_ArchivesEvent verifies that delivering to +// a database target marks the delivery delivered and archives +// the full event into a separate per-webhook archive file. +func TestDeliverDatabase_ArchivesEvent(t *testing.T) { + t.Parallel() + + dataDir := t.TempDir() + dbMgr := database.NewTestWebhookDBManager(dataDir) + + e := delivery.NewTestEngineWithDB( + nil, dbMgr, + archiveTestLogger(), + &http.Client{Timeout: 5 * time.Second}, + 1, + ) + + webhookDB := testWebhookDB(t) + event := seedEvent(t, webhookDB, `{"archived":true}`) + d := seedDatabaseTargetDelivery(t, webhookDB, event, "") + + e.ExportDeliverDatabase(webhookDB, d) + + var updated database.Delivery + + require.NoError(t, webhookDB.First( + &updated, "id = ?", d.ID, + ).Error) + assert.Equal(t, + database.DeliveryStatusDelivered, updated.Status, + "database target should mark the delivery delivered", + ) + + archivePath := filepath.Join( + dataDir, + fmt.Sprintf("archive-%s.db", event.WebhookID), + ) + assert.FileExists(t, archivePath) + + rdb := openArchiveDBForRead(t, archivePath) + + var rows []delivery.ExportArchivedEvent + + require.NoError(t, rdb.Find(&rows).Error) + require.Len(t, rows, 1) + assert.Equal(t, event.ID, rows[0].EventID) + assert.Equal(t, event.WebhookID, rows[0].WebhookID) + assert.Equal(t, event.Method, rows[0].Method) + assert.JSONEq(t, `{"archived":true}`, rows[0].Body) +} + +func TestArchiveWriter_WritesRow(t *testing.T) { + t.Parallel() + + path := filepath.Join(t.TempDir(), "archive-wh.db") + w := delivery.NewExportArchiveWriter( + path, archiveTestLogger(), 0, + ) + + row := delivery.ExportArchivedEvent{ + EventID: "ev-1", + WebhookID: "wh-1", + EntrypointID: "ep-1", + Method: "POST", + Headers: `{"X":"Y"}`, + Body: `{"hello":"world"}`, + ContentType: "application/json", + } + + require.NoError(t, w.Write(row, 0)) + assert.FileExists(t, path) + + var got []delivery.ExportArchivedEvent + + require.NoError(t, w.DB().Find(&got).Error) + require.Len(t, got, 1) + assert.Equal(t, "ev-1", got[0].EventID) + assert.Equal(t, "wh-1", got[0].WebhookID) + assert.Equal(t, "ep-1", got[0].EntrypointID) + assert.Equal(t, row.Method, got[0].Method) + assert.Equal(t, row.ContentType, got[0].ContentType) + assert.JSONEq(t, `{"hello":"world"}`, got[0].Body) + assert.False(t, got[0].ArchivedAt.IsZero()) +} + +func TestArchiveWriter_RecreatesAfterRemoval( + t *testing.T, +) { + t.Parallel() + + path := filepath.Join(t.TempDir(), "archive-wh.db") + w := delivery.NewExportArchiveWriter( + path, archiveTestLogger(), 0, + ) + + require.NoError(t, w.Write( + delivery.ExportArchivedEvent{EventID: "a"}, 0, + )) + assert.FileExists(t, path) + + // The operator moves the archive away while the handle is + // still open. + removeArchiveFiles(t, path) + require.NoFileExists(t, path) + + // The next write recreates the file with a fresh schema and + // only the new row. + require.NoError(t, w.Write( + delivery.ExportArchivedEvent{EventID: "b"}, 0, + )) + assert.FileExists(t, path) + + var got []delivery.ExportArchivedEvent + + require.NoError(t, w.DB().Find(&got).Error) + require.Len(t, got, 1) + assert.Equal(t, "b", got[0].EventID) +} + +func TestArchiveWriter_ReopenDebounce(t *testing.T) { + t.Parallel() + + // A generous debounce keeps the two rapid writes inside + // the window even on a heavily loaded test machine. + path := filepath.Join(t.TempDir(), "archive-wh.db") + w := delivery.NewExportArchiveWriter( + path, archiveTestLogger(), 2*time.Second, + ) + + require.NoError(t, w.Write( + delivery.ExportArchivedEvent{EventID: "a"}, 0, + )) + require.NoError(t, w.Write( + delivery.ExportArchivedEvent{EventID: "b"}, 0, + )) + + // Two writes inside the debounce window trigger only the + // initial open — no extra close/reopen. + assert.Equal(t, 1, w.Reopens()) + + time.Sleep(2100 * time.Millisecond) + + require.NoError(t, w.Write( + delivery.ExportArchivedEvent{EventID: "c"}, 0, + )) + + // A write after the window elapses closes and reopens once. + assert.Equal(t, 2, w.Reopens()) +} + +func TestArchiveWriter_ExpiryPrune(t *testing.T) { + t.Parallel() + + path := filepath.Join(t.TempDir(), "archive-wh.db") + w := delivery.NewExportArchiveWriter( + path, archiveTestLogger(), 0, + ) + + require.NoError(t, w.Open(0)) + + old := delivery.ExportArchivedEvent{ + EventID: "old", + ArchivedAt: time.Now().Add(-2 * time.Hour), + } + fresh := delivery.ExportArchivedEvent{ + EventID: "fresh", + ArchivedAt: time.Now(), + } + + require.NoError(t, w.DB().Create(&old).Error) + require.NoError(t, w.DB().Create(&fresh).Error) + + // Reopening with a one-hour expiry prunes the old row. + require.NoError(t, w.Reopen(time.Hour)) + + var got []delivery.ExportArchivedEvent + + require.NoError(t, w.DB().Find(&got).Error) + require.Len(t, got, 1) + assert.Equal(t, "fresh", got[0].EventID) +} + +func TestParseArchiveExpiry(t *testing.T) { + t.Parallel() + + cases := []struct { + name string + in string + want time.Duration + wantErr bool + }{ + {"empty config", "", 0, false}, + {"explicit never", `{"expiry":"never"}`, 0, false}, + {"empty expiry", `{"expiry":""}`, 0, false}, + {"duration", `{"expiry":"1h"}`, time.Hour, false}, + {"unparseable", `{"expiry":"nonsense"}`, 0, true}, + {"zero duration", `{"expiry":"0s"}`, 0, true}, + {"negative duration", `{"expiry":"-5h"}`, 0, true}, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + got, err := delivery.ExportParseArchiveExpiry(tc.in) + if tc.wantErr { + require.Error(t, err) + + return + } + + require.NoError(t, err) + assert.Equal(t, tc.want, got) + }) + } +} + +// seedDatabaseTargetDelivery seeds a pending delivery for a +// database target with the given config JSON and returns the +// in-memory delivery the target handler is invoked with. +func seedDatabaseTargetDelivery( + t *testing.T, + webhookDB *gorm.DB, + event database.Event, + config string, +) *database.Delivery { + t.Helper() + + dlv := seedDelivery( + t, webhookDB, event.ID, uuid.New().String(), + database.DeliveryStatusPending, + ) + + d := &database.Delivery{ + EventID: event.ID, + TargetID: dlv.TargetID, + Status: database.DeliveryStatusPending, + Event: event, + Target: database.Target{ + Name: "test-db", + Type: database.TargetTypeDatabase, + Config: config, + }, + } + d.ID = dlv.ID + + return d +} + +// TestDeliverDatabase_ArchiveFailureFailsDelivery verifies that +// an archive error (here: an unparseable expiry in the target +// config) fails the delivery loudly: the attempt is recorded as +// failed with the error and the delivery is marked failed, not +// delivered. +func TestDeliverDatabase_ArchiveFailureFailsDelivery( + t *testing.T, +) { + t.Parallel() + + dataDir := t.TempDir() + + e := delivery.NewTestEngineWithDB( + nil, database.NewTestWebhookDBManager(dataDir), + archiveTestLogger(), + &http.Client{Timeout: 5 * time.Second}, + 1, + ) + + webhookDB := testWebhookDB(t) + event := seedEvent(t, webhookDB, `{"archived":false}`) + d := seedDatabaseTargetDelivery( + t, webhookDB, event, `{"expiry":"nonsense"}`, + ) + + e.ExportDeliverDatabase(webhookDB, d) + + var updated database.Delivery + + require.NoError(t, webhookDB.First( + &updated, "id = ?", d.ID, + ).Error) + assert.Equal(t, + database.DeliveryStatusFailed, updated.Status, + "archive failure must mark the delivery failed", + ) + + var results []database.DeliveryResult + + require.NoError(t, webhookDB.Where( + "delivery_id = ?", d.ID, + ).Find(&results).Error) + require.Len(t, results, 1) + assert.False(t, + results[0].Success, + "the attempt must be recorded as failed", + ) + assert.Contains(t, + results[0].Error, "nonsense", + "the archive error must be recorded on the attempt", + ) + + assert.NoFileExists(t, + filepath.Join( + dataDir, + fmt.Sprintf("archive-%s.db", event.WebhookID), + ), + "no archive file should exist for a failed config", + ) +} + +func TestValidateArchiveExpiry(t *testing.T) { + t.Parallel() + + valid := []string{"", "never", "1h", "720h", "30m"} + for _, in := range valid { + require.NoError(t, + delivery.ValidateArchiveExpiry(in), + "expiry %q should be accepted", in, + ) + } + + invalid := []string{"nonsense", "7d", "-5h", "0s", "0"} + for _, in := range invalid { + require.Error(t, + delivery.ValidateArchiveExpiry(in), + "expiry %q should be rejected", in, + ) + } +} diff --git a/internal/handlers/export_test.go b/internal/handlers/export_test.go index 831ed96..cff6c09 100644 --- a/internal/handlers/export_test.go +++ b/internal/handlers/export_test.go @@ -22,3 +22,13 @@ func (s *Handlers) BuildSlackTargetConfigForTest( ) (string, error) { return s.buildSlackTargetConfig(w, r, targetURL) } + +// BuildDatabaseTargetConfigForTest exposes +// buildDatabaseTargetConfig for use in the handlers_test +// package. +func (s *Handlers) BuildDatabaseTargetConfigForTest( + w http.ResponseWriter, + expiry string, +) (string, error) { + return s.buildDatabaseTargetConfig(w, expiry) +} diff --git a/internal/handlers/handlers_test.go b/internal/handlers/handlers_test.go index ec2791b..8aae54a 100644 --- a/internal/handlers/handlers_test.go +++ b/internal/handlers/handlers_test.go @@ -186,3 +186,57 @@ func TestRenderTemplate(t *testing.T) { t, http.StatusInternalServerError, w.Code, ) } + +func TestBuildDatabaseTargetConfig_Valid(t *testing.T) { + t.Parallel() + + var h *handlers.Handlers + + app := newTestApp(t, &h) + app.RequireStart() + + t.Cleanup(app.RequireStop) + + // Empty expiry: the keep-forever default, empty config. + w := httptest.NewRecorder() + cfg, err := h.BuildDatabaseTargetConfigForTest(w, "") + require.NoError(t, err) + assert.Empty(t, cfg) + + // Explicit never is stored as config. + w = httptest.NewRecorder() + cfg, err = h.BuildDatabaseTargetConfigForTest(w, "never") + require.NoError(t, err) + assert.JSONEq(t, `{"expiry":"never"}`, cfg) + + // A positive duration is stored as config. + w = httptest.NewRecorder() + cfg, err = h.BuildDatabaseTargetConfigForTest(w, "720h") + require.NoError(t, err) + assert.JSONEq(t, `{"expiry":"720h"}`, cfg) +} + +func TestBuildDatabaseTargetConfig_RejectsBadExpiry( + t *testing.T, +) { + t.Parallel() + + var h *handlers.Handlers + + app := newTestApp(t, &h) + app.RequireStart() + + t.Cleanup(app.RequireStop) + + for _, bad := range []string{"nonsense", "7d", "-5h"} { + w := httptest.NewRecorder() + cfg, err := h.BuildDatabaseTargetConfigForTest(w, bad) + + require.Error(t, err, "expiry %q", bad) + assert.Empty(t, cfg) + assert.Equal( + t, http.StatusBadRequest, w.Code, + "expiry %q should be rejected with 400", bad, + ) + } +} diff --git a/internal/handlers/source_management.go b/internal/handlers/source_management.go index 917307a..7d43ad5 100644 --- a/internal/handlers/source_management.go +++ b/internal/handlers/source_management.go @@ -5,6 +5,7 @@ import ( "errors" "net/http" "strconv" + "strings" "github.com/go-chi/chi" "github.com/google/uuid" @@ -815,6 +816,7 @@ func (h *Handlers) processTargetCreate( targetType := database.TargetType(r.FormValue("type")) targetURL := r.FormValue("url") maxRetriesStr := r.FormValue("max_retries") + expiry := r.FormValue("expiry") if name == "" { http.Error( @@ -834,7 +836,7 @@ func (h *Handlers) processTargetCreate( } configJSON, err := h.buildTargetConfig( - w, r, targetType, targetURL, + w, r, targetType, targetURL, expiry, ) if err != nil { return @@ -892,18 +894,22 @@ func parseNonNegativeInt(s string) int { } // buildTargetConfig builds the JSON config string for a target. +// The expiry form value is read by the caller (which bounds the +// request body) and applies to database targets only. func (h *Handlers) buildTargetConfig( w http.ResponseWriter, r *http.Request, targetType database.TargetType, - targetURL string, + targetURL, expiry string, ) (string, error) { switch targetType { case database.TargetTypeHTTP: return h.buildHTTPTargetConfig(w, r, targetURL) case database.TargetTypeSlack: return h.buildSlackTargetConfig(w, r, targetURL) - case database.TargetTypeDatabase, database.TargetTypeLog: + case database.TargetTypeDatabase: + return h.buildDatabaseTargetConfig(w, expiry) + case database.TargetTypeLog: return "", nil default: http.Error( @@ -1013,6 +1019,47 @@ func (h *Handlers) buildSlackTargetConfig( return string(configBytes), nil } +// buildDatabaseTargetConfig builds config JSON for a database +// (archive) target. The optional expiry (a form value read by +// the caller, which bounds the request body) is validated here, +// at creation time, so an unparseable value is rejected with a +// 400 instead of failing every subsequent delivery. An empty +// expiry yields an empty config (the keep-forever default). +func (h *Handlers) buildDatabaseTargetConfig( + w http.ResponseWriter, + expiry string, +) (string, error) { + expiry = strings.TrimSpace(expiry) + if expiry == "" { + return "", nil + } + + err := delivery.ValidateArchiveExpiry(expiry) + if err != nil { + http.Error( + w, + "Invalid archive expiry: "+err.Error(), + http.StatusBadRequest, + ) + + return "", err + } + + cfg := map[string]any{"expiry": expiry} + + configBytes, err := json.Marshal(cfg) + if err != nil { + http.Error( + w, "Internal server error", + http.StatusInternalServerError, + ) + + return "", err + } + + return string(configBytes), nil +} + // HandleEntrypointDelete handles deleting an entrypoint. func (h *Handlers) HandleEntrypointDelete() http.HandlerFunc { return h.deleteChildResource( diff --git a/templates/source_detail.html b/templates/source_detail.html index 5f20957..5b9944a 100644 --- a/templates/source_detail.html +++ b/templates/source_detail.html @@ -113,6 +113,10 @@

Slack or Mattermost incoming webhook URL. Payloads are pretty-printed in code blocks.

+
+ +

Archive expiry: "never" (default) keeps rows forever, or a duration like "720h" prunes older rows.

+