Compare commits
4
Commits
e1e9ba85ea
...
81d758d756
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
81d758d756 | ||
|
|
608c3b21c8 | ||
|
|
a0bbfca28d | ||
|
|
e0b211f960 |
@@ -2051,7 +2051,7 @@ rescans the database anyway).
|
||||
| ----------- | -------- |
|
||||
| **Closed** | Normal operation. Deliveries flow through. Consecutive failures are counted. |
|
||||
| **Open** | Target appears down. Deliveries are skipped and rescheduled for after the cooldown. |
|
||||
| **Half-Open** | Cooldown expired. One probe delivery is allowed to test if the target has recovered. |
|
||||
| **Half-Open** | Cooldown expired. One probe delivery is allowed to test if the target has recovered. Other deliveries are rescheduled for one whole cooldown later. |
|
||||
|
||||
**Transitions:**
|
||||
|
||||
@@ -2089,7 +2089,9 @@ operations), and log targets (stdout) do not use circuit breakers.
|
||||
When a circuit is open and a new delivery arrives, the engine marks the
|
||||
delivery as `retrying` and schedules a retry timer for after the
|
||||
remaining cooldown period. This ensures no deliveries are lost — they're
|
||||
just delayed until the target is healthy again.
|
||||
just delayed until the target is healthy again. A delivery already in
|
||||
`retrying` keeps that status without another database write each time
|
||||
the breaker turns it away.
|
||||
|
||||
### Metrics
|
||||
|
||||
@@ -2103,7 +2105,7 @@ arriving and being stored, they are just not getting anywhere.
|
||||
| Metric | Type | Meaning |
|
||||
| ------ | ---- | ------- |
|
||||
| `webhooker_events_received_total` | counter | Events received and durably stored. Compare against the delivery counters on one dashboard |
|
||||
| `webhooker_delivery_attempts_total` | counter | Delivery attempts actually dispatched to a target. A delivery an open circuit breaker refused is not one: it is counted as a retry instead |
|
||||
| `webhooker_delivery_attempts_total` | counter | Delivery attempts actually dispatched to a target. A delivery a circuit breaker refused is not one: it is counted as a retry instead, but only when the refusal moves it into `retrying` |
|
||||
| `webhooker_deliveries_succeeded_total` | counter | Deliveries that reached `delivered` |
|
||||
| `webhooker_deliveries_failed_total` | counter | Deliveries that failed terminally and will not be retried |
|
||||
| `webhooker_delivery_retries_total` | counter | Deliveries put back into `retrying` |
|
||||
|
||||
@@ -76,12 +76,20 @@ func (cb *CircuitBreaker) Allow() bool {
|
||||
}
|
||||
}
|
||||
|
||||
// CooldownRemaining returns how much time is left before
|
||||
// an open circuit transitions to half-open.
|
||||
// CooldownRemaining returns how long a delivery that Allow refused
|
||||
// should wait before it is tried again. Closed, it returns zero.
|
||||
// Open, it returns what is left of the cooldown, or zero once that
|
||||
// has passed. Half-open, it returns the whole cooldown: the one
|
||||
// probe delivery is still in flight, and if it fails the circuit
|
||||
// reopens for that long.
|
||||
func (cb *CircuitBreaker) CooldownRemaining() time.Duration {
|
||||
cb.mu.Lock()
|
||||
defer cb.mu.Unlock()
|
||||
|
||||
if cb.state == CircuitHalfOpen {
|
||||
return cb.cooldown
|
||||
}
|
||||
|
||||
if cb.state != CircuitOpen {
|
||||
return 0
|
||||
}
|
||||
|
||||
@@ -267,7 +267,7 @@ func TestCircuitBreaker_CooldownRemaining_ClosedReturnsZero(
|
||||
)
|
||||
}
|
||||
|
||||
func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero(
|
||||
func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsCooldown(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
@@ -282,9 +282,11 @@ func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero(
|
||||
|
||||
require.True(t, cb.Allow())
|
||||
|
||||
assert.Equal(t, time.Duration(0),
|
||||
// The cooldown newShortCooldownCB gives the breaker.
|
||||
assert.Equal(t, 50*time.Millisecond,
|
||||
cb.CooldownRemaining(),
|
||||
"half-open circuit should have zero cooldown remaining",
|
||||
"a delivery refused while half-open should wait "+
|
||||
"a whole cooldown",
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
+54
-18
@@ -728,9 +728,7 @@ func (e *Engine) recoverSingleRetry(
|
||||
// webhook on one bad read would be a far larger fault than
|
||||
// the strand it is meant to clear.
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
e.failMissingTargetRetry(
|
||||
webhookDB, webhookID, d,
|
||||
)
|
||||
e.failMissingTarget(webhookDB, webhookID, d)
|
||||
|
||||
return
|
||||
}
|
||||
@@ -1133,9 +1131,7 @@ func (e *Engine) sweepSingleRetry(
|
||||
// Deleted is terminal, unreadable is not; see
|
||||
// recoverSingleRetry.
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
e.failMissingTargetRetry(
|
||||
webhookDB, webhookID, d,
|
||||
)
|
||||
e.failMissingTarget(webhookDB, webhookID, d)
|
||||
|
||||
return
|
||||
}
|
||||
@@ -1249,19 +1245,19 @@ func (e *Engine) failUnretryableRetry(
|
||||
e.failDelivery(webhookDB, d, target.Type, reason)
|
||||
}
|
||||
|
||||
// failMissingTargetRetry terminally fails an orphaned retrying
|
||||
// delivery whose target row is gone. Both restart recovery and the
|
||||
// periodic sweep call it, so the transition exists once.
|
||||
// failMissingTarget terminally fails a recovered delivery, pending or
|
||||
// retrying, whose target row is gone. Restart recovery and the periodic
|
||||
// sweep call it for both statuses, so the transition exists once.
|
||||
//
|
||||
// Until it existed both paths logged the failed lookup and returned,
|
||||
// which left the delivery retrying for the life of the database and
|
||||
// Until it existed those paths logged the failed lookup and moved on,
|
||||
// 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
|
||||
// with a recorded reason is the treatment the other orphaned-retry
|
||||
// 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
|
||||
// action, not a system fault.
|
||||
func (e *Engine) failMissingTargetRetry(
|
||||
func (e *Engine) failMissingTarget(
|
||||
webhookDB *gorm.DB,
|
||||
webhookID string,
|
||||
d *database.Delivery,
|
||||
@@ -1274,13 +1270,37 @@ func (e *Engine) failMissingTargetRetry(
|
||||
|
||||
defer e.inflight.release(d.ID)
|
||||
|
||||
// The batch was read before ownership was taken, and a worker may
|
||||
// have settled the delivery and let it go in between. Only a row
|
||||
// still in the status the batch read is failed.
|
||||
row, err := e.loadDelivery(webhookDB, d.ID)
|
||||
if err != nil {
|
||||
e.log.Error(
|
||||
"failed to load delivery",
|
||||
"delivery_id", d.ID,
|
||||
"error", err,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if row.Status != d.Status {
|
||||
e.log.Debug(
|
||||
"delivery already handled, not failed",
|
||||
"delivery_id", d.ID,
|
||||
"status", row.Status,
|
||||
)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
targetType, reason := e.missingTargetReason(d.TargetID)
|
||||
|
||||
e.log.Warn(
|
||||
"failing orphaned retrying delivery: "+
|
||||
"its target no longer exists",
|
||||
"failing recovered delivery: its target no longer exists",
|
||||
"webhook_id", webhookID,
|
||||
"delivery_id", d.ID,
|
||||
"status", d.Status,
|
||||
"target_id", d.TargetID,
|
||||
"target_type", targetType,
|
||||
)
|
||||
@@ -1314,15 +1334,14 @@ func (e *Engine) missingTargetReason(
|
||||
if err != nil {
|
||||
return "", fmt.Sprintf(
|
||||
"target %s no longer exists; the delivery "+
|
||||
"cannot be retried and has been failed "+
|
||||
"terminally",
|
||||
"has been failed terminally",
|
||||
targetID,
|
||||
)
|
||||
}
|
||||
|
||||
return target.Type, fmt.Sprintf(
|
||||
"target %q (type %s) was deleted; the delivery "+
|
||||
"cannot be retried and has been failed terminally",
|
||||
"has been failed terminally",
|
||||
target.Name, target.Type,
|
||||
)
|
||||
}
|
||||
@@ -2021,14 +2040,31 @@ func (e *Engine) sendRecoveredDeliveries(
|
||||
|
||||
target, ok := targetMap[deliveries[i].TargetID]
|
||||
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(
|
||||
"target not found for delivery",
|
||||
"failed to load target for recovered delivery",
|
||||
"delivery_id", deliveries[i].ID,
|
||||
"target_id", deliveries[i].TargetID,
|
||||
"error", err,
|
||||
)
|
||||
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
if !e.takeForRedispatch(
|
||||
webhookDB, deliveries[i].ID,
|
||||
|
||||
@@ -17,6 +17,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/driver/sqlite"
|
||||
@@ -24,6 +25,7 @@ import (
|
||||
_ "modernc.org/sqlite"
|
||||
"sneak.berlin/go/webhooker/internal/database"
|
||||
"sneak.berlin/go/webhooker/internal/delivery"
|
||||
"sneak.berlin/go/webhooker/internal/metrics"
|
||||
)
|
||||
|
||||
// testContentType is the event content type used in tests.
|
||||
@@ -894,6 +896,100 @@ func TestDeliverHTTP_CircuitBreakerBlocks(t *testing.T) {
|
||||
)
|
||||
}
|
||||
|
||||
// recordingScheduler keeps the delay of every retry it is asked to
|
||||
// schedule, and schedules nothing.
|
||||
type recordingScheduler struct {
|
||||
delays []time.Duration
|
||||
}
|
||||
|
||||
func (s *recordingScheduler) ScheduleRetry(
|
||||
_ delivery.Task, delay time.Duration,
|
||||
) {
|
||||
s.delays = append(s.delays, delay)
|
||||
}
|
||||
|
||||
// TestDeliverHTTP_HalfOpenBreakerDelaysQueuedTasks proves that while a
|
||||
// half-open breaker's one probe delivery is in flight, every other task
|
||||
// for the target is put back with a whole cooldown as its delay rather
|
||||
// than none, and that its status is written the first time the breaker
|
||||
// turns it away and not on each pass after that.
|
||||
func TestDeliverHTTP_HalfOpenBreakerDelaysQueuedTasks(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db := testWebhookDB(t)
|
||||
e := testEngine(t, 1)
|
||||
|
||||
// Every write of retrying moves the retry counter, so on a registry
|
||||
// this test owns the counter is the number of those writes.
|
||||
reg := prometheus.NewRegistry()
|
||||
e.ExportSetMetrics(metrics.New(reg))
|
||||
|
||||
targetID := uuid.New().String()
|
||||
cb := newShortCooldownCB(t)
|
||||
e.ExportSetCircuitBreaker(targetID, cb)
|
||||
|
||||
for range delivery.ExportDefaultFailureThreshold {
|
||||
cb.RecordFailure()
|
||||
}
|
||||
|
||||
time.Sleep(60 * time.Millisecond)
|
||||
|
||||
require.True(t, cb.Allow(), "the probe delivery should go through")
|
||||
require.Equal(t, delivery.CircuitHalfOpen, cb.State())
|
||||
|
||||
cfg := newHTTPTargetConfig(
|
||||
"http://will-not-be-called.invalid",
|
||||
)
|
||||
sched := &recordingScheduler{}
|
||||
|
||||
const queued, passes = 3, 4
|
||||
|
||||
for range queued {
|
||||
event := seedEvent(t, db, `{"cb":"half-open"}`)
|
||||
dlv := seedDelivery(
|
||||
t, db, event.ID, targetID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
for range passes {
|
||||
// Each pass starts from the stored row, as a retry does.
|
||||
var row database.Delivery
|
||||
|
||||
require.NoError(t, db.First(
|
||||
&row, "id = ?", dlv.ID,
|
||||
).Error)
|
||||
|
||||
fix := buildHTTPFixture(
|
||||
row, event, targetID,
|
||||
"test-cb-half-open", cfg, 5, 1,
|
||||
)
|
||||
|
||||
e.ExportDeliverHTTPWithScheduler(
|
||||
context.TODO(), db, fix.Delivery, fix.Task, sched,
|
||||
)
|
||||
}
|
||||
|
||||
assertDeliveryStatus(t, db, dlv.ID,
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
}
|
||||
|
||||
require.Len(t, sched.delays, queued*passes)
|
||||
|
||||
for _, delay := range sched.delays {
|
||||
// The cooldown newShortCooldownCB gives the breaker.
|
||||
assert.Equal(t, 50*time.Millisecond, delay,
|
||||
"a task turned away while half-open should wait "+
|
||||
"a whole cooldown",
|
||||
)
|
||||
}
|
||||
|
||||
assert.InDelta(t, float64(queued),
|
||||
mCounter(t, reg, mRetries, mTypeHTTP), 0,
|
||||
"status should be written once per task, not once per pass",
|
||||
)
|
||||
}
|
||||
|
||||
func TestGetCircuitBreaker_CreatesOnDemand(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
@@ -101,6 +101,19 @@ func (e *Engine) ExportDeliverHTTP(
|
||||
e.httpTarget.Deliver(ctx, webhookDB, d, task, e)
|
||||
}
|
||||
|
||||
// ExportDeliverHTTPWithScheduler delivers via the http target, handing
|
||||
// any retry to sched instead of the engine, so a test can see the
|
||||
// delay each retry is given.
|
||||
func (e *Engine) ExportDeliverHTTPWithScheduler(
|
||||
ctx context.Context,
|
||||
webhookDB *gorm.DB,
|
||||
d *database.Delivery,
|
||||
task *Task,
|
||||
sched Scheduler,
|
||||
) {
|
||||
e.httpTarget.Deliver(ctx, webhookDB, d, task, sched)
|
||||
}
|
||||
|
||||
// ExportDeliverDatabase delivers via the database target.
|
||||
func (e *Engine) ExportDeliverDatabase(
|
||||
webhookDB *gorm.DB, d *database.Delivery,
|
||||
@@ -179,6 +192,14 @@ func (e *Engine) ExportGetCircuitBreaker(
|
||||
return e.httpTarget.getCircuitBreaker(targetID)
|
||||
}
|
||||
|
||||
// ExportSetCircuitBreaker makes cb the http target's circuit breaker
|
||||
// for targetID, so a test can use one with a short cooldown.
|
||||
func (e *Engine) ExportSetCircuitBreaker(
|
||||
targetID string, cb *CircuitBreaker,
|
||||
) {
|
||||
e.httpTarget.circuitBreakers.Store(targetID, cb)
|
||||
}
|
||||
|
||||
// ExportParseHTTPConfig exposes parseHTTPConfig.
|
||||
func (e *Engine) ExportParseHTTPConfig(
|
||||
configJSON string,
|
||||
@@ -321,6 +342,31 @@ func (e *Engine) ExportRecoverRetryingDeliveries(
|
||||
e.recoverRetryingDeliveries(webhookDB, webhookID)
|
||||
}
|
||||
|
||||
// ExportFailMissingTarget exposes failMissingTarget, so a test can hand
|
||||
// it a delivery as a batch read it earlier.
|
||||
func (e *Engine) ExportFailMissingTarget(
|
||||
webhookDB *gorm.DB,
|
||||
webhookID string,
|
||||
d *database.Delivery,
|
||||
) {
|
||||
e.failMissingTarget(webhookDB, webhookID, d)
|
||||
}
|
||||
|
||||
// ExportSendRecoveredDeliveries exposes sendRecoveredDeliveries, so a
|
||||
// test can hand it a target map that lacks a delivery's target.
|
||||
func (e *Engine) ExportSendRecoveredDeliveries(
|
||||
ctx context.Context,
|
||||
webhookDB *gorm.DB,
|
||||
deliveries []database.Delivery,
|
||||
webhookID string,
|
||||
targetMap map[string]database.Target,
|
||||
settled map[string]struct{},
|
||||
) {
|
||||
e.sendRecoveredDeliveries(
|
||||
ctx, webhookDB, deliveries, webhookID, targetMap, settled,
|
||||
)
|
||||
}
|
||||
|
||||
// ExportDeliveryCh returns the delivery channel.
|
||||
func (e *Engine) ExportDeliveryCh() chan Task {
|
||||
return e.deliveryCh
|
||||
|
||||
@@ -412,9 +412,10 @@ func TestDeliveryMetrics_BreakerBlockedIsNotAnAttempt(
|
||||
|
||||
s.Engine.ExportProcessRetryTask(context.TODO(), &blocked)
|
||||
|
||||
// The breaker refused it: rescheduled, so the retry counter
|
||||
// moved, but nothing was attempted or timed.
|
||||
assert.InDelta(t, retriesBefore+1,
|
||||
// The breaker refused it: rescheduled without rewriting the
|
||||
// retrying status it already had, so the retry counter did not
|
||||
// move, and nothing was attempted or timed.
|
||||
assert.InDelta(t, retriesBefore,
|
||||
mCounter(t, reg, mRetries, mTypeHTTP), 0)
|
||||
assert.InDelta(t, threshold,
|
||||
mCounter(t, reg, mAttempts, mTypeHTTP), 0)
|
||||
|
||||
@@ -197,10 +197,14 @@ func (c *httpCore) circuitBreakerBlock(
|
||||
"cooldown_remaining", remaining,
|
||||
)
|
||||
|
||||
// A delivery already at retrying is left as it is, so a task
|
||||
// the breaker keeps turning away writes nothing each time.
|
||||
if d.Status != database.DeliveryStatusRetrying {
|
||||
c.eng.settleStatus(
|
||||
webhookDB, d, d.Target.Type,
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
}
|
||||
|
||||
retryTask := *task
|
||||
sched.ScheduleRetry(retryTask, remaining)
|
||||
|
||||
@@ -18,16 +18,17 @@ import (
|
||||
// https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed
|
||||
// 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
|
||||
// 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
|
||||
// a target whose type was written by a build that knew a type this one
|
||||
// does not.
|
||||
const tUnknownType = database.TargetType("pubsub")
|
||||
|
||||
// tSeedDeletedTarget creates a target, a retrying delivery against it
|
||||
// with one recorded failed attempt, and then deletes the target the
|
||||
// way the source page does.
|
||||
// tSeedDeletedTarget creates a target, a delivery against it at the
|
||||
// given status with one recorded failed attempt, and then deletes the
|
||||
// target the way the source page does.
|
||||
//
|
||||
// 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
|
||||
@@ -36,6 +37,7 @@ func tSeedDeletedTarget(
|
||||
t *testing.T,
|
||||
s iSetup,
|
||||
name, url string,
|
||||
status database.DeliveryStatus,
|
||||
) string {
|
||||
t.Helper()
|
||||
|
||||
@@ -51,8 +53,7 @@ func tSeedDeletedTarget(
|
||||
)
|
||||
|
||||
d := iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, targetID,
|
||||
database.DeliveryStatusRetrying,
|
||||
t, s.WebhookDB, event.ID, targetID, status,
|
||||
)
|
||||
|
||||
iSeedFailedResult(t, s.WebhookDB, d.ID)
|
||||
@@ -173,6 +174,7 @@ func TestRecoverSingleRetry_TargetDeleted(t *testing.T) {
|
||||
|
||||
deliveryID := tSeedDeletedTarget(
|
||||
t, s, "gone-on-recovery", "http://example.com/hook",
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
|
||||
s.Engine.ExportRecoverWebhookDeliveries(
|
||||
@@ -210,6 +212,7 @@ func TestSweepSingleRetry_TargetDeleted(t *testing.T) {
|
||||
|
||||
deliveryID := tSeedDeletedTarget(
|
||||
t, s, "gone-on-sweep", "http://example.com/hook",
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
|
||||
// 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
|
||||
// path to the same rule as the existing one: no target row, and so no
|
||||
// TestFailMissingTarget_WritesNoTargetRow holds the new terminal path
|
||||
// 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
|
||||
// database. See https://git.eeqj.de/sneak/webhooker/issues/206.
|
||||
func TestFailMissingTargetRetry_WritesNoTargetRow(
|
||||
func TestFailMissingTarget_WritesNoTargetRow(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
@@ -297,6 +300,7 @@ func TestFailMissingTargetRetry_WritesNoTargetRow(
|
||||
|
||||
deliveryID := tSeedDeletedTarget(
|
||||
t, s, "credential-bearing", hookURL,
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
|
||||
s.Engine.ExportSweepWebhookRetries(
|
||||
@@ -529,3 +533,260 @@ func TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone(
|
||||
|
||||
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",
|
||||
)
|
||||
}
|
||||
|
||||
// TestFailMissingTarget_LeavesASettledDeliveryAlone: the recovery paths
|
||||
// read their batch before taking ownership, and a worker may send a
|
||||
// delivery and let it go in between. The terminal write goes by the row
|
||||
// as it is now, not as the batch read it.
|
||||
func TestFailMissingTarget_LeavesASettledDeliveryAlone(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
s := newISetup(t)
|
||||
|
||||
deliveryID := tSeedDeletedTarget(
|
||||
t, s, "gone-after-sending", "http://example.com/hook",
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
var batch database.Delivery
|
||||
|
||||
require.NoError(t, s.WebhookDB.First(
|
||||
&batch, "id = ?", deliveryID,
|
||||
).Error)
|
||||
|
||||
// A worker settles the delivery after the batch was read.
|
||||
require.NoError(t, s.WebhookDB.Model(&database.Delivery{}).
|
||||
Where("id = ?", deliveryID).
|
||||
Update("status", database.DeliveryStatusDelivered).Error)
|
||||
|
||||
s.Engine.ExportFailMissingTarget(
|
||||
s.WebhookDB, s.WebhookID, &batch,
|
||||
)
|
||||
|
||||
iAssertStatus(
|
||||
t, s.WebhookDB, deliveryID,
|
||||
database.DeliveryStatusDelivered,
|
||||
)
|
||||
|
||||
assert.Len(t, iResults(t, s.WebhookDB, deliveryID), 1,
|
||||
"a delivery settled after the batch read was then failed",
|
||||
)
|
||||
assert.Zero(t, s.Engine.ExportInflightHeld())
|
||||
}
|
||||
|
||||
// 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 queued by the
|
||||
// first sweep and not again by the second.
|
||||
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)
|
||||
|
||||
tasks := fDrain(s.Engine)
|
||||
require.Len(t, tasks, 1,
|
||||
"the first sweep did not queue the healthy delivery",
|
||||
)
|
||||
assert.Equal(t, healthy.ID, tasks[0].DeliveryID)
|
||||
|
||||
s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID)
|
||||
|
||||
assert.Empty(t, fDrain(s.Engine),
|
||||
"the second sweep queued a delivery again",
|
||||
)
|
||||
|
||||
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")
|
||||
}
|
||||
|
||||
// TestSendRecoveredDeliveries_TargetMissingFromMap: the batch's target
|
||||
// map is empty when its query failed, so every delivery in the batch is
|
||||
// looked up on its own. A healthy one is sent to the target that lookup
|
||||
// finds.
|
||||
func TestSendRecoveredDeliveries_TargetMissingFromMap(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
s := newISetup(t)
|
||||
|
||||
targetID := uuid.New().String()
|
||||
|
||||
iCreateTarget(
|
||||
t, s.MainDB, targetID, s.WebhookID, "found-on-lookup",
|
||||
database.TargetTypeLog, "", 0,
|
||||
)
|
||||
|
||||
event := iSeedEvent(
|
||||
t, s.WebhookDB, s.WebhookID, `{"map":"empty"}`,
|
||||
)
|
||||
|
||||
d := iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, targetID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
s.Engine.ExportSendRecoveredDeliveries(
|
||||
context.Background(), s.WebhookDB,
|
||||
[]database.Delivery{d}, s.WebhookID,
|
||||
map[string]database.Target{}, nil,
|
||||
)
|
||||
|
||||
tasks := fDrain(s.Engine)
|
||||
require.Len(t, tasks, 1,
|
||||
"the healthy delivery was not queued exactly once",
|
||||
)
|
||||
assert.Equal(t, d.ID, tasks[0].DeliveryID)
|
||||
assert.Equal(t, targetID, tasks[0].TargetID)
|
||||
assert.Equal(t, database.TargetTypeLog, tasks[0].TargetType)
|
||||
|
||||
iAssertStatus(
|
||||
t, s.WebhookDB, d.ID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
}
|
||||
|
||||
// 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