package delivery_test import ( "bytes" "compress/gzip" "crypto/rand" "encoding/base64" "encoding/json" "fmt" "io" "os" "path/filepath" "runtime" "runtime/debug" "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, archiveTestLogger(), ) 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, archiveTestLogger(), ) 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, archiveTestLogger(), ) 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 output in memory whole: exporting a 16 MiB archive grows the // heap by less than half of that. The bodies are random bytes in // base64, which gzip shrinks by only a quarter, so an export that read // every row before writing, or built the JSON or the gzipped file // before writing it, would hold at least 12 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) // Base64 makes four characters of every three bytes. random := make([]byte, bodySize/4*3) for range rows { _, _ = rand.Read(random) require.NoError(t, w.Write(delivery.ExportArchivedEvent{ Body: base64.StdEncoding.EncodeToString(random), }, 0)) } export, err := delivery.OpenArchiveExport( t.Context(), path, archiveTestLogger(), ) 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), ), ) }