check / check (push) Successful in 3m35s
Restart recovery could find a just-written delivery pending, send it and release it before the receiver's Notify queued the same delivery. Notify's claim then succeeded on the released id, and the worker sent it again because the new-task path never read the delivery's row. Before sending a new task the worker now reads the delivery's status by primary key and skips the task unless the row still says pending, as the retry path already does for retrying. Nothing else can change the row while the worker owns the delivery. A row left pending by a failed bookkeeping write is still sent again. loadRetryDelivery is renamed loadDelivery now that both paths use it. Model: opus-5-5
496 lines
12 KiB
Go
496 lines
12 KiB
Go
package delivery_test
|
|
|
|
import (
|
|
"context"
|
|
"path/filepath"
|
|
"strings"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
"sneak.berlin/go/webhooker/internal/database"
|
|
"sneak.berlin/go/webhooker/internal/delivery"
|
|
)
|
|
|
|
// These tests pin the rule that decides whether a delivery may be
|
|
// handed back to a worker: the engine re-dispatches only what it does
|
|
// not already own. Age alone is not that rule — a healthy delivery
|
|
// waiting in a 10000-deep channel is old and must not be re-sent. See
|
|
// https://git.eeqj.de/sneak/webhooker/issues/256.
|
|
|
|
// fSweepSetup seeds the main database with the webhook row the sweep
|
|
// enumerates, and returns the setup.
|
|
func fSweepSetup(
|
|
t *testing.T, targetID, name string,
|
|
) iSetup {
|
|
t.Helper()
|
|
|
|
s := newISetup(t)
|
|
|
|
iCreateTarget(t, s.MainDB, targetID,
|
|
s.WebhookID, name,
|
|
database.TargetTypeLog, "", 0,
|
|
)
|
|
|
|
require.NoError(t, s.MainDB.Create(&database.Webhook{
|
|
BaseModel: database.BaseModel{ID: s.WebhookID},
|
|
UserID: uuid.New().String(),
|
|
Name: name,
|
|
}).Error)
|
|
|
|
return s
|
|
}
|
|
|
|
// fDrain collects every task the engine has queued.
|
|
//
|
|
// Every caller drives the dispatch paths synchronously and has already
|
|
// waited for them to return, so anything they queued is in the channel
|
|
// by now. The short grace covers nothing but scheduler jitter, and is
|
|
// kept small because one of these tests runs the drain forty times.
|
|
func fDrain(e *delivery.Engine) []delivery.Task {
|
|
var out []delivery.Task
|
|
|
|
for {
|
|
select {
|
|
case task := <-e.ExportDeliveryCh():
|
|
out = append(out, task)
|
|
case task := <-e.ExportRetryCh():
|
|
out = append(out, task)
|
|
case <-time.After(25 * time.Millisecond):
|
|
return out
|
|
}
|
|
}
|
|
}
|
|
|
|
// TestArchiveHandleIsWAL closes the last gap in the durability
|
|
// evidence: the main and per-webhook tiers each assert their journal
|
|
// mode on a live handle, and the archive tier gets its settings from
|
|
// the same code path but nothing checked the running file.
|
|
func TestArchiveHandleIsWAL(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
w := delivery.NewExportArchiveWriter(
|
|
filepath.Join(t.TempDir(), "archive-wal.db"),
|
|
archiveTestLogger(), 0,
|
|
)
|
|
|
|
require.NoError(t, w.Open(0))
|
|
|
|
var mode string
|
|
|
|
row := w.DB().Raw("pragma journal_mode").Row()
|
|
require.NoError(t, row.Scan(&mode))
|
|
assert.Equal(t, "wal", strings.ToLower(mode))
|
|
|
|
var busy string
|
|
|
|
row = w.DB().Raw("pragma busy_timeout").Row()
|
|
require.NoError(t, row.Scan(&busy))
|
|
assert.Equal(t, "10000", busy)
|
|
}
|
|
|
|
// TestSweepLeavesAQueuedDeliveryAlone is the case the age bound cannot
|
|
// see. The delivery is queued and untouched, so its row is arbitrarily
|
|
// old and still perfectly healthy; only ownership distinguishes it
|
|
// from a stranded one.
|
|
func TestSweepLeavesAQueuedDeliveryAlone(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
targetID := uuid.New().String()
|
|
s := fSweepSetup(t, targetID, "queued")
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"queued":true}`,
|
|
)
|
|
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
rAgePending(t, s.WebhookDB, d.ID)
|
|
|
|
// Queued exactly as the receiver queues it, and never dequeued:
|
|
// no workers are running in this engine.
|
|
s.Engine.Notify([]delivery.Task{{
|
|
DeliveryID: d.ID,
|
|
EventID: event.ID,
|
|
WebhookID: s.WebhookID,
|
|
TargetID: targetID,
|
|
}})
|
|
|
|
require.Equal(t, 1, s.Engine.ExportInflightHeld())
|
|
|
|
s.Engine.ExportSweepWebhookRetries(
|
|
context.Background(), s.WebhookID,
|
|
)
|
|
|
|
tasks := fDrain(s.Engine)
|
|
assert.Len(
|
|
t, tasks, 1,
|
|
"the sweep must not queue a delivery that is "+
|
|
"already waiting for a worker",
|
|
)
|
|
}
|
|
|
|
// TestRecoveryAndSweepDoNotDoubleDispatch drives the two entry points
|
|
// the engine starts concurrently against one aged pending row. Before
|
|
// ownership they both dispatched it.
|
|
func TestRecoveryAndSweepDoNotDoubleDispatch(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
targetID := uuid.New().String()
|
|
s := fSweepSetup(t, targetID, "racing")
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"racing":true}`,
|
|
)
|
|
|
|
ctx := context.Background()
|
|
|
|
for range 40 {
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
rAgePending(t, s.WebhookDB, d.ID)
|
|
|
|
var wg sync.WaitGroup
|
|
|
|
wg.Go(func() {
|
|
s.Engine.ExportRecoverPendingDeliveries(
|
|
ctx, s.WebhookDB, s.WebhookID,
|
|
)
|
|
})
|
|
wg.Go(func() {
|
|
s.Engine.ExportSweepWebhookRetries(
|
|
ctx, s.WebhookID,
|
|
)
|
|
})
|
|
wg.Wait()
|
|
|
|
tasks := fDrain(s.Engine)
|
|
require.Len(
|
|
t, tasks, 1,
|
|
"delivery %s dispatched %d times",
|
|
d.ID, len(tasks),
|
|
)
|
|
|
|
// No worker runs in this engine, so the reference the winner
|
|
// took is never released and earlier iterations' deliveries
|
|
// stay owned — which is itself the property under test, since
|
|
// both paths see them on every subsequent pass.
|
|
}
|
|
}
|
|
|
|
// TestConcurrentClaimsOfOneDeliveryYieldOneOwner exercises the
|
|
// exclusion directly, rather than arguing it from a SQL predicate.
|
|
func TestConcurrentClaimsOfOneDeliveryYieldOneOwner(
|
|
t *testing.T,
|
|
) {
|
|
t.Parallel()
|
|
|
|
eng := newISetup(t).Engine
|
|
deliveryID := uuid.New().String()
|
|
|
|
var (
|
|
wg sync.WaitGroup
|
|
mu sync.Mutex
|
|
won int
|
|
)
|
|
|
|
for range 64 {
|
|
wg.Go(func() {
|
|
if eng.ExportRetainDelivery(deliveryID) {
|
|
mu.Lock()
|
|
won++
|
|
mu.Unlock()
|
|
}
|
|
})
|
|
}
|
|
|
|
wg.Wait()
|
|
|
|
assert.Equal(t, 1, won)
|
|
assert.Equal(t, 1, eng.ExportInflightHeld())
|
|
}
|
|
|
|
// TestOwnershipIsReleasedAfterDelivery guards the other direction: a
|
|
// leaked reference hides a delivery from every sweep for the life of
|
|
// the process.
|
|
func TestOwnershipIsReleasedAfterDelivery(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
s := newISetup(t)
|
|
targetID := uuid.New().String()
|
|
|
|
iCreateTarget(t, s.MainDB, targetID,
|
|
s.WebhookID, "released",
|
|
database.TargetTypeLog, "", 0,
|
|
)
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"released":true}`,
|
|
)
|
|
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
|
|
s.Engine.ExportStart()
|
|
|
|
defer func() {
|
|
require.NoError(
|
|
t, s.Engine.ExportStop(context.Background()),
|
|
)
|
|
}()
|
|
|
|
body := `{"released":true}`
|
|
|
|
s.Engine.Notify([]delivery.Task{{
|
|
DeliveryID: d.ID,
|
|
EventID: event.ID,
|
|
WebhookID: s.WebhookID,
|
|
TargetID: targetID,
|
|
TargetName: "released",
|
|
TargetType: database.TargetTypeLog,
|
|
Body: &body,
|
|
EntrypointID: event.EntrypointID,
|
|
}})
|
|
|
|
iWaitForDelivered(t, s.WebhookDB, d.ID)
|
|
|
|
assert.Eventually(
|
|
t,
|
|
func() bool {
|
|
return s.Engine.ExportInflightHeld() == 0
|
|
},
|
|
2*time.Second, 20*time.Millisecond,
|
|
"the delivery stayed owned after it was delivered",
|
|
)
|
|
}
|
|
|
|
// 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
|
|
// successful result, and re-sending it is the same duplicate.
|
|
func TestRetryingRecoverySkipsASuccessfulResult(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
targetID := uuid.New().String()
|
|
s := fSweepSetup(t, targetID, "retry-settled")
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"retry":true}`,
|
|
)
|
|
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusRetrying,
|
|
)
|
|
rSeedResult(t, s.WebhookDB, d.ID, 1, false)
|
|
rSeedResult(t, s.WebhookDB, d.ID, 2, true)
|
|
|
|
s.Engine.ExportRecoverRetryingDeliveries(
|
|
s.WebhookDB, s.WebhookID,
|
|
)
|
|
|
|
assert.Empty(
|
|
t, fDrain(s.Engine),
|
|
"a retrying delivery holding a successful result "+
|
|
"must not be sent again",
|
|
)
|
|
|
|
iAssertStatus(
|
|
t, s.WebhookDB, d.ID,
|
|
database.DeliveryStatusDelivered,
|
|
)
|
|
}
|
|
|
|
// TestRetryingSweepSkipsASuccessfulResult is the same rule on the
|
|
// periodic sweep's retrying arm.
|
|
func TestRetryingSweepSkipsASuccessfulResult(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
targetID := uuid.New().String()
|
|
s := fSweepSetup(t, targetID, "retry-swept")
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"swept":true}`,
|
|
)
|
|
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusRetrying,
|
|
)
|
|
rSeedResult(t, s.WebhookDB, d.ID, 1, false)
|
|
rSeedResult(t, s.WebhookDB, d.ID, 2, true)
|
|
|
|
s.Engine.ExportSweepWebhookRetries(
|
|
context.Background(), s.WebhookID,
|
|
)
|
|
|
|
assert.Empty(t, fDrain(s.Engine))
|
|
|
|
iAssertStatus(
|
|
t, s.WebhookDB, d.ID,
|
|
database.DeliveryStatusDelivered,
|
|
)
|
|
|
|
var attempts int64
|
|
|
|
require.NoError(t, s.WebhookDB.
|
|
Model(&database.DeliveryResult{}).
|
|
Where("delivery_id = ?", d.ID).
|
|
Count(&attempts).Error)
|
|
assert.Equal(
|
|
t, int64(2), attempts,
|
|
"settling must not invent an attempt",
|
|
)
|
|
}
|
|
|
|
// TestScheduledRetryIsNotSweptDuringBackoff closes the window between
|
|
// a target scheduling a retry and the timer firing. The row says
|
|
// retrying and nothing is running, which is exactly what an orphaned
|
|
// retry looks like from the database.
|
|
func TestScheduledRetryIsNotSweptDuringBackoff(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
targetID := uuid.New().String()
|
|
s := fSweepSetup(t, targetID, "backoff")
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"backoff":true}`,
|
|
)
|
|
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusRetrying,
|
|
)
|
|
|
|
s.Engine.ExportScheduleRetry(delivery.Task{
|
|
DeliveryID: d.ID,
|
|
EventID: event.ID,
|
|
WebhookID: s.WebhookID,
|
|
TargetID: targetID,
|
|
AttemptNum: 2,
|
|
}, time.Hour)
|
|
|
|
require.Equal(t, 1, s.Engine.ExportInflightHeld())
|
|
|
|
s.Engine.ExportSweepWebhookRetries(
|
|
context.Background(), s.WebhookID,
|
|
)
|
|
|
|
assert.Empty(
|
|
t, fDrain(s.Engine),
|
|
"the sweep must not duplicate a retry that is "+
|
|
"already scheduled",
|
|
)
|
|
}
|
|
|
|
// TestRedispatchStampsTheRow pins the cadence control: a stranded
|
|
// delivery that has just been handed out is not selected again by the
|
|
// next tick a minute later.
|
|
func TestRedispatchStampsTheRow(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
targetID := uuid.New().String()
|
|
s := fSweepSetup(t, targetID, "stamped")
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"stamped":true}`,
|
|
)
|
|
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
rAgePending(t, s.WebhookDB, d.ID)
|
|
|
|
ctx := context.Background()
|
|
|
|
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
|
|
require.Len(t, fDrain(s.Engine), 1)
|
|
|
|
var row database.Delivery
|
|
|
|
require.NoError(t, s.WebhookDB.
|
|
First(&row, "id = ?", d.ID).Error)
|
|
assert.WithinDuration(
|
|
t, time.Now(), row.UpdatedAt, time.Minute,
|
|
"a re-dispatched delivery must be stamped so the "+
|
|
"next tick does not select it again",
|
|
)
|
|
}
|