Files
webhooker/internal/delivery/engine.go
clawbot a6a306d810
All checks were successful
check / check (push) Successful in 2m47s
Evict archive writers on deletion and sweep idle archives (closes #89)
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.

Each of those guards is pinned by a test that fails when the guard is
removed. The adopt-during-sweep window -- a delivery claiming the
sweep's own registry entry while that sweep is still running -- is
driven directly against the registry, because a delivery placed between
two sweeps never reaches the release path at all. The requirement that
an idle archive ends the sweep closed is asserted both on a writer
proven to hold an open handle beforehand and, end to end, on a
delivery-owned entry the sweep keeps, rather than on an entry the sweep
has already released and which therefore reports "not open" either
way. The "never" expiry short circuit is checked against an archive
file that has never been migrated, so any open of it would be
observable as a created table.
2026-08-09 05:22:49 +00:00

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
}
}
}