All checks were successful
check / check (push) Successful in 2m56s
The per-webhook archiveWriter registry in the database delivery target was never evicted, so a deleted webhook's writer -- and any archive file handle open within its debounce window -- lingered for the process lifetime. Separately, expiry pruning ran only when an archive was (re)opened, and reopens only happen on writes, so an archive belonging to a webhook that stopped receiving events kept its expired rows forever. Eviction: a new one-method delivery.WebhookEvictor interface (kept separate from Notifier: archiving lifecycle is not notification) is implemented by the Engine and injected into the handlers. Deleting a webhook, or deleting its last database target, drops the writer from the registry and closes its handle under the writer's own mutex, so eviction can never race an in-flight write. An evicted writer refuses further writes rather than reopening a file nothing holds. The archive file is deliberately left on disk: it is long-term storage an operator may want to keep or move away, and destroying it as a side effect of deleting a webhook would be unrecoverable. Idle sweep: a new ArchiveSweeper, modelled on the event RetentionReaper (fx lifecycle hooks, cancellable context, WaitGroup, ticker loop), prunes archives whose database target declares a positive expiry. It reuses the existing RETENTION_SWEEP_INTERVAL rather than adding a config key. It never creates an archive -- a missing file is skipped, and the reopen uses SQLite mode=rw so the file cannot be conjured even if it disappears mid-sweep -- routes the prune through the per-webhook writer so its mutex orders the sweep against concurrent writes, and leaves the archive closed so the move-the-file-away workflow keeps working. A failure for one webhook is logged and the sweep continues. Archives with no expiry or the expiry "never" are untouched. The sweep loop's context is rooted at context.Background(), not at the fx OnStart hook context. The hook context carries fx's 15 second start timeout, so a loop derived from it is cancelled three quarters of an hour before the first tick under the default one-hour interval, giving a sweeper that never sweeps. OnStop still cancels the loop and waits on the WaitGroup, so shutdown is unchanged. The sweep also never leaves a registry entry behind. Reaching the writer through the ordinary create-and-cache accessor would let a sweep that raced a webhook deletion re-insert a writer for a webhook that no longer exists, which nothing would ever evict again -- the very leak this change closes. An entry the sweep has to create is marked sweep-owned and released when the prune finishes, unless a delivery claimed it meanwhile, in which case it belongs to the registry and an eviction can still reach it. A writer evicted underneath a sweep is an ordinary interleaving and is logged at debug, not error.
1061 lines
21 KiB
Go
1061 lines
21 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/logger"
|
|
)
|
|
|
|
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
|
|
|
|
// 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
|
|
}
|
|
|
|
// 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
|
|
|
|
// 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,
|
|
}
|
|
|
|
e.initTargets(&http.Client{
|
|
Timeout: httpClientTimeout,
|
|
Transport: NewSSRFSafeTransport(),
|
|
})
|
|
|
|
lc.Append(fx.Hook{
|
|
OnStart: func(ctx context.Context) error {
|
|
e.start(ctx)
|
|
|
|
return nil
|
|
},
|
|
OnStop: func(_ context.Context) error {
|
|
e.stop()
|
|
|
|
return nil
|
|
},
|
|
})
|
|
|
|
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,
|
|
)
|
|
}
|
|
})
|
|
}
|
|
|
|
func (e *Engine) start(ctx context.Context) {
|
|
ctx, cancel := context.WithCancel(ctx)
|
|
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.log.Info(
|
|
"delivery engine started",
|
|
"workers", e.workers,
|
|
)
|
|
}
|
|
|
|
func (e *Engine) stop() {
|
|
e.log.Info("delivery engine stopping")
|
|
e.cancel()
|
|
e.wg.Wait()
|
|
e.log.Info("delivery engine stopped")
|
|
}
|
|
|
|
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
|
|
// they are skipped.
|
|
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 {
|
|
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),
|
|
)
|
|
|
|
targetMap := e.loadTargetMap(deliveries)
|
|
|
|
e.sendRecoveredDeliveries(
|
|
ctx, deliveries, webhookID, targetMap,
|
|
)
|
|
}
|
|
|
|
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],
|
|
)
|
|
}
|
|
}
|
|
|
|
// sweepSingleRetry re-enqueues an orphaned retrying delivery
|
|
// whose backoff window has elapsed, delegating the backoff
|
|
// decision to the delivery's target. Targets that do not own
|
|
// durable retries are skipped.
|
|
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 {
|
|
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:
|
|
}
|
|
}
|
|
|
|
// 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.updateDeliveryStatus(
|
|
webhookDB, d, database.DeliveryStatusFailed,
|
|
)
|
|
|
|
return
|
|
}
|
|
|
|
target.Deliver(ctx, webhookDB, d, task, e)
|
|
}
|
|
|
|
// recordResult persists a DeliveryResult row describing a
|
|
// single attempt. It is a cross-target helper the targets
|
|
// call.
|
|
func (e *Engine) recordResult(
|
|
webhookDB *gorm.DB,
|
|
d *database.Delivery,
|
|
attemptNum int,
|
|
success bool,
|
|
statusCode int,
|
|
respBody, errMsg string,
|
|
durationMs int64,
|
|
) {
|
|
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 {
|
|
e.log.Error(
|
|
"failed to record delivery result",
|
|
"delivery_id", d.ID,
|
|
"error", err,
|
|
)
|
|
}
|
|
}
|
|
|
|
// updateDeliveryStatus persists a new status for a delivery.
|
|
// It is a cross-target helper the targets call.
|
|
func (e *Engine) updateDeliveryStatus(
|
|
webhookDB *gorm.DB,
|
|
d *database.Delivery,
|
|
status database.DeliveryStatus,
|
|
) {
|
|
err := webhookDB.Model(d).
|
|
Update("status", status).Error
|
|
if err != nil {
|
|
e.log.Error(
|
|
"failed to update delivery status",
|
|
"delivery_id", d.ID,
|
|
"status", status,
|
|
"error", 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)
|
|
}
|
|
|
|
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
|
|
}
|
|
|
|
func (e *Engine) sendRecoveredDeliveries(
|
|
ctx context.Context,
|
|
deliveries []database.Delivery,
|
|
webhookID string,
|
|
targetMap map[string]database.Target,
|
|
) {
|
|
for i := range deliveries {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
default:
|
|
}
|
|
|
|
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
|
|
}
|
|
|
|
task := buildRecoveryTask(
|
|
&deliveries[i], webhookID,
|
|
&deliveries[i].Event, &target, 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
|
|
}
|
|
}
|
|
}
|