check / check (push) Waiting to run
Each database target on the webhook page has a Download button that streams its archive as gzipped JSON, archive-WEBHOOKNAME-TARGETNAME-TIME.json.gz, with names made safe by delivery.ArchiveFileName's function. The export reads one consistent snapshot through one cursor in a read-only transaction, so archive writes carry on, and holds the rename lock only while it reads the stored names and opens the file. It extends its write deadline as it writes, so a large archive downloads for as long as the client reads; a failure after the response has started aborts the connection so the browser marks the download failed. The request limit is now the service's own middleware, which no longer writes a 504 over a response already started. Model: opus-5-5
437 lines
11 KiB
Go
437 lines
11 KiB
Go
package delivery
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"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 is the one ArchivePath gives for the webhook and the target as
|
|
// the main database names 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,
|
|
)
|
|
}
|
|
|
|
w := newArchiveWriter(
|
|
ArchivePath(t.eng.dbManager, &target.Webhook, &target),
|
|
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)
|
|
}
|