check / check (push) Waiting to run
Each database target on the webhook page has a Download button that streams its archive as gzipped JSON, archive-WEBHOOKNAME-TARGETNAME-TIME.json.gz, with names made safe by delivery.ArchiveFileName's function. The export reads one consistent snapshot through one cursor in a read-only transaction, so archive writes carry on, and holds the rename lock only while it reads the stored names and opens the file. It extends its write deadline as it writes, so a large archive downloads for as long as the client reads; a failure after the response has started aborts the connection so the browser marks the download failed. The request limit is now the service's own middleware, which no longer writes a 504 over a response already started. Model: opus-5-5
413 lines
12 KiB
Go
413 lines
12 KiB
Go
package delivery_test
|
|
|
|
import (
|
|
"bufio"
|
|
"bytes"
|
|
"compress/gzip"
|
|
"crypto/rand"
|
|
"encoding/base64"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"path/filepath"
|
|
"runtime"
|
|
"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. 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, err := delivery.OpenArchiveExport(
|
|
t.Context(), path, archiveTestLogger(),
|
|
)
|
|
require.NoError(t, err)
|
|
|
|
defer func() { require.NoError(t, export.Close()) }()
|
|
|
|
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),
|
|
),
|
|
)
|
|
}
|