Make SQLite durable under concurrent readers and stop re-delivering stranded webhooks (closes #256)
All checks were successful
check / check (push) Successful in 3m3s
All checks were successful
check / check (push) Successful in 3m3s
An operator running `sqlite3 <db> .dump` against their own per-webhook database wedged it: inbound webhooks rejected with HTTP 500, delivered webhooks stranded at `pending`, and every one of them POSTed a second time on the next restart while the event log recorded a single attempt. Durability. Every SQLite file — main, per-webhook, and archive — now opens through one path, `internal/database/sqlite_open.go`, in WAL journal mode with a 10-second busy timeout, `BEGIN IMMEDIATE` transactions, and a bounded connection pool. WAL is what stops a reader blocking writers at all. `_txlock=immediate` is what stops a `COMMIT` failing while its transaction stays open on a pooled connection, which is how four `database is locked` errors became 593 `cannot start a transaction within a transaction`. `cache=shared` is gone, because under it an in-process conflict is SQLITE_LOCKED, which the busy handler does not retry. The busy timeout is applied before journal_mode: the driver runs DSN pragmas in order on every new connection, and `PRAGMA journal_mode` takes a lock, so the reverse order leaves the one pragma that can block uncovered by the handler meant to cover it. Eligibility. `internal/delivery/inflight.go` holds the set of deliveries the engine owns — taken when a task is queued, when a target schedules a retry, and by every recovery path before it re-dispatches; dropped when the worker that ran the task returns. Recovery and both sweep arms re-dispatch only what the set does not hold. Nothing decides that from a row's age: a delivery waiting in a 10000-deep channel is arbitrarily old and perfectly healthy, and reasoning from age re-sends it. `takeForRedispatch` is the single gate every re-dispatch goes through — ownership first, then a conditional update confirming the row is still in the status the batch read. Bookkeeping. `recordResult` and `updateDeliveryStatus` return their errors instead of logging and dropping them, and a caller whose bookkeeping write failed writes nothing at all: the delivery keeps whichever non-terminal status it already held, and the sweeps recover it. Every recovery path — pending and retrying alike — first settles any delivery that already holds a successful `DeliveryResult` rather than sending it again. Recovery continues each delivery's own attempt numbering instead of restarting at 1. The sweep gains a `pending`-with-age-bound arm, so a stranded delivery no longer waits for a restart. Docs. WAL produces `-wal`/`-shm` sidecars, so the backup and restore procedures in README.md are corrected against measurement: both documented procedures were re-run against a live instance, a `-wal` left by a crash carries data the `.db` alone does not, and an archive file normally holds its rows in a `-wal` rather than in the `.db`.
This commit is contained in:
@@ -41,6 +41,31 @@ const (
|
||||
// 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
|
||||
@@ -157,6 +182,12 @@ type Engine struct {
|
||||
// dbTarget is retained so the engine can reach the archive
|
||||
// writer registry for webhook eviction and the idle sweep.
|
||||
dbTarget *databaseTarget
|
||||
|
||||
// inflight is the set of deliveries this engine currently owns.
|
||||
// Recovery and the sweeps re-dispatch only what it does not
|
||||
// hold. Held by value: its zero value works, so no constructor
|
||||
// can leave it out. See inflight.go.
|
||||
inflight inflightSet
|
||||
}
|
||||
|
||||
// New creates and registers the delivery engine with the
|
||||
@@ -189,12 +220,27 @@ func New(
|
||||
// are ready.
|
||||
func (e *Engine) Notify(tasks []Task) {
|
||||
for i := range tasks {
|
||||
// Owned before it is queued, and until the worker that runs
|
||||
// it returns. A task can sit in a 10000-deep channel for a
|
||||
// long time on a healthy system, and nothing may re-send it
|
||||
// while it waits. See inflight.go.
|
||||
if !e.inflight.retainIdle(tasks[i].DeliveryID) {
|
||||
e.log.Warn(
|
||||
"delivery already in flight, not queued again",
|
||||
"delivery_id", tasks[i].DeliveryID,
|
||||
"event_id", tasks[i].EventID,
|
||||
)
|
||||
|
||||
continue
|
||||
}
|
||||
|
||||
select {
|
||||
case e.deliveryCh <- tasks[i]:
|
||||
default:
|
||||
e.inflight.release(tasks[i].DeliveryID)
|
||||
e.log.Warn(
|
||||
"delivery channel full, "+
|
||||
"task will be recovered on restart",
|
||||
"task will be recovered by the sweep",
|
||||
"delivery_id", tasks[i].DeliveryID,
|
||||
"event_id", tasks[i].EventID,
|
||||
)
|
||||
@@ -229,10 +275,20 @@ func (e *Engine) ScheduleRetry(
|
||||
"next_attempt", task.AttemptNum,
|
||||
)
|
||||
|
||||
// The reference is taken here rather than when the timer fires,
|
||||
// so the delivery stays owned across the whole backoff window.
|
||||
// Its caller is a target inside Deliver, so the engine already
|
||||
// owns it; this second reference is what keeps that ownership
|
||||
// alive after the worker returns and the row sits at retrying
|
||||
// with nothing running. Without it the sweep finds the row
|
||||
// orphaned and sends it again.
|
||||
e.inflight.retain(task.DeliveryID)
|
||||
|
||||
time.AfterFunc(delay, func() {
|
||||
select {
|
||||
case e.retryCh <- task:
|
||||
default:
|
||||
e.inflight.release(task.DeliveryID)
|
||||
e.log.Warn(
|
||||
"retry channel full, delivery "+
|
||||
"will be recovered by periodic sweep",
|
||||
@@ -332,13 +388,35 @@ func (e *Engine) worker(ctx context.Context) {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case task := <-e.deliveryCh:
|
||||
e.processNewTask(ctx, &task)
|
||||
e.runTask(ctx, &task, e.processNewTask)
|
||||
case task := <-e.retryCh:
|
||||
e.processRetryTask(ctx, &task)
|
||||
e.runTask(ctx, &task, e.processRetryTask)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// runTask runs one task and then drops the reference the queueing
|
||||
// side took on its delivery.
|
||||
//
|
||||
// The release is deferred rather than written after the call because
|
||||
// every early return inside the processing paths must drop it too: a
|
||||
// delivery whose database could not be opened is one the engine has
|
||||
// stopped working on, and leaving it owned would hide it from the
|
||||
// sweep forever.
|
||||
//
|
||||
// Ownership does not necessarily end here. A target that scheduled a
|
||||
// retry took its own reference before this one is dropped, so the
|
||||
// delivery stays owned through the backoff window.
|
||||
func (e *Engine) runTask(
|
||||
ctx context.Context,
|
||||
task *Task,
|
||||
run func(context.Context, *Task),
|
||||
) {
|
||||
defer e.inflight.release(task.DeliveryID)
|
||||
|
||||
run(ctx, task)
|
||||
}
|
||||
|
||||
func (e *Engine) recoverPending(ctx context.Context) {
|
||||
defer e.wg.Done()
|
||||
|
||||
@@ -526,7 +604,16 @@ func (e *Engine) recoverRetryingDeliveries(
|
||||
return
|
||||
}
|
||||
|
||||
settled := e.reconcileDelivered(
|
||||
webhookDB, webhookID, retrying,
|
||||
e.loadTargetMap(retrying),
|
||||
)
|
||||
|
||||
for i := range retrying {
|
||||
if _, ok := settled[retrying[i].ID]; ok {
|
||||
continue
|
||||
}
|
||||
|
||||
e.recoverSingleRetry(
|
||||
webhookDB, webhookID, &retrying[i],
|
||||
)
|
||||
@@ -588,6 +675,10 @@ func (e *Engine) recoverSingleRetry(
|
||||
d, webhookID, &event, &target, attemptNum+1,
|
||||
)
|
||||
|
||||
if !e.rescheduleRecovered(webhookDB, task, remaining) {
|
||||
return
|
||||
}
|
||||
|
||||
e.log.Info(
|
||||
"recovering retrying delivery",
|
||||
"webhook_id", webhookID,
|
||||
@@ -595,8 +686,6 @@ func (e *Engine) recoverSingleRetry(
|
||||
"attempt", attemptNum,
|
||||
"remaining_backoff", remaining,
|
||||
)
|
||||
|
||||
e.ScheduleRetry(task, remaining)
|
||||
}
|
||||
|
||||
func (e *Engine) recoverPendingDeliveries(
|
||||
@@ -606,12 +695,14 @@ func (e *Engine) recoverPendingDeliveries(
|
||||
) {
|
||||
var deliveries []database.Delivery
|
||||
|
||||
// No Preload: event bodies are read one at a time in
|
||||
// sendRecoveredDeliveries, and only for the deliveries actually
|
||||
// being sent.
|
||||
result := webhookDB.
|
||||
Where(
|
||||
"status = ?",
|
||||
database.DeliveryStatusPending,
|
||||
).
|
||||
Preload("Event").
|
||||
Find(&deliveries)
|
||||
|
||||
if result.Error != nil {
|
||||
@@ -634,11 +725,137 @@ func (e *Engine) recoverPendingDeliveries(
|
||||
"count", len(deliveries),
|
||||
)
|
||||
|
||||
e.recoverPendingBatch(
|
||||
ctx, webhookDB, webhookID, deliveries,
|
||||
)
|
||||
}
|
||||
|
||||
// recoverPendingBatch settles every delivery in the batch that was
|
||||
// already delivered, and re-dispatches only the rest. Both the
|
||||
// restart-time recovery and the periodic sweep go through it, so a
|
||||
// pending delivery is treated the same however it was found.
|
||||
func (e *Engine) recoverPendingBatch(
|
||||
ctx context.Context,
|
||||
webhookDB *gorm.DB,
|
||||
webhookID string,
|
||||
deliveries []database.Delivery,
|
||||
) {
|
||||
targetMap := e.loadTargetMap(deliveries)
|
||||
|
||||
e.sendRecoveredDeliveries(
|
||||
ctx, deliveries, webhookID, targetMap,
|
||||
settled := e.reconcileDelivered(
|
||||
webhookDB, webhookID, deliveries, targetMap,
|
||||
)
|
||||
|
||||
e.sendRecoveredDeliveries(
|
||||
ctx, webhookDB, deliveries, webhookID,
|
||||
targetMap, settled,
|
||||
)
|
||||
}
|
||||
|
||||
// reconcileDelivered finds the deliveries in a recovered batch that
|
||||
// already have a successful DeliveryResult, marks them delivered, and
|
||||
// returns their ids so the caller does not send them a second time.
|
||||
//
|
||||
// This is the state the engine previously had no way to represent. A
|
||||
// delivery is left in a non-terminal state by a failed bookkeeping
|
||||
// write, and that covers two different histories: nothing was ever
|
||||
// sent, or the send reached the receiver and only the status write
|
||||
// failed. Re-sending was the sole option, so every stranded row
|
||||
// produced a duplicate at the receiver and an event log that recorded
|
||||
// one attempt for two POSTs. A successful result row distinguishes
|
||||
// them: it is written before the status, so its presence means the
|
||||
// wire I/O happened and was recorded, and all that is missing is the
|
||||
// status.
|
||||
//
|
||||
// Every recovery path runs this, not only the pending one. A delivery
|
||||
// abandoned at retrying can hold a successful result just as a pending
|
||||
// one can — a second attempt that reached the receiver and whose status
|
||||
// write then failed sits at retrying with success recorded — and
|
||||
// re-sending it is the same duplicate.
|
||||
//
|
||||
// Deliveries whose result row itself never landed are not in the
|
||||
// returned set and are re-sent, recorded as the further attempt they
|
||||
// are. That is honest at-least-once delivery rather than a silent
|
||||
// duplicate.
|
||||
func (e *Engine) reconcileDelivered(
|
||||
webhookDB *gorm.DB,
|
||||
webhookID string,
|
||||
deliveries []database.Delivery,
|
||||
targetMap map[string]database.Target,
|
||||
) map[string]struct{} {
|
||||
settled := make(map[string]struct{})
|
||||
|
||||
if len(deliveries) == 0 {
|
||||
return settled
|
||||
}
|
||||
|
||||
ids := make([]string, 0, len(deliveries))
|
||||
for i := range deliveries {
|
||||
ids = append(ids, deliveries[i].ID)
|
||||
}
|
||||
|
||||
var deliveredIDs []string
|
||||
|
||||
err := webhookDB.
|
||||
Model(&database.DeliveryResult{}).
|
||||
Where(
|
||||
"delivery_id IN ? AND success = ?", ids, true,
|
||||
).
|
||||
Distinct().
|
||||
Pluck("delivery_id", &deliveredIDs).Error
|
||||
if err != nil {
|
||||
// Every delivery stays out of the settled set, so the batch
|
||||
// is re-sent exactly as it was before this check existed.
|
||||
// That is the safe direction: a duplicate delivery beats
|
||||
// declaring a delivery successful on a query that failed.
|
||||
e.log.Error(
|
||||
"failed to query successful delivery results; "+
|
||||
"pending deliveries will be re-sent",
|
||||
"webhook_id", webhookID,
|
||||
"error", err,
|
||||
)
|
||||
|
||||
return settled
|
||||
}
|
||||
|
||||
for _, id := range deliveredIDs {
|
||||
settled[id] = struct{}{}
|
||||
}
|
||||
|
||||
if len(settled) == 0 {
|
||||
return settled
|
||||
}
|
||||
|
||||
e.log.Info(
|
||||
"settling recovered deliveries that already succeeded",
|
||||
"webhook_id", webhookID,
|
||||
"count", len(settled),
|
||||
)
|
||||
|
||||
for i := range deliveries {
|
||||
if _, ok := settled[deliveries[i].ID]; !ok {
|
||||
continue
|
||||
}
|
||||
|
||||
// A delivery the engine is working on right now settles
|
||||
// itself; writing over it from here would race that worker.
|
||||
if !e.inflight.retainIdle(deliveries[i].ID) {
|
||||
delete(settled, deliveries[i].ID)
|
||||
|
||||
continue
|
||||
}
|
||||
|
||||
e.settleStatus(
|
||||
webhookDB,
|
||||
&deliveries[i],
|
||||
targetMap[deliveries[i].TargetID].Type,
|
||||
database.DeliveryStatusDelivered,
|
||||
)
|
||||
|
||||
e.inflight.release(deliveries[i].ID)
|
||||
}
|
||||
|
||||
return settled
|
||||
}
|
||||
|
||||
func (e *Engine) retrySweep(ctx context.Context) {
|
||||
@@ -722,6 +939,11 @@ func (e *Engine) sweepWebhookRetries(
|
||||
return
|
||||
}
|
||||
|
||||
settled := e.reconcileDelivered(
|
||||
webhookDB, webhookID, retrying,
|
||||
e.loadTargetMap(retrying),
|
||||
)
|
||||
|
||||
for i := range retrying {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
@@ -729,10 +951,70 @@ func (e *Engine) sweepWebhookRetries(
|
||||
default:
|
||||
}
|
||||
|
||||
if _, ok := settled[retrying[i].ID]; ok {
|
||||
continue
|
||||
}
|
||||
|
||||
e.sweepSingleRetry(
|
||||
webhookDB, webhookID, &retrying[i],
|
||||
)
|
||||
}
|
||||
|
||||
e.sweepWebhookPending(ctx, webhookDB, webhookID)
|
||||
}
|
||||
|
||||
// sweepWebhookPending recovers deliveries stranded at pending.
|
||||
//
|
||||
// A delivery is created pending and leaves that state only when its
|
||||
// outcome is written, so a pending row the engine does not own is one
|
||||
// whose bookkeeping write failed — the state that used to sit there
|
||||
// until a restart, and then produce a duplicate at the receiver. The
|
||||
// sweep gives it the same reconcile-then-dispatch treatment restart
|
||||
// recovery gets, so it costs a minute rather than an operator
|
||||
// noticing.
|
||||
//
|
||||
// What keeps the sweep off live work is ownership, checked per
|
||||
// delivery in takeForRedispatch, not the age bound in this query.
|
||||
// A delivery waiting in deliveryCh is pending and arbitrarily old —
|
||||
// the channel holds 10000 tasks and 10 workers drain it — so
|
||||
// reasoning from the row's age alone re-sends it. See inflight.go.
|
||||
func (e *Engine) sweepWebhookPending(
|
||||
ctx context.Context,
|
||||
webhookDB *gorm.DB,
|
||||
webhookID string,
|
||||
) {
|
||||
var pending []database.Delivery
|
||||
|
||||
err := webhookDB.
|
||||
Where(
|
||||
"status = ? AND updated_at < ?",
|
||||
database.DeliveryStatusPending,
|
||||
time.Now().Add(-pendingSweepMinAge),
|
||||
).
|
||||
Limit(pendingSweepBatch).
|
||||
Find(&pending).Error
|
||||
if err != nil {
|
||||
e.log.Error(
|
||||
"retry sweep: "+
|
||||
"failed to query pending deliveries",
|
||||
"webhook_id", webhookID,
|
||||
"error", err,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if len(pending) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
e.log.Info(
|
||||
"retry sweep: recovering stranded pending deliveries",
|
||||
"webhook_id", webhookID,
|
||||
"count", len(pending),
|
||||
)
|
||||
|
||||
e.recoverPendingBatch(ctx, webhookDB, webhookID, pending)
|
||||
}
|
||||
|
||||
// sweepSingleRetry re-enqueues an orphaned retrying delivery
|
||||
@@ -789,17 +1071,20 @@ func (e *Engine) sweepSingleRetry(
|
||||
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:
|
||||
if !e.redispatch(
|
||||
e.retryCh, webhookDB, task,
|
||||
database.DeliveryStatusRetrying,
|
||||
) {
|
||||
return
|
||||
}
|
||||
|
||||
e.log.Info(
|
||||
"retry sweep: "+
|
||||
"recovered orphaned retrying delivery",
|
||||
"delivery_id", d.ID,
|
||||
"webhook_id", webhookID,
|
||||
"attempt", attemptNum+1,
|
||||
)
|
||||
}
|
||||
|
||||
// failUnretryableRetry terminally fails an orphaned retrying
|
||||
@@ -822,6 +1107,16 @@ func (e *Engine) failUnretryableRetry(
|
||||
d *database.Delivery,
|
||||
target *database.Target,
|
||||
) {
|
||||
// Terminal, and reached from the recovery paths, so it takes
|
||||
// ownership like every other write they make: a delivery the
|
||||
// engine is still attempting must not be failed underneath the
|
||||
// worker running it.
|
||||
if !e.inflight.retainIdle(d.ID) {
|
||||
return
|
||||
}
|
||||
|
||||
defer e.inflight.release(d.ID)
|
||||
|
||||
e.log.Warn(
|
||||
"failing orphaned retrying delivery: target "+
|
||||
"type no longer supports retries",
|
||||
@@ -839,7 +1134,7 @@ func (e *Engine) failUnretryableRetry(
|
||||
target.Type,
|
||||
)
|
||||
|
||||
e.recordResult(
|
||||
err := e.recordResult(
|
||||
webhookDB,
|
||||
d,
|
||||
e.countAttempts(webhookDB, d.ID)+1,
|
||||
@@ -849,6 +1144,11 @@ func (e *Engine) failUnretryableRetry(
|
||||
reason,
|
||||
0,
|
||||
)
|
||||
if err != nil {
|
||||
e.bookkeepingFailed(d, err)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
// The type is passed rather than assigned onto d: the delivery
|
||||
// is loaded here without its target relation, and populating
|
||||
@@ -856,7 +1156,7 @@ func (e *Engine) failUnretryableRetry(
|
||||
// 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(
|
||||
e.settleStatus(
|
||||
webhookDB, d, target.Type,
|
||||
database.DeliveryStatusFailed,
|
||||
)
|
||||
@@ -878,7 +1178,7 @@ func (e *Engine) processDelivery(
|
||||
"type", d.Target.Type,
|
||||
)
|
||||
|
||||
e.updateDeliveryStatus(
|
||||
e.settleStatus(
|
||||
webhookDB, d, d.Target.Type,
|
||||
database.DeliveryStatusFailed,
|
||||
)
|
||||
@@ -910,6 +1210,14 @@ func (e *Engine) observeAttempt(
|
||||
// recordResult persists a DeliveryResult row describing a
|
||||
// single attempt. It is a cross-target helper the targets
|
||||
// call.
|
||||
//
|
||||
// It returns its error rather than swallowing it. A DeliveryResult
|
||||
// row is the only record that an attempt happened at all, so a
|
||||
// caller that ignored a failed write would go on to mark the
|
||||
// delivery delivered — leaving the event log claiming one attempt
|
||||
// for a receiver that got two. Every caller must instead stop
|
||||
// advancing the delivery's status and let it stay in the
|
||||
// non-terminal state it already holds; see bookkeepingFailed.
|
||||
func (e *Engine) recordResult(
|
||||
webhookDB *gorm.DB,
|
||||
d *database.Delivery,
|
||||
@@ -918,7 +1226,7 @@ func (e *Engine) recordResult(
|
||||
statusCode int,
|
||||
respBody, errMsg string,
|
||||
durationMs int64,
|
||||
) {
|
||||
) error {
|
||||
result := &database.DeliveryResult{
|
||||
DeliveryID: d.ID,
|
||||
AttemptNum: attemptNum,
|
||||
@@ -931,12 +1239,44 @@ func (e *Engine) recordResult(
|
||||
|
||||
err := webhookDB.Create(result).Error
|
||||
if err != nil {
|
||||
e.log.Error(
|
||||
"failed to record delivery result",
|
||||
"delivery_id", d.ID,
|
||||
"error", err,
|
||||
return fmt.Errorf(
|
||||
"recording delivery result for %s: %w", d.ID, err,
|
||||
)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// bookkeepingFailed reports that a delivery's own record of what
|
||||
// happened could not be written, and deliberately writes nothing in
|
||||
// response.
|
||||
//
|
||||
// Leaving the row alone is the whole point. A delivery is created
|
||||
// pending and only ever leaves that state through
|
||||
// updateDeliveryStatus, so a delivery whose bookkeeping write failed
|
||||
// is still pending or retrying — the two non-terminal states, per
|
||||
// DeliveryStatus.Terminal — and both are swept and recovered. Writing
|
||||
// anything here would need the very database that just refused a
|
||||
// write, and would be one more thing to fail; not writing cannot.
|
||||
//
|
||||
// The cost is honest at-least-once behaviour: a send that reached the
|
||||
// receiver but whose result row did not land is attempted again, and
|
||||
// recorded as the further attempt it is. What no longer happens is the
|
||||
// silent duplicate — a second POST the event log denies ever
|
||||
// occurred. Where the result row *did* land and only the status write
|
||||
// failed, reconcileDelivered settles the row without re-sending.
|
||||
func (e *Engine) bookkeepingFailed(
|
||||
d *database.Delivery, err error,
|
||||
) {
|
||||
e.log.Error(
|
||||
"delivery bookkeeping write failed; leaving delivery "+
|
||||
"in a recoverable state",
|
||||
"delivery_id", d.ID,
|
||||
"event_id", d.EventID,
|
||||
"target_id", d.TargetID,
|
||||
"status", d.Status,
|
||||
"error", err,
|
||||
)
|
||||
}
|
||||
|
||||
// updateDeliveryStatus persists a new status for a delivery.
|
||||
@@ -952,26 +1292,51 @@ func (e *Engine) recordResult(
|
||||
//
|
||||
// The counter moves only after the row is written, so a transition
|
||||
// the database rejected is not claimed as an outcome that happened.
|
||||
// For the same reason the error is returned rather than logged and
|
||||
// dropped: a delivery whose status write failed has not reached that
|
||||
// status, and its caller must not act as though it had.
|
||||
func (e *Engine) updateDeliveryStatus(
|
||||
webhookDB *gorm.DB,
|
||||
d *database.Delivery,
|
||||
targetType database.TargetType,
|
||||
status database.DeliveryStatus,
|
||||
) {
|
||||
) error {
|
||||
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 fmt.Errorf(
|
||||
"updating delivery %s to status %s: %w",
|
||||
d.ID, status, err,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
e.mtr.DeliveryStatusChanged(targetType, status)
|
||||
// An empty type means the target row is gone — a delivery being
|
||||
// settled long after its target was deleted. The row still has to
|
||||
// be settled, but the counter is left alone rather than given a
|
||||
// series labelled with the empty string.
|
||||
if targetType != "" {
|
||||
e.mtr.DeliveryStatusChanged(targetType, status)
|
||||
}
|
||||
|
||||
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 {
|
||||
@@ -1065,6 +1430,184 @@ func (e *Engine) countAttempts(
|
||||
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) {
|
||||
@@ -1132,6 +1675,10 @@ func buildRecoveryTask(
|
||||
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))
|
||||
@@ -1168,12 +1715,33 @@ func (e *Engine) loadTargetMap(
|
||||
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():
|
||||
@@ -1181,6 +1749,10 @@ func (e *Engine) sendRecoveredDeliveries(
|
||||
default:
|
||||
}
|
||||
|
||||
if _, ok := settled[deliveries[i].ID]; ok {
|
||||
continue
|
||||
}
|
||||
|
||||
target, ok := targetMap[deliveries[i].TargetID]
|
||||
if !ok {
|
||||
e.log.Error(
|
||||
@@ -1192,22 +1764,40 @@ func (e *Engine) sendRecoveredDeliveries(
|
||||
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,
|
||||
&deliveries[i].Event, &target, 1,
|
||||
&deliveries[i], webhookID, &event, &target,
|
||||
attempts[deliveries[i].ID]+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
|
||||
}
|
||||
e.queueRecovered(e.deliveryCh, task)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user