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
This commit is contained in:
@@ -1270,6 +1270,30 @@ func (e *Engine) failMissingTarget(
|
|||||||
|
|
||||||
defer e.inflight.release(d.ID)
|
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)
|
targetType, reason := e.missingTargetReason(d.TargetID)
|
||||||
|
|
||||||
e.log.Warn(
|
e.log.Warn(
|
||||||
|
|||||||
@@ -342,6 +342,16 @@ func (e *Engine) ExportRecoverRetryingDeliveries(
|
|||||||
e.recoverRetryingDeliveries(webhookDB, webhookID)
|
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.
|
// ExportDeliveryCh returns the delivery channel.
|
||||||
func (e *Engine) ExportDeliveryCh() chan Task {
|
func (e *Engine) ExportDeliveryCh() chan Task {
|
||||||
return e.deliveryCh
|
return e.deliveryCh
|
||||||
|
|||||||
@@ -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
|
// TestSweepPending_TargetDeleted sweeps twice over a batch that also
|
||||||
// holds a healthy stranded delivery. The one whose target is gone is
|
// 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) {
|
func TestSweepPending_TargetDeleted(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
@@ -628,13 +671,20 @@ func TestSweepPending_TargetDeleted(t *testing.T) {
|
|||||||
|
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
|
|
||||||
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
|
|
||||||
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
|
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
|
||||||
|
|
||||||
tasks := fDrain(s.Engine)
|
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)
|
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(
|
iAssertStatus(
|
||||||
t, s.WebhookDB, deliveryID,
|
t, s.WebhookDB, deliveryID,
|
||||||
database.DeliveryStatusFailed,
|
database.DeliveryStatusFailed,
|
||||||
|
|||||||
Reference in New Issue
Block a user