From a55d4208a589d8b17050a1620608439921bc685f Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Fri, 2 Oct 2026 12:13:02 +0000 Subject: [PATCH] Add a Download button that exports a database target's archive as gzipped JSON (closes #374) Each database target on the webhook page links to /hook/ID/targets/TARGETID/download, which streams the target's archive as archive-WEBHOOK-TARGET-YYYYMMDDTHHMMSSZ.json.gz: the webhook, the target, exported_at, and archived_events, one object per row keyed by column, a body that is not valid UTF-8 in base64 with body_encoding. The export reads one row at a time on its own connection inside a read-only transaction: one snapshot, and no write lock. The handler holds the rename lock only while it reads the names and opens the file. A missing archive exports empty and is not created. Model: opus-5-5 --- README.md | 26 ++ internal/delivery/target_database.go | 16 +- internal/delivery/target_database_export.go | 270 +++++++++++++ .../delivery/target_database_export_test.go | 373 ++++++++++++++++++ internal/handlers/handlers.go | 3 +- internal/handlers/target_download.go | 82 ++++ internal/handlers/target_download_test.go | 185 +++++++++ internal/handlers/target_edit_test.go | 8 +- internal/server/archive_download_test.go | 69 ++++ internal/server/routes.go | 4 + templates/source_detail.html | 3 + 11 files changed, 1026 insertions(+), 13 deletions(-) create mode 100644 internal/delivery/target_database_export.go create mode 100644 internal/delivery/target_database_export_test.go create mode 100644 internal/handlers/target_download.go create mode 100644 internal/handlers/target_download_test.go create mode 100644 internal/server/archive_download_test.go diff --git a/README.md b/README.md index 227ffa3..1314e73 100644 --- a/README.md +++ b/README.md @@ -1992,6 +1992,30 @@ Because each `database` target has its own archive file, a target's webhook with different expiries keep two archives, each pruned on its own schedule. +Each `database` target on the webhook page has a **Download** button, +which returns its archive as one gzipped JSON file, +`archive-{webhook_name}-{target_name}-{YYYYMMDDTHHMMSSZ}.json.gz`, the +names made safe as above and the time in UTC. The file holds one +object: `webhook` and `target`, each an `id` and a `name`; +`exported_at`; and `archived_events`, one object per archived row with +every column, keyed by column name. A body that is not valid UTF-8 is +written in base64, with `"body_encoding": "base64"` beside it. An +archive that does not exist yet, or was moved away, downloads with an +empty `archived_events`; the download never creates the file. + +The download streams: each row is read and written out compressed +before the next is read, so neither the archive nor the JSON is held in +memory. It reads on a connection of its own, inside one read-only +transaction, so the file holds the archive as it stood when the +download started, and archive writes go on meanwhile, since under WAL a +reader never blocks a writer. While it runs, the `-wal` cannot be +checkpointed past what it reads, so a long download lets the `-wal` +grow. It finds the file by the stored names under the lock that webhook +edits, target edits and target creation hold, and lets go once the file +is open: a rename during the download moves the file without affecting +it. Like every request it is cut off after 60 seconds, which leaves the +file incomplete and failing to decompress. + Deleting a webhook releases its archives: the delivery engine's cached archive writers are dropped and their file handles closed, so nothing lingers after the webhook is gone. The archive **files themselves are @@ -2854,6 +2878,7 @@ returns to the page that was asked for. | `POST` | `/hook/{id}/targets` | Add target to webhook | | `GET` | `/hook/{id}/targets/{targetID}/edit` | Edit target form. The one page that renders a target's destination URL and header values in full, rather than masked | | `POST` | `/hook/{id}/targets/{targetID}/edit` | Edit target submission | +| `GET` | `/hook/{id}/targets/{targetID}/download` | Download a `database` target's archive as one gzipped JSON file. See [Database Architecture](#database-architecture) | | `POST` | `/hook/{id}/targets/{targetID}/delete` | Delete a target | | `POST` | `/hook/{id}/targets/{targetID}/toggle` | Enable or disable a target | @@ -2934,6 +2959,7 @@ webhooker/ │ │ ├── target_slack.go # Slack/Mattermost incoming-webhook target │ │ ├── target_database.go # Database archive target │ │ ├── target_database_archive.go # Archive file lifecycle and pruning +│ │ ├── target_database_export.go # Archive download as gzipped JSON │ │ ├── target_log.go # Log target (stdout) │ │ ├── target_config_view.go # Masked target config for templates │ │ ├── archive_sweeper.go # Periodic pruning of idle archives diff --git a/internal/delivery/target_database.go b/internal/delivery/target_database.go index 3fb52db..483e317 100644 --- a/internal/delivery/target_database.go +++ b/internal/delivery/target_database.go @@ -3,7 +3,6 @@ package delivery import ( "context" "fmt" - "path/filepath" "strings" "sync" "time" @@ -277,10 +276,9 @@ func (t *databaseTarget) releaseSweepWriter( } // newWriter builds the writer for a database target's archive. The -// file lives beside the webhook's event database in the data -// directory and is named for the webhook and the target as the main -// database has them now; from then on only rename changes the name -// the writer uses. It does not touch the archive file. +// file is the one ArchivePath gives for the webhook and the target as +// the main database names them now; from then on only rename changes +// the name the writer uses. It does not touch the archive file. func (t *databaseTarget) newWriter( targetID string, ) (*archiveWriter, error) { @@ -299,12 +297,10 @@ func (t *databaseTarget) newWriter( ) } - dir := filepath.Dir(t.eng.dbManager.DBPath(target.WebhookID)) - name := ArchiveFileName( - target.Webhook.Name, target.Name, target.ID, + w := newArchiveWriter( + ArchivePath(t.eng.dbManager, &target.Webhook, &target), + t.eng.log, ) - - w := newArchiveWriter(filepath.Join(dir, name), t.eng.log) w.webhookID = target.WebhookID return w, nil diff --git a/internal/delivery/target_database_export.go b/internal/delivery/target_database_export.go new file mode 100644 index 0000000..d781327 --- /dev/null +++ b/internal/delivery/target_database_export.go @@ -0,0 +1,270 @@ +package delivery + +import ( + "compress/gzip" + "context" + "database/sql" + "encoding/base64" + "encoding/json" + "errors" + "fmt" + "io" + "path/filepath" + "time" + "unicode/utf8" + + "sneak.berlin/go/webhooker/internal/database" +) + +// archiveTableQuery counts the archive's table: 0 when the archive +// writer has created the file but not yet the table in it. +const archiveTableQuery = "SELECT count(*) FROM sqlite_master " + + "WHERE type = 'table' AND name = 'archived_events'" + +// archiveNextRowQuery reads every column of the first archived row +// after a given id. An export reads the archive a row at a time this +// way rather than through one cursor, because the Scan check in +// internal/gormlog accepts only a Scan straight on QueryRowContext's +// result. +const archiveNextRowQuery = "SELECT id, event_id, webhook_id, " + + "entrypoint_id, method, headers, body, content_type, archived_at " + + "FROM archived_events WHERE id > ? ORDER BY id LIMIT 1" + +// ArchivePath returns where a database target's archive file is: in +// the data directory, beside the webhook's event database, under the +// name ArchiveFileName gives it. +func ArchivePath( + dbMgr *database.WebhookDBManager, + webhook *database.Webhook, + target *database.Target, +) string { + return filepath.Join( + filepath.Dir(dbMgr.DBPath(webhook.ID)), + ArchiveFileName(webhook.Name, target.Name, target.ID), + ) +} + +// ArchiveExportFileName returns the name a database target's archive +// downloads under: +// archive-WEBHOOKNAME-TARGETNAME-YYYYMMDDTHHMMSSZ.json.gz, the names +// made safe as in ArchiveFileName and the time in UTC. +func ArchiveExportFileName( + webhookName, targetName string, at time.Time, +) string { + return "archive-" + archiveNamePart(webhookName) + "-" + + archiveNamePart(targetName) + "-" + + at.UTC().Format("20060102T150405Z") + ".json.gz" +} + +// ArchiveExport is a database target's archive opened for download. +// It reads the file on its own connection, inside one read-only +// transaction, so it writes out the archive as it stood when +// OpenArchiveExport returned. +// +// Archives are in WAL mode, where a reader works from a snapshot and +// never blocks a writer: archive writes go on while an export is open, +// and the export does not see them. SQLite cannot checkpoint the -wal +// past an open snapshot, so the -wal grows until the export is closed. +type ArchiveExport struct { + db *sql.DB + tx *sql.Tx + + // empty is true when there is nothing to read: no file, or a file + // without the archive's table yet. + empty bool +} + +// exportedName is how an export names its webhook and its target. +type exportedName struct { + ID string `json:"id"` + Name string `json:"name"` +} + +// OpenArchiveExport opens the archive file at path for export and +// takes the snapshot the export reads. It never creates the file: with +// no file at path, the export has no rows. +// +// Once it has returned, the file is open, so a rename or a move of it +// does not affect the export, which reads the same file under its new +// name. +func OpenArchiveExport( + ctx context.Context, path string, +) (*ArchiveExport, error) { + if !fileExists(path) { + return &ArchiveExport{empty: true}, nil + } + + db, err := database.OpenSQLite(path, archiveModeExisting) + if err != nil { + return nil, fmt.Errorf("opening archive %s: %w", path, err) + } + + // ReadOnly makes the driver begin a deferred transaction in place + // of the BEGIN IMMEDIATE the connection string asks for, so the + // export never takes the archive's write lock. + tx, err := db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true}) + if err != nil { + _ = db.Close() + + return nil, fmt.Errorf("reading archive %s: %w", path, err) + } + + // The transaction's first read is what takes the snapshot. + var tables int + + err = tx.QueryRowContext(ctx, archiveTableQuery).Scan(&tables) + if err != nil { + _ = tx.Rollback() + _ = db.Close() + + return nil, fmt.Errorf("reading archive %s: %w", path, err) + } + + return &ArchiveExport{db: db, tx: tx, empty: tables == 0}, nil +} + +// WriteGzipJSON writes the export to w as one gzipped JSON object: +// webhook and target, each an id and a name; exported_at; and +// archived_events, one object per archived row, keyed by column name. +// A body that is not valid UTF-8 cannot be a JSON string, so it is +// written in base64, with "body_encoding": "base64" beside it. +// +// Each row is written out before the next is read, so neither the +// archive nor its JSON is ever held in memory whole. After an error +// the gzip stream is left unfinished, so what was written does not +// decompress as a whole file. +func (x *ArchiveExport) WriteGzipJSON( + ctx context.Context, + w io.Writer, + webhook *database.Webhook, + target *database.Target, + exportedAt time.Time, +) error { + head, err := json.Marshal(map[string]any{ + "webhook": exportedName{ID: webhook.ID, Name: webhook.Name}, + "target": exportedName{ID: target.ID, Name: target.Name}, + "exported_at": exportedAt.UTC(), + }) + if err != nil { + return fmt.Errorf("encoding archive export: %w", err) + } + + zw := gzip.NewWriter(w) + + err = x.writeJSON(ctx, zw, head) + if err != nil { + return fmt.Errorf("writing archive export: %w", err) + } + + return zw.Close() +} + +// Close ends the export's transaction and closes its connection. +func (x *ArchiveExport) Close() error { + if x.db == nil { + return nil + } + + _ = x.tx.Rollback() + + return x.db.Close() +} + +// writeJSON writes head with archived_events added as its last key, +// the rows going into it one at a time. +func (x *ArchiveExport) writeJSON( + ctx context.Context, w io.Writer, head []byte, +) error { + // head goes out without its closing brace, so that + // archived_events can follow it. + _, err := w.Write(head[:len(head)-1]) + if err != nil { + return err + } + + _, err = io.WriteString(w, `,"archived_events":[`) + if err != nil { + return err + } + + err = x.writeRows(ctx, w) + if err != nil { + return err + } + + _, err = io.WriteString(w, "\n]}\n") + + return err +} + +// writeRows writes each archived row to w, oldest first, one per line, +// separated by commas. +func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error { + if x.empty { + return nil + } + + next, err := x.tx.PrepareContext(ctx, archiveNextRowQuery) + if err != nil { + return err + } + + defer func() { _ = next.Close() }() + + var ev archivedEvent + + for sep := "\n"; ; sep = ",\n" { + err = next.QueryRowContext(ctx, ev.ID).Scan( + &ev.ID, &ev.EventID, &ev.WebhookID, &ev.EntrypointID, + &ev.Method, &ev.Headers, &ev.Body, &ev.ContentType, + &ev.ArchivedAt, + ) + if errors.Is(err, sql.ErrNoRows) { + return nil + } + + if err != nil { + return err + } + + _, err = io.WriteString(w, sep) + if err != nil { + return err + } + + err = writeRow(w, &ev) + if err != nil { + return err + } + } +} + +// writeRow writes an archived row to w as a JSON object keyed by +// column name, its body in base64 when it is not valid UTF-8. +func writeRow(w io.Writer, ev *archivedEvent) error { + row := map[string]any{ + "id": ev.ID, + "event_id": ev.EventID, + "webhook_id": ev.WebhookID, + "entrypoint_id": ev.EntrypointID, + "method": ev.Method, + "headers": ev.Headers, + "body": ev.Body, + "content_type": ev.ContentType, + "archived_at": ev.ArchivedAt.UTC(), + } + + if !utf8.ValidString(ev.Body) { + row["body"] = base64.StdEncoding.EncodeToString([]byte(ev.Body)) + row["body_encoding"] = "base64" + } + + line, err := json.Marshal(row) + if err != nil { + return err + } + + _, err = w.Write(line) + + return err +} diff --git a/internal/delivery/target_database_export_test.go b/internal/delivery/target_database_export_test.go new file mode 100644 index 0000000..95f8383 --- /dev/null +++ b/internal/delivery/target_database_export_test.go @@ -0,0 +1,373 @@ +package delivery_test + +import ( + "bytes" + "compress/gzip" + "encoding/base64" + "encoding/json" + "fmt" + "io" + "os" + "path/filepath" + "runtime" + "runtime/debug" + "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 webhook and the target the export tests' archives belong to. +const ( + exportWebhookID = "wh-export" + exportWebhookName = "Orders (EU)" + exportTargetID = "tgt-export" + exportTargetName = "Long-term archive" +) + +const ( + // binaryBody is a body that is not valid UTF-8. + binaryBody = "\xff\xfe\x00\x01binary\x80" + + // openedEventID is the event the snapshot tests archive before + // they open the export. + openedEventID = "opened" +) + +// writeExportTo writes export to w as the archive of the export tests' +// webhook and target, exported at 2026-10-02T12:03:04Z. +func writeExportTo( + t *testing.T, export *delivery.ArchiveExport, w io.Writer, +) error { + t.Helper() + + return export.WriteGzipJSON( + t.Context(), w, + &database.Webhook{ + BaseModel: database.BaseModel{ID: exportWebhookID}, + Name: exportWebhookName, + }, + &database.Target{ + BaseModel: database.BaseModel{ID: exportTargetID}, + Name: exportTargetName, + }, + time.Date(2026, 10, 2, 12, 3, 4, 0, time.UTC), + ) +} + +// exportArchive runs a whole export of the archive at path and returns +// its JSON, decompressed and parsed. +func exportArchive(t *testing.T, path string) map[string]any { + t.Helper() + + export, err := delivery.OpenArchiveExport(t.Context(), path) + require.NoError(t, err) + + defer func() { require.NoError(t, export.Close()) }() + + return writeExport(t, export) +} + +// writeExport writes an opened export and returns its JSON, +// decompressed and parsed. Reading to the end makes the gzip reader +// check that the stream was finished. +func writeExport( + t *testing.T, export *delivery.ArchiveExport, +) map[string]any { + t.Helper() + + var buf bytes.Buffer + + require.NoError(t, writeExportTo(t, export, &buf)) + + zr, err := gzip.NewReader(&buf) + require.NoError(t, err) + + raw, err := io.ReadAll(zr) + require.NoError(t, err) + + var got map[string]any + + require.NoError(t, json.Unmarshal(raw, &got)) + + return got +} + +// exportedEvents returns an export's archived_events. +func exportedEvents(t *testing.T, got map[string]any) []map[string]any { + t.Helper() + + list, ok := got["archived_events"].([]any) + require.True(t, ok, "archived_events must be an array: %v", got) + + events := make([]map[string]any, len(list)) + + for i, v := range list { + events[i], ok = v.(map[string]any) + require.True(t, ok, "an archived event must be an object: %v", v) + } + + return events +} + +// exportedEventIDs returns the event_id of each of an export's +// archived_events. +func exportedEventIDs(t *testing.T, got map[string]any) []string { + t.Helper() + + events := exportedEvents(t, got) + ids := make([]string, 0, len(events)) + + for _, ev := range events { + ids = append(ids, fmt.Sprint(ev["event_id"])) + } + + return ids +} + +// TestArchiveExport_MatchesStoredRows proves an export holds the +// webhook, the target, the time, and every column of every stored +// row: a body that is valid UTF-8 as a string, and one that is not in +// base64, marked as such. +func TestArchiveExport_MatchesStoredRows(t *testing.T) { + t.Parallel() + + path := filepath.Join(t.TempDir(), "archive.db") + w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0) + bodies := []string{`{"order":1}`, "plain text", "", binaryBody} + + for i, body := range bodies { + require.NoError(t, w.Write(delivery.ExportArchivedEvent{ + EventID: fmt.Sprintf("ev-%d", i), + WebhookID: exportWebhookID, + EntrypointID: "ep-1", + Method: "POST", + Headers: `{"X-Test":["yes"]}`, + Body: body, + ContentType: testContentType, + }, 0)) + } + + var stored []delivery.ExportArchivedEvent + + require.NoError(t, openArchiveDBForRead(t, path). + Order("id").Find(&stored).Error) + + got := exportArchive(t, path) + + assert.Equal(t, + map[string]any{"id": exportWebhookID, "name": exportWebhookName}, + got["webhook"], + ) + assert.Equal(t, + map[string]any{"id": exportTargetID, "name": exportTargetName}, + got["target"], + ) + assert.Equal(t, "2026-10-02T12:03:04Z", got["exported_at"]) + + events := exportedEvents(t, got) + require.Len(t, events, len(bodies)) + + for i, row := range stored { + assertExportedRow(t, row, events[i]) + } +} + +// assertExportedRow checks that ev, from an export, holds every column +// of the stored row. +func assertExportedRow( + t *testing.T, row delivery.ExportArchivedEvent, ev map[string]any, +) { + t.Helper() + + archivedAt, err := time.Parse( + time.RFC3339Nano, fmt.Sprint(ev["archived_at"]), + ) + require.NoError(t, err) + assert.True(t, archivedAt.Equal(row.ArchivedAt)) + + assert.EqualValues(t, row.ID, ev["id"]) + assert.Equal(t, row.EventID, ev["event_id"]) + assert.Equal(t, row.WebhookID, ev["webhook_id"]) + assert.Equal(t, row.EntrypointID, ev["entrypoint_id"]) + assert.Equal(t, row.Method, ev["method"]) + assert.Equal(t, row.Headers, ev["headers"]) + assert.Equal(t, row.ContentType, ev["content_type"]) + + if row.Body != binaryBody { + assert.Equal(t, row.Body, ev["body"]) + assert.Len(t, ev, 9, "the nine columns and nothing else: %v", ev) + + return + } + + body, err := base64.StdEncoding.DecodeString(fmt.Sprint(ev["body"])) + require.NoError(t, err) + assert.Equal(t, binaryBody, string(body)) + assert.Equal(t, "base64", ev["body_encoding"]) + assert.Len(t, ev, 10, "the nine columns and body_encoding: %v", ev) +} + +// TestArchiveExport_Empty proves an archive with nothing in it exports +// as an empty archived_events: no file, which the export must not +// create; a file the archive writer has not yet put its table in; and +// a table with no rows. +func TestArchiveExport_Empty(t *testing.T) { + t.Parallel() + + dir := t.TempDir() + missing := filepath.Join(dir, "missing.db") + noTable := filepath.Join(dir, "no-table.db") + noRows := filepath.Join(dir, "no-rows.db") + + require.NoError(t, os.WriteFile(noTable, nil, 0o600)) + require.NoError(t, + delivery.NewExportArchiveWriter(noRows, archiveTestLogger(), 0). + Open(0), + ) + + for _, path := range []string{missing, noTable, noRows} { + assert.Empty(t, exportedEvents(t, exportArchive(t, path)), path) + } + + for _, suffix := range archiveFileSuffixes() { + assert.NoFileExists(t, missing+suffix) + } +} + +// TestArchiveExport_ReadsOneSnapshot proves an export writes the +// archive as it was when it was opened, and holds up no archive +// write: a row written while the export is open is stored, and is not +// in the export. A write held up for the whole busy timeout would +// fail. +func TestArchiveExport_ReadsOneSnapshot(t *testing.T) { + t.Parallel() + + path := filepath.Join(t.TempDir(), "archive.db") + w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0) + + require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: openedEventID}, 0)) + + export, err := delivery.OpenArchiveExport(t.Context(), path) + require.NoError(t, err) + + defer func() { require.NoError(t, export.Close()) }() + + require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "during"}, 0)) + + assert.Equal(t, + []string{openedEventID}, exportedEventIDs(t, writeExport(t, export)), + ) + + var stored int64 + + require.NoError(t, openArchiveDBForRead(t, path). + Model(&delivery.ExportArchivedEvent{}).Count(&stored).Error) + assert.Equal(t, int64(2), stored) +} + +// TestArchiveExport_SurvivesRename proves that renaming the archive +// while an export of it is open, as renaming its webhook or target +// does, leaves the export reading the same file. +func TestArchiveExport_SurvivesRename(t *testing.T) { + t.Parallel() + + path := filepath.Join(t.TempDir(), "archive-old.db") + w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0) + + require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: openedEventID}, 0)) + + export, err := delivery.OpenArchiveExport(t.Context(), path) + require.NoError(t, err) + + defer func() { require.NoError(t, export.Close()) }() + + require.NoError(t, w.Rename("archive-new.db")) + require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "after"}, 0)) + require.NoFileExists(t, path) + + assert.Equal(t, + []string{openedEventID}, exportedEventIDs(t, writeExport(t, export)), + ) +} + +// heapPeak is an io.Writer that discards what it is given and records +// the largest heap it saw at a write. +type heapPeak struct { + max uint64 +} + +func (p *heapPeak) Write(b []byte) (int, error) { + var m runtime.MemStats + + runtime.ReadMemStats(&m) + p.max = max(p.max, m.HeapAlloc) + + return len(b), nil +} + +// TestArchiveExport_Streams proves an export holds neither the archive +// nor its JSON in memory whole: exporting a 16 MiB archive grows the +// heap by less than half of that. An export that read every row before +// writing, or built the JSON before writing it, would hold all 16 MiB +// at a write. +// +//nolint:paralleltest // It measures the heap, which tests share. +func TestArchiveExport_Streams(t *testing.T) { + const ( + rows = 64 + bodySize = 256 << 10 + limit = rows * bodySize / 2 + ) + + path := filepath.Join(t.TempDir(), "archive.db") + w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0) + body := strings.Repeat("x", bodySize) + + for range rows { + require.NoError(t, w.Write(delivery.ExportArchivedEvent{Body: body}, 0)) + } + + export, err := delivery.OpenArchiveExport(t.Context(), path) + require.NoError(t, err) + + defer func() { require.NoError(t, export.Close()) }() + + // A low GC target collects garbage soon after it is made, so the + // heap at each write is close to what the export is holding. + defer debug.SetGCPercent(debug.SetGCPercent(10)) + + runtime.GC() + + var start runtime.MemStats + + runtime.ReadMemStats(&start) + + peak := &heapPeak{} + + require.NoError(t, writeExportTo(t, export, peak)) + assert.Less(t, peak.max, start.HeapAlloc+limit, + "heap at the start %d, at its peak %d", start.HeapAlloc, peak.max, + ) +} + +// TestArchiveExportFileName proves the download is named for the +// webhook and the target, with the names made safe as for the archive +// file, and the export time in UTC. +func TestArchiveExportFileName(t *testing.T) { + t.Parallel() + + cest := time.FixedZone("CEST", int((2 * time.Hour).Seconds())) + + assert.Equal(t, + "archive-orders-eu-long-term-archive-20261002T120304Z.json.gz", + delivery.ArchiveExportFileName( + exportWebhookName, exportTargetName, + time.Date(2026, 10, 2, 14, 3, 4, 0, cest), + ), + ) +} diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index f47ab1b..22c50c1 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -97,7 +97,8 @@ type Handlers struct { // names through the archive rename, the save and any move back. // Interleaved, one could rename an archive between another's // rename and save, leaving the file named for one edit and the - // stored names from the other. + // stored names from the other. An archive download holds it while + // it reads the stored names and opens the file they give. renameMu sync.Mutex // dummyVerifications counts the equivalent-cost verifications diff --git a/internal/handlers/target_download.go b/internal/handlers/target_download.go new file mode 100644 index 0000000..5d8e18c --- /dev/null +++ b/internal/handlers/target_download.go @@ -0,0 +1,82 @@ +package handlers + +import ( + "net/http" + "time" + + "sneak.berlin/go/webhooker/internal/database" + "sneak.berlin/go/webhooker/internal/delivery" +) + +// HandleTargetDownload serves a database target's archive as one +// gzipped JSON file, named for the webhook, the target and the time; +// see delivery.ArchiveExport.WriteGzipJSON for what it holds. Other +// target types have no archive and are a 404. +func (h *Handlers) HandleTargetDownload() http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + webhook, target, export, ok := h.openTargetArchive(w, r) + if !ok { + return + } + + defer func() { _ = export.Close() }() + + now := time.Now() + + w.Header().Set("Content-Type", "application/gzip") + w.Header().Set( + "Content-Disposition", + `attachment; filename="`+delivery.ArchiveExportFileName( + webhook.Name, target.Name, now, + )+`"`, + ) + + err := export.WriteGzipJSON(r.Context(), w, &webhook, target, now) + if err != nil { + // The 200 has gone out. The client is left with a file + // that does not decompress; the log is the record. + h.log.Error( + "failed to export archive", + "target_id", target.ID, + "error", err, + ) + } + } +} + +// openTargetArchive opens the archive of the request's database target +// for export. It reports false once it has written the response. +// +// 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 +// the one those names give. It lets go before the export is streamed: +// once the file is open, a rename does not affect the export. +func (h *Handlers) openTargetArchive( + w http.ResponseWriter, + r *http.Request, +) (database.Webhook, *database.Target, *delivery.ArchiveExport, bool) { + h.renameMu.Lock() + defer h.renameMu.Unlock() + + webhook, target, ok := h.ownedTarget(w, r) + if !ok { + return database.Webhook{}, nil, nil, false + } + + if target.Type != database.TargetTypeDatabase { + h.renderError(w, r, http.StatusNotFound) + + return database.Webhook{}, nil, nil, false + } + + export, err := delivery.OpenArchiveExport( + r.Context(), delivery.ArchivePath(h.dbMgr, &webhook, target), + ) + if err != nil { + h.serverError(w, r, "failed to open archive for export", err) + + return database.Webhook{}, nil, nil, false + } + + return webhook, target, export, true +} diff --git a/internal/handlers/target_download_test.go b/internal/handlers/target_download_test.go new file mode 100644 index 0000000..6286f43 --- /dev/null +++ b/internal/handlers/target_download_test.go @@ -0,0 +1,185 @@ +package handlers_test + +import ( + "compress/gzip" + "encoding/json" + "net/http" + "net/http/httptest" + "net/url" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "sneak.berlin/go/webhooker/internal/database" +) + +// downloadPath is the archive download route of a target. +func downloadPath(webhookID, targetID string) string { + return "/hook/" + webhookID + "/targets/" + targetID + "/download" +} + +// renameTarget submits the edit form renaming a target to Renamed. +func renameTarget( + env *sourceTestEnv, webhookID, targetID string, +) *httptest.ResponseRecorder { + form := url.Values{} + form.Set("name", "Renamed") + + return submitTargetEdit(env, webhookID, targetID, form) +} + +// TestHandleTargetDownload proves a database target's archive +// downloads as a gzipped JSON attachment named for the webhook, the +// target and the time, here with no archive file yet, so with no +// rows; and that a target of another type has no download. +func TestHandleTargetDownload(t *testing.T) { + t.Parallel() + + env := setupSourceTest(t) + wh := seedWebhookWithRetention(t, env.db, 7) + archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase) + logTarget := seedTarget(t, env.db, wh.ID, database.TargetTypeLog) + + w := serveTarget( + env, http.MethodGet, downloadPath(wh.ID, archive.ID), nil, + ) + require.Equal(t, http.StatusOK, w.Code, w.Body.String()) + assert.Equal(t, "application/gzip", w.Header().Get("Content-Type")) + assert.Regexp(t, + `^attachment; filename="archive-seeded-t-database-`+ + `\d{8}T\d{6}Z\.json\.gz"$`, + w.Header().Get("Content-Disposition"), + ) + + zr, err := gzip.NewReader(w.Body) + require.NoError(t, err) + + var got map[string]json.RawMessage + + require.NoError(t, json.NewDecoder(zr).Decode(&got)) + assert.JSONEq(t, + `{"id":"`+archive.ID+`","name":"t-database"}`, + string(got["target"]), + ) + assert.JSONEq(t, `[]`, string(got["archived_events"])) + + w = serveTarget( + env, http.MethodGet, downloadPath(wh.ID, logTarget.ID), nil, + ) + assert.Equal(t, http.StatusNotFound, w.Code) +} + +// TestHandleTargetDownload_WaitsForRename proves a download reads the +// target's names and opens its archive under the lock a rename holds: +// started while an edit is renaming the archive, it waits, and is +// named for the target's new name. +func TestHandleTargetDownload_WaitsForRename(t *testing.T) { + t.Parallel() + + env := setupSourceTest(t) + wh := seedWebhookWithRetention(t, env.db, 7) + archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase) + + renaming, release := env.archives.BlockNextRename() + edited := make(chan *httptest.ResponseRecorder, 1) + + go func() { + edited <- renameTarget(env, wh.ID, archive.ID) + }() + + <-renaming + + downloaded := make(chan *httptest.ResponseRecorder, 1) + + go func() { + downloaded <- serveTarget( + env, http.MethodGet, downloadPath(wh.ID, archive.ID), nil, + ) + }() + + select { + case <-downloaded: + release() + t.Fatal("the download did not wait for the rename") + case <-time.After(100 * time.Millisecond): + } + + release() + require.Equal(t, http.StatusSeeOther, (<-edited).Code) + + w := <-downloaded + require.Equal(t, http.StatusOK, w.Code) + assert.Contains(t, + w.Header().Get("Content-Disposition"), "archive-seeded-renamed-", + ) +} + +// stalledWriter is a response writer whose first write waits until +// resume is closed, closing writing when it starts to wait. +type stalledWriter struct { + *httptest.ResponseRecorder + + once sync.Once + writing chan struct{} + resume chan struct{} +} + +func (s *stalledWriter) Write(b []byte) (int, error) { + s.once.Do(func() { + close(s.writing) + <-s.resume + }) + + return s.ResponseRecorder.Write(b) +} + +// TestHandleTargetDownload_StreamsWithoutTheLock proves a download +// lets go of the rename lock once its archive is open: while the +// download is stalled writing, an edit can still rename the target. +func TestHandleTargetDownload_StreamsWithoutTheLock(t *testing.T) { + t.Parallel() + + env := setupSourceTest(t) + wh := seedWebhookWithRetention(t, env.db, 7) + archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase) + + req := httptest.NewRequestWithContext( + t.Context(), http.MethodGet, downloadPath(wh.ID, archive.ID), nil, + ) + for _, c := range env.cookies { + req.AddCookie(c) + } + + sw := &stalledWriter{ + ResponseRecorder: httptest.NewRecorder(), + writing: make(chan struct{}), + resume: make(chan struct{}), + } + downloaded := make(chan struct{}) + + go func() { + targetRouter(env).ServeHTTP(sw, req) + close(downloaded) + }() + + <-sw.writing + + edited := make(chan *httptest.ResponseRecorder, 1) + + go func() { + edited <- renameTarget(env, wh.ID, archive.ID) + }() + + select { + case w := <-edited: + assert.Equal(t, http.StatusSeeOther, w.Code) + case <-time.After(10 * time.Second): + t.Error("the rename waited for the download") + } + + close(sw.resume) + <-downloaded + assert.Equal(t, http.StatusOK, sw.Code) +} diff --git a/internal/handlers/target_edit_test.go b/internal/handlers/target_edit_test.go index be24c0e..942d765 100644 --- a/internal/handlers/target_edit_test.go +++ b/internal/handlers/target_edit_test.go @@ -37,14 +37,18 @@ const ( editAuthHeader = "Authorization: Bearer " + editBearerSecret ) -// targetRouter mounts the target create and edit routes on a chi -// router so the handlers see the URL parameters they read. +// targetRouter mounts the target create, edit and download routes on +// a chi router so the handlers see the URL parameters they read. func targetRouter(env *sourceTestEnv) *chi.Mux { router := chi.NewRouter() router.Post( "/hook/{sourceID}/targets", env.handlers.HandleTargetCreate(), ) + router.Get( + "/hook/{sourceID}/targets/{targetID}/download", + env.handlers.HandleTargetDownload(), + ) router.Get( "/hook/{sourceID}/targets/{targetID}/edit", env.handlers.HandleTargetEdit(), diff --git a/internal/server/archive_download_test.go b/internal/server/archive_download_test.go new file mode 100644 index 0000000..6e2239a --- /dev/null +++ b/internal/server/archive_download_test.go @@ -0,0 +1,69 @@ +package server_test + +import ( + "compress/gzip" + "encoding/json" + "net/http" + "regexp" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/gorm/clause" + "sneak.berlin/go/webhooker/internal/database" +) + +// TestHook_DownloadArchive follows the Download link the webhook page +// shows for a database target, and only for it, and gets the archive +// as a gzipped JSON file. Signed out, the link leads to the login page. +func TestHook_DownloadArchive(t *testing.T) { + t.Parallel() + + env := newTestEnv(t) + + userID, _ := env.seedUser(t, "archivist", "somepassword") + cookies := env.authCookies(t, userID, "archivist") + wh := env.seedWebhook(t, userID) + env.seedTarget(t, wh.ID) + + archive := &database.Target{ + WebhookID: wh.ID, + Name: "kept", + Type: database.TargetTypeDatabase, + Active: true, + } + require.NoError(t, + env.db.DB().Omit(clause.Associations).Create(archive).Error, + ) + + page := env.get("/hook/"+wh.ID, cookies) + require.Equal(t, http.StatusOK, page.Code) + + links := regexp.MustCompile( + `href="(/hook/[^/"]+/targets/[^/"]+/download)"`, + ).FindAllStringSubmatch(page.Body.String(), -1) + require.Len(t, links, 1, "only the database target has a Download") + + link := links[0][1] + assert.Equal(t, + "/hook/"+wh.ID+"/targets/"+archive.ID+"/download", link, + ) + + w := env.get(link, cookies) + require.Equal(t, http.StatusOK, w.Code) + assert.Equal(t, "application/gzip", w.Header().Get("Content-Type")) + + zr, err := gzip.NewReader(w.Body) + require.NoError(t, err) + + var got map[string]json.RawMessage + + require.NoError(t, json.NewDecoder(zr).Decode(&got)) + assert.JSONEq(t, + `{"id":"`+wh.ID+`","name":"routed"}`, string(got["webhook"]), + ) + + w = env.get(link, nil) + assert.Equal(t, http.StatusSeeOther, w.Code) + assert.Contains(t, w.Header().Get("Location"), "/pages/login") +} diff --git a/internal/server/routes.go b/internal/server/routes.go index 24c63d0..e8bab01 100644 --- a/internal/server/routes.go +++ b/internal/server/routes.go @@ -312,6 +312,10 @@ func (s *Server) setupSourceRoutes() { "/targets/{targetID}/edit", s.h.HandleTargetEditSubmit(), ) + r.Get( + "/targets/{targetID}/download", + s.h.HandleTargetDownload(), + ) r.Post( "/targets/{targetID}/delete", s.h.HandleTargetDelete(), diff --git a/templates/source_detail.html b/templates/source_detail.html index 2a89922..05ff02e 100644 --- a/templates/source_detail.html +++ b/templates/source_detail.html @@ -157,6 +157,9 @@ {{else}} Inactive {{end}} + {{if eq .Type "database"}} + Download + {{end}} Edit