Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
83740b1de1 |
@@ -438,6 +438,31 @@ func (e *Engine) processNewTask(
|
||||
return
|
||||
}
|
||||
|
||||
// Restart recovery can send and release this delivery before the
|
||||
// receiver's Notify queues it. Ownership cannot refuse a delivery
|
||||
// nobody holds, so the row decides whether it still needs sending.
|
||||
row, err := e.loadDelivery(webhookDB, task.DeliveryID)
|
||||
if err != nil {
|
||||
e.log.Error(
|
||||
"failed to load delivery",
|
||||
"delivery_id", task.DeliveryID,
|
||||
"error", err,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if row.Status != database.DeliveryStatusPending {
|
||||
e.log.Info(
|
||||
"delivery already handled, not sent again",
|
||||
"delivery_id", task.DeliveryID,
|
||||
"event_id", task.EventID,
|
||||
"status", row.Status,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
event := buildEventFromTask(task)
|
||||
|
||||
event, err = e.hydrateEvent(
|
||||
@@ -482,7 +507,7 @@ func (e *Engine) processRetryTask(
|
||||
return
|
||||
}
|
||||
|
||||
d, err := e.loadRetryDelivery(
|
||||
d, err := e.loadDelivery(
|
||||
webhookDB, task.DeliveryID,
|
||||
)
|
||||
if err != nil {
|
||||
@@ -1643,7 +1668,7 @@ func (e *Engine) hydrateEvent(
|
||||
return event, nil
|
||||
}
|
||||
|
||||
func (e *Engine) loadRetryDelivery(
|
||||
func (e *Engine) loadDelivery(
|
||||
webhookDB *gorm.DB, deliveryID string,
|
||||
) (*database.Delivery, error) {
|
||||
var d database.Delivery
|
||||
|
||||
@@ -273,6 +273,73 @@ func TestOwnershipIsReleasedAfterDelivery(t *testing.T) {
|
||||
)
|
||||
}
|
||||
|
||||
// TestNotifyAfterRecoveryDoesNotSendAgain is the startup race of
|
||||
// https://git.eeqj.de/sneak/webhooker/issues/299. The receiver has
|
||||
// written a delivery, restart recovery finds it pending, sends it and
|
||||
// releases it, and only then does the receiver's Notify for it arrive.
|
||||
// Nothing owns the delivery by then, so Notify takes it.
|
||||
func TestNotifyAfterRecoveryDoesNotSendAgain(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
targetID := uuid.New().String()
|
||||
s := fSweepSetup(t, targetID, "recovered")
|
||||
|
||||
event := iSeedEvent(
|
||||
t, s.WebhookDB, s.WebhookID, `{"recovered":true}`,
|
||||
)
|
||||
|
||||
d := iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, targetID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
s.Engine.ExportStart()
|
||||
|
||||
defer func() {
|
||||
require.NoError(
|
||||
t, s.Engine.ExportStop(context.Background()),
|
||||
)
|
||||
}()
|
||||
|
||||
// Restart recovery sends the delivery and lets it go.
|
||||
iWaitForDelivered(t, s.WebhookDB, d.ID)
|
||||
require.Eventually(
|
||||
t,
|
||||
func() bool {
|
||||
return s.Engine.ExportInflightHeld() == 0
|
||||
},
|
||||
5*time.Second, 20*time.Millisecond,
|
||||
)
|
||||
|
||||
body := event.Body
|
||||
|
||||
s.Engine.Notify([]delivery.Task{{
|
||||
DeliveryID: d.ID,
|
||||
EventID: event.ID,
|
||||
WebhookID: s.WebhookID,
|
||||
TargetID: targetID,
|
||||
TargetName: "recovered",
|
||||
TargetType: database.TargetTypeLog,
|
||||
Body: &body,
|
||||
EntrypointID: event.EntrypointID,
|
||||
}})
|
||||
|
||||
// Notify took the delivery, and a worker releases it once it has
|
||||
// run the task.
|
||||
require.Eventually(
|
||||
t,
|
||||
func() bool {
|
||||
return s.Engine.ExportInflightHeld() == 0
|
||||
},
|
||||
5*time.Second, 20*time.Millisecond,
|
||||
)
|
||||
|
||||
assert.Len(
|
||||
t, iResults(t, s.WebhookDB, d.ID), 1,
|
||||
"the delivery was sent a second time",
|
||||
)
|
||||
}
|
||||
|
||||
// TestRetryingRecoverySkipsASuccessfulResult is the retrying-side twin
|
||||
// of the pending reconcile. A second attempt that reached the receiver
|
||||
// and whose status write then failed sits at retrying holding a
|
||||
|
||||
Reference in New Issue
Block a user