All checks were successful
check / check (push) Successful in 3m38s
An operator running `sqlite3 <db> .dump` against their own per-webhook database wedged it: 60 of 60 inbound webhooks rejected with HTTP 500, 206 delivered webhooks stranded at `pending`, and every one of them POSTed a second time on the next restart while the event log recorded a single attempt. Durability. Every SQLite file — main, per-webhook, and archive — now opens through one path, `internal/database/sqlite_open.go`, in WAL journal mode with a 10-second busy timeout, `BEGIN IMMEDIATE` transactions, and a bounded connection pool. WAL is what stops a reader blocking writers at all. `_txlock=immediate` is what stops a `COMMIT` failing while its transaction stays open on a pooled connection, which is how four `database is locked` errors became 593 `cannot start a transaction within a transaction`: a deferred transaction that upgrades to a write lock mid-flight gets SQLITE_BUSY without the busy handler being consulted. `cache=shared` is gone, because under it an in-process conflict is SQLITE_LOCKED, which the busy handler does not retry. Delivery. `recordResult` and `updateDeliveryStatus` return their errors instead of logging and dropping them, and a caller whose bookkeeping write failed writes nothing at all — the delivery keeps whichever non-terminal status it already held, and both sweeps recover it. Recovery and the sweep now reconcile before re-sending: a pending delivery that already holds a successful `DeliveryResult` is marked delivered rather than sent again, which is the state that did not previously exist. A delivery handed back out is claimed by compare-and-set so successive sweeps cannot send it repeatedly, and it continues its own attempt numbering instead of restarting at 1. The sweep gains a `pending`-with-age-bound arm, so a stranded delivery no longer waits for a restart. Docs. WAL produces `-wal`/`-shm` sidecars, so the backup and restore procedures in README.md are corrected: both documented procedures were re-run against a live instance, and a `-wal` left by a crash carries data the `.db` alone does not. Verified by reproducing the failure on unmodified `next` first — 6 targets, 60 events at 5/s, a concurrent `.dump` reader — which gave 38 HTTP 500s and 112 duplicate POSTs at the sinks across a restart. Both arms of the matched pair now show 0 inbound 500s, 0 engine write errors, and 0 new requests at the sinks after a restart, counted by payload.
1576 lines
38 KiB
Go
1576 lines
38 KiB
Go
// Package delivery manages asynchronous event delivery
|
|
// to configured targets.
|
|
package delivery
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log/slog"
|
|
"net/http"
|
|
"sync"
|
|
"time"
|
|
|
|
"go.uber.org/fx"
|
|
"gorm.io/gorm"
|
|
"sneak.berlin/go/webhooker/internal/database"
|
|
"sneak.berlin/go/webhooker/internal/lifecycle"
|
|
"sneak.berlin/go/webhooker/internal/logger"
|
|
"sneak.berlin/go/webhooker/internal/metrics"
|
|
)
|
|
|
|
const (
|
|
// deliveryChannelSize is the buffer size for the delivery
|
|
// channel. New Tasks from the webhook handler are sent
|
|
// here. Workers drain this channel. Sized large enough
|
|
// that the webhook handler should never block under
|
|
// normal load.
|
|
deliveryChannelSize = 10000
|
|
|
|
// retryChannelSize is the buffer size for the retry
|
|
// channel. Timer-fired retries are sent here for
|
|
// processing by workers.
|
|
retryChannelSize = 10000
|
|
|
|
// defaultWorkers is the number of worker goroutines in
|
|
// the delivery engine pool. At most this many deliveries
|
|
// are in-flight at any time, preventing goroutine
|
|
// explosions regardless of queue depth.
|
|
defaultWorkers = 10
|
|
|
|
// retrySweepInterval is how often the periodic retry
|
|
// sweep runs.
|
|
retrySweepInterval = 60 * time.Second
|
|
|
|
// pendingSweepMinAge is how long a delivery must have sat at
|
|
// pending before the sweep treats it as stranded rather than as
|
|
// in flight.
|
|
//
|
|
// A delivery is pending from the moment it is created until its
|
|
// outcome is written, which includes the whole time a worker
|
|
// spends on it, so the bound has to clear the longest a live
|
|
// attempt can take: httpClientTimeout plus queueing behind the
|
|
// other deliveries in front of it. Five minutes is far above
|
|
// that, and still recovers a stranded delivery in minutes rather
|
|
// than at the next restart.
|
|
pendingSweepMinAge = 5 * time.Minute
|
|
|
|
// pendingSweepBatch bounds how many stranded pending deliveries
|
|
// one sweep of one webhook re-dispatches. The sweep runs every
|
|
// retrySweepInterval, so a larger backlog drains across
|
|
// successive sweeps instead of arriving as one burst against a
|
|
// database that was already struggling to accept writes.
|
|
pendingSweepBatch = 500
|
|
|
|
// MaxInlineBodySize is the maximum event body size that
|
|
// will be carried inline in a Task through the channel.
|
|
// Bodies at or above this size are left nil and fetched
|
|
// from the per-webhook database on demand.
|
|
MaxInlineBodySize = 16 * 1024
|
|
|
|
// httpClientTimeout is the timeout for outbound HTTP
|
|
// requests.
|
|
httpClientTimeout = 30 * time.Second
|
|
|
|
// maxBodyLog is the maximum response body length to
|
|
// store in DeliveryResult.
|
|
maxBodyLog = 4096
|
|
|
|
// maxBackoffShift caps the exponential backoff shift to
|
|
// avoid integer overflow in the 1<<shift expression.
|
|
maxBackoffShift = 30
|
|
|
|
// httpSuccessMin is the lower bound (inclusive) of the
|
|
// HTTP success status code range.
|
|
httpSuccessMin = 200
|
|
|
|
// httpSuccessMax is the upper bound (exclusive) of the
|
|
// HTTP success status code range.
|
|
httpSuccessMax = 300
|
|
)
|
|
|
|
// Task contains everything needed to deliver an event to a
|
|
// single target.
|
|
type Task struct {
|
|
DeliveryID string
|
|
EventID string
|
|
WebhookID string
|
|
EntrypointID string
|
|
|
|
TargetID string
|
|
TargetName string
|
|
TargetType database.TargetType
|
|
TargetConfig string
|
|
MaxRetries int
|
|
|
|
Method string
|
|
Headers string
|
|
ContentType string
|
|
Body *string
|
|
|
|
AttemptNum int
|
|
}
|
|
|
|
// Notifier is the interface for notifying the delivery
|
|
// engine about new deliveries.
|
|
type Notifier interface {
|
|
Notify(tasks []Task)
|
|
}
|
|
|
|
// WebhookEvictor releases the delivery engine's per-webhook
|
|
// state for a webhook that no longer needs it — currently the
|
|
// cached archive writer of the database target, whose open
|
|
// file handle would otherwise outlive the webhook.
|
|
//
|
|
// It is deliberately separate from Notifier and deliberately
|
|
// one method wide: archiving lifecycle is not notification, and
|
|
// a single-method interface keeps the handlers package free of
|
|
// any dependency on the engine's internals while staying
|
|
// trivially fakeable in tests.
|
|
//
|
|
// EvictWebhook never deletes an archive file. It is idempotent
|
|
// and is a no-op for a webhook with no engine state.
|
|
type WebhookEvictor interface {
|
|
EvictWebhook(webhookID string)
|
|
}
|
|
|
|
// EngineParams are the fx dependencies for the delivery
|
|
// engine.
|
|
type EngineParams struct {
|
|
fx.In
|
|
|
|
DB *database.Database
|
|
DBManager *database.WebhookDBManager
|
|
Logger *logger.Logger
|
|
SSRFGuard *Guard
|
|
}
|
|
|
|
// Engine processes queued deliveries in the background
|
|
// using a bounded worker pool architecture. It owns only
|
|
// the cross-target machinery: the worker pool, the queue
|
|
// and retry channels, restart recovery, and the persistence
|
|
// helpers and Scheduler that individual targets rely on.
|
|
// Each target type owns its own delivery, including retries,
|
|
// backoff, and circuit breaking.
|
|
type Engine struct {
|
|
database *database.Database
|
|
dbManager *database.WebhookDBManager
|
|
log *slog.Logger
|
|
cancel context.CancelFunc
|
|
wg sync.WaitGroup
|
|
deliveryCh chan Task
|
|
retryCh chan Task
|
|
workers int
|
|
|
|
// mtr is the delivery metric set. Production wires the
|
|
// process-wide one; a test can substitute a set registered on
|
|
// a private registry so its assertions are not disturbed by
|
|
// deliveries other tests are making at the same time.
|
|
mtr *metrics.Set
|
|
|
|
// targets maps each target type to its implementation.
|
|
targets map[database.TargetType]Target
|
|
|
|
// httpTarget is retained so tests can reach the HTTP
|
|
// target's shared client and circuit breakers.
|
|
httpTarget *httpTarget
|
|
|
|
// dbTarget is retained so the engine can reach the archive
|
|
// writer registry for webhook eviction and the idle sweep.
|
|
dbTarget *databaseTarget
|
|
}
|
|
|
|
// New creates and registers the delivery engine with the
|
|
// fx lifecycle.
|
|
func New(
|
|
lc fx.Lifecycle,
|
|
params EngineParams,
|
|
) *Engine {
|
|
e := &Engine{
|
|
database: params.DB,
|
|
dbManager: params.DBManager,
|
|
log: params.Logger.Get(),
|
|
deliveryCh: make(chan Task, deliveryChannelSize),
|
|
retryCh: make(chan Task, retryChannelSize),
|
|
workers: defaultWorkers,
|
|
mtr: metrics.Default(),
|
|
}
|
|
|
|
e.initTargets(&http.Client{
|
|
Timeout: httpClientTimeout,
|
|
Transport: params.SSRFGuard.NewSSRFSafeTransport(),
|
|
})
|
|
|
|
e.registerHooks(lc)
|
|
|
|
return e
|
|
}
|
|
|
|
// Notify signals the delivery engine that new deliveries
|
|
// are ready.
|
|
func (e *Engine) Notify(tasks []Task) {
|
|
for i := range tasks {
|
|
select {
|
|
case e.deliveryCh <- tasks[i]:
|
|
default:
|
|
e.log.Warn(
|
|
"delivery channel full, "+
|
|
"task will be recovered on restart",
|
|
"delivery_id", tasks[i].DeliveryID,
|
|
"event_id", tasks[i].EventID,
|
|
)
|
|
}
|
|
}
|
|
}
|
|
|
|
// EvictWebhook implements WebhookEvictor. It releases the
|
|
// engine's per-webhook archiving state: the database target's
|
|
// cached archive writer is dropped from the registry and its
|
|
// file handle closed. The archive file itself is left on disk
|
|
// — it is long-term storage the operator owns.
|
|
func (e *Engine) EvictWebhook(webhookID string) {
|
|
if e.dbTarget == nil {
|
|
return
|
|
}
|
|
|
|
e.dbTarget.evict(webhookID)
|
|
}
|
|
|
|
// ScheduleRetry schedules a task to be re-enqueued onto the
|
|
// retry channel after delay. It implements the Scheduler
|
|
// interface the targets use to own their durable retries.
|
|
func (e *Engine) ScheduleRetry(
|
|
task Task, delay time.Duration,
|
|
) {
|
|
e.log.Debug(
|
|
"scheduling delivery retry",
|
|
"webhook_id", task.WebhookID,
|
|
"delivery_id", task.DeliveryID,
|
|
"delay", delay,
|
|
"next_attempt", task.AttemptNum,
|
|
)
|
|
|
|
time.AfterFunc(delay, func() {
|
|
select {
|
|
case e.retryCh <- task:
|
|
default:
|
|
e.log.Warn(
|
|
"retry channel full, delivery "+
|
|
"will be recovered by periodic sweep",
|
|
"delivery_id", task.DeliveryID,
|
|
"webhook_id", task.WebhookID,
|
|
)
|
|
}
|
|
})
|
|
}
|
|
|
|
// registerHooks wires the engine's start and stop into the fx
|
|
// lifecycle. The start hook's context is deliberately ignored
|
|
// (see start for why the worker pool must not inherit it); the
|
|
// stop hook's context is honoured (see stop).
|
|
func (e *Engine) registerHooks(lc fx.Lifecycle) {
|
|
lc.Append(fx.Hook{
|
|
//nolint:contextcheck // Not inheriting the hook context
|
|
// is the point: see start.
|
|
OnStart: func(_ context.Context) error {
|
|
e.start()
|
|
|
|
return nil
|
|
},
|
|
OnStop: func(ctx context.Context) error {
|
|
return e.stop(ctx)
|
|
},
|
|
})
|
|
}
|
|
|
|
// start launches the worker pool, restart recovery, and the
|
|
// periodic retry sweep.
|
|
//
|
|
// Their context is derived from context.Background(), NOT from
|
|
// the fx OnStart hook context. The hook context carries fx's
|
|
// start timeout (15s by default) and is cancelled once the start
|
|
// phase completes, so goroutines derived from it stop a few
|
|
// seconds into the process: every worker would return and the
|
|
// engine would silently stop delivering webhooks entirely. A
|
|
// long-lived goroutine must outlive the startup phase, so its
|
|
// lifetime is bounded by OnStop instead: stop cancels this
|
|
// context and waits on the WaitGroup.
|
|
func (e *Engine) start() {
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
e.cancel = cancel
|
|
|
|
for range e.workers {
|
|
e.wg.Add(1)
|
|
|
|
go e.worker(ctx)
|
|
}
|
|
|
|
e.wg.Add(1)
|
|
|
|
go e.recoverPending(ctx)
|
|
|
|
e.wg.Add(1)
|
|
|
|
go e.retrySweep(ctx)
|
|
|
|
e.wg.Add(1)
|
|
|
|
go e.queueDepthSampler(ctx)
|
|
|
|
e.log.Info(
|
|
"delivery engine started",
|
|
"workers", e.workers,
|
|
)
|
|
}
|
|
|
|
// stop cancels the worker pool's context and waits for the pool
|
|
// to drain, bounded by the stop hook's context: a wedged worker
|
|
// must not hang the process past fx's stop timeout.
|
|
func (e *Engine) stop(ctx context.Context) error {
|
|
e.log.Info("delivery engine stopping")
|
|
|
|
if e.cancel != nil {
|
|
e.cancel()
|
|
}
|
|
|
|
err := lifecycle.WaitForShutdown(
|
|
ctx, e.log, "delivery engine", &e.wg,
|
|
)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
e.log.Info("delivery engine stopped")
|
|
|
|
return nil
|
|
}
|
|
|
|
func (e *Engine) worker(ctx context.Context) {
|
|
defer e.wg.Done()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case task := <-e.deliveryCh:
|
|
e.processNewTask(ctx, &task)
|
|
case task := <-e.retryCh:
|
|
e.processRetryTask(ctx, &task)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (e *Engine) recoverPending(ctx context.Context) {
|
|
defer e.wg.Done()
|
|
|
|
e.recoverInFlight(ctx)
|
|
}
|
|
|
|
func (e *Engine) processNewTask(
|
|
ctx context.Context, task *Task,
|
|
) {
|
|
webhookDB, err := e.dbManager.GetDB(task.WebhookID)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to get webhook database",
|
|
"webhook_id", task.WebhookID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
event := buildEventFromTask(task)
|
|
|
|
event, err = e.resolveEventBody(
|
|
webhookDB, event, task,
|
|
)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to fetch event body from database",
|
|
"event_id", task.EventID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
target := buildTargetFromTask(task)
|
|
|
|
d := &database.Delivery{
|
|
EventID: task.EventID,
|
|
TargetID: task.TargetID,
|
|
Status: database.DeliveryStatusPending,
|
|
Event: event,
|
|
Target: target,
|
|
}
|
|
d.ID = task.DeliveryID
|
|
|
|
e.processDelivery(ctx, webhookDB, d, task)
|
|
}
|
|
|
|
func (e *Engine) processRetryTask(
|
|
ctx context.Context, task *Task,
|
|
) {
|
|
webhookDB, err := e.dbManager.GetDB(task.WebhookID)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to get webhook database for retry",
|
|
"webhook_id", task.WebhookID,
|
|
"delivery_id", task.DeliveryID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
d, err := e.loadRetryDelivery(
|
|
webhookDB, task.DeliveryID,
|
|
)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to load delivery for retry",
|
|
"delivery_id", task.DeliveryID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
if d.Status != database.DeliveryStatusRetrying {
|
|
e.log.Debug(
|
|
"skipping retry for delivery "+
|
|
"no longer in retrying status",
|
|
"delivery_id", d.ID,
|
|
"status", d.Status,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
event := buildEventFromTask(task)
|
|
|
|
event, err = e.resolveEventBody(
|
|
webhookDB, event, task,
|
|
)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to fetch event body for retry",
|
|
"event_id", task.EventID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
target := buildTargetFromTask(task)
|
|
d.EventID = task.EventID
|
|
d.TargetID = task.TargetID
|
|
d.Event = event
|
|
d.Target = target
|
|
|
|
e.processDelivery(ctx, webhookDB, d, task)
|
|
}
|
|
|
|
func (e *Engine) recoverInFlight(ctx context.Context) {
|
|
var webhookIDs []string
|
|
|
|
err := e.database.DB().
|
|
Model(&database.Webhook{}).
|
|
Pluck("id", &webhookIDs).Error
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to query webhook IDs for recovery",
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
for _, webhookID := range webhookIDs {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
default:
|
|
}
|
|
|
|
if !e.dbManager.DBExists(webhookID) {
|
|
continue
|
|
}
|
|
|
|
e.recoverWebhookDeliveries(ctx, webhookID)
|
|
}
|
|
}
|
|
|
|
func (e *Engine) recoverWebhookDeliveries(
|
|
ctx context.Context, webhookID string,
|
|
) {
|
|
webhookDB, err := e.dbManager.GetDB(webhookID)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to get webhook database for recovery",
|
|
"webhook_id", webhookID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
e.recoverPendingDeliveries(
|
|
ctx, webhookDB, webhookID,
|
|
)
|
|
|
|
e.recoverRetryingDeliveries(
|
|
webhookDB, webhookID,
|
|
)
|
|
}
|
|
|
|
func (e *Engine) recoverRetryingDeliveries(
|
|
webhookDB *gorm.DB, webhookID string,
|
|
) {
|
|
var retrying []database.Delivery
|
|
|
|
err := webhookDB.
|
|
Where(
|
|
"status = ?",
|
|
database.DeliveryStatusRetrying,
|
|
).
|
|
Find(&retrying).Error
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to query retrying deliveries "+
|
|
"for recovery",
|
|
"webhook_id", webhookID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
for i := range retrying {
|
|
e.recoverSingleRetry(
|
|
webhookDB, webhookID, &retrying[i],
|
|
)
|
|
}
|
|
}
|
|
|
|
// recoverSingleRetry hands an orphaned retrying delivery back
|
|
// to its target to recompute the remaining backoff, then
|
|
// reschedules it. Targets that do not own durable retries
|
|
// (fire-and-forget) never produce retrying deliveries, so a
|
|
// delivery found in that state has had its target's type
|
|
// changed underneath it and is terminally failed.
|
|
func (e *Engine) recoverSingleRetry(
|
|
webhookDB *gorm.DB,
|
|
webhookID string,
|
|
d *database.Delivery,
|
|
) {
|
|
target, err := e.loadTarget(d.TargetID)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to load target for retrying "+
|
|
"delivery recovery",
|
|
"delivery_id", d.ID,
|
|
"target_id", d.TargetID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
rs, ok := e.targets[target.Type].(rescheduler)
|
|
if !ok {
|
|
e.failUnretryableRetry(
|
|
webhookDB, webhookID, d, &target,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
attemptNum := e.countAttempts(webhookDB, d.ID)
|
|
remaining := rs.remainingBackoff(
|
|
webhookDB, d.ID, attemptNum,
|
|
)
|
|
|
|
event, err := e.loadEvent(webhookDB, d.EventID)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to load event for retrying "+
|
|
"delivery recovery",
|
|
"delivery_id", d.ID,
|
|
"event_id", d.EventID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
task := buildRecoveryTask(
|
|
d, webhookID, &event, &target, attemptNum+1,
|
|
)
|
|
|
|
e.log.Info(
|
|
"recovering retrying delivery",
|
|
"webhook_id", webhookID,
|
|
"delivery_id", d.ID,
|
|
"attempt", attemptNum,
|
|
"remaining_backoff", remaining,
|
|
)
|
|
|
|
e.ScheduleRetry(task, remaining)
|
|
}
|
|
|
|
func (e *Engine) recoverPendingDeliveries(
|
|
ctx context.Context,
|
|
webhookDB *gorm.DB,
|
|
webhookID string,
|
|
) {
|
|
var deliveries []database.Delivery
|
|
|
|
result := webhookDB.
|
|
Where(
|
|
"status = ?",
|
|
database.DeliveryStatusPending,
|
|
).
|
|
Preload("Event").
|
|
Find(&deliveries)
|
|
|
|
if result.Error != nil {
|
|
e.log.Error(
|
|
"failed to query pending deliveries",
|
|
"webhook_id", webhookID,
|
|
"error", result.Error,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
if len(deliveries) == 0 {
|
|
return
|
|
}
|
|
|
|
e.log.Info(
|
|
"recovering pending deliveries",
|
|
"webhook_id", webhookID,
|
|
"count", len(deliveries),
|
|
)
|
|
|
|
e.recoverPendingBatch(
|
|
ctx, webhookDB, webhookID, deliveries,
|
|
)
|
|
}
|
|
|
|
// recoverPendingBatch settles every delivery in the batch that was
|
|
// already delivered, and re-dispatches only the rest. Both the
|
|
// restart-time recovery and the periodic sweep go through it, so a
|
|
// pending delivery is treated the same however it was found.
|
|
func (e *Engine) recoverPendingBatch(
|
|
ctx context.Context,
|
|
webhookDB *gorm.DB,
|
|
webhookID string,
|
|
deliveries []database.Delivery,
|
|
) {
|
|
targetMap := e.loadTargetMap(deliveries)
|
|
|
|
settled := e.reconcileDelivered(
|
|
webhookDB, webhookID, deliveries, targetMap,
|
|
)
|
|
|
|
e.sendRecoveredDeliveries(
|
|
ctx, webhookDB, deliveries, webhookID,
|
|
targetMap, settled,
|
|
)
|
|
}
|
|
|
|
// reconcileDelivered finds the deliveries in a pending batch that
|
|
// already have a successful DeliveryResult, marks them delivered, and
|
|
// returns their ids so the caller does not send them a second time.
|
|
//
|
|
// This is the state the engine previously had no way to represent. A
|
|
// delivery is left pending by a failed bookkeeping write, and that
|
|
// covers two different histories: nothing was ever sent, or the send
|
|
// reached the receiver and only the status write failed. Re-sending
|
|
// was the sole option, so every stranded row produced a duplicate at
|
|
// the receiver and an event log that recorded one attempt for two
|
|
// POSTs. A successful result row distinguishes them: it is written
|
|
// before the status, so its presence means the wire I/O happened and
|
|
// was recorded, and all that is missing is the status.
|
|
//
|
|
// Deliveries whose result row itself never landed are not in the
|
|
// returned set and are re-sent, recorded as the further attempt they
|
|
// are. That is honest at-least-once delivery rather than a silent
|
|
// duplicate.
|
|
func (e *Engine) reconcileDelivered(
|
|
webhookDB *gorm.DB,
|
|
webhookID string,
|
|
deliveries []database.Delivery,
|
|
targetMap map[string]database.Target,
|
|
) map[string]struct{} {
|
|
settled := make(map[string]struct{})
|
|
|
|
if len(deliveries) == 0 {
|
|
return settled
|
|
}
|
|
|
|
ids := make([]string, 0, len(deliveries))
|
|
for i := range deliveries {
|
|
ids = append(ids, deliveries[i].ID)
|
|
}
|
|
|
|
var deliveredIDs []string
|
|
|
|
err := webhookDB.
|
|
Model(&database.DeliveryResult{}).
|
|
Where(
|
|
"delivery_id IN ? AND success = ?", ids, true,
|
|
).
|
|
Distinct().
|
|
Pluck("delivery_id", &deliveredIDs).Error
|
|
if err != nil {
|
|
// Every delivery stays out of the settled set, so the batch
|
|
// is re-sent exactly as it was before this check existed.
|
|
// That is the safe direction: a duplicate delivery beats
|
|
// declaring a delivery successful on a query that failed.
|
|
e.log.Error(
|
|
"failed to query successful delivery results; "+
|
|
"pending deliveries will be re-sent",
|
|
"webhook_id", webhookID,
|
|
"error", err,
|
|
)
|
|
|
|
return settled
|
|
}
|
|
|
|
for _, id := range deliveredIDs {
|
|
settled[id] = struct{}{}
|
|
}
|
|
|
|
if len(settled) == 0 {
|
|
return settled
|
|
}
|
|
|
|
e.log.Info(
|
|
"settling pending deliveries that already succeeded",
|
|
"webhook_id", webhookID,
|
|
"count", len(settled),
|
|
)
|
|
|
|
for i := range deliveries {
|
|
if _, ok := settled[deliveries[i].ID]; !ok {
|
|
continue
|
|
}
|
|
|
|
e.settleStatus(
|
|
webhookDB,
|
|
&deliveries[i],
|
|
targetMap[deliveries[i].TargetID].Type,
|
|
database.DeliveryStatusDelivered,
|
|
)
|
|
}
|
|
|
|
return settled
|
|
}
|
|
|
|
func (e *Engine) retrySweep(ctx context.Context) {
|
|
defer e.wg.Done()
|
|
|
|
ticker := time.NewTicker(retrySweepInterval)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
e.sweepOrphanedRetries(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (e *Engine) sweepOrphanedRetries(
|
|
ctx context.Context,
|
|
) {
|
|
var webhookIDs []string
|
|
|
|
err := e.database.DB().
|
|
Model(&database.Webhook{}).
|
|
Pluck("id", &webhookIDs).Error
|
|
if err != nil {
|
|
e.log.Error(
|
|
"retry sweep: failed to query webhook IDs",
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
for _, webhookID := range webhookIDs {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
default:
|
|
}
|
|
|
|
if !e.dbManager.DBExists(webhookID) {
|
|
continue
|
|
}
|
|
|
|
e.sweepWebhookRetries(ctx, webhookID)
|
|
}
|
|
}
|
|
|
|
func (e *Engine) sweepWebhookRetries(
|
|
ctx context.Context, webhookID string,
|
|
) {
|
|
webhookDB, err := e.dbManager.GetDB(webhookID)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"retry sweep: failed to get webhook database",
|
|
"webhook_id", webhookID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
var retrying []database.Delivery
|
|
|
|
err = webhookDB.
|
|
Where(
|
|
"status = ?",
|
|
database.DeliveryStatusRetrying,
|
|
).
|
|
Find(&retrying).Error
|
|
if err != nil {
|
|
e.log.Error(
|
|
"retry sweep: "+
|
|
"failed to query retrying deliveries",
|
|
"webhook_id", webhookID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
for i := range retrying {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
default:
|
|
}
|
|
|
|
e.sweepSingleRetry(
|
|
webhookDB, webhookID, &retrying[i],
|
|
)
|
|
}
|
|
|
|
e.sweepWebhookPending(ctx, webhookDB, webhookID)
|
|
}
|
|
|
|
// sweepWebhookPending recovers deliveries stranded at pending.
|
|
//
|
|
// A delivery is created pending and leaves that state only when its
|
|
// outcome is written, so a pending row older than the age bound is one
|
|
// whose bookkeeping write failed — the state that used to sit there
|
|
// until a restart, and then produce a duplicate at the receiver. The
|
|
// sweep gives it the same reconcile-then-dispatch treatment restart
|
|
// recovery gets, so it costs a minute rather than an operator
|
|
// noticing.
|
|
//
|
|
// The age bound is what keeps the sweep off deliveries the workers
|
|
// still hold: a delivery in flight is pending too, and re-dispatching
|
|
// one would race the worker that owns it. It is measured on updated_at
|
|
// rather than created_at because claimPending stamps that column when
|
|
// a delivery is handed out, which is what stops the next sweep, a
|
|
// minute later, from sending the same delivery again while the first
|
|
// attempt is still running.
|
|
func (e *Engine) sweepWebhookPending(
|
|
ctx context.Context,
|
|
webhookDB *gorm.DB,
|
|
webhookID string,
|
|
) {
|
|
var pending []database.Delivery
|
|
|
|
err := webhookDB.
|
|
Where(
|
|
"status = ? AND updated_at < ?",
|
|
database.DeliveryStatusPending,
|
|
time.Now().Add(-pendingSweepMinAge),
|
|
).
|
|
Preload("Event").
|
|
Limit(pendingSweepBatch).
|
|
Find(&pending).Error
|
|
if err != nil {
|
|
e.log.Error(
|
|
"retry sweep: "+
|
|
"failed to query pending deliveries",
|
|
"webhook_id", webhookID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
if len(pending) == 0 {
|
|
return
|
|
}
|
|
|
|
e.log.Info(
|
|
"retry sweep: recovering stranded pending deliveries",
|
|
"webhook_id", webhookID,
|
|
"count", len(pending),
|
|
)
|
|
|
|
e.recoverPendingBatch(ctx, webhookDB, webhookID, pending)
|
|
}
|
|
|
|
// sweepSingleRetry re-enqueues an orphaned retrying delivery
|
|
// whose backoff window has elapsed, delegating the backoff
|
|
// decision to the delivery's target. A delivery whose target
|
|
// no longer owns durable retries is terminally failed.
|
|
func (e *Engine) sweepSingleRetry(
|
|
webhookDB *gorm.DB,
|
|
webhookID string,
|
|
d *database.Delivery,
|
|
) {
|
|
target, err := e.loadTarget(d.TargetID)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"retry sweep: failed to load target",
|
|
"delivery_id", d.ID,
|
|
"target_id", d.TargetID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
rs, ok := e.targets[target.Type].(rescheduler)
|
|
if !ok {
|
|
e.failUnretryableRetry(
|
|
webhookDB, webhookID, d, &target,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
attemptNum := e.countAttempts(webhookDB, d.ID)
|
|
|
|
if !rs.backoffElapsed(
|
|
webhookDB, d.ID, attemptNum,
|
|
) {
|
|
return
|
|
}
|
|
|
|
event, err := e.loadEvent(webhookDB, d.EventID)
|
|
if err != nil {
|
|
e.log.Error(
|
|
"retry sweep: failed to load event",
|
|
"delivery_id", d.ID,
|
|
"event_id", d.EventID,
|
|
"error", err,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
task := buildRecoveryTask(
|
|
d, webhookID, &event, &target, attemptNum+1,
|
|
)
|
|
|
|
select {
|
|
case e.retryCh <- task:
|
|
e.log.Info(
|
|
"retry sweep: "+
|
|
"recovered orphaned retrying delivery",
|
|
"delivery_id", d.ID,
|
|
"webhook_id", webhookID,
|
|
"attempt", attemptNum+1,
|
|
)
|
|
default:
|
|
}
|
|
}
|
|
|
|
// failUnretryableRetry terminally fails an orphaned retrying
|
|
// delivery whose target type no longer supports retries. Both
|
|
// restart recovery and the periodic sweep call it, so the
|
|
// terminal transition exists once.
|
|
//
|
|
// This is only reachable when a target's type has been changed
|
|
// out from under an in-flight retrying delivery (or the type is
|
|
// unknown to the registry): fire-and-forget targets never set
|
|
// status retrying themselves. Re-dispatching under the new type
|
|
// would be a delivery the operator never asked for, and leaving
|
|
// the row retrying strands it forever, so the delivery is
|
|
// failed with a recorded reason. The event stays stored, but
|
|
// nothing redelivers it today. Logged at warn, not error: this
|
|
// is operator-caused state, not a system fault.
|
|
func (e *Engine) failUnretryableRetry(
|
|
webhookDB *gorm.DB,
|
|
webhookID string,
|
|
d *database.Delivery,
|
|
target *database.Target,
|
|
) {
|
|
e.log.Warn(
|
|
"failing orphaned retrying delivery: target "+
|
|
"type no longer supports retries",
|
|
"webhook_id", webhookID,
|
|
"delivery_id", d.ID,
|
|
"target_id", target.ID,
|
|
"target_name", target.Name,
|
|
"target_type", target.Type,
|
|
)
|
|
|
|
reason := fmt.Sprintf(
|
|
"target type %q does not support retries; "+
|
|
"delivery was left retrying by a previous "+
|
|
"target type and has been failed terminally",
|
|
target.Type,
|
|
)
|
|
|
|
err := e.recordResult(
|
|
webhookDB,
|
|
d,
|
|
e.countAttempts(webhookDB, d.ID)+1,
|
|
false,
|
|
0,
|
|
"",
|
|
reason,
|
|
0,
|
|
)
|
|
if err != nil {
|
|
e.bookkeepingFailed(d, err)
|
|
|
|
return
|
|
}
|
|
|
|
// The type is passed rather than assigned onto d: the delivery
|
|
// is loaded here without its target relation, and populating
|
|
// d.Target would make GORM's SaveBeforeAssociations upsert the
|
|
// whole target row — plaintext config, which for a slack target
|
|
// is the credential — into the per-webhook event database. See
|
|
// https://git.eeqj.de/sneak/webhooker/issues/206.
|
|
e.settleStatus(
|
|
webhookDB, d, target.Type,
|
|
database.DeliveryStatusFailed,
|
|
)
|
|
}
|
|
|
|
// processDelivery dispatches a delivery to the target that
|
|
// owns its type. Unknown target types fail the delivery.
|
|
func (e *Engine) processDelivery(
|
|
ctx context.Context,
|
|
webhookDB *gorm.DB,
|
|
d *database.Delivery,
|
|
task *Task,
|
|
) {
|
|
target, ok := e.targets[d.Target.Type]
|
|
if !ok {
|
|
e.log.Error(
|
|
"unknown target type",
|
|
"target_id", d.TargetID,
|
|
"type", d.Target.Type,
|
|
)
|
|
|
|
e.settleStatus(
|
|
webhookDB, d, d.Target.Type,
|
|
database.DeliveryStatusFailed,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
target.Deliver(ctx, webhookDB, d, task, e)
|
|
}
|
|
|
|
// observeAttempt counts one delivery attempt that was actually
|
|
// dispatched to a target, and records how long it took.
|
|
//
|
|
// It is called from the dispatch paths rather than from around
|
|
// Target.Deliver, because Deliver is also entered for deliveries
|
|
// that never reach the wire: a delivery an open circuit breaker
|
|
// refuses sends nothing, records no DeliveryResult, and is
|
|
// rescheduled. Counting those would climb the attempts counter with
|
|
// no traffic behind it and fill the duration histogram with
|
|
// microsecond samples, which would make the delivery-duration
|
|
// quantiles improve during exactly the outage they exist to reveal.
|
|
func (e *Engine) observeAttempt(
|
|
t database.TargetType, dur time.Duration,
|
|
) {
|
|
e.mtr.DeliveryAttempted(t)
|
|
e.mtr.ObserveDeliveryDuration(t, dur)
|
|
}
|
|
|
|
// recordResult persists a DeliveryResult row describing a
|
|
// single attempt. It is a cross-target helper the targets
|
|
// call.
|
|
//
|
|
// It returns its error rather than swallowing it. A DeliveryResult
|
|
// row is the only record that an attempt happened at all, so a
|
|
// caller that ignored a failed write would go on to mark the
|
|
// delivery delivered — leaving the event log claiming one attempt
|
|
// for a receiver that got two. Every caller must instead stop
|
|
// advancing the delivery's status and let it stay in the
|
|
// non-terminal state it already holds; see bookkeepingFailed.
|
|
func (e *Engine) recordResult(
|
|
webhookDB *gorm.DB,
|
|
d *database.Delivery,
|
|
attemptNum int,
|
|
success bool,
|
|
statusCode int,
|
|
respBody, errMsg string,
|
|
durationMs int64,
|
|
) error {
|
|
result := &database.DeliveryResult{
|
|
DeliveryID: d.ID,
|
|
AttemptNum: attemptNum,
|
|
Success: success,
|
|
StatusCode: statusCode,
|
|
ResponseBody: truncate(respBody, maxBodyLog),
|
|
Error: errMsg,
|
|
Duration: durationMs,
|
|
}
|
|
|
|
err := webhookDB.Create(result).Error
|
|
if err != nil {
|
|
return fmt.Errorf(
|
|
"recording delivery result for %s: %w", d.ID, err,
|
|
)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// bookkeepingFailed reports that a delivery's own record of what
|
|
// happened could not be written, and deliberately writes nothing in
|
|
// response.
|
|
//
|
|
// Leaving the row alone is the whole point. A delivery is created
|
|
// pending and only ever leaves that state through
|
|
// updateDeliveryStatus, so a delivery whose bookkeeping write failed
|
|
// is still pending or retrying — the two non-terminal states, per
|
|
// DeliveryStatus.Terminal — and both are swept and recovered. Writing
|
|
// anything here would need the very database that just refused a
|
|
// write, and would be one more thing to fail; not writing cannot.
|
|
//
|
|
// The cost is honest at-least-once behaviour: a send that reached the
|
|
// receiver but whose result row did not land is attempted again, and
|
|
// recorded as the further attempt it is. What no longer happens is the
|
|
// silent duplicate — a second POST the event log denies ever
|
|
// occurred. Where the result row *did* land and only the status write
|
|
// failed, reconcileDelivered settles the row without re-sending.
|
|
func (e *Engine) bookkeepingFailed(
|
|
d *database.Delivery, err error,
|
|
) {
|
|
e.log.Error(
|
|
"delivery bookkeeping write failed; leaving delivery "+
|
|
"in a recoverable state",
|
|
"delivery_id", d.ID,
|
|
"event_id", d.EventID,
|
|
"target_id", d.TargetID,
|
|
"status", d.Status,
|
|
"error", err,
|
|
)
|
|
}
|
|
|
|
// updateDeliveryStatus persists a new status for a delivery.
|
|
// It is a cross-target helper the targets call, and therefore the
|
|
// single point where a delivery's outcome — delivered, terminally
|
|
// failed, or put back into retry — is counted.
|
|
//
|
|
// The target type is a parameter rather than read off d.Target
|
|
// because one caller — failUnretryableRetry — deliberately holds a
|
|
// delivery loaded without its target relation, and must keep it that
|
|
// way: a populated d.Target makes GORM upsert the target row, config
|
|
// and all, into the per-webhook database.
|
|
//
|
|
// The counter moves only after the row is written, so a transition
|
|
// the database rejected is not claimed as an outcome that happened.
|
|
// For the same reason the error is returned rather than logged and
|
|
// dropped: a delivery whose status write failed has not reached that
|
|
// status, and its caller must not act as though it had.
|
|
func (e *Engine) updateDeliveryStatus(
|
|
webhookDB *gorm.DB,
|
|
d *database.Delivery,
|
|
targetType database.TargetType,
|
|
status database.DeliveryStatus,
|
|
) error {
|
|
err := webhookDB.Model(d).
|
|
Update("status", status).Error
|
|
if err != nil {
|
|
return fmt.Errorf(
|
|
"updating delivery %s to status %s: %w",
|
|
d.ID, status, err,
|
|
)
|
|
}
|
|
|
|
e.mtr.DeliveryStatusChanged(targetType, status)
|
|
|
|
return nil
|
|
}
|
|
|
|
// settleStatus moves a delivery to its outcome status and reports a
|
|
// failed write through bookkeepingFailed, which leaves the row
|
|
// recoverable. It exists so the target call sites read as one
|
|
// statement rather than four lines of identical error handling.
|
|
func (e *Engine) settleStatus(
|
|
webhookDB *gorm.DB,
|
|
d *database.Delivery,
|
|
targetType database.TargetType,
|
|
status database.DeliveryStatus,
|
|
) {
|
|
err := e.updateDeliveryStatus(
|
|
webhookDB, d, targetType, status,
|
|
)
|
|
if err != nil {
|
|
e.bookkeepingFailed(d, err)
|
|
}
|
|
}
|
|
|
|
func truncate(s string, maxLen int) string {
|
|
if len(s) <= maxLen {
|
|
return s
|
|
}
|
|
|
|
return s[:maxLen]
|
|
}
|
|
|
|
// --- Helper functions ---
|
|
|
|
func buildEventFromTask(task *Task) database.Event {
|
|
event := database.Event{
|
|
EntrypointID: task.EntrypointID,
|
|
Method: task.Method,
|
|
Headers: task.Headers,
|
|
ContentType: task.ContentType,
|
|
}
|
|
|
|
event.ID = task.EventID
|
|
event.WebhookID = task.WebhookID
|
|
|
|
return event
|
|
}
|
|
|
|
func buildTargetFromTask(task *Task) database.Target {
|
|
target := database.Target{
|
|
Name: task.TargetName,
|
|
Type: task.TargetType,
|
|
Config: task.TargetConfig,
|
|
MaxRetries: task.MaxRetries,
|
|
}
|
|
|
|
target.ID = task.TargetID
|
|
|
|
return target
|
|
}
|
|
|
|
func (e *Engine) resolveEventBody(
|
|
webhookDB *gorm.DB,
|
|
event database.Event,
|
|
task *Task,
|
|
) (database.Event, error) {
|
|
if task.Body != nil {
|
|
event.Body = *task.Body
|
|
|
|
return event, nil
|
|
}
|
|
|
|
var dbEvent database.Event
|
|
|
|
err := webhookDB.Select("body").
|
|
First(&dbEvent, "id = ?", task.EventID).Error
|
|
if err != nil {
|
|
return event, fmt.Errorf(
|
|
"fetching event body: %w", err,
|
|
)
|
|
}
|
|
|
|
event.Body = dbEvent.Body
|
|
|
|
return event, nil
|
|
}
|
|
|
|
func (e *Engine) loadRetryDelivery(
|
|
webhookDB *gorm.DB, deliveryID string,
|
|
) (*database.Delivery, error) {
|
|
var d database.Delivery
|
|
|
|
err := webhookDB.Select("id", "status").
|
|
First(&d, "id = ?", deliveryID).Error
|
|
if err != nil {
|
|
return nil, fmt.Errorf(
|
|
"loading delivery: %w", err,
|
|
)
|
|
}
|
|
|
|
return &d, nil
|
|
}
|
|
|
|
func (e *Engine) countAttempts(
|
|
webhookDB *gorm.DB, deliveryID string,
|
|
) int {
|
|
var resultCount int64
|
|
|
|
webhookDB.Model(&database.DeliveryResult{}).
|
|
Where("delivery_id = ?", deliveryID).
|
|
Count(&resultCount)
|
|
|
|
return int(resultCount)
|
|
}
|
|
|
|
// claimPending takes ownership of a pending delivery before it is
|
|
// re-dispatched, and reports whether the claim succeeded.
|
|
//
|
|
// The claim is a compare-and-set on the status: it takes effect only
|
|
// while the delivery is still pending, so a worker that settled the
|
|
// delivery between the query and here wins and nothing is re-sent.
|
|
// Stamping updated_at is the claim itself — the sweep selects on that
|
|
// column, so a delivery handed out now is out of the sweep's reach for
|
|
// a further pendingSweepMinAge, rather than being sent again on every
|
|
// sweep for as long as the attempt takes.
|
|
//
|
|
// A claim that cannot be written means the database is refusing
|
|
// writes, which is the condition that stranded this delivery in the
|
|
// first place. Not sending is then the right answer: the attempt
|
|
// could not be recorded either, and an unrecordable send is exactly
|
|
// the duplicate this issue is about.
|
|
func (e *Engine) claimPending(
|
|
webhookDB *gorm.DB, d *database.Delivery,
|
|
) bool {
|
|
res := webhookDB.
|
|
Model(&database.Delivery{}).
|
|
Where(
|
|
"id = ? AND status = ?",
|
|
d.ID, database.DeliveryStatusPending,
|
|
).
|
|
UpdateColumn("updated_at", time.Now())
|
|
if res.Error != nil {
|
|
e.log.Error(
|
|
"failed to claim pending delivery for recovery; "+
|
|
"leaving it for the next sweep",
|
|
"delivery_id", d.ID,
|
|
"error", res.Error,
|
|
)
|
|
|
|
return false
|
|
}
|
|
|
|
return res.RowsAffected == 1
|
|
}
|
|
|
|
// countAttemptsBatch counts the recorded attempts of every delivery
|
|
// in a batch with one grouped query, keyed by delivery id. Deliveries
|
|
// with no attempts are simply absent from the result, which reads back
|
|
// as the zero this caller wants.
|
|
//
|
|
// One query rather than one per delivery: this runs on the recovery
|
|
// path, which is a burst of writes against a database that has just
|
|
// been under enough contention to strand these rows in the first
|
|
// place. See https://git.eeqj.de/sneak/webhooker/issues/256.
|
|
func (e *Engine) countAttemptsBatch(
|
|
webhookDB *gorm.DB, deliveries []database.Delivery,
|
|
) map[string]int {
|
|
counts := make(map[string]int, len(deliveries))
|
|
|
|
if len(deliveries) == 0 {
|
|
return counts
|
|
}
|
|
|
|
ids := make([]string, 0, len(deliveries))
|
|
for i := range deliveries {
|
|
ids = append(ids, deliveries[i].ID)
|
|
}
|
|
|
|
// One delivery id per recorded attempt, tallied here rather than
|
|
// grouped in SQL: internal/gormlog forbids (*gorm.DB).Scan, which
|
|
// a GROUP BY into a struct would need, and an attempt row per
|
|
// delivery is bounded by the target's MaxRetries.
|
|
var attemptIDs []string
|
|
|
|
err := webhookDB.
|
|
Model(&database.DeliveryResult{}).
|
|
Where("delivery_id IN ?", ids).
|
|
Pluck("delivery_id", &attemptIDs).Error
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to count delivery attempts for recovery",
|
|
"error", err,
|
|
)
|
|
|
|
return counts
|
|
}
|
|
|
|
for _, id := range attemptIDs {
|
|
counts[id]++
|
|
}
|
|
|
|
return counts
|
|
}
|
|
|
|
func (e *Engine) loadEvent(
|
|
webhookDB *gorm.DB, eventID string,
|
|
) (database.Event, error) {
|
|
var event database.Event
|
|
|
|
err := webhookDB.
|
|
First(&event, "id = ?", eventID).Error
|
|
if err != nil {
|
|
return event, fmt.Errorf(
|
|
"loading event: %w", err,
|
|
)
|
|
}
|
|
|
|
return event, nil
|
|
}
|
|
|
|
func (e *Engine) loadTarget(
|
|
targetID string,
|
|
) (database.Target, error) {
|
|
var target database.Target
|
|
|
|
err := e.database.DB().
|
|
First(&target, "id = ?", targetID).Error
|
|
if err != nil {
|
|
return target, fmt.Errorf(
|
|
"loading target: %w", err,
|
|
)
|
|
}
|
|
|
|
return target, nil
|
|
}
|
|
|
|
func buildRecoveryTask(
|
|
d *database.Delivery,
|
|
webhookID string,
|
|
event *database.Event,
|
|
target *database.Target,
|
|
attemptNum int,
|
|
) Task {
|
|
var bodyPtr *string
|
|
|
|
if len(event.Body) < MaxInlineBodySize {
|
|
bodyStr := event.Body
|
|
bodyPtr = &bodyStr
|
|
}
|
|
|
|
return Task{
|
|
DeliveryID: d.ID,
|
|
EventID: d.EventID,
|
|
WebhookID: webhookID,
|
|
EntrypointID: event.EntrypointID,
|
|
TargetID: target.ID,
|
|
TargetName: target.Name,
|
|
TargetType: target.Type,
|
|
TargetConfig: target.Config,
|
|
MaxRetries: target.MaxRetries,
|
|
Method: event.Method,
|
|
Headers: event.Headers,
|
|
ContentType: event.ContentType,
|
|
Body: bodyPtr,
|
|
AttemptNum: attemptNum,
|
|
}
|
|
}
|
|
|
|
func (e *Engine) loadTargetMap(
|
|
deliveries []database.Delivery,
|
|
) map[string]database.Target {
|
|
seen := make(map[string]bool)
|
|
|
|
targetIDs := make([]string, 0, len(deliveries))
|
|
|
|
for _, d := range deliveries {
|
|
if !seen[d.TargetID] {
|
|
targetIDs = append(targetIDs, d.TargetID)
|
|
seen[d.TargetID] = true
|
|
}
|
|
}
|
|
|
|
var targets []database.Target
|
|
|
|
err := e.database.DB().
|
|
Where("id IN ?", targetIDs).
|
|
Find(&targets).Error
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to load targets from main DB",
|
|
"error", err,
|
|
)
|
|
|
|
return nil
|
|
}
|
|
|
|
targetMap := make(
|
|
map[string]database.Target, len(targets),
|
|
)
|
|
|
|
for _, t := range targets {
|
|
targetMap[t.ID] = t
|
|
}
|
|
|
|
return targetMap
|
|
}
|
|
|
|
// sendRecoveredDeliveries re-dispatches pending deliveries, skipping
|
|
// the ids in settled — those already reached their receiver and have
|
|
// been marked delivered by reconcileDelivered.
|
|
func (e *Engine) sendRecoveredDeliveries(
|
|
ctx context.Context,
|
|
webhookDB *gorm.DB,
|
|
deliveries []database.Delivery,
|
|
webhookID string,
|
|
targetMap map[string]database.Target,
|
|
settled map[string]struct{},
|
|
) {
|
|
// The attempt number continues each delivery's own history
|
|
// rather than restarting at 1. A recovered delivery may already
|
|
// have recorded attempts, and numbering the next one 1 again
|
|
// both collides in the event log and hands the retry path a
|
|
// backoff computed from the wrong attempt.
|
|
attempts := e.countAttemptsBatch(webhookDB, deliveries)
|
|
|
|
for i := range deliveries {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
default:
|
|
}
|
|
|
|
if _, ok := settled[deliveries[i].ID]; ok {
|
|
continue
|
|
}
|
|
|
|
target, ok := targetMap[deliveries[i].TargetID]
|
|
if !ok {
|
|
e.log.Error(
|
|
"target not found for delivery",
|
|
"delivery_id", deliveries[i].ID,
|
|
"target_id", deliveries[i].TargetID,
|
|
)
|
|
|
|
continue
|
|
}
|
|
|
|
if !e.claimPending(webhookDB, &deliveries[i]) {
|
|
continue
|
|
}
|
|
|
|
task := buildRecoveryTask(
|
|
&deliveries[i], webhookID,
|
|
&deliveries[i].Event, &target,
|
|
attempts[deliveries[i].ID]+1,
|
|
)
|
|
|
|
select {
|
|
case e.deliveryCh <- task:
|
|
default:
|
|
e.log.Warn(
|
|
"delivery channel full during "+
|
|
"recovery, remaining deliveries "+
|
|
"will be recovered on next restart",
|
|
"delivery_id", deliveries[i].ID,
|
|
)
|
|
|
|
return
|
|
}
|
|
}
|
|
}
|