From 3f690e8d10b6f654871418bf19455d5fbf64f111 Mon Sep 17 00:00:00 2001 From: sneak Date: Tue, 29 Sep 2026 01:39:40 +0000 Subject: [PATCH] Do not resend a delivery that restart recovery already sent (closes #299) 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 --- internal/delivery/engine.go | 29 ++++++++++++- internal/delivery/inflight_test.go | 67 ++++++++++++++++++++++++++++++ 2 files changed, 94 insertions(+), 2 deletions(-) diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index e023918..1143cbf 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -438,6 +438,31 @@ func (e *Engine) processNewTask( 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, err = e.hydrateEvent( @@ -482,7 +507,7 @@ func (e *Engine) processRetryTask( return } - d, err := e.loadRetryDelivery( + d, err := e.loadDelivery( webhookDB, task.DeliveryID, ) if err != nil { @@ -1643,7 +1668,7 @@ func (e *Engine) hydrateEvent( return event, nil } -func (e *Engine) loadRetryDelivery( +func (e *Engine) loadDelivery( webhookDB *gorm.DB, deliveryID string, ) (*database.Delivery, error) { var d database.Delivery diff --git a/internal/delivery/inflight_test.go b/internal/delivery/inflight_test.go index 1c86003..08a0f69 100644 --- a/internal/delivery/inflight_test.go +++ b/internal/delivery/inflight_test.go @@ -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 // of the pending reconcile. A second attempt that reached the receiver // and whose status write then failed sits at retrying holding a -- 2.54.0