Files
webhooker/internal/delivery/engine.go
sneak 7cc2e201ab
All checks were successful
check / check (push) Successful in 3m4s
Add an egress CIDR allowlist to the SSRF guard (closes #204)
The SSRF blocklist had no escape hatch, so the thing webhooker is
mostly for — taking a public webhook and forwarding it to something
on your own network — could not be configured at all. Every private
address, Docker sibling and loopback service was permanently
unreachable as a delivery destination.

ALLOWED_EGRESS_CIDRS (default empty) names blocks that delivery
targets may reach despite the default blocklist. It is an allowlist
and only ever adds destinations: there is no boolean, and no value
disables SSRF protection wholesale. Empty, it adds nothing and the
guard permits and refuses the same addresses it did before, save
the two spellings named below.

A small set of addresses is refused before the allowlist is
consulted, so no supplied CIDR opens one — not the exact address,
not a supernet, not 0.0.0.0/0 or ::/0. alwaysBlockedNetworks in
internal/delivery/ssrf.go is the authoritative list and states the
membership criterion in full; it is deliberately not copied here,
because a copy drifts out of date. In short: the provider fixes the
address, so a host route for it collides with nothing the operator
runs, and reaching it discloses credentials or user data. A
publicly routable address never qualifies however well it meets
both — nothing in this set can be reopened, so blocking one here
would leave the operator no escape hatch at all. Those belong in
blockedNetworks, which an allowlist can override.

Every entry is already inside the default blocklist, which is what
makes it unconditional rather than newly blocked, with two
exceptions that are the one behaviour change visible when the
allowlist is unset: ::a9fe:a9fe and 64:ff9b::a9fe:a9fe, the
IPv4-compatible and NAT64 spellings of 169.254.169.254, were
reachable before and are refused now. net.IPNet.Contains normalises
only the IPv4-mapped form via To4(), so 169.254.0.0/16 never
matched those two. The ten it does cover now report a metadata
error rather than the generic private-range one.

The policy now lives in one function, Guard.checkIP, which both
target-creation validation and the delivery dialer call. The two
paths previously decided separately, which is how they came to
disagree about a destination. The guard is built once from config
and injected via fx into both the handlers and the delivery engine,
so there is a single instance and a single answer.

A set-but-unparseable value aborts startup naming the variable,
reusing the existing envPrefixList parser. A non-empty list is
logged at startup with the blocks spelled out, not counted, so the
hole is visible in the log of any deployment that has one.

Tests: an allowlisted loopback CIDR both validates and delivers to a
live server (and the same URL still fails without the allowlist); a
private address outside the listed block stays refused on both
paths; every unconditionally blocked address stays refused on both
paths under an allowlist that covers it, and the set itself is
pinned entry by entry; public addresses are unaffected either way;
and config coverage for parsing, startup abort, and the warning's
contents.
2026-08-20 08:09:12 +00:00

1214 lines
26 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
// 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),
)
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. 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,
)
e.recordResult(
webhookDB,
d,
e.countAttempts(webhookDB, d.ID)+1,
false,
0,
"",
reason,
0,
)
// 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.updateDeliveryStatus(
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.updateDeliveryStatus(
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.
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, 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.
func (e *Engine) updateDeliveryStatus(
webhookDB *gorm.DB,
d *database.Delivery,
targetType database.TargetType,
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,
)
return
}
e.mtr.DeliveryStatusChanged(targetType, status)
}
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
}
}
}