package delivery_test import ( "bufio" "bytes" "compress/gzip" "crypto/rand" "encoding/base64" "encoding/json" "fmt" "io" "os" "path/filepath" "runtime" "strings" "sync" "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" ) // binaryBody is a body that is not valid UTF-8. const binaryBody = "\xff\xfe\x00\x01binary\x80" // 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), ) } // listExport lists the archive at path for export, as the archive of a // target whose names do not change. func listExport(t *testing.T, path string) *delivery.ArchiveExport { t.Helper() return newExport(t, path, &sync.Mutex{}, func() (string, error) { return path, nil }) } // newExport lists the archive at path for export, to find each file // again under the path currentPath gives, holding lock while it does. // Nothing else takes lock while it lists, so it does not hold lock. func newExport( t *testing.T, path string, lock sync.Locker, currentPath func() (string, error), ) *delivery.ArchiveExport { t.Helper() export, err := delivery.NewArchiveExport( path, lock, currentPath, archiveTestLogger(), ) require.NoError(t, err) return export } // 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() return writeExport(t, listExport(t, path)) } // 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) } } // unlockHook is a sync.Locker that runs fn each time it is unlocked. An // export unlocks its lock right after it opens a file. type unlockHook struct { sync.Mutex fn func() } func (u *unlockHook) Unlock() { u.Mutex.Unlock() u.fn() } // TestArchiveExport_ReadsOneSnapshot proves an export writes a file // out as it was when the export opened it, and holds up no archive // write: a row written after the export was listed but before the file // was opened is in the export, and one written while the file is open // is stored, and is not. 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: "listed"}, 0)) opened := &unlockHook{fn: func() { require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "during"}, 0)) }} export := newExport(t, path, opened, func() (string, error) { return path, nil }) require.NoError(t, w.Write(delivery.ExportArchivedEvent{EventID: "before-open"}, 0)) assert.Equal(t, []string{"listed", "before-open"}, exportedEventIDs(t, writeExport(t, export)), ) var stored int64 require.NoError(t, openArchiveDBForRead(t, path). Model(&delivery.ExportArchivedEvent{}).Count(&stored).Error) assert.Equal(t, int64(3), stored) } // TestArchiveExport_FindsFilesAfterRename proves that renaming the // archive after an export has listed it, as renaming its webhook or // target does, loses no file: the export finds each file again by its // period under the new name. A file moved away by then is skipped. func TestArchiveExport_FindsFilesAfterRename(t *testing.T) { t.Parallel() dir := t.TempDir() path := filepath.Join(dir, "archive-old.db") w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0) for _, period := range []string{"", dayPeriod, hourPeriod} { require.NoError(t, w.WritePeriod( delivery.ExportArchivedEvent{EventID: "in-" + period}, 0, period, )) } current := path export := newExport(t, path, &sync.Mutex{}, func() (string, error) { return current, nil }) require.NoError(t, w.Rename("archive-new.db")) current = filepath.Join(dir, "archive-new.db") removeArchiveFiles(t, periodPath(current, dayPeriod)) assert.Equal(t, []string{"in-", "in-" + hourPeriod}, exportedEventIDs(t, writeExport(t, export)), ) } // 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 was listed 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 := listExport(t, path) 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") } // openFilesPeak is an io.Writer that discards what it is given and // records the most archive files in dir the process had open at any // write, as /proc/self/fd lists the files a process has open. type openFilesPeak struct { dir string max int } func (p *openFilesPeak) Write(b []byte) (int, error) { fds, err := os.ReadDir("/proc/self/fd") if err != nil { return 0, err } open := map[string]bool{} for _, fd := range fds { file, err := os.Readlink(filepath.Join("/proc/self/fd", fd.Name())) if err == nil && filepath.Dir(file) == p.dir && strings.HasSuffix(file, ".db") { open[file] = true } } p.max = max(p.max, len(open)) return len(b), nil } // TestArchiveExport_OneFileOpenAtATime exports a target with a file for // each of 24 hours and proves the export never had more than one of // them open, and had one open while it wrote. Each file holds a row of // 48 KiB of random base64, which gzip shrinks little, so the export // writes output while it reads each file. func TestArchiveExport_OneFileOpenAtATime(t *testing.T) { t.Parallel() if runtime.GOOS != "linux" { t.Skip("only Linux lists a process's open files in /proc/self/fd") } // Readlink gives each open file's path with no symbolic link in it. dir, err := filepath.EvalSymlinks(t.TempDir()) require.NoError(t, err) path := filepath.Join(dir, "archive-wh.db") w := delivery.NewExportArchiveWriter(path, archiveTestLogger(), 0) random := make([]byte, 36<<10) for hour := range 24 { _, _ = rand.Read(random) require.NoError(t, w.WritePeriod(delivery.ExportArchivedEvent{ Body: base64.StdEncoding.EncodeToString(random), }, 0, fmt.Sprintf("2026-10-01-%02d", hour))) } // The writer's own handle on the last file is not the export's. w.Evict() peak := &openFilesPeak{dir: dir} require.NoError(t, writeExportTo(t, listExport(t, path), peak)) assert.Equal(t, 1, peak.max) } // heapPeak is an io.Writer that discards what it is given and records // the largest heap it saw at a write. It collects garbage before each // reading, so the heap it reads is what is still held. type heapPeak struct { max uint64 } func (p *heapPeak) Write(b []byte) (int, error) { var m runtime.MemStats runtime.GC() runtime.ReadMemStats(&m) p.max = max(p.max, m.HeapAlloc) return len(b), nil } // exportHeapGrowth exports an archive of rows random bodies, each // bodySize bytes of base64, and returns how far the heap rose above // where it stood when the export began, at its highest. func exportHeapGrowth(t *testing.T, rows, bodySize int) uint64 { t.Helper() 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 := listExport(t, path) runtime.GC() var start runtime.MemStats runtime.ReadMemStats(&start) // Through a buffer, the heap is read once per 8 KiB of output // rather than at each of gzip's small writes, which takes far // longer. peak := &heapPeak{max: start.HeapAlloc} buffered := bufio.NewWriterSize(peak, 8<<10) require.NoError(t, writeExportTo(t, export, buffered)) require.NoError(t, buffered.Flush()) return peak.max - start.HeapAlloc } // TestArchiveExport_Streams proves an export holds neither the archive // nor its output in memory whole: exporting 384 KiB more of archive // raises the heap's peak by less than half of that. The export's own // memory, mostly gzip's compressor, is the same for both archives, so // it cancels out. 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 raise the peak by at least three quarters of the difference. // // The smaller archive has two rows so that its export, too, writes // out more than the 8 KiB buffer in exportHeapGrowth before it ends: // the heap must be read while the export's own memory is held. // //nolint:paralleltest // It measures the heap, which tests share. func TestArchiveExport_Streams(t *testing.T) { const ( bodySize = 16 << 10 smallRows = 2 largeRows = smallRows + 24 limit = (largeRows - smallRows) * bodySize / 2 ) small := exportHeapGrowth(t, smallRows, bodySize) large := exportHeapGrowth(t, largeRows, bodySize) assert.Less(t, large, small+limit, "the heap rose by %d for %d rows and by %d for %d rows", small, smallRows, large, largeRows, ) } // 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), ), ) }