// Package delivery manages asynchronous event delivery // to configured targets. package delivery import ( "context" "errors" "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 // pendingSweepMinAge is how long a delivery must have sat // untouched at pending before the sweep will look at it. // // It is not what keeps the sweep off live work — inflightSet is, // and it is exact. This bound sets the re-dispatch cadence for a // delivery that really is stranded: without it, a delivery the // database will not let the engine settle would be re-sent on // every 60-second tick. // // It is nonetheless set clear of the longest legitimate attempt, // so that the two guards do not both have to be right. That // length is MaxTargetTimeoutSeconds (300s), the per-target // timeout the target form accepts — not httpClientTimeout, which // is merely the default. Fifteen minutes leaves a margin of // three times the ceiling rather than the zero margin the two // equal values would have given. pendingSweepMinAge = 15 * time.Minute // pendingSweepBatch bounds how many stranded pending deliveries // one sweep of one webhook re-dispatches. The sweep runs every // retrySweepInterval, so a larger backlog drains across // successive sweeps instead of arriving as one burst against a // database that was already struggling to accept writes. pendingSweepBatch = 500 // 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< 0 { return database.AddTotals(tx, database.Totals{Failures: 1}) } return nil } // settleStatus moves a delivery to its outcome status and reports a // failed write through bookkeepingFailed, which leaves the row // recoverable. It exists so the target call sites read as one // statement rather than four lines of identical error handling. func (e *Engine) settleStatus( webhookDB *gorm.DB, d *database.Delivery, targetType database.TargetType, status database.DeliveryStatus, ) { err := e.updateDeliveryStatus( webhookDB, d, targetType, status, ) if err != nil { e.bookkeepingFailed(d, err) } } func truncate(s string, maxLen int) string { if len(s) <= maxLen { return s } return s[:maxLen] } // --- Helper functions --- // buildEventFromTask reconstructs the event a Task describes, as far // as the Task itself goes. The fields it cannot fill — the body when // it was too large to inline, and the receipt time, which no Task // carries — come from the stored row in hydrateEvent, which every // caller of this function runs next. 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 } // hydrateEvent fills in the event fields a Task does not carry, by // reading the stored event row. // // CreatedAt is the event's receipt time and lives only in that row. // The Slack target renders it into every message it sends, so an // unhydrated event puts the zero time in front of a human on every // notification the product delivers. See // https://git.eeqj.de/sneak/webhooker/issues/257. // // The body comes from the same row when the Task did not inline it, // which is the case for a body at or above MaxInlineBodySize. // // A read failure is fatal to the delivery only when the body depended // on it. When the Task inlined the body, the delivery has everything // it needs to be sent and goes ahead with the timestamp unset: the row // can be gone under a retention reap while a queued delivery still // holds its body, and dropping a deliverable event to protect one // metadata field would be a worse failure than the one it prevents. func (e *Engine) hydrateEvent( webhookDB *gorm.DB, event database.Event, task *Task, ) (database.Event, error) { columns := []string{"created_at"} if task.Body == nil { columns = append(columns, "body") } var dbEvent database.Event err := webhookDB.Select(columns). First(&dbEvent, "id = ?", task.EventID).Error if err != nil { if task.Body == nil { return event, fmt.Errorf( "fetching event body: %w", err, ) } e.log.Warn( "could not read the stored event; delivering "+ "the inlined body without its receipt time", "event_id", task.EventID, "delivery_id", task.DeliveryID, "error", err, ) event.Body = *task.Body return event, nil } event.CreatedAt = dbEvent.CreatedAt if task.Body != nil { event.Body = *task.Body } else { event.Body = dbEvent.Body } return event, nil } func (e *Engine) loadDelivery( 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) } // takeForRedispatch decides whether a recovered delivery may be sent // again, and takes it if so. It is the single gate every re-dispatch // path goes through, and it asks two separate questions in order. // // First, does the engine already own this delivery? Ownership is // exact and mutually exclusive, so a delivery queued, being attempted, // or waiting out a retry backoff is refused here, and two dispatchers // racing for the same delivery cannot both win. See inflight.go. // // Second, is the row still in the status that made it eligible? The // batch was read some time ago and a worker may have settled a row // since. The check is a conditional update rather than a read so the // answer cannot go stale between asking and acting. // // Stamping updated_at is the same statement, and it is a cadence // control rather than a claim: the pending sweep selects on that // column, so a delivery handed out now is not selected again on the // next tick a minute later but after pendingSweepMinAge. A delivery // the database refuses to settle is therefore retried on that // interval instead of every tick. // // A failed write is a refusal. It means the database is not accepting // writes, which is the condition that stranded the delivery in the // first place; an attempt that cannot be recorded is exactly the // unlogged duplicate this is all here to prevent. // // The caller must release ownership if it then fails to queue the // task. func (e *Engine) takeForRedispatch( webhookDB *gorm.DB, deliveryID string, eligible database.DeliveryStatus, ) bool { if !e.inflight.retainIdle(deliveryID) { return false } res := webhookDB. Model(&database.Delivery{}). Where( "id = ? AND status = ?", deliveryID, eligible, ). UpdateColumn("updated_at", time.Now()) if res.Error != nil { e.log.Error( "failed to mark delivery for re-dispatch; "+ "leaving it for a later sweep", "delivery_id", deliveryID, "error", res.Error, ) e.inflight.release(deliveryID) return false } if res.RowsAffected != 1 { // Settled underneath us between the query and here. e.inflight.release(deliveryID) return false } return true } // queueRecovered puts an owned delivery's task on a worker channel, // dropping the ownership the gate took if it does not fit. func (e *Engine) queueRecovered( ch chan<- Task, task Task, ) bool { select { case ch <- task: return true default: e.inflight.release(task.DeliveryID) e.log.Warn( "worker channel full during recovery; "+ "delivery will be recovered by a later sweep", "delivery_id", task.DeliveryID, "webhook_id", task.WebhookID, ) return false } } // redispatch hands a recovered delivery to a worker channel through // takeForRedispatch, and reports whether the task was queued. func (e *Engine) redispatch( ch chan<- Task, webhookDB *gorm.DB, task Task, eligible database.DeliveryStatus, ) bool { if !e.takeForRedispatch( webhookDB, task.DeliveryID, eligible, ) { return false } return e.queueRecovered(ch, task) } // rescheduleRecovered hands an orphaned retrying delivery back to the // retry timer, through the same gate. It reports whether the delivery // was rescheduled. // // The reference taken by the gate is dropped as soon as ScheduleRetry // has taken its own, which it does before returning: what keeps the // delivery owned through the backoff window is ScheduleRetry's // reference, not this one. func (e *Engine) rescheduleRecovered( webhookDB *gorm.DB, task Task, delay time.Duration, ) bool { if !e.takeForRedispatch( webhookDB, task.DeliveryID, database.DeliveryStatusRetrying, ) { return false } defer e.inflight.release(task.DeliveryID) e.ScheduleRetry(task, delay) return true } // countAttemptsBatch counts the recorded attempts of every delivery // in a batch with one grouped query, keyed by delivery id. Deliveries // with no attempts are simply absent from the result, which reads back // as the zero this caller wants. // // One query rather than one per delivery: this runs on the recovery // path, which is a burst of writes against a database that has just // been under enough contention to strand these rows in the first // place. See https://git.eeqj.de/sneak/webhooker/issues/256. func (e *Engine) countAttemptsBatch( webhookDB *gorm.DB, deliveries []database.Delivery, ) map[string]int { counts := make(map[string]int, len(deliveries)) if len(deliveries) == 0 { return counts } ids := make([]string, 0, len(deliveries)) for i := range deliveries { ids = append(ids, deliveries[i].ID) } // One delivery id per recorded attempt, tallied here rather than // grouped in SQL: internal/gormlog forbids (*gorm.DB).Scan, which // a GROUP BY into a struct would need, and an attempt row per // delivery is bounded by the target's MaxRetries. var attemptIDs []string err := webhookDB. Model(&database.DeliveryResult{}). Where("delivery_id IN ?", ids). Pluck("delivery_id", &attemptIDs).Error if err != nil { e.log.Error( "failed to count delivery attempts for recovery", "error", err, ) return counts } for _, id := range attemptIDs { counts[id]++ } return counts } 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 { if len(deliveries) == 0 { return nil } 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 } // sendRecoveredDeliveries re-dispatches pending deliveries, skipping // the ids in settled — those already reached their receiver and have // been marked delivered by reconcileDelivered. // // The skip and takeForRedispatch's status check answer different // questions and neither replaces the other. This one is "did this // delivery already succeed", which is what settles the row to // delivered instead of sending it, and which is the only thing that // keeps the retrying paths from terminally failing a delivery that // reconcile just settled. The status check is "is the row still what // the batch query said it was", which catches a worker settling it to // anything at all in between. func (e *Engine) sendRecoveredDeliveries( ctx context.Context, webhookDB *gorm.DB, deliveries []database.Delivery, webhookID string, targetMap map[string]database.Target, settled map[string]struct{}, ) { // The attempt number continues each delivery's own history // rather than restarting at 1. A recovered delivery may already // have recorded attempts, and numbering the next one 1 again // both collides in the event log and hands the retry path a // backoff computed from the wrong attempt. attempts := e.countAttemptsBatch(webhookDB, deliveries) for i := range deliveries { select { case <-ctx.Done(): return default: } if _, ok := settled[deliveries[i].ID]; ok { continue } target, ok := targetMap[deliveries[i].TargetID] if !ok { // A missing entry does not mean the target is gone: the // map is also empty when its query failed. Only a lookup // that finds no row ends the delivery; any other error // leaves it pending for the next sweep. See // recoverSingleRetry. var err error target, err = e.loadTarget(deliveries[i].TargetID) if errors.Is(err, gorm.ErrRecordNotFound) { e.failMissingTarget(webhookDB, webhookID, &deliveries[i]) continue } if err != nil { e.log.Error( "failed to load target for recovered delivery", "delivery_id", deliveries[i].ID, "target_id", deliveries[i].TargetID, "error", err, ) continue } } if !e.takeForRedispatch( webhookDB, deliveries[i].ID, database.DeliveryStatusPending, ) { continue } // The body is read here, one delivery at a time and only for // deliveries that are actually being sent, rather than // preloaded across the whole batch. A batch is up to // pendingSweepBatch rows at up to the 1 MB ingest cap, and // most of a sweep's batch is refused by the gate above — so // preloading would hold hundreds of megabytes per webhook per // tick to build tasks it then discards. event, err := e.loadEvent( webhookDB, deliveries[i].EventID, ) if err != nil { e.log.Error( "failed to load event for recovered delivery", "delivery_id", deliveries[i].ID, "event_id", deliveries[i].EventID, "error", err, ) e.inflight.release(deliveries[i].ID) continue } task := buildRecoveryTask( &deliveries[i], webhookID, &event, &target, attempts[deliveries[i].ID]+1, ) e.queueRecovered(e.deliveryCh, task) } }