Deploy: main into prod #343
@@ -438,6 +438,31 @@ func (e *Engine) processNewTask(
|
|||||||
return
|
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 := buildEventFromTask(task)
|
||||||
|
|
||||||
event, err = e.hydrateEvent(
|
event, err = e.hydrateEvent(
|
||||||
@@ -482,7 +507,7 @@ func (e *Engine) processRetryTask(
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
d, err := e.loadRetryDelivery(
|
d, err := e.loadDelivery(
|
||||||
webhookDB, task.DeliveryID,
|
webhookDB, task.DeliveryID,
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -1643,7 +1668,7 @@ func (e *Engine) hydrateEvent(
|
|||||||
return event, nil
|
return event, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (e *Engine) loadRetryDelivery(
|
func (e *Engine) loadDelivery(
|
||||||
webhookDB *gorm.DB, deliveryID string,
|
webhookDB *gorm.DB, deliveryID string,
|
||||||
) (*database.Delivery, error) {
|
) (*database.Delivery, error) {
|
||||||
var d database.Delivery
|
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
|
// TestRetryingRecoverySkipsASuccessfulResult is the retrying-side twin
|
||||||
// of the pending reconcile. A second attempt that reached the receiver
|
// of the pending reconcile. A second attempt that reached the receiver
|
||||||
// and whose status write then failed sits at retrying holding a
|
// and whose status write then failed sits at retrying holding a
|
||||||
|
|||||||
Reference in New Issue
Block a user