Do not resend a delivery that restart recovery already sent (closes #299)
check / check (push) Successful in 3m39s

Restart recovery could find a just-written delivery pending, send it
and release it before the receiver's Notify queued the same delivery.
Notify's claim then succeeded on the released id, and the worker sent
it again because the new-task path never read the delivery's row.

Before sending a new task the worker now reads the delivery's status
by primary key and skips the task unless the row still says pending,
as the retry path already does for retrying. Nothing else can change
the row while the worker owns the delivery. A row left pending by a
failed bookkeeping write is still sent again.

loadRetryDelivery is renamed loadDelivery now that both paths use it.

Model: opus-5-5
This commit is contained in:
2026-09-29 01:39:40 +00:00
parent aeeeca5ea1
commit 3f690e8d10
2 changed files with 94 additions and 2 deletions
+27 -2
View File
@@ -438,6 +438,31 @@ func (e *Engine) processNewTask(
return return
} }
// Restart recovery can send and release this delivery before the
// receiver's Notify queues it. Ownership cannot refuse a delivery
// nobody holds, so the row decides whether it still needs sending.
row, err := e.loadDelivery(webhookDB, task.DeliveryID)
if err != nil {
e.log.Error(
"failed to load delivery",
"delivery_id", task.DeliveryID,
"error", err,
)
return
}
if row.Status != database.DeliveryStatusPending {
e.log.Info(
"delivery already handled, not sent again",
"delivery_id", task.DeliveryID,
"event_id", task.EventID,
"status", row.Status,
)
return
}
event := buildEventFromTask(task) event := buildEventFromTask(task)
event, err = e.hydrateEvent( event, err = e.hydrateEvent(
@@ -482,7 +507,7 @@ func (e *Engine) processRetryTask(
return return
} }
d, err := e.loadRetryDelivery( d, err := e.loadDelivery(
webhookDB, task.DeliveryID, webhookDB, task.DeliveryID,
) )
if err != nil { if err != nil {
@@ -1643,7 +1668,7 @@ func (e *Engine) hydrateEvent(
return event, nil return event, nil
} }
func (e *Engine) loadRetryDelivery( func (e *Engine) loadDelivery(
webhookDB *gorm.DB, deliveryID string, webhookDB *gorm.DB, deliveryID string,
) (*database.Delivery, error) { ) (*database.Delivery, error) {
var d database.Delivery var d database.Delivery
+67
View File
@@ -273,6 +273,73 @@ func TestOwnershipIsReleasedAfterDelivery(t *testing.T) {
) )
} }
// TestNotifyAfterRecoveryDoesNotSendAgain is the startup race of
// https://git.eeqj.de/sneak/webhooker/issues/299. The receiver has
// written a delivery, restart recovery finds it pending, sends it and
// releases it, and only then does the receiver's Notify for it arrive.
// Nothing owns the delivery by then, so Notify takes it.
func TestNotifyAfterRecoveryDoesNotSendAgain(t *testing.T) {
t.Parallel()
targetID := uuid.New().String()
s := fSweepSetup(t, targetID, "recovered")
event := iSeedEvent(
t, s.WebhookDB, s.WebhookID, `{"recovered":true}`,
)
d := iSeedDelivery(
t, s.WebhookDB, event.ID, targetID,
database.DeliveryStatusPending,
)
s.Engine.ExportStart()
defer func() {
require.NoError(
t, s.Engine.ExportStop(context.Background()),
)
}()
// Restart recovery sends the delivery and lets it go.
iWaitForDelivered(t, s.WebhookDB, d.ID)
require.Eventually(
t,
func() bool {
return s.Engine.ExportInflightHeld() == 0
},
5*time.Second, 20*time.Millisecond,
)
body := event.Body
s.Engine.Notify([]delivery.Task{{
DeliveryID: d.ID,
EventID: event.ID,
WebhookID: s.WebhookID,
TargetID: targetID,
TargetName: "recovered",
TargetType: database.TargetTypeLog,
Body: &body,
EntrypointID: event.EntrypointID,
}})
// Notify took the delivery, and a worker releases it once it has
// run the task.
require.Eventually(
t,
func() bool {
return s.Engine.ExportInflightHeld() == 0
},
5*time.Second, 20*time.Millisecond,
)
assert.Len(
t, iResults(t, s.WebhookDB, d.ID), 1,
"the delivery was sent a second time",
)
}
// TestRetryingRecoverySkipsASuccessfulResult is the retrying-side twin // TestRetryingRecoverySkipsASuccessfulResult is the retrying-side twin
// of the pending reconcile. A second attempt that reached the receiver // of the pending reconcile. A second attempt that reached the receiver
// and whose status write then failed sits at retrying holding a // and whose status write then failed sits at retrying holding a