check / check (push) Waiting to run
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
271 lines
7.3 KiB
Go
271 lines
7.3 KiB
Go
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
|
|
}
|