diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index 67bb261..c96dda7 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -1270,6 +1270,30 @@ func (e *Engine) failMissingTarget( defer e.inflight.release(d.ID) + // The batch was read before ownership was taken, and a worker may + // have settled the delivery and let it go in between. Only a row + // still in the status the batch read is failed. + row, err := e.loadDelivery(webhookDB, d.ID) + if err != nil { + e.log.Error( + "failed to load delivery", + "delivery_id", d.ID, + "error", err, + ) + + return + } + + if row.Status != d.Status { + e.log.Debug( + "delivery already handled, not failed", + "delivery_id", d.ID, + "status", row.Status, + ) + + return + } + targetType, reason := e.missingTargetReason(d.TargetID) e.log.Warn( diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index 37c2695..b3974c8 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -342,6 +342,16 @@ func (e *Engine) ExportRecoverRetryingDeliveries( e.recoverRetryingDeliveries(webhookDB, webhookID) } +// ExportFailMissingTarget exposes failMissingTarget, so a test can hand +// it a delivery as a batch read it earlier. +func (e *Engine) ExportFailMissingTarget( + webhookDB *gorm.DB, + webhookID string, + d *database.Delivery, +) { + e.failMissingTarget(webhookDB, webhookID, d) +} + // ExportDeliveryCh returns the delivery channel. func (e *Engine) ExportDeliveryCh() chan Task { return e.deliveryCh diff --git a/internal/delivery/terminal_state_test.go b/internal/delivery/terminal_state_test.go index c4054b6..d448271 100644 --- a/internal/delivery/terminal_state_test.go +++ b/internal/delivery/terminal_state_test.go @@ -601,9 +601,52 @@ func TestRecoverPending_TargetDeleted_LeavesAnOwnedDeliveryAlone( ) } +// TestFailMissingTarget_LeavesASettledDeliveryAlone: the recovery paths +// read their batch before taking ownership, and a worker may send a +// delivery and let it go in between. The terminal write goes by the row +// as it is now, not as the batch read it. +func TestFailMissingTarget_LeavesASettledDeliveryAlone( + t *testing.T, +) { + t.Parallel() + + s := newISetup(t) + + deliveryID := tSeedDeletedTarget( + t, s, "gone-after-sending", "http://example.com/hook", + database.DeliveryStatusPending, + ) + + var batch database.Delivery + + require.NoError(t, s.WebhookDB.First( + &batch, "id = ?", deliveryID, + ).Error) + + // A worker settles the delivery after the batch was read. + require.NoError(t, s.WebhookDB.Model(&database.Delivery{}). + Where("id = ?", deliveryID). + Update("status", database.DeliveryStatusDelivered).Error) + + s.Engine.ExportFailMissingTarget( + s.WebhookDB, s.WebhookID, &batch, + ) + + iAssertStatus( + t, s.WebhookDB, deliveryID, + database.DeliveryStatusDelivered, + ) + + assert.Len(t, iResults(t, s.WebhookDB, deliveryID), 1, + "a delivery settled after the batch read was then failed", + ) + assert.Zero(t, s.Engine.ExportInflightHeld()) +} + // TestSweepPending_TargetDeleted sweeps twice over a batch that also // holds a healthy stranded delivery. The one whose target is gone is -// failed once and then left alone; the healthy one is sent exactly once. +// failed once and then left alone; the healthy one is queued by the +// first sweep and not again by the second. func TestSweepPending_TargetDeleted(t *testing.T) { t.Parallel() @@ -628,13 +671,20 @@ func TestSweepPending_TargetDeleted(t *testing.T) { ctx := context.Background() - s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID) s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID) tasks := fDrain(s.Engine) - require.Len(t, tasks, 1) + require.Len(t, tasks, 1, + "the first sweep did not queue the healthy delivery", + ) assert.Equal(t, healthy.ID, tasks[0].DeliveryID) + s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID) + + assert.Empty(t, fDrain(s.Engine), + "the second sweep queued a delivery again", + ) + iAssertStatus( t, s.WebhookDB, deliveryID, database.DeliveryStatusFailed,