From 608c3b21c82c781b36c3567c4d0b1a58ee8943c3 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 04:13:05 +0000 Subject: [PATCH] Re-read a delivery before failing it for a missing target failMissingTarget failed the delivery as the batch had read it, so a delivery a worker sent and let go between the batch read and the ownership check could end failed after a successful attempt. It now reads the row once it owns the delivery and fails it only if the status is still the one the batch read, as processNewTask does. A row that cannot be read is left alone. TestSweepPending_TargetDeleted now checks that the first sweep already queues the healthy delivery and the second does not queue it again. Model: opus-5-5 --- internal/delivery/engine.go | 24 ++++++++++ internal/delivery/export_test.go | 10 +++++ internal/delivery/terminal_state_test.go | 56 ++++++++++++++++++++++-- 3 files changed, 87 insertions(+), 3 deletions(-) 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,