check / check (push) Successful in 4m45s
The receiver dropped the query string of every request it received, so a sender's URL parameters were silently lost. Each event now keeps it, as sent, in a new `raw_query` column of the per-webhook `events` table; a resubmitted copy carries its original's. The event log and the event's page show it in the shared request block, the event log leaving out one over 32 KiB with a link, as for headers. The archive and log targets carry it. HTTP targets gain "Pass the query string on to this target", off by default: on, deliveries, replays and resubmits append it to the target URL, joined with `&` to one already there. The access log still hides it. Model: opus-5-5
388 lines
10 KiB
Go
388 lines
10 KiB
Go
package delivery
|
|
|
|
import (
|
|
"compress/gzip"
|
|
"context"
|
|
"database/sql"
|
|
"encoding/base64"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"io/fs"
|
|
"log/slog"
|
|
"os"
|
|
"path/filepath"
|
|
"sync"
|
|
"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 listed for download. It
|
|
// opens one of the target's files at a time, only when its rows are
|
|
// about to be written out, and closes it before it opens the next, so
|
|
// an export holds at most one file open however many the target has.
|
|
//
|
|
// Each file is read on its own connection inside one read-only
|
|
// transaction, so its rows are written out as the file stood when it
|
|
// was opened. Archives are in WAL mode, where a reader works from a
|
|
// snapshot and never blocks a writer: archive writes go on while a file
|
|
// is open, and the export does not see them. SQLite cannot checkpoint
|
|
// a -wal past an open snapshot, so the open file's -wal grows until the
|
|
// export has written that file out.
|
|
type ArchiveExport struct {
|
|
// periods are the periods of the target's files when the export
|
|
// was listed, "" for the file without one, in the order
|
|
// archiveFiles lists them.
|
|
periods []string
|
|
|
|
// lock is held while currentPath is called and a file is opened,
|
|
// so that a rename, which holds it too, cannot move the file in
|
|
// between.
|
|
lock sync.Locker
|
|
|
|
// currentPath returns the path ArchivePath gives the target under
|
|
// the names stored for it now, which a rename may have changed since
|
|
// the export was listed.
|
|
currentPath func() (string, error)
|
|
|
|
log *slog.Logger
|
|
}
|
|
|
|
// exportFile is one archive file opened for an export.
|
|
type exportFile struct {
|
|
db *sql.DB
|
|
tx *gorm.DB
|
|
|
|
// period is the period in the file's name, "" for none.
|
|
period string
|
|
|
|
// empty is true for 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"`
|
|
}
|
|
|
|
// NewArchiveExport lists a database target's archive files for export,
|
|
// given the path ArchivePath gives it (see archiveFiles). It opens none
|
|
// of them. Its caller holds lock, which every rename of the target's
|
|
// files runs under, from reading the names path is made of until it
|
|
// returns, so the files it lists are the ones those names give.
|
|
//
|
|
// WriteGzipJSON, called without lock held, finds each file again by its
|
|
// period under the path currentPath gives, holding lock while it does
|
|
// and while it opens the file, so a rename during the export loses no
|
|
// file. A file that is gone by then, emptied by the sweep or moved
|
|
// away, is skipped. The export never creates a file: with no files, it
|
|
// has no rows.
|
|
func NewArchiveExport(
|
|
path string,
|
|
lock sync.Locker,
|
|
currentPath func() (string, error),
|
|
log *slog.Logger,
|
|
) (*ArchiveExport, error) {
|
|
files, err := archiveFiles(path)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
x := &ArchiveExport{lock: lock, currentPath: currentPath, log: log}
|
|
|
|
for _, file := range files {
|
|
x.periods = append(x.periods, file.period)
|
|
}
|
|
|
|
return x, nil
|
|
}
|
|
|
|
// openExportFile opens one archive file for an export and takes its
|
|
// snapshot. The transaction lasts as long as ctx does.
|
|
func openExportFile(
|
|
ctx context.Context, file archiveFile, log *slog.Logger,
|
|
) (*exportFile, error) {
|
|
db, err := database.OpenSQLite(file.path, archiveModeExisting)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("opening archive %s: %w", file.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", file.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", file.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", file.path, err)
|
|
}
|
|
|
|
return &exportFile{
|
|
db: db, tx: tx, period: file.period, 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,
|
|
// the files in the order archiveFiles lists them. A row from a file
|
|
// named for a period has "period" beside its columns. 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, and each file is
|
|
// closed once its rows are written, before the next is opened. When it
|
|
// returns, no file is open. 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()
|
|
}
|
|
|
|
// openFile finds the target's archive file for period under the path
|
|
// currentPath gives now, and opens it for the export, holding x.lock
|
|
// for both. For a file that is gone, the error wraps fs.ErrNotExist.
|
|
func (x *ArchiveExport) openFile(
|
|
ctx context.Context, period string,
|
|
) (*exportFile, error) {
|
|
x.lock.Lock()
|
|
defer x.lock.Unlock()
|
|
|
|
path, err := x.currentPath()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("finding archive file: %w", err)
|
|
}
|
|
|
|
file := archiveFile{path: archivePeriodPath(path, period), period: period}
|
|
|
|
_, err = os.Stat(file.path)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return openExportFile(ctx, file, x.log)
|
|
}
|
|
|
|
// close ends the file's transaction and closes its connection.
|
|
func (f *exportFile) close() error {
|
|
_ = f.tx.Rollback()
|
|
|
|
return f.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 the archived rows of each file to w, one per line,
|
|
// separated by commas, opening each file in turn and closing it once
|
|
// its rows are written.
|
|
func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error {
|
|
sep := "\n"
|
|
|
|
for _, period := range x.periods {
|
|
f, err := x.openFile(ctx, period)
|
|
if errors.Is(err, fs.ErrNotExist) {
|
|
continue
|
|
}
|
|
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
sep, err = f.writeRows(ctx, w, sep)
|
|
|
|
err = errors.Join(err, f.close())
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// writeRows writes the file's archived rows to w, oldest first, the
|
|
// first after sep and each other after ",\n". It returns what goes
|
|
// before the next row: sep again when the file had no rows.
|
|
func (f *exportFile) writeRows(
|
|
ctx context.Context, w io.Writer, sep string,
|
|
) (string, error) {
|
|
if f.empty {
|
|
return sep, nil
|
|
}
|
|
|
|
rows, err := f.tx.WithContext(ctx).
|
|
Model(&archivedEvent{}).Order("id").Rows()
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
|
|
defer func() { _ = rows.Close() }()
|
|
|
|
for ; rows.Next(); sep = ",\n" {
|
|
var ev archivedEvent
|
|
|
|
err = f.tx.ScanRows(rows, &ev)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
|
|
_, err = io.WriteString(w, sep)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
|
|
err = writeRow(w, &ev, f.period)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
}
|
|
|
|
return sep, 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, with
|
|
// the period of its file beside them unless that is "".
|
|
func writeRow(w io.Writer, ev *archivedEvent, period string) error {
|
|
row := map[string]any{
|
|
"id": ev.ID,
|
|
"event_id": ev.EventID,
|
|
"webhook_id": ev.WebhookID,
|
|
"entrypoint_id": ev.EntrypointID,
|
|
"method": ev.Method,
|
|
"raw_query": ev.RawQuery,
|
|
"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"
|
|
}
|
|
|
|
if period != "" {
|
|
row["period"] = period
|
|
}
|
|
|
|
line, err := json.Marshal(row)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
_, err = w.Write(line)
|
|
|
|
return err
|
|
}
|