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
276 lines
7.1 KiB
Go
276 lines
7.1 KiB
Go
package delivery
|
|
|
|
import (
|
|
"compress/gzip"
|
|
"context"
|
|
"database/sql"
|
|
"encoding/base64"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"path/filepath"
|
|
"time"
|
|
"unicode/utf8"
|
|
|
|
"gorm.io/driver/sqlite"
|
|
"gorm.io/gorm"
|
|
"sneak.berlin/go/webhooker/internal/database"
|
|
"sneak.berlin/go/webhooker/internal/gormlog"
|
|
)
|
|
|
|
// 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'"
|
|
|
|
// 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 *gorm.DB
|
|
|
|
// 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.
|
|
//
|
|
// The transaction lasts as long as ctx does, so ctx must last for the
|
|
// whole export.
|
|
func OpenArchiveExport(
|
|
ctx context.Context, path string, log *slog.Logger,
|
|
) (*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)
|
|
}
|
|
|
|
gdb, err := gorm.Open(
|
|
sqlite.Dialector{Conn: db}, &gorm.Config{
|
|
// Never leave this at GORM's default. See
|
|
// internal/gormlog.
|
|
Logger: gormlog.New(log),
|
|
},
|
|
)
|
|
if err != nil {
|
|
_ = db.Close()
|
|
|
|
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 := gdb.WithContext(ctx).Begin(&sql.TxOptions{ReadOnly: true})
|
|
if tx.Error != nil {
|
|
_ = db.Close()
|
|
|
|
return nil, fmt.Errorf("reading archive %s: %w", path, tx.Error)
|
|
}
|
|
|
|
// The transaction's first read is what takes the snapshot.
|
|
var tables int
|
|
|
|
err = tx.Raw(archiveTableQuery).Row().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
|
|
}
|
|
|
|
rows, err := x.tx.WithContext(ctx).
|
|
Model(&archivedEvent{}).Order("id").Rows()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
defer func() { _ = rows.Close() }()
|
|
|
|
for sep := "\n"; rows.Next(); sep = ",\n" {
|
|
var ev archivedEvent
|
|
|
|
err = x.tx.ScanRows(rows, &ev)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
_, err = io.WriteString(w, sep)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
err = writeRow(w, &ev)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return rows.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
|
|
}
|