Files
webhooker/internal/delivery/target_database_export.go
T
clawbot 8b617efa63
check / check (push) Successful in 3m27s
Rotate a database target's archive monthly, daily or hourly (closes #379)
A database target's rotation setting (none, monthly, daily or hourly) puts the UTC period of each event's receive time in its archive file name, so each file holds exactly its period's events. It is on the new webhook page and both target forms, and shown in the target list. Renames move every one of a target's files and move them all back if one fails. The sweep prunes one file at a time under the target's lock and deletes a rotated file it leaves empty. Download opens one file at a time, oldest period first, finding each again under the target's current name. The target list names the current file and totals all of them.

Model: opus-5-5
2026-10-03 04:18:37 +02:00

387 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,
"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
}