check / check (push) Waiting to run
A database target's rotation (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, the add target form and the target edit form, and shown in the target list. Renames move every one of a target's files and move them back if one fails. The sweep prunes every file, one at a time under the target's lock, and deletes a rotated file it leaves empty. Download lists the files, then opens one at a time, oldest first, finding each again under the target's current names, and gives each row its period. The target list names the current file and totals the size of all of them. Model: opus-5-5
689 lines
18 KiB
Go
689 lines
18 KiB
Go
package delivery
|
|
|
|
import (
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io/fs"
|
|
"log/slog"
|
|
"os"
|
|
"path/filepath"
|
|
"sync"
|
|
"time"
|
|
|
|
"gorm.io/driver/sqlite"
|
|
"gorm.io/gorm"
|
|
"sneak.berlin/go/webhooker/internal/database"
|
|
"sneak.berlin/go/webhooker/internal/gormlog"
|
|
)
|
|
|
|
// archiveExpiryNever is the expiry sentinel (and default) that
|
|
// disables pruning so archived rows are kept forever.
|
|
const archiveExpiryNever = "never"
|
|
|
|
// archiveReopenDebounce bounds how often an archive file is
|
|
// closed and reopened. After each write the handle is closed
|
|
// and reopened so an operator can move the file away for
|
|
// offline archiving, but never more than once per this window.
|
|
const archiveReopenDebounce = time.Second
|
|
|
|
const (
|
|
// archiveModeCreate is the SQLite URI mode used by the write
|
|
// path: open the archive file, creating it if missing, so a
|
|
// first write (or a write after the operator moved the file
|
|
// away) recreates it.
|
|
archiveModeCreate = database.SQLiteModeCreate
|
|
|
|
// archiveModeExisting is the SQLite URI mode used by the idle
|
|
// sweep: open read-write but never create. A sweep must never
|
|
// conjure an empty archive file for a webhook that has a
|
|
// database target but has never received an event.
|
|
archiveModeExisting = database.SQLiteModeExisting
|
|
)
|
|
|
|
var (
|
|
// errArchiveMissingWebhookID is returned when an event to
|
|
// archive has no webhook id to record in its archive row.
|
|
errArchiveMissingWebhookID = errors.New(
|
|
"cannot archive event without a webhook id",
|
|
)
|
|
|
|
// errArchiveNoDataDir is returned when the database target
|
|
// has no webhook database manager and so cannot locate the
|
|
// data directory for archive files.
|
|
errArchiveNoDataDir = errors.New(
|
|
"database target has no data directory",
|
|
)
|
|
|
|
// errArchiveExpiryNotPositive is returned when a
|
|
// user-supplied archive expiry parses as a duration but is
|
|
// zero or negative; "never" is the way to disable pruning.
|
|
errArchiveExpiryNotPositive = errors.New(
|
|
"expiry must be a positive duration or \"never\"",
|
|
)
|
|
|
|
// errArchiveWriterEvicted is returned when a writer that has
|
|
// been evicted (its target or its webhook was deleted) is used
|
|
// again. An evicted writer is no longer in the registry, so
|
|
// reopening its file would leak a handle nothing owns.
|
|
errArchiveWriterEvicted = errors.New(
|
|
"archive writer has been evicted",
|
|
)
|
|
|
|
// ErrArchiveNameTaken is returned when an archive cannot be
|
|
// renamed because a file already has the new name. That file may
|
|
// be an archive with rows of its own, so it is never replaced.
|
|
ErrArchiveNameTaken = errors.New(
|
|
"a file already has the archive's new name",
|
|
)
|
|
)
|
|
|
|
// databaseTargetConfig is the optional per-target JSON config
|
|
// for a database (archive) target.
|
|
type databaseTargetConfig struct {
|
|
// Expiry is a Go duration (e.g. "720h") after which
|
|
// archived rows are pruned, or "never" (the default) to
|
|
// keep them forever.
|
|
Expiry string `json:"expiry"`
|
|
|
|
// Rotation is none (the default), monthly, daily or hourly: see
|
|
// archivePeriod.
|
|
Rotation string `json:"rotation"`
|
|
}
|
|
|
|
// archivedEvent is one fully captured webhook event stored in a
|
|
// database target's archive for long-term retention. It is a
|
|
// self-contained copy — independent of the per-webhook event
|
|
// database, which may prune events under its own retention.
|
|
type archivedEvent struct {
|
|
ID uint `gorm:"primaryKey;autoIncrement"`
|
|
EventID string `gorm:"index"`
|
|
WebhookID string
|
|
EntrypointID string
|
|
Method string
|
|
Headers string
|
|
Body string
|
|
ContentType string
|
|
|
|
// ArchivedAt is when the row was archived and is the age
|
|
// basis for expiry pruning.
|
|
ArchivedAt time.Time `gorm:"index"`
|
|
}
|
|
|
|
// parseArchiveExpiry reads the optional expiry from a database
|
|
// target's config JSON. An empty config, an empty expiry, or
|
|
// the literal "never" all mean keep forever, returned as a zero
|
|
// duration. Any other value must parse as a positive Go
|
|
// duration; a set-but-invalid value (unparseable, zero, or
|
|
// negative) is an error rather than a silent default, matching
|
|
// ValidateArchiveExpiry at target creation.
|
|
func parseArchiveExpiry(
|
|
configJSON string,
|
|
) (time.Duration, error) {
|
|
if configJSON == "" {
|
|
return 0, nil
|
|
}
|
|
|
|
var cfg databaseTargetConfig
|
|
|
|
err := json.Unmarshal([]byte(configJSON), &cfg)
|
|
if err != nil {
|
|
return 0, fmt.Errorf(
|
|
"parsing database target config: %w", err,
|
|
)
|
|
}
|
|
|
|
if cfg.Expiry == "" || cfg.Expiry == archiveExpiryNever {
|
|
return 0, nil
|
|
}
|
|
|
|
dur, err := time.ParseDuration(cfg.Expiry)
|
|
if err != nil {
|
|
return 0, fmt.Errorf(
|
|
"parsing archive expiry %q: %w", cfg.Expiry, err,
|
|
)
|
|
}
|
|
|
|
if dur <= 0 {
|
|
return 0, fmt.Errorf(
|
|
"%w: %q", errArchiveExpiryNotPositive, cfg.Expiry,
|
|
)
|
|
}
|
|
|
|
return dur, nil
|
|
}
|
|
|
|
// ValidateArchiveExpiry checks a user-supplied archive expiry
|
|
// for a database target at configuration time. Valid values are
|
|
// empty, "never" (both meaning keep forever), or a positive Go
|
|
// duration such as "720h". Anything else is an error, so a bad
|
|
// expiry is rejected when the target is created rather than
|
|
// failing every subsequent delivery.
|
|
func ValidateArchiveExpiry(expiry string) error {
|
|
if expiry == "" || expiry == archiveExpiryNever {
|
|
return nil
|
|
}
|
|
|
|
dur, err := time.ParseDuration(expiry)
|
|
if err != nil {
|
|
return fmt.Errorf(
|
|
"expiry must be %q or a Go duration "+
|
|
"such as \"720h\": %w",
|
|
archiveExpiryNever, err,
|
|
)
|
|
}
|
|
|
|
if dur <= 0 {
|
|
return fmt.Errorf(
|
|
"%w: %q", errArchiveExpiryNotPositive, expiry,
|
|
)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// archiveWriter owns one database target's archive SQLite files.
|
|
// It serialises writes, and after each write closes and reopens
|
|
// the file (debounced to at most once per debounce window) so
|
|
// an operator can move the file away for offline archiving. The
|
|
// next write recreates a moved or removed file, because the
|
|
// file is opened create-if-missing and its schema is migrated
|
|
// on every open.
|
|
type archiveWriter struct {
|
|
mu sync.Mutex
|
|
|
|
// path is the target's archive file as ArchivePath names it. A
|
|
// target that rotates writes to the files archivePeriodPath names
|
|
// for path and a period instead.
|
|
path string
|
|
|
|
// current is the file db is open on.
|
|
current string
|
|
|
|
log *slog.Logger
|
|
debounce time.Duration
|
|
db *gorm.DB
|
|
lastReopen time.Time
|
|
reopens int
|
|
|
|
// now is the clock the reopen debounce is measured on. It is
|
|
// time.Now outside tests.
|
|
now func() time.Time
|
|
|
|
// evicted marks a writer that has been removed from the
|
|
// registry. Its handle is closed and it must never open the
|
|
// file again: nothing holds it any more, so a reopen would
|
|
// leak the handle for the process lifetime.
|
|
evicted bool
|
|
|
|
// webhookID is the webhook the archive's target belongs to,
|
|
// so deleting the webhook can find its writers. It is set
|
|
// when the writer is created and never changes.
|
|
webhookID string
|
|
|
|
// sweepOwned marks a registry entry that the idle sweep
|
|
// created because no writer was cached for the target. The
|
|
// sweep removes such an entry again when it is done, so a
|
|
// sweep can never leave — or resurrect — a registry entry
|
|
// for a target that has been deleted. A delivery that adopts
|
|
// the writer clears the flag, handing the entry to the
|
|
// registry proper.
|
|
//
|
|
// Unlike every other field here it is guarded by
|
|
// databaseTarget.mu, not by this writer's mu: it describes the
|
|
// registry entry rather than the file.
|
|
sweepOwned bool
|
|
}
|
|
|
|
// newArchiveWriter builds an archiveWriter for a file path with
|
|
// the default reopen debounce.
|
|
func newArchiveWriter(
|
|
path string, log *slog.Logger,
|
|
) *archiveWriter {
|
|
return &archiveWriter{
|
|
path: path,
|
|
log: log,
|
|
debounce: archiveReopenDebounce,
|
|
now: time.Now,
|
|
}
|
|
}
|
|
|
|
// write appends the event as a row to the archive file for period
|
|
// (see archivePeriodPath), then applies the debounced close/reopen.
|
|
// When period names a different file from the one open, the open one
|
|
// is closed first. It recreates the archive file if it was moved or
|
|
// removed since the last open. A positive expiry prunes rows older
|
|
// than it on each (re)open.
|
|
func (w *archiveWriter) write(
|
|
row archivedEvent, expiry time.Duration, period string,
|
|
) error {
|
|
w.mu.Lock()
|
|
defer w.mu.Unlock()
|
|
|
|
if w.evicted {
|
|
return fmt.Errorf(
|
|
"%w: %s", errArchiveWriterEvicted, w.path,
|
|
)
|
|
}
|
|
|
|
file := archivePeriodPath(w.path, period)
|
|
|
|
if w.db == nil || w.current != file || !fileExists(file) {
|
|
err := w.reopen(file, expiry)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
row.ArchivedAt = time.Now()
|
|
|
|
err := w.db.Create(&row).Error
|
|
if err != nil {
|
|
return fmt.Errorf(
|
|
"archiving event to %s: %w", file, err,
|
|
)
|
|
}
|
|
|
|
if w.now().Sub(w.lastReopen) >= w.debounce {
|
|
return w.reopen(file, expiry)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// open opens (creating if missing) an archive file, migrates
|
|
// its schema, records the reopen time, and prunes expired rows
|
|
// when expiry is positive.
|
|
func (w *archiveWriter) open(file string, expiry time.Duration) error {
|
|
return w.openMode(file, archiveModeCreate, expiry)
|
|
}
|
|
|
|
// openMode opens an archive file with the given SQLite URI
|
|
// mode, migrates its schema, records the reopen time, and
|
|
// prunes expired rows when expiry is positive. The write path
|
|
// passes archiveModeCreate so a missing file is recreated; the
|
|
// idle sweep passes archiveModeExisting so a missing file is an
|
|
// error rather than a newly conjured empty archive.
|
|
func (w *archiveWriter) openMode(
|
|
file, mode string, expiry time.Duration,
|
|
) error {
|
|
// Opened through database.OpenSQLite so an archive file carries
|
|
// the same WAL journaling, busy timeout, immediate-transaction
|
|
// locking, and pool bounds as every other database file. See
|
|
// internal/database/sqlite_open.go.
|
|
sqlDB, err := database.OpenSQLite(file, mode)
|
|
if err != nil {
|
|
return fmt.Errorf(
|
|
"opening archive database %s: %w", file, err,
|
|
)
|
|
}
|
|
|
|
gdb, err := gorm.Open(
|
|
sqlite.Dialector{Conn: sqlDB}, &gorm.Config{
|
|
// Never leave this at GORM's default. See
|
|
// internal/gormlog.
|
|
Logger: gormlog.New(w.log),
|
|
},
|
|
)
|
|
if err != nil {
|
|
_ = sqlDB.Close()
|
|
|
|
return fmt.Errorf(
|
|
"connecting to archive database %s: %w",
|
|
file, err,
|
|
)
|
|
}
|
|
|
|
err = gdb.AutoMigrate(&archivedEvent{})
|
|
if err != nil {
|
|
_ = sqlDB.Close()
|
|
|
|
return fmt.Errorf(
|
|
"migrating archive database %s: %w", file, err,
|
|
)
|
|
}
|
|
|
|
w.db = gdb
|
|
w.current = file
|
|
w.lastReopen = w.now()
|
|
w.reopens++
|
|
|
|
if expiry > 0 {
|
|
w.prune(expiry)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// reopen closes any open handle and opens file afresh. The
|
|
// fresh open recreates the file if it was moved away.
|
|
func (w *archiveWriter) reopen(file string, expiry time.Duration) error {
|
|
w.close()
|
|
|
|
return w.open(file, expiry)
|
|
}
|
|
|
|
// close closes the underlying handle, if any.
|
|
func (w *archiveWriter) close() {
|
|
if w.db == nil {
|
|
return
|
|
}
|
|
|
|
sqlDB, err := w.db.DB()
|
|
if err == nil {
|
|
_ = sqlDB.Close()
|
|
}
|
|
|
|
w.db = nil
|
|
}
|
|
|
|
// sweepExpired prunes the target's archive files, which may have
|
|
// gone idle, with no write to trigger the usual on-reopen prune. It
|
|
// lists the files under the writer's own mutex, then takes the mutex
|
|
// again for one file at a time, so a write waits for at most one
|
|
// file's prune, and each prune is ordered against concurrent writes
|
|
// rather than reaching around them to the file.
|
|
//
|
|
// It never creates an archive file: it prunes only the files
|
|
// archiveFiles lists, skips one that is gone by the time it is
|
|
// reached (moved away, or renamed since the listing), and opens each
|
|
// with archiveModeExisting so SQLite itself refuses to create one if
|
|
// the file disappears between the check and the open. A file named
|
|
// for a period that the prune leaves empty is deleted.
|
|
//
|
|
// The archive is left CLOSED afterwards. An idle archive holding
|
|
// no handle is what keeps the operator's move-the-file-away
|
|
// workflow working; the next write reopens (and recreates) the
|
|
// file as it always has.
|
|
func (w *archiveWriter) sweepExpired(expiry time.Duration) error {
|
|
w.mu.Lock()
|
|
files, err := archiveFiles(w.path)
|
|
w.mu.Unlock()
|
|
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
var errs []error
|
|
|
|
for _, file := range files {
|
|
err = w.sweepFile(file, expiry)
|
|
if errors.Is(err, errArchiveWriterEvicted) {
|
|
return err
|
|
}
|
|
|
|
if err != nil {
|
|
errs = append(errs, err)
|
|
}
|
|
}
|
|
|
|
return errors.Join(errs...)
|
|
}
|
|
|
|
// sweepFile prunes one of the target's archive files for sweepExpired,
|
|
// holding w.mu while it does. It skips a file that is gone, and deletes
|
|
// the file, with its -wal and -shm, when it is named for a period and
|
|
// the prune leaves it empty.
|
|
func (w *archiveWriter) sweepFile(
|
|
file archiveFile, expiry time.Duration,
|
|
) error {
|
|
w.mu.Lock()
|
|
defer w.mu.Unlock()
|
|
|
|
if w.evicted {
|
|
return fmt.Errorf(
|
|
"%w: %s", errArchiveWriterEvicted, w.path,
|
|
)
|
|
}
|
|
|
|
if !fileExists(file.path) {
|
|
return nil
|
|
}
|
|
|
|
// Drop any live handle first so the prune runs against a
|
|
// freshly opened file, matching the write path's semantics.
|
|
w.close()
|
|
|
|
err := w.openMode(file.path, archiveModeExisting, expiry)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if file.period == "" {
|
|
w.close()
|
|
|
|
return nil
|
|
}
|
|
|
|
var rows int64
|
|
|
|
err = w.db.Model(&archivedEvent{}).Count(&rows).Error
|
|
|
|
w.close()
|
|
|
|
if err != nil {
|
|
return fmt.Errorf(
|
|
"counting rows in archive %s: %w", file.path, err,
|
|
)
|
|
}
|
|
|
|
if rows > 0 {
|
|
return nil
|
|
}
|
|
|
|
for _, suffix := range []string{"", "-wal", "-shm"} {
|
|
err = os.Remove(file.path + suffix)
|
|
if err != nil && !errors.Is(err, fs.ErrNotExist) {
|
|
return fmt.Errorf("deleting empty archive file: %w", err)
|
|
}
|
|
}
|
|
|
|
w.log.Info("deleted empty archive file", "path", file.path)
|
|
|
|
return nil
|
|
}
|
|
|
|
// rename gives every one of the target's archive files the new
|
|
// name, keeping the period in the name of each (see
|
|
// archivePeriodPath), and the writer uses the files under that name
|
|
// from now on. The handle is closed first, which folds the -wal into
|
|
// the .db; any -wal or -shm still beside a file (left by a crash) is
|
|
// moved with it, because SQLite finds them by name. A target with no
|
|
// files is not an error: the operator may have moved them away, and
|
|
// the next write creates its file under the new name.
|
|
//
|
|
// If a file already has one of the new names, nothing is moved and
|
|
// the error is ErrArchiveNameTaken. If one file fails to move, those
|
|
// already moved are moved back before the error is returned, so the
|
|
// archive is never split across two names.
|
|
func (w *archiveWriter) rename(name string) error {
|
|
w.mu.Lock()
|
|
defer w.mu.Unlock()
|
|
|
|
if w.evicted {
|
|
return fmt.Errorf(
|
|
"%w: %s", errArchiveWriterEvicted, w.path,
|
|
)
|
|
}
|
|
|
|
path := filepath.Join(filepath.Dir(w.path), name)
|
|
if path == w.path {
|
|
return nil
|
|
}
|
|
|
|
files, err := archiveFiles(w.path)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// from[i] moves to to[i].
|
|
var from, to []string
|
|
|
|
for _, file := range files {
|
|
renamed := archivePeriodPath(path, file.period)
|
|
|
|
for _, suffix := range []string{"", "-wal", "-shm"} {
|
|
from = append(from, file.path+suffix)
|
|
to = append(to, renamed+suffix)
|
|
}
|
|
}
|
|
|
|
for _, taken := range to {
|
|
if fileExists(taken) {
|
|
return fmt.Errorf(
|
|
"%w: %s", ErrArchiveNameTaken, filepath.Base(taken),
|
|
)
|
|
}
|
|
}
|
|
|
|
w.close()
|
|
|
|
for i := range from {
|
|
err = os.Rename(from[i], to[i])
|
|
if err == nil || errors.Is(err, fs.ErrNotExist) {
|
|
continue
|
|
}
|
|
|
|
for j := range i {
|
|
backErr := os.Rename(to[j], from[j])
|
|
if backErr != nil && !errors.Is(backErr, fs.ErrNotExist) {
|
|
w.log.Error(
|
|
"failed to move archive file back",
|
|
"from", to[j],
|
|
"to", from[j],
|
|
"error", backErr,
|
|
)
|
|
}
|
|
}
|
|
|
|
return fmt.Errorf(
|
|
"renaming archive %s to %s: %w", from[i], to[i], err,
|
|
)
|
|
}
|
|
|
|
w.path = path
|
|
|
|
return nil
|
|
}
|
|
|
|
// evict closes the writer's handle and marks it unusable. It is
|
|
// called when the writer leaves the registry, because its target
|
|
// or its webhook was deleted, or at shutdown. The archive FILE is
|
|
// deliberately left on disk: it is long-term storage an operator
|
|
// may still want.
|
|
func (w *archiveWriter) evict() {
|
|
w.mu.Lock()
|
|
defer w.mu.Unlock()
|
|
|
|
w.evicted = true
|
|
|
|
w.close()
|
|
}
|
|
|
|
// prune deletes archived rows older than expiry, measured from
|
|
// each row's archived time. It runs on every (re)open, so a
|
|
// steadily written archive is swept by its own write traffic. An
|
|
// archive that goes idle receives no further reopens, which is
|
|
// why ArchiveSweeper exists to drive sweepExpired on a timer.
|
|
// Failures are logged, not fatal: a prune error must not stop
|
|
// archiving.
|
|
func (w *archiveWriter) prune(expiry time.Duration) {
|
|
cutoff := time.Now().Add(-expiry)
|
|
|
|
res := w.db.Where("archived_at < ?", cutoff).
|
|
Delete(&archivedEvent{})
|
|
if res.Error != nil {
|
|
w.log.Error(
|
|
"failed to prune expired archive rows",
|
|
"path", w.current,
|
|
"error", res.Error,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
if res.RowsAffected > 0 {
|
|
w.log.Info(
|
|
"pruned expired archive rows",
|
|
"path", w.current,
|
|
"rows_deleted", res.RowsAffected,
|
|
)
|
|
}
|
|
}
|
|
|
|
// ArchiveFileInfo is what the metadata of a database target's archive
|
|
// files says about them.
|
|
type ArchiveFileInfo struct {
|
|
// Files counts the files.
|
|
Files int
|
|
|
|
// Size is the bytes on disk of the files and their -wal together.
|
|
Size int64
|
|
|
|
// Written is when a file or a -wal was last modified, whichever is
|
|
// latest: a write lands in the -wal first.
|
|
Written time.Time
|
|
}
|
|
|
|
// StatArchive reads the metadata of a database target's archive
|
|
// files, given the path ArchivePath gives it (see archiveFiles), and
|
|
// of their -wal, without opening them. With no files, which is so
|
|
// before the first write and after the operator moved them away, the
|
|
// error wraps fs.ErrNotExist.
|
|
func StatArchive(path string) (ArchiveFileInfo, error) {
|
|
files, err := archiveFiles(path)
|
|
if err != nil {
|
|
return ArchiveFileInfo{}, err
|
|
}
|
|
|
|
var info ArchiveFileInfo
|
|
|
|
for _, file := range files {
|
|
db, err := os.Stat(file.path)
|
|
if errors.Is(err, fs.ErrNotExist) {
|
|
continue
|
|
}
|
|
|
|
if err != nil {
|
|
return ArchiveFileInfo{}, err
|
|
}
|
|
|
|
info.Files++
|
|
info.Size += db.Size()
|
|
|
|
if db.ModTime().After(info.Written) {
|
|
info.Written = db.ModTime()
|
|
}
|
|
|
|
wal, err := os.Stat(file.path + "-wal")
|
|
if errors.Is(err, fs.ErrNotExist) {
|
|
continue
|
|
}
|
|
|
|
if err != nil {
|
|
return ArchiveFileInfo{}, err
|
|
}
|
|
|
|
info.Size += wal.Size()
|
|
|
|
if wal.ModTime().After(info.Written) {
|
|
info.Written = wal.ModTime()
|
|
}
|
|
}
|
|
|
|
if info.Files == 0 {
|
|
return ArchiveFileInfo{}, fmt.Errorf(
|
|
"no archive file for %s: %w", path, fs.ErrNotExist,
|
|
)
|
|
}
|
|
|
|
return info, nil
|
|
}
|
|
|
|
// fileExists reports whether a path currently exists.
|
|
func fileExists(path string) bool {
|
|
_, err := os.Stat(path)
|
|
|
|
return err == nil
|
|
}
|