check / check (push) Waiting to run
Each database target now writes its own archive file, archive-WEBHOOKNAME-TARGETNAME-TARGETID.db, named by delivery.ArchiveFileName, in place of one archive per webhook keyed on its UUID. Renaming a webhook or a target renames its archive files (with any -wal and -shm) before the new name is saved, never over an existing file, and moves every one back if a rename or the save fails. The webhook edit, the target edit and target creation share one lock so no two interleave. Deleting a webhook or target leaves its files on disk. Nothing looks for the old archive-WEBHOOKID.db files. The README gives the naming and the recovery steps. Model: opus-5-5
441 lines
11 KiB
Go
441 lines
11 KiB
Go
package delivery
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"path/filepath"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"gorm.io/gorm"
|
|
"sneak.berlin/go/webhooker/internal/database"
|
|
)
|
|
|
|
// archiveNameMaxLen is how many characters of a webhook or target
|
|
// name an archive file name keeps.
|
|
const archiveNameMaxLen = 40
|
|
|
|
// databaseTarget is a no-retry target that archives the full
|
|
// inbound event into the target's own archive SQLite file, separate
|
|
// from the per-webhook event database. The event is already
|
|
// persisted in the per-webhook event DB by the time delivery runs;
|
|
// the database target additionally writes a durable long-term copy
|
|
// into the file ArchiveFileName names and then records a single
|
|
// attempt whose outcome reflects whether the archive write
|
|
// succeeded. See archiveWriter for the close/reopen, auto-recreate,
|
|
// and expiry semantics.
|
|
type databaseTarget struct {
|
|
eng *Engine
|
|
|
|
// writers holds one archive writer per database target, keyed
|
|
// by target ID.
|
|
mu sync.Mutex
|
|
writers map[string]*archiveWriter
|
|
}
|
|
|
|
// ArchiveFileName returns the file name of a database target's
|
|
// archive: archive-WEBHOOKNAME-TARGETNAME-TARGETID.db, with both
|
|
// names passed through archiveNamePart. The target ID keeps the
|
|
// name unique when two targets' names come out the same.
|
|
func ArchiveFileName(webhookName, targetName, targetID string) string {
|
|
return "archive-" + archiveNamePart(webhookName) + "-" +
|
|
archiveNamePart(targetName) + "-" + targetID + ".db"
|
|
}
|
|
|
|
// archiveNamePart makes a webhook or target name safe to put in a
|
|
// file name. It is lowercased; ASCII letters and digits are kept,
|
|
// every other run of characters becomes a single "-", and no "-" is
|
|
// left at either end. It is cut to archiveNameMaxLen characters, and
|
|
// a name with nothing left is "unnamed".
|
|
func archiveNamePart(name string) string {
|
|
var b strings.Builder
|
|
|
|
dash := false
|
|
|
|
for _, r := range strings.ToLower(name) {
|
|
if (r < 'a' || r > 'z') && (r < '0' || r > '9') {
|
|
dash = b.Len() > 0
|
|
|
|
continue
|
|
}
|
|
|
|
if dash {
|
|
b.WriteByte('-')
|
|
|
|
dash = false
|
|
}
|
|
|
|
b.WriteRune(r)
|
|
}
|
|
|
|
part := b.String()
|
|
if len(part) > archiveNameMaxLen {
|
|
part = strings.TrimRight(part[:archiveNameMaxLen], "-")
|
|
}
|
|
|
|
if part == "" {
|
|
return "unnamed"
|
|
}
|
|
|
|
return part
|
|
}
|
|
|
|
// Deliver implements Target. It archives the event, then
|
|
// records one successful attempt and marks the delivery
|
|
// delivered. An archiving error fails the delivery: the
|
|
// attempt is recorded as failed with the error and the
|
|
// delivery is marked failed, so a target that could not do
|
|
// its one job (archiving) never reports success. The target
|
|
// does not retry; the event remains durably stored in the
|
|
// per-webhook event database.
|
|
func (t *databaseTarget) Deliver(
|
|
_ context.Context,
|
|
webhookDB *gorm.DB,
|
|
d *database.Delivery,
|
|
_ *Task,
|
|
_ Scheduler,
|
|
) {
|
|
start := time.Now()
|
|
|
|
err := t.archive(d)
|
|
|
|
elapsed := time.Since(start)
|
|
|
|
t.eng.observeAttempt(d.Target.Type, elapsed)
|
|
|
|
if err != nil {
|
|
t.eng.log.Error(
|
|
"failed to archive event to database target",
|
|
"delivery_id", d.ID,
|
|
"event_id", d.EventID,
|
|
"error", err,
|
|
)
|
|
|
|
recErr := t.eng.recordResult(
|
|
webhookDB, d, 1, false, 0, "",
|
|
err.Error(), elapsed.Milliseconds(),
|
|
)
|
|
if recErr != nil {
|
|
t.eng.bookkeepingFailed(d, recErr)
|
|
|
|
return
|
|
}
|
|
|
|
t.eng.settleStatus(
|
|
webhookDB, d, d.Target.Type,
|
|
database.DeliveryStatusFailed,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
recErr := t.eng.recordResult(
|
|
webhookDB, d, 1, true, 0, "", "",
|
|
elapsed.Milliseconds(),
|
|
)
|
|
if recErr != nil {
|
|
t.eng.bookkeepingFailed(d, recErr)
|
|
|
|
return
|
|
}
|
|
|
|
t.eng.settleStatus(
|
|
webhookDB, d, d.Target.Type,
|
|
database.DeliveryStatusDelivered,
|
|
)
|
|
}
|
|
|
|
// archive writes the full event as a row into the target's
|
|
// archive database, honouring the optional per-target expiry
|
|
// parsed from the target config JSON.
|
|
func (t *databaseTarget) archive(d *database.Delivery) error {
|
|
webhookID := d.Event.WebhookID
|
|
if webhookID == "" {
|
|
return errArchiveMissingWebhookID
|
|
}
|
|
|
|
expiry, err := parseArchiveExpiry(d.Target.Config)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
w, err := t.writerFor(d.TargetID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
row := archivedEvent{
|
|
EventID: d.Event.ID,
|
|
WebhookID: webhookID,
|
|
EntrypointID: d.Event.EntrypointID,
|
|
Method: d.Event.Method,
|
|
Headers: d.Event.Headers,
|
|
Body: d.Event.Body,
|
|
ContentType: d.Event.ContentType,
|
|
}
|
|
|
|
return w.write(row, expiry)
|
|
}
|
|
|
|
// writerFor returns the archive writer for a database target,
|
|
// creating and caching it on first use. Each target has one writer
|
|
// so its close/reopen debounce state is shared across concurrent
|
|
// deliveries, and so a rename and the idle sweep take the same lock
|
|
// as its writes.
|
|
func (t *databaseTarget) writerFor(
|
|
targetID string,
|
|
) (*archiveWriter, error) {
|
|
t.mu.Lock()
|
|
defer t.mu.Unlock()
|
|
|
|
w, ok := t.writers[targetID]
|
|
if !ok {
|
|
var err error
|
|
|
|
w, err = t.newWriter(targetID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
if t.writers == nil {
|
|
t.writers = make(map[string]*archiveWriter)
|
|
}
|
|
|
|
t.writers[targetID] = w
|
|
}
|
|
|
|
// A delivery claims the entry: even if the idle sweep created
|
|
// it moments ago, it now belongs to the registry proper and
|
|
// the sweep must leave it in place when it finishes.
|
|
w.sweepOwned = false
|
|
|
|
return w, nil
|
|
}
|
|
|
|
// sweepWriterFor returns the archive writer the idle sweep should
|
|
// prune a target's archive through, together with whether the sweep
|
|
// itself created the registry entry.
|
|
//
|
|
// The sweep must route its prune through the registered writer so
|
|
// the writer's mutex orders it against concurrent writes, but it
|
|
// must never leave a registry entry behind: a sweep that ran
|
|
// concurrently with the target's deletion would otherwise
|
|
// re-create an entry that nothing will ever evict again, which is
|
|
// exactly the leak eviction exists to prevent. An entry the sweep
|
|
// creates is therefore marked sweep-owned and handed back to
|
|
// releaseSweepWriter when the sweep is done.
|
|
func (t *databaseTarget) sweepWriterFor(
|
|
targetID string,
|
|
) (*archiveWriter, bool, error) {
|
|
t.mu.Lock()
|
|
defer t.mu.Unlock()
|
|
|
|
w, ok := t.writers[targetID]
|
|
if ok {
|
|
return w, false, nil
|
|
}
|
|
|
|
w, err := t.newWriter(targetID)
|
|
if err != nil {
|
|
return nil, false, err
|
|
}
|
|
|
|
if t.writers == nil {
|
|
t.writers = make(map[string]*archiveWriter)
|
|
}
|
|
|
|
w.sweepOwned = true
|
|
t.writers[targetID] = w
|
|
|
|
return w, true, nil
|
|
}
|
|
|
|
// releaseSweepWriter drops a registry entry that the idle sweep
|
|
// created, so a sweep leaves the registry exactly as it found it.
|
|
//
|
|
// The entry is removed only if it is still the very writer the
|
|
// sweep installed and no delivery has claimed it in the meantime
|
|
// (writerFor clears sweepOwned when it hands a writer to the
|
|
// write path). Both conditions are evaluated under the registry
|
|
// lock, so an eviction that raced the sweep — which removes the
|
|
// entry outright — simply finds nothing left to do here, and a
|
|
// delivery that adopted the writer keeps a registered, evictable
|
|
// one.
|
|
func (t *databaseTarget) releaseSweepWriter(
|
|
targetID string, w *archiveWriter,
|
|
) {
|
|
t.mu.Lock()
|
|
defer t.mu.Unlock()
|
|
|
|
cur, ok := t.writers[targetID]
|
|
if !ok || cur != w || !cur.sweepOwned {
|
|
return
|
|
}
|
|
|
|
delete(t.writers, targetID)
|
|
}
|
|
|
|
// newWriter builds the writer for a database target's archive. The
|
|
// file lives beside the webhook's event database in the data
|
|
// directory and is named for the webhook and the target as the main
|
|
// database has them now; from then on only rename changes the name
|
|
// the writer uses. It does not touch the archive file.
|
|
func (t *databaseTarget) newWriter(
|
|
targetID string,
|
|
) (*archiveWriter, error) {
|
|
if t.eng.dbManager == nil {
|
|
return nil, errArchiveNoDataDir
|
|
}
|
|
|
|
var target database.Target
|
|
|
|
err := t.eng.database.DB().
|
|
Preload("Webhook").
|
|
First(&target, "id = ?", targetID).Error
|
|
if err != nil {
|
|
return nil, fmt.Errorf(
|
|
"loading database target %s: %w", targetID, err,
|
|
)
|
|
}
|
|
|
|
dir := filepath.Dir(t.eng.dbManager.DBPath(target.WebhookID))
|
|
name := ArchiveFileName(
|
|
target.Webhook.Name, target.Name, target.ID,
|
|
)
|
|
|
|
w := newArchiveWriter(filepath.Join(dir, name), t.eng.log)
|
|
w.webhookID = target.WebhookID
|
|
|
|
return w, nil
|
|
}
|
|
|
|
// rename moves a database target's archive file to the name for
|
|
// webhookName and targetName. It goes through the target's writer,
|
|
// so the move holds the lock that writes and the idle sweep take,
|
|
// and later writes use the new name.
|
|
//
|
|
// The writer is created if there is none, and it stays cached. The
|
|
// handlers rename before they save the new name, so until the save
|
|
// the main database still has the old one; a delivery in that window
|
|
// must find this writer rather than build one from the old name.
|
|
func (t *databaseTarget) rename(
|
|
targetID, webhookName, targetName string,
|
|
) error {
|
|
w, err := t.writerFor(targetID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return w.rename(ArchiveFileName(webhookName, targetName, targetID))
|
|
}
|
|
|
|
// evict drops a database target's archive writer from the registry
|
|
// and closes its handle, so a deleted target does not leave a
|
|
// writer (and an open archive handle within its debounce window)
|
|
// alive for the process lifetime.
|
|
//
|
|
// The map entry is removed under the registry lock, which is
|
|
// then released before the handle is closed under the writer's
|
|
// own lock: that ordering keeps the registry available to other
|
|
// targets while an in-flight write on this one drains, and
|
|
// closing under the writer's lock means eviction can never race
|
|
// a write.
|
|
//
|
|
// Eviction is idempotent and silent for a target with no writer,
|
|
// which is the common case: only a database target that has
|
|
// received an event or been renamed has one. It never deletes the
|
|
// archive file.
|
|
func (t *databaseTarget) evict(targetID string) {
|
|
t.mu.Lock()
|
|
|
|
w, ok := t.writers[targetID]
|
|
if ok {
|
|
delete(t.writers, targetID)
|
|
}
|
|
|
|
t.mu.Unlock()
|
|
|
|
if !ok {
|
|
return
|
|
}
|
|
|
|
w.evict()
|
|
|
|
t.eng.log.Info(
|
|
"evicted archive writer",
|
|
"target_id", targetID,
|
|
"path", w.path,
|
|
)
|
|
}
|
|
|
|
// evictWebhook evicts, exactly as evict does, the writer of every
|
|
// database target of a webhook.
|
|
func (t *databaseTarget) evictWebhook(webhookID string) {
|
|
t.mu.Lock()
|
|
|
|
var gone []*archiveWriter
|
|
|
|
for targetID, w := range t.writers {
|
|
if w.webhookID == webhookID {
|
|
delete(t.writers, targetID)
|
|
|
|
gone = append(gone, w)
|
|
}
|
|
}
|
|
|
|
t.mu.Unlock()
|
|
|
|
for _, w := range gone {
|
|
w.evict()
|
|
|
|
t.eng.log.Info(
|
|
"evicted archive writer",
|
|
"webhook_id", webhookID,
|
|
"path", w.path,
|
|
)
|
|
}
|
|
}
|
|
|
|
// evictAll evicts every cached archive writer, exactly as evict
|
|
// does for one target. The engine calls it at shutdown, once its
|
|
// workers have returned. Closing the last handle on an archive
|
|
// moves the contents of its -wal into the .db and removes the
|
|
// -wal, so a clean stop leaves each archive as a single file.
|
|
func (t *databaseTarget) evictAll() {
|
|
t.mu.Lock()
|
|
|
|
writers := t.writers
|
|
t.writers = nil
|
|
|
|
t.mu.Unlock()
|
|
|
|
for _, w := range writers {
|
|
w.evict()
|
|
}
|
|
}
|
|
|
|
// sweepArchive prunes one database target's archive of rows older
|
|
// than expiry, without requiring a write. A missing archive file is
|
|
// left missing (see sweepExpired), so a sweep never creates an
|
|
// archive for a target that has never received an event.
|
|
//
|
|
// It also never leaves a registry entry behind: an entry it had
|
|
// to create to reach the writer's mutex is released again once
|
|
// the prune is done, so a sweep racing a target deletion cannot
|
|
// resurrect the writer the eviction just dropped.
|
|
func (t *databaseTarget) sweepArchive(
|
|
targetID string, expiry time.Duration,
|
|
) error {
|
|
w, created, err := t.sweepWriterFor(targetID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if created {
|
|
defer t.releaseSweepWriter(targetID, w)
|
|
}
|
|
|
|
return w.sweepExpired(expiry)
|
|
}
|