From a0bbfca28d01357304cc55345cd0c24c598fd067 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 1/3] Fail a pending delivery whose target was deleted (closes #293) Restart recovery and the pending sweep skipped a pending delivery whose target was missing from the batch's target map, and the sweep did so again every minute for the life of the database. A miss now asks loadTarget: no row fails the delivery with a recorded reason, through the ownership-gated function the retrying paths already use, renamed failMissingTarget with its log line and reason text made to fit both statuses. Any other error leaves the delivery pending, because the map is also empty when its query failed. Model: opus-5-5 --- internal/delivery/engine.go | 58 +++++--- internal/delivery/terminal_state_test.go | 182 +++++++++++++++++++++-- 2 files changed, 208 insertions(+), 32 deletions(-) diff --git a/internal/delivery/engine.go b/internal/delivery/engine.go index 1143cbf..67bb261 100644 --- a/internal/delivery/engine.go +++ b/internal/delivery/engine.go @@ -728,9 +728,7 @@ func (e *Engine) recoverSingleRetry( // webhook on one bad read would be a far larger fault than // the strand it is meant to clear. if errors.Is(err, gorm.ErrRecordNotFound) { - e.failMissingTargetRetry( - webhookDB, webhookID, d, - ) + e.failMissingTarget(webhookDB, webhookID, d) return } @@ -1133,9 +1131,7 @@ func (e *Engine) sweepSingleRetry( // Deleted is terminal, unreadable is not; see // recoverSingleRetry. if errors.Is(err, gorm.ErrRecordNotFound) { - e.failMissingTargetRetry( - webhookDB, webhookID, d, - ) + e.failMissingTarget(webhookDB, webhookID, d) return } @@ -1249,19 +1245,19 @@ func (e *Engine) failUnretryableRetry( e.failDelivery(webhookDB, d, target.Type, reason) } -// failMissingTargetRetry terminally fails an orphaned retrying -// delivery whose target row is gone. Both restart recovery and the -// periodic sweep call it, so the transition exists once. +// failMissingTarget terminally fails a recovered delivery, pending or +// retrying, whose target row is gone. Restart recovery and the periodic +// sweep call it for both statuses, so the transition exists once. // -// Until it existed both paths logged the failed lookup and returned, -// which left the delivery retrying for the life of the database and +// Until it existed those paths logged the failed lookup and moved on, +// which left the delivery where it was for the life of the database and // the sweep repeating the same error every minute forever. Failing it // with a recorded reason is the treatment the other orphaned-retry // cases already get, so all of them read alike in the event log. // // Logged at warn rather than error: a deleted target is an operator // action, not a system fault. -func (e *Engine) failMissingTargetRetry( +func (e *Engine) failMissingTarget( webhookDB *gorm.DB, webhookID string, d *database.Delivery, @@ -1277,10 +1273,10 @@ func (e *Engine) failMissingTargetRetry( targetType, reason := e.missingTargetReason(d.TargetID) e.log.Warn( - "failing orphaned retrying delivery: "+ - "its target no longer exists", + "failing recovered delivery: its target no longer exists", "webhook_id", webhookID, "delivery_id", d.ID, + "status", d.Status, "target_id", d.TargetID, "target_type", targetType, ) @@ -1314,15 +1310,14 @@ func (e *Engine) missingTargetReason( if err != nil { return "", fmt.Sprintf( "target %s no longer exists; the delivery "+ - "cannot be retried and has been failed "+ - "terminally", + "has been failed terminally", targetID, ) } return target.Type, fmt.Sprintf( "target %q (type %s) was deleted; the delivery "+ - "cannot be retried and has been failed terminally", + "has been failed terminally", target.Name, target.Type, ) } @@ -2021,13 +2016,30 @@ func (e *Engine) sendRecoveredDeliveries( target, ok := targetMap[deliveries[i].TargetID] if !ok { - e.log.Error( - "target not found for delivery", - "delivery_id", deliveries[i].ID, - "target_id", deliveries[i].TargetID, - ) + // A missing entry does not mean the target is gone: the + // map is also empty when its query failed. Only a lookup + // that finds no row ends the delivery; any other error + // leaves it pending for the next sweep. See + // recoverSingleRetry. + var err error - continue + target, err = e.loadTarget(deliveries[i].TargetID) + if errors.Is(err, gorm.ErrRecordNotFound) { + e.failMissingTarget(webhookDB, webhookID, &deliveries[i]) + + continue + } + + if err != nil { + e.log.Error( + "failed to load target for recovered delivery", + "delivery_id", deliveries[i].ID, + "target_id", deliveries[i].TargetID, + "error", err, + ) + + continue + } } if !e.takeForRedispatch( diff --git a/internal/delivery/terminal_state_test.go b/internal/delivery/terminal_state_test.go index 382a6ff..c4054b6 100644 --- a/internal/delivery/terminal_state_test.go +++ b/internal/delivery/terminal_state_test.go @@ -18,16 +18,17 @@ import ( // https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed // with nothing in its event log to say why, and a retrying delivery // whose target was deleted, which used to keep sending and then never -// terminalise. +// terminalise. Section 4 is the same deleted-target gap for a pending +// delivery: https://git.eeqj.de/sneak/webhooker/issues/293. // tUnknownType is a target type no build implements. It stands in for // a target whose type was written by a build that knew a type this one // does not. const tUnknownType = database.TargetType("pubsub") -// tSeedDeletedTarget creates a target, a retrying delivery against it -// with one recorded failed attempt, and then deletes the target the -// way the source page does. +// tSeedDeletedTarget creates a target, a delivery against it at the +// given status with one recorded failed attempt, and then deletes the +// target the way the source page does. // // It asserts the delete is soft, because that is the whole reason the // engine could not tell a deleted target from a target id that never @@ -36,6 +37,7 @@ func tSeedDeletedTarget( t *testing.T, s iSetup, name, url string, + status database.DeliveryStatus, ) string { t.Helper() @@ -51,8 +53,7 @@ func tSeedDeletedTarget( ) d := iSeedDelivery( - t, s.WebhookDB, event.ID, targetID, - database.DeliveryStatusRetrying, + t, s.WebhookDB, event.ID, targetID, status, ) iSeedFailedResult(t, s.WebhookDB, d.ID) @@ -173,6 +174,7 @@ func TestRecoverSingleRetry_TargetDeleted(t *testing.T) { deliveryID := tSeedDeletedTarget( t, s, "gone-on-recovery", "http://example.com/hook", + database.DeliveryStatusRetrying, ) s.Engine.ExportRecoverWebhookDeliveries( @@ -210,6 +212,7 @@ func TestSweepSingleRetry_TargetDeleted(t *testing.T) { deliveryID := tSeedDeletedTarget( t, s, "gone-on-sweep", "http://example.com/hook", + database.DeliveryStatusRetrying, ) // Twice, because the bug was an error the sweep repeated every @@ -278,11 +281,11 @@ func TestSweepSingleRetry_TargetNeverExisted(t *testing.T) { ) } -// TestFailMissingTargetRetry_WritesNoTargetRow holds the new terminal -// path to the same rule as the existing one: no target row, and so no +// TestFailMissingTarget_WritesNoTargetRow holds the new terminal path +// to the same rule as the existing one: no target row, and so no // plaintext target config, may be written into the per-webhook event // database. See https://git.eeqj.de/sneak/webhooker/issues/206. -func TestFailMissingTargetRetry_WritesNoTargetRow( +func TestFailMissingTarget_WritesNoTargetRow( t *testing.T, ) { t.Parallel() @@ -297,6 +300,7 @@ func TestFailMissingTargetRetry_WritesNoTargetRow( deliveryID := tSeedDeletedTarget( t, s, "credential-bearing", hookURL, + database.DeliveryStatusRetrying, ) s.Engine.ExportSweepWebhookRetries( @@ -529,3 +533,163 @@ func TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone( assert.Zero(t, s.Engine.ExportInflightHeld()) } + +// --- 4. A pending delivery whose target is gone --- + +func TestRecoverPending_TargetDeleted(t *testing.T) { + t.Parallel() + + s := newISetup(t) + + deliveryID := tSeedDeletedTarget( + t, s, "gone-while-pending", "http://example.com/hook", + database.DeliveryStatusPending, + ) + + s.Engine.ExportRecoverWebhookDeliveries( + context.Background(), s.WebhookID, + ) + + iAssertStatus( + t, s.WebhookDB, deliveryID, + database.DeliveryStatusFailed, + ) + + last := tLastResult(t, s, deliveryID, 2) + + assert.False(t, last.Success) + assert.Equal(t, 2, last.AttemptNum) + assert.Contains(t, last.Error, "gone-while-pending") + assert.Contains(t, last.Error, "was deleted") + + assert.Empty(t, fDrain(s.Engine), + "a delivery whose target is gone was sent", + ) + assert.Zero(t, s.Engine.ExportInflightHeld(), + "the terminal path leaked its ownership reference", + ) +} + +// TestRecoverPending_TargetDeleted_LeavesAnOwnedDeliveryAlone: the +// terminal write takes ownership like every other recovery write, so a +// delivery the engine still holds is not failed underneath its worker. +func TestRecoverPending_TargetDeleted_LeavesAnOwnedDeliveryAlone( + t *testing.T, +) { + t.Parallel() + + s := newISetup(t) + + deliveryID := tSeedDeletedTarget( + t, s, "gone-but-owned", "http://example.com/hook", + database.DeliveryStatusPending, + ) + + require.True(t, s.Engine.ExportRetainDelivery(deliveryID)) + + s.Engine.ExportRecoverWebhookDeliveries( + context.Background(), s.WebhookID, + ) + + iAssertStatus( + t, s.WebhookDB, deliveryID, + database.DeliveryStatusPending, + ) + + assert.Len(t, iResults(t, s.WebhookDB, deliveryID), 1, + "a delivery the engine owns was failed underneath it", + ) +} + +// 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. +func TestSweepPending_TargetDeleted(t *testing.T) { + t.Parallel() + + liveTargetID := uuid.New().String() + s := fSweepSetup(t, liveTargetID, "still-there") + + deliveryID := tSeedDeletedTarget( + t, s, "gone-on-pending-sweep", "http://example.com/hook", + database.DeliveryStatusPending, + ) + rAgePending(t, s.WebhookDB, deliveryID) + + event := iSeedEvent( + t, s.WebhookDB, s.WebhookID, `{"target":"live"}`, + ) + + healthy := iSeedDelivery( + t, s.WebhookDB, event.ID, liveTargetID, + database.DeliveryStatusPending, + ) + rAgePending(t, s.WebhookDB, healthy.ID) + + ctx := context.Background() + + s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID) + s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID) + + tasks := fDrain(s.Engine) + require.Len(t, tasks, 1) + assert.Equal(t, healthy.ID, tasks[0].DeliveryID) + + iAssertStatus( + t, s.WebhookDB, deliveryID, + database.DeliveryStatusFailed, + ) + + last := tLastResult(t, s, deliveryID, 2) + + assert.Contains(t, last.Error, "gone-on-pending-sweep") + assert.Contains(t, last.Error, "was deleted") +} + +// TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone: a failed +// read of the main database is not a deleted target. Restart recovery +// holds every pending delivery of the webhook in one batch, so failing +// on this would fail all of them. +func TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone( + t *testing.T, +) { + t.Parallel() + + s := newISetup(t) + + targetID := uuid.New().String() + + iCreateTarget( + t, s.MainDB, targetID, s.WebhookID, "healthy", + database.TargetTypeLog, "", 0, + ) + + event := iSeedEvent( + t, s.WebhookDB, s.WebhookID, `{"still":"pending"}`, + ) + + d := iSeedDelivery( + t, s.WebhookDB, event.ID, targetID, + database.DeliveryStatusPending, + ) + + sqlDB, err := s.MainDB.DB() + require.NoError(t, err) + require.NoError(t, sqlDB.Close()) + + s.Engine.ExportRecoverPendingDeliveries( + context.Background(), s.WebhookDB, s.WebhookID, + ) + + iAssertStatus( + t, s.WebhookDB, d.ID, + database.DeliveryStatusPending, + ) + + assert.Empty(t, iResults(t, s.WebhookDB, d.ID), + "an unreadable main database produced a terminal "+ + "failure row", + ) + assert.Empty(t, fDrain(s.Engine)) + assert.Zero(t, s.Engine.ExportInflightHeld()) +} -- 2.54.0 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 2/3] 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, -- 2.54.0 From 81d758d75659e78ec7257e7376ccd39c1db9cb9a Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 04:15:01 +0000 Subject: [PATCH 3/3] Test a pending delivery whose target is found only on lookup When the batch's target map lacks a delivery's target, as it does for every delivery when the batch query fails, sendRecoveredDeliveries looks the target up on its own. The new test hands it an empty map and checks the healthy delivery is queued once, with its real target id and type. Model: opus-5-5 --- internal/delivery/export_test.go | 15 ++++++++ internal/delivery/terminal_state_test.go | 47 ++++++++++++++++++++++++ 2 files changed, 62 insertions(+) diff --git a/internal/delivery/export_test.go b/internal/delivery/export_test.go index b3974c8..9dd2531 100644 --- a/internal/delivery/export_test.go +++ b/internal/delivery/export_test.go @@ -352,6 +352,21 @@ func (e *Engine) ExportFailMissingTarget( e.failMissingTarget(webhookDB, webhookID, d) } +// ExportSendRecoveredDeliveries exposes sendRecoveredDeliveries, so a +// test can hand it a target map that lacks a delivery's target. +func (e *Engine) ExportSendRecoveredDeliveries( + ctx context.Context, + webhookDB *gorm.DB, + deliveries []database.Delivery, + webhookID string, + targetMap map[string]database.Target, + settled map[string]struct{}, +) { + e.sendRecoveredDeliveries( + ctx, webhookDB, deliveries, webhookID, targetMap, settled, + ) +} + // 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 d448271..8eda81a 100644 --- a/internal/delivery/terminal_state_test.go +++ b/internal/delivery/terminal_state_test.go @@ -696,6 +696,53 @@ func TestSweepPending_TargetDeleted(t *testing.T) { assert.Contains(t, last.Error, "was deleted") } +// TestSendRecoveredDeliveries_TargetMissingFromMap: the batch's target +// map is empty when its query failed, so every delivery in the batch is +// looked up on its own. A healthy one is sent to the target that lookup +// finds. +func TestSendRecoveredDeliveries_TargetMissingFromMap( + t *testing.T, +) { + t.Parallel() + + s := newISetup(t) + + targetID := uuid.New().String() + + iCreateTarget( + t, s.MainDB, targetID, s.WebhookID, "found-on-lookup", + database.TargetTypeLog, "", 0, + ) + + event := iSeedEvent( + t, s.WebhookDB, s.WebhookID, `{"map":"empty"}`, + ) + + d := iSeedDelivery( + t, s.WebhookDB, event.ID, targetID, + database.DeliveryStatusPending, + ) + + s.Engine.ExportSendRecoveredDeliveries( + context.Background(), s.WebhookDB, + []database.Delivery{d}, s.WebhookID, + map[string]database.Target{}, nil, + ) + + tasks := fDrain(s.Engine) + require.Len(t, tasks, 1, + "the healthy delivery was not queued exactly once", + ) + assert.Equal(t, d.ID, tasks[0].DeliveryID) + assert.Equal(t, targetID, tasks[0].TargetID) + assert.Equal(t, database.TargetTypeLog, tasks[0].TargetType) + + iAssertStatus( + t, s.WebhookDB, d.ID, + database.DeliveryStatusPending, + ) +} + // TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone: a failed // read of the main database is not a deleted target. Restart recovery // holds every pending delivery of the webhook in one batch, so failing -- 2.54.0