Fail a pending delivery whose target was deleted (closes #293)
check / check (push) Successful in 4m21s
check / check (push) Successful in 4m21s
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
This commit is contained in:
+30
-18
@@ -728,9 +728,7 @@ func (e *Engine) recoverSingleRetry(
|
|||||||
// webhook on one bad read would be a far larger fault than
|
// webhook on one bad read would be a far larger fault than
|
||||||
// the strand it is meant to clear.
|
// the strand it is meant to clear.
|
||||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||||
e.failMissingTargetRetry(
|
e.failMissingTarget(webhookDB, webhookID, d)
|
||||||
webhookDB, webhookID, d,
|
|
||||||
)
|
|
||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -1133,9 +1131,7 @@ func (e *Engine) sweepSingleRetry(
|
|||||||
// Deleted is terminal, unreadable is not; see
|
// Deleted is terminal, unreadable is not; see
|
||||||
// recoverSingleRetry.
|
// recoverSingleRetry.
|
||||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||||
e.failMissingTargetRetry(
|
e.failMissingTarget(webhookDB, webhookID, d)
|
||||||
webhookDB, webhookID, d,
|
|
||||||
)
|
|
||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -1249,19 +1245,19 @@ func (e *Engine) failUnretryableRetry(
|
|||||||
e.failDelivery(webhookDB, d, target.Type, reason)
|
e.failDelivery(webhookDB, d, target.Type, reason)
|
||||||
}
|
}
|
||||||
|
|
||||||
// failMissingTargetRetry terminally fails an orphaned retrying
|
// failMissingTarget terminally fails a recovered delivery, pending or
|
||||||
// delivery whose target row is gone. Both restart recovery and the
|
// retrying, whose target row is gone. Restart recovery and the periodic
|
||||||
// periodic sweep call it, so the transition exists once.
|
// sweep call it for both statuses, so the transition exists once.
|
||||||
//
|
//
|
||||||
// Until it existed both paths logged the failed lookup and returned,
|
// Until it existed those paths logged the failed lookup and moved on,
|
||||||
// which left the delivery retrying for the life of the database and
|
// 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
|
// the sweep repeating the same error every minute forever. Failing it
|
||||||
// with a recorded reason is the treatment the other orphaned-retry
|
// with a recorded reason is the treatment the other orphaned-retry
|
||||||
// cases already get, so all of them read alike in the event log.
|
// 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
|
// Logged at warn rather than error: a deleted target is an operator
|
||||||
// action, not a system fault.
|
// action, not a system fault.
|
||||||
func (e *Engine) failMissingTargetRetry(
|
func (e *Engine) failMissingTarget(
|
||||||
webhookDB *gorm.DB,
|
webhookDB *gorm.DB,
|
||||||
webhookID string,
|
webhookID string,
|
||||||
d *database.Delivery,
|
d *database.Delivery,
|
||||||
@@ -1277,10 +1273,10 @@ func (e *Engine) failMissingTargetRetry(
|
|||||||
targetType, reason := e.missingTargetReason(d.TargetID)
|
targetType, reason := e.missingTargetReason(d.TargetID)
|
||||||
|
|
||||||
e.log.Warn(
|
e.log.Warn(
|
||||||
"failing orphaned retrying delivery: "+
|
"failing recovered delivery: its target no longer exists",
|
||||||
"its target no longer exists",
|
|
||||||
"webhook_id", webhookID,
|
"webhook_id", webhookID,
|
||||||
"delivery_id", d.ID,
|
"delivery_id", d.ID,
|
||||||
|
"status", d.Status,
|
||||||
"target_id", d.TargetID,
|
"target_id", d.TargetID,
|
||||||
"target_type", targetType,
|
"target_type", targetType,
|
||||||
)
|
)
|
||||||
@@ -1314,15 +1310,14 @@ func (e *Engine) missingTargetReason(
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return "", fmt.Sprintf(
|
return "", fmt.Sprintf(
|
||||||
"target %s no longer exists; the delivery "+
|
"target %s no longer exists; the delivery "+
|
||||||
"cannot be retried and has been failed "+
|
"has been failed terminally",
|
||||||
"terminally",
|
|
||||||
targetID,
|
targetID,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
return target.Type, fmt.Sprintf(
|
return target.Type, fmt.Sprintf(
|
||||||
"target %q (type %s) was deleted; the delivery "+
|
"target %q (type %s) was deleted; the delivery "+
|
||||||
"cannot be retried and has been failed terminally",
|
"has been failed terminally",
|
||||||
target.Name, target.Type,
|
target.Name, target.Type,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
@@ -2021,14 +2016,31 @@ func (e *Engine) sendRecoveredDeliveries(
|
|||||||
|
|
||||||
target, ok := targetMap[deliveries[i].TargetID]
|
target, ok := targetMap[deliveries[i].TargetID]
|
||||||
if !ok {
|
if !ok {
|
||||||
|
// 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
|
||||||
|
|
||||||
|
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(
|
e.log.Error(
|
||||||
"target not found for delivery",
|
"failed to load target for recovered delivery",
|
||||||
"delivery_id", deliveries[i].ID,
|
"delivery_id", deliveries[i].ID,
|
||||||
"target_id", deliveries[i].TargetID,
|
"target_id", deliveries[i].TargetID,
|
||||||
|
"error", err,
|
||||||
)
|
)
|
||||||
|
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
if !e.takeForRedispatch(
|
if !e.takeForRedispatch(
|
||||||
webhookDB, deliveries[i].ID,
|
webhookDB, deliveries[i].ID,
|
||||||
|
|||||||
@@ -18,16 +18,17 @@ import (
|
|||||||
// https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed
|
// https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed
|
||||||
// with nothing in its event log to say why, and a retrying delivery
|
// 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
|
// 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
|
// 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
|
// a target whose type was written by a build that knew a type this one
|
||||||
// does not.
|
// does not.
|
||||||
const tUnknownType = database.TargetType("pubsub")
|
const tUnknownType = database.TargetType("pubsub")
|
||||||
|
|
||||||
// tSeedDeletedTarget creates a target, a retrying delivery against it
|
// tSeedDeletedTarget creates a target, a delivery against it at the
|
||||||
// with one recorded failed attempt, and then deletes the target the
|
// given status with one recorded failed attempt, and then deletes the
|
||||||
// way the source page does.
|
// target the way the source page does.
|
||||||
//
|
//
|
||||||
// It asserts the delete is soft, because that is the whole reason the
|
// 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
|
// engine could not tell a deleted target from a target id that never
|
||||||
@@ -36,6 +37,7 @@ func tSeedDeletedTarget(
|
|||||||
t *testing.T,
|
t *testing.T,
|
||||||
s iSetup,
|
s iSetup,
|
||||||
name, url string,
|
name, url string,
|
||||||
|
status database.DeliveryStatus,
|
||||||
) string {
|
) string {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
@@ -51,8 +53,7 @@ func tSeedDeletedTarget(
|
|||||||
)
|
)
|
||||||
|
|
||||||
d := iSeedDelivery(
|
d := iSeedDelivery(
|
||||||
t, s.WebhookDB, event.ID, targetID,
|
t, s.WebhookDB, event.ID, targetID, status,
|
||||||
database.DeliveryStatusRetrying,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
iSeedFailedResult(t, s.WebhookDB, d.ID)
|
iSeedFailedResult(t, s.WebhookDB, d.ID)
|
||||||
@@ -173,6 +174,7 @@ func TestRecoverSingleRetry_TargetDeleted(t *testing.T) {
|
|||||||
|
|
||||||
deliveryID := tSeedDeletedTarget(
|
deliveryID := tSeedDeletedTarget(
|
||||||
t, s, "gone-on-recovery", "http://example.com/hook",
|
t, s, "gone-on-recovery", "http://example.com/hook",
|
||||||
|
database.DeliveryStatusRetrying,
|
||||||
)
|
)
|
||||||
|
|
||||||
s.Engine.ExportRecoverWebhookDeliveries(
|
s.Engine.ExportRecoverWebhookDeliveries(
|
||||||
@@ -210,6 +212,7 @@ func TestSweepSingleRetry_TargetDeleted(t *testing.T) {
|
|||||||
|
|
||||||
deliveryID := tSeedDeletedTarget(
|
deliveryID := tSeedDeletedTarget(
|
||||||
t, s, "gone-on-sweep", "http://example.com/hook",
|
t, s, "gone-on-sweep", "http://example.com/hook",
|
||||||
|
database.DeliveryStatusRetrying,
|
||||||
)
|
)
|
||||||
|
|
||||||
// Twice, because the bug was an error the sweep repeated every
|
// 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
|
// TestFailMissingTarget_WritesNoTargetRow holds the new terminal path
|
||||||
// path to the same rule as the existing one: no target row, and so no
|
// 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
|
// plaintext target config, may be written into the per-webhook event
|
||||||
// database. See https://git.eeqj.de/sneak/webhooker/issues/206.
|
// database. See https://git.eeqj.de/sneak/webhooker/issues/206.
|
||||||
func TestFailMissingTargetRetry_WritesNoTargetRow(
|
func TestFailMissingTarget_WritesNoTargetRow(
|
||||||
t *testing.T,
|
t *testing.T,
|
||||||
) {
|
) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
@@ -297,6 +300,7 @@ func TestFailMissingTargetRetry_WritesNoTargetRow(
|
|||||||
|
|
||||||
deliveryID := tSeedDeletedTarget(
|
deliveryID := tSeedDeletedTarget(
|
||||||
t, s, "credential-bearing", hookURL,
|
t, s, "credential-bearing", hookURL,
|
||||||
|
database.DeliveryStatusRetrying,
|
||||||
)
|
)
|
||||||
|
|
||||||
s.Engine.ExportSweepWebhookRetries(
|
s.Engine.ExportSweepWebhookRetries(
|
||||||
@@ -529,3 +533,163 @@ func TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone(
|
|||||||
|
|
||||||
assert.Zero(t, s.Engine.ExportInflightHeld())
|
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())
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user