All checks were successful
check / check (push) Successful in 5s
The delivery engine worker pool and the retention reaper both rooted their goroutines in the fx OnStart hook context, which fx cancels 15s into startup. Both now use context.WithCancel(context.Background()), bounded by OnStop.
1048 lines
20 KiB
Go
1048 lines
20 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)
|
|
}
|
|
|
|
// 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
|
|
}
|
|
|
|
// 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(),
|
|
})
|
|
|
|
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,
|
|
)
|
|
}
|
|
}
|
|
}
|
|
|
|
// 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.
|
|
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(_ context.Context) error {
|
|
e.stop()
|
|
|
|
return nil
|
|
},
|
|
})
|
|
}
|
|
|
|
// 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.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
|
|
}
|
|
}
|
|
}
|