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