package delivery import ( "database/sql" "encoding/json" "errors" "fmt" "log/slog" "os" "sync" "time" "gorm.io/driver/sqlite" "gorm.io/gorm" ) // 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 var ( // errArchiveMissingWebhookID is returned when an event to // archive has no webhook id to key its archive file on. 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\"", ) ) // 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"` } // archivedEvent is one fully captured webhook event stored in a // per-webhook archive database 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 is parsed as a Go duration. 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, nil } 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 per-webhook archive SQLite file. 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 string log *slog.Logger debounce time.Duration db *gorm.DB lastReopen time.Time reopens int } // 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, } } // write appends the event as a row, then applies the debounced // close/reopen. 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, ) error { w.mu.Lock() defer w.mu.Unlock() if w.db == nil || !fileExists(w.path) { err := w.reopen(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", w.path, err, ) } if time.Since(w.lastReopen) >= w.debounce { return w.reopen(expiry) } return nil } // open opens (creating if missing) the archive file, migrates // its schema, records the reopen time, and prunes expired rows // when expiry is positive. func (w *archiveWriter) open(expiry time.Duration) error { dbURL := fmt.Sprintf("file:%s?mode=rwc", w.path) sqlDB, err := sql.Open("sqlite", dbURL) if err != nil { return fmt.Errorf( "opening archive database %s: %w", w.path, err, ) } gdb, err := gorm.Open( sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, ) if err != nil { _ = sqlDB.Close() return fmt.Errorf( "connecting to archive database %s: %w", w.path, err, ) } err = gdb.AutoMigrate(&archivedEvent{}) if err != nil { _ = sqlDB.Close() return fmt.Errorf( "migrating archive database %s: %w", w.path, err, ) } w.db = gdb w.lastReopen = time.Now() w.reopens++ if expiry > 0 { w.prune(expiry) } return nil } // reopen closes any open handle and opens the file afresh. The // fresh open recreates the file if it was moved away. func (w *archiveWriter) reopen(expiry time.Duration) error { w.close() return w.open(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 } // prune deletes archived rows older than expiry, measured from // each row's archived time. It runs on every (re)open, and // because the file is reopened after writes this keeps the // archive swept without a separate background sweeper. 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.path, "error", res.Error, ) return } if res.RowsAffected > 0 { w.log.Info( "pruned expired archive rows", "path", w.path, "rows_deleted", res.RowsAffected, ) } } // fileExists reports whether a path currently exists. func fileExists(path string) bool { _, err := os.Stat(path) return err == nil }