check / check (push) Waiting to run
withRetry's check on a failed result write could be removed with every test still passing, and so could six other error checks in target_http.go. Each now has a test that fails without it: the circuit breaker learning a failed send whose result went unrecorded, the result write for an invalid config, building the request, reading the response body, decoding the target config, and decoding the stored inbound headers. The two backoff lookups' error checks stay unpinned: without them a failed lookup leaves a zero time, which gives the same answer, so no test can tell the difference. Model: opus-5-5
466 lines
11 KiB
Go
466 lines
11 KiB
Go
package delivery_test
|
|
|
|
import (
|
|
"context"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
"gorm.io/gorm"
|
|
"sneak.berlin/go/webhooker/internal/database"
|
|
"sneak.berlin/go/webhooker/internal/delivery"
|
|
)
|
|
|
|
// These tests cover the delivery half of
|
|
// https://git.eeqj.de/sneak/webhooker/issues/256: a delivery that
|
|
// reached its receiver but whose bookkeeping write failed used to be
|
|
// left at pending and re-sent on the next restart, giving the receiver
|
|
// a second copy while the event log recorded one attempt.
|
|
|
|
// rSeedResult records a DeliveryResult against a delivery, standing in
|
|
// for the attempt row the send path writes before the status.
|
|
func rSeedResult(
|
|
t *testing.T,
|
|
db *gorm.DB,
|
|
deliveryID string,
|
|
attemptNum int,
|
|
success bool,
|
|
) {
|
|
t.Helper()
|
|
|
|
require.NoError(t, db.Create(&database.DeliveryResult{
|
|
DeliveryID: deliveryID,
|
|
AttemptNum: attemptNum,
|
|
Success: success,
|
|
}).Error)
|
|
}
|
|
|
|
// rAgePending backdates a delivery past the sweep's age bound, which is
|
|
// what separates a stranded delivery from one a worker still holds.
|
|
func rAgePending(
|
|
t *testing.T, db *gorm.DB, deliveryID string,
|
|
) {
|
|
t.Helper()
|
|
|
|
old := time.Now().Add(
|
|
-2 * delivery.ExportPendingSweepMinAge,
|
|
)
|
|
|
|
require.NoError(t, db.Model(&database.Delivery{}).
|
|
Where("id = ?", deliveryID).
|
|
UpdateColumn("updated_at", old).Error)
|
|
}
|
|
|
|
func TestRecoverySkipsPendingWithSuccessfulResult(
|
|
t *testing.T,
|
|
) {
|
|
t.Parallel()
|
|
|
|
s := newISetup(t)
|
|
targetID := uuid.New().String()
|
|
|
|
iCreateTarget(t, s.MainDB, targetID,
|
|
s.WebhookID, "already-delivered",
|
|
database.TargetTypeLog, "", 0,
|
|
)
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"delivered":true}`,
|
|
)
|
|
|
|
// The delivery whose send succeeded and whose result row landed:
|
|
// only the status write failed, so it sits at pending.
|
|
done := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
rSeedResult(t, s.WebhookDB, done.ID, 1, true)
|
|
|
|
// A delivery that was genuinely never attempted.
|
|
fresh := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
|
|
s.Engine.ExportRecoverPendingDeliveries(
|
|
context.Background(), s.WebhookDB, s.WebhookID,
|
|
)
|
|
|
|
select {
|
|
case task := <-s.Engine.ExportDeliveryCh():
|
|
assert.Equal(
|
|
t, fresh.ID, task.DeliveryID,
|
|
"only the unattempted delivery may be re-sent",
|
|
)
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("expected the unattempted delivery")
|
|
}
|
|
|
|
select {
|
|
case task := <-s.Engine.ExportDeliveryCh():
|
|
t.Fatalf(
|
|
"re-sent an already delivered delivery: %s",
|
|
task.DeliveryID,
|
|
)
|
|
case <-time.After(200 * time.Millisecond):
|
|
}
|
|
|
|
// It is settled rather than merely skipped: leaving it pending
|
|
// would strand it again on the next sweep.
|
|
iAssertStatus(
|
|
t, s.WebhookDB, done.ID,
|
|
database.DeliveryStatusDelivered,
|
|
)
|
|
}
|
|
|
|
// TestRecoveryContinuesTheAttemptNumbering pins the audit trail: a
|
|
// recovered delivery that already recorded two attempts is re-sent as
|
|
// attempt three, not as attempt one again.
|
|
func TestRecoveryContinuesTheAttemptNumbering(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
s := newISetup(t)
|
|
targetID := uuid.New().String()
|
|
|
|
iCreateTarget(t, s.MainDB, targetID,
|
|
s.WebhookID, "numbering",
|
|
database.TargetTypeLog, "", 0,
|
|
)
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"numbering":true}`,
|
|
)
|
|
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
|
|
rSeedResult(t, s.WebhookDB, d.ID, 1, false)
|
|
rSeedResult(t, s.WebhookDB, d.ID, 2, false)
|
|
|
|
s.Engine.ExportRecoverPendingDeliveries(
|
|
context.Background(), s.WebhookDB, s.WebhookID,
|
|
)
|
|
|
|
select {
|
|
case task := <-s.Engine.ExportDeliveryCh():
|
|
assert.Equal(t, d.ID, task.DeliveryID)
|
|
assert.Equal(t, 3, task.AttemptNum)
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("expected the delivery to be recovered")
|
|
}
|
|
}
|
|
|
|
// TestSweepRecoversStrandedPending is the half that removes the
|
|
// restart requirement: a delivery left at pending is picked up by the
|
|
// periodic sweep.
|
|
func TestSweepRecoversStrandedPending(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
targetID := uuid.New().String()
|
|
s := fSweepSetup(t, targetID, "stranded")
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"stranded":true}`,
|
|
)
|
|
|
|
stranded := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
rAgePending(t, s.WebhookDB, stranded.ID)
|
|
|
|
// A delivery a worker may still be holding: young, and therefore
|
|
// none of the sweep's business.
|
|
inFlight := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
|
|
s.Engine.ExportSweepWebhookRetries(
|
|
context.Background(), s.WebhookID,
|
|
)
|
|
|
|
select {
|
|
case task := <-s.Engine.ExportDeliveryCh():
|
|
assert.Equal(t, stranded.ID, task.DeliveryID)
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("expected the stranded delivery")
|
|
}
|
|
|
|
select {
|
|
case task := <-s.Engine.ExportDeliveryCh():
|
|
t.Fatalf(
|
|
"swept an in-flight delivery: %s",
|
|
task.DeliveryID,
|
|
)
|
|
case <-time.After(200 * time.Millisecond):
|
|
}
|
|
|
|
iAssertStatus(
|
|
t, s.WebhookDB, inFlight.ID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
}
|
|
|
|
// TestSweepClaimsAStrandedDeliveryOnlyOnce guards the repeat the sweep
|
|
// would otherwise be: the row stays pending for as long as the attempt
|
|
// runs, and a sweep a minute later must not send it a second time.
|
|
func TestSweepClaimsAStrandedDeliveryOnlyOnce(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
targetID := uuid.New().String()
|
|
s := fSweepSetup(t, targetID, "claimed")
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"claimed":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)
|
|
|
|
select {
|
|
case task := <-s.Engine.ExportDeliveryCh():
|
|
assert.Equal(t, d.ID, task.DeliveryID)
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("expected the stranded delivery")
|
|
}
|
|
|
|
// The delivery is still pending — nothing has run it yet — but
|
|
// the claim must keep the next sweep off it.
|
|
iAssertStatus(
|
|
t, s.WebhookDB, d.ID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
|
|
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
|
|
|
|
select {
|
|
case task := <-s.Engine.ExportDeliveryCh():
|
|
t.Fatalf(
|
|
"sent a claimed delivery again: %s",
|
|
task.DeliveryID,
|
|
)
|
|
case <-time.After(200 * time.Millisecond):
|
|
}
|
|
}
|
|
|
|
// TestSweepSettlesStrandedPendingWithoutResending is the sweep's own
|
|
// version of the reconcile: a stranded delivery holding a successful
|
|
// result is settled where it stands, and the receiver hears nothing.
|
|
func TestSweepSettlesStrandedPendingWithoutResending(
|
|
t *testing.T,
|
|
) {
|
|
t.Parallel()
|
|
|
|
targetID := uuid.New().String()
|
|
s := fSweepSetup(t, targetID, "settled")
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"settled":true}`,
|
|
)
|
|
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
rSeedResult(t, s.WebhookDB, d.ID, 1, true)
|
|
rAgePending(t, s.WebhookDB, d.ID)
|
|
|
|
s.Engine.ExportSweepWebhookRetries(
|
|
context.Background(), s.WebhookID,
|
|
)
|
|
|
|
select {
|
|
case task := <-s.Engine.ExportDeliveryCh():
|
|
t.Fatalf(
|
|
"re-sent a delivery that already succeeded: %s",
|
|
task.DeliveryID,
|
|
)
|
|
case <-time.After(200 * time.Millisecond):
|
|
}
|
|
|
|
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(1), attempts,
|
|
"settling must not invent an attempt",
|
|
)
|
|
}
|
|
|
|
// TestFailedResultWriteLeavesDeliveryRecoverable is the rule the
|
|
// targets now follow: a bookkeeping write that fails must not advance
|
|
// the status, because pending and retrying are the states the sweeps
|
|
// recover and delivered is a claim the database refused to record.
|
|
func TestFailedResultWriteLeavesDeliveryRecoverable(
|
|
t *testing.T,
|
|
) {
|
|
t.Parallel()
|
|
|
|
s := newISetup(t)
|
|
targetID := uuid.New().String()
|
|
|
|
var hits atomic.Int64
|
|
|
|
ts := httptest.NewServer(http.HandlerFunc(
|
|
func(w http.ResponseWriter, _ *http.Request) {
|
|
hits.Add(1)
|
|
w.WriteHeader(http.StatusOK)
|
|
},
|
|
))
|
|
defer ts.Close()
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"unwritable":true}`,
|
|
)
|
|
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
|
|
// Drop the table the attempt row goes in, so the send succeeds
|
|
// and only the bookkeeping write fails.
|
|
require.NoError(
|
|
t,
|
|
s.WebhookDB.Exec("drop table delivery_results").Error,
|
|
)
|
|
|
|
full := &database.Delivery{
|
|
EventID: event.ID,
|
|
TargetID: targetID,
|
|
Status: database.DeliveryStatusPending,
|
|
Event: event,
|
|
Target: database.Target{
|
|
Name: "unwritable",
|
|
Type: database.TargetTypeHTTP,
|
|
Config: iHTTPConfig(ts.URL),
|
|
},
|
|
}
|
|
full.ID = d.ID
|
|
|
|
s.Engine.ExportDeliverHTTP(
|
|
context.Background(), s.WebhookDB, full,
|
|
&delivery.Task{DeliveryID: d.ID, AttemptNum: 1},
|
|
)
|
|
|
|
assert.Equal(
|
|
t, int64(1), hits.Load(),
|
|
"the send itself must still happen",
|
|
)
|
|
|
|
iAssertStatus(
|
|
t, s.WebhookDB, d.ID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
}
|
|
|
|
// TestFailedResultWriteWithRetriesLeavesDeliveryRecoverable is the same
|
|
// rule for a target with retries: whatever the receiver answered, the
|
|
// delivery stays pending and no retry is scheduled. The circuit breaker
|
|
// still learns the answer, because it describes the target's health,
|
|
// not the database's.
|
|
func TestFailedResultWriteWithRetriesLeavesDeliveryRecoverable(
|
|
t *testing.T,
|
|
) {
|
|
t.Parallel()
|
|
|
|
tests := []struct {
|
|
name string
|
|
answer int
|
|
wantBreaker delivery.CircuitState
|
|
}{
|
|
{"send succeeded", http.StatusOK, delivery.CircuitClosed},
|
|
{"send failed", http.StatusBadGateway, delivery.CircuitOpen},
|
|
}
|
|
|
|
for _, tc := range tests {
|
|
t.Run(tc.name, func(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
s := newISetup(t)
|
|
targetID := uuid.New().String()
|
|
|
|
ts := httptest.NewServer(http.HandlerFunc(
|
|
func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(tc.answer)
|
|
},
|
|
))
|
|
defer ts.Close()
|
|
|
|
event := iSeedEvent(
|
|
t, s.WebhookDB, s.WebhookID, `{"unwritable":true}`,
|
|
)
|
|
|
|
d := iSeedDelivery(
|
|
t, s.WebhookDB, event.ID, targetID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
|
|
require.NoError(
|
|
t,
|
|
s.WebhookDB.Exec("drop table delivery_results").Error,
|
|
)
|
|
|
|
// A single failure opens this breaker.
|
|
cb := delivery.NewTestCircuitBreaker(1, time.Minute)
|
|
s.Engine.ExportSetCircuitBreaker(targetID, cb)
|
|
|
|
full := &database.Delivery{
|
|
EventID: event.ID,
|
|
TargetID: targetID,
|
|
Status: database.DeliveryStatusPending,
|
|
Event: event,
|
|
Target: database.Target{
|
|
Name: "unwritable",
|
|
Type: database.TargetTypeHTTP,
|
|
Config: iHTTPConfig(ts.URL),
|
|
MaxRetries: 3,
|
|
},
|
|
}
|
|
full.ID = d.ID
|
|
|
|
sched := &recordingScheduler{}
|
|
|
|
s.Engine.ExportDeliverHTTPWithScheduler(
|
|
context.Background(), s.WebhookDB, full,
|
|
&delivery.Task{
|
|
DeliveryID: d.ID,
|
|
TargetID: targetID,
|
|
AttemptNum: 1,
|
|
},
|
|
sched,
|
|
)
|
|
|
|
iAssertStatus(
|
|
t, s.WebhookDB, d.ID,
|
|
database.DeliveryStatusPending,
|
|
)
|
|
assert.Empty(t, sched.delays, "no retry may be scheduled")
|
|
assert.Equal(t, tc.wantBreaker, cb.State())
|
|
})
|
|
}
|
|
}
|