Files
webhooker/internal/delivery/engine.go
clawbot c4022c0834
All checks were successful
check / check (push) Successful in 3m12s
Correct release-blocking documentation inaccuracies (closes #141)
Four defects found by the integration review, each of which would have
made the README or the release notes untrue at the moment of tagging.

RETENTION_SWEEP_INTERVAL was absent from the README env table while two
other passages referred to it as documented. Enumerated every variable
read by internal/config from the source (12 in total) rather than by
eye; that was the only one missing.

TODO.md omitted five of the units landed in this milestone (#64, #79,
#90, #113, #118), two of them credential-exposure fixes, which are
precisely the entries a reader of the release notes wants to find. The
list is now derived from git log origin/main..origin/next.

The README sold Replay in the present tense as a core capability while
no redelivery code exists anywhere in the tree, and the roadmap entry
for it had been dropped without it being implemented. Every mention of
replay or redelivery in the README is now either marked planned or
already under a Planned heading: the Rationale item, the Use Cases
bullet, the Event model's field description, and the delivery-semantics
passage on target-type edits, which told an operator that a terminally
failed delivery was recoverable by hand when nothing can recover it.
The comment on failUnretryableRetry made the same claim and is
corrected with it. The roadmap entry is back in TODO.md.

The production JS asset shipped a console.log on load. The rest of the
file and every other shipped asset were checked; that was the only one
(alpine.min.js is vendored and untouched).

No behavioural change: the only non-comment, non-documentation edit is
the deleted console.log.
2026-08-12 11:01:29 +00:00

1144 lines
24 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(),
})
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.
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 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,
)
e.updateDeliveryStatus(
webhookDB, d, 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, 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
}
}
}