4 Commits
Author SHA1 Message Date
clawbot 81d758d756 Test a pending delivery whose target is found only on lookup
check / check (push) Successful in 4m31s
When the batch's target map lacks a delivery's target, as it does
for every delivery when the batch query fails, sendRecoveredDeliveries
looks the target up on its own. The new test hands it an empty map
and checks the healthy delivery is queued once, with its real target
id and type.

Model: opus-5-5
2026-09-29 04:15:01 +00:00
clawbot 608c3b21c8 Re-read a delivery before failing it for a missing target
failMissingTarget failed the delivery as the batch had read it, so a
delivery a worker sent and let go between the batch read and the
ownership check could end failed after a successful attempt. It now
reads the row once it owns the delivery and fails it only if the
status is still the one the batch read, as processNewTask does. A row
that cannot be read is left alone.

TestSweepPending_TargetDeleted now checks that the first sweep already
queues the healthy delivery and the second does not queue it again.

Model: opus-5-5
2026-09-29 04:13:05 +00:00
clawbot a0bbfca28d Fail a pending delivery whose target was deleted (closes #293)
Restart recovery and the pending sweep skipped a pending delivery
whose target was missing from the batch's target map, and the sweep
did so again every minute for the life of the database. A miss now
asks loadTarget: no row fails the delivery with a recorded reason,
through the ownership-gated function the retrying paths already use,
renamed failMissingTarget with its log line and reason text made to
fit both statuses. Any other error leaves the delivery pending,
because the map is also empty when its query failed.

Model: opus-5-5
2026-09-29 04:13:05 +00:00
clawbot e0b211f960 Make deliveries refused while half-open wait a cooldown (closes #306)
check / check (push) Successful in 5m0s
While the breaker was half-open, Allow refused every delivery but the
probe and CooldownRemaining returned zero, so each queued task for the
target went straight back onto the retry channel and rewrote its status
on every pass until the probe finished.

CooldownRemaining now returns the whole cooldown while half-open, so a
refused delivery waits that long. A refused delivery already at
retrying is not written again, so the retry counter now moves only
when a refusal moves a delivery into retrying.

Model: opus-5-5
2026-09-29 05:48:20 +02:00
8 changed files with 211 additions and 15 deletions
+5 -3
View File
@@ -2051,7 +2051,7 @@ rescans the database anyway).
| ----------- | -------- | | ----------- | -------- |
| **Closed** | Normal operation. Deliveries flow through. Consecutive failures are counted. | | **Closed** | Normal operation. Deliveries flow through. Consecutive failures are counted. |
| **Open** | Target appears down. Deliveries are skipped and rescheduled for after the cooldown. | | **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:** **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 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 delivery as `retrying` and schedules a retry timer for after the
remaining cooldown period. This ensures no deliveries are lost — they're 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 ### Metrics
@@ -2103,7 +2105,7 @@ arriving and being stored, they are just not getting anywhere.
| Metric | Type | Meaning | | Metric | Type | Meaning |
| ------ | ---- | ------- | | ------ | ---- | ------- |
| `webhooker_events_received_total` | counter | Events received and durably stored. Compare against the delivery counters on one dashboard | | `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_succeeded_total` | counter | Deliveries that reached `delivered` |
| `webhooker_deliveries_failed_total` | counter | Deliveries that failed terminally and will not be retried | | `webhooker_deliveries_failed_total` | counter | Deliveries that failed terminally and will not be retried |
| `webhooker_delivery_retries_total` | counter | Deliveries put back into `retrying` | | `webhooker_delivery_retries_total` | counter | Deliveries put back into `retrying` |
+10 -2
View File
@@ -76,12 +76,20 @@ func (cb *CircuitBreaker) Allow() bool {
} }
} }
// CooldownRemaining returns how much time is left before // CooldownRemaining returns how long a delivery that Allow refused
// an open circuit transitions to half-open. // 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 { func (cb *CircuitBreaker) CooldownRemaining() time.Duration {
cb.mu.Lock() cb.mu.Lock()
defer cb.mu.Unlock() defer cb.mu.Unlock()
if cb.state == CircuitHalfOpen {
return cb.cooldown
}
if cb.state != CircuitOpen { if cb.state != CircuitOpen {
return 0 return 0
} }
+5 -3
View File
@@ -267,7 +267,7 @@ func TestCircuitBreaker_CooldownRemaining_ClosedReturnsZero(
) )
} }
func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero( func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsCooldown(
t *testing.T, t *testing.T,
) { ) {
t.Parallel() t.Parallel()
@@ -282,9 +282,11 @@ func TestCircuitBreaker_CooldownRemaining_HalfOpenReturnsZero(
require.True(t, cb.Allow()) 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(), cb.CooldownRemaining(),
"half-open circuit should have zero cooldown remaining", "a delivery refused while half-open should wait "+
"a whole cooldown",
) )
} }
+96
View File
@@ -17,6 +17,7 @@ import (
"time" "time"
"github.com/google/uuid" "github.com/google/uuid"
"github.com/prometheus/client_golang/prometheus"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require" "github.com/stretchr/testify/require"
"gorm.io/driver/sqlite" "gorm.io/driver/sqlite"
@@ -24,6 +25,7 @@ import (
_ "modernc.org/sqlite" _ "modernc.org/sqlite"
"sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery" "sneak.berlin/go/webhooker/internal/delivery"
"sneak.berlin/go/webhooker/internal/metrics"
) )
// testContentType is the event content type used in tests. // 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) { func TestGetCircuitBreaker_CreatesOnDemand(t *testing.T) {
t.Parallel() t.Parallel()
+36
View File
@@ -101,6 +101,19 @@ func (e *Engine) ExportDeliverHTTP(
e.httpTarget.Deliver(ctx, webhookDB, d, task, e) 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. // ExportDeliverDatabase delivers via the database target.
func (e *Engine) ExportDeliverDatabase( func (e *Engine) ExportDeliverDatabase(
webhookDB *gorm.DB, d *database.Delivery, webhookDB *gorm.DB, d *database.Delivery,
@@ -179,6 +192,14 @@ func (e *Engine) ExportGetCircuitBreaker(
return e.httpTarget.getCircuitBreaker(targetID) 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. // ExportParseHTTPConfig exposes parseHTTPConfig.
func (e *Engine) ExportParseHTTPConfig( func (e *Engine) ExportParseHTTPConfig(
configJSON string, configJSON string,
@@ -331,6 +352,21 @@ func (e *Engine) ExportFailMissingTarget(
e.failMissingTarget(webhookDB, webhookID, d) 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. // ExportDeliveryCh returns the delivery channel.
func (e *Engine) ExportDeliveryCh() chan Task { func (e *Engine) ExportDeliveryCh() chan Task {
return e.deliveryCh return e.deliveryCh
+4 -3
View File
@@ -412,9 +412,10 @@ func TestDeliveryMetrics_BreakerBlockedIsNotAnAttempt(
s.Engine.ExportProcessRetryTask(context.TODO(), &blocked) s.Engine.ExportProcessRetryTask(context.TODO(), &blocked)
// The breaker refused it: rescheduled, so the retry counter // The breaker refused it: rescheduled without rewriting the
// moved, but nothing was attempted or timed. // retrying status it already had, so the retry counter did not
assert.InDelta(t, retriesBefore+1, // move, and nothing was attempted or timed.
assert.InDelta(t, retriesBefore,
mCounter(t, reg, mRetries, mTypeHTTP), 0) mCounter(t, reg, mRetries, mTypeHTTP), 0)
assert.InDelta(t, threshold, assert.InDelta(t, threshold,
mCounter(t, reg, mAttempts, mTypeHTTP), 0) mCounter(t, reg, mAttempts, mTypeHTTP), 0)
+8 -4
View File
@@ -197,10 +197,14 @@ func (c *httpCore) circuitBreakerBlock(
"cooldown_remaining", remaining, "cooldown_remaining", remaining,
) )
c.eng.settleStatus( // A delivery already at retrying is left as it is, so a task
webhookDB, d, d.Target.Type, // the breaker keeps turning away writes nothing each time.
database.DeliveryStatusRetrying, if d.Status != database.DeliveryStatusRetrying {
) c.eng.settleStatus(
webhookDB, d, d.Target.Type,
database.DeliveryStatusRetrying,
)
}
retryTask := *task retryTask := *task
sched.ScheduleRetry(retryTask, remaining) sched.ScheduleRetry(retryTask, remaining)
+47
View File
@@ -696,6 +696,53 @@ func TestSweepPending_TargetDeleted(t *testing.T) {
assert.Contains(t, last.Error, "was deleted") 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 // TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone: a failed
// read of the main database is not a deleted target. Restart recovery // read of the main database is not a deleted target. Restart recovery
// holds every pending delivery of the webhook in one batch, so failing // holds every pending delivery of the webhook in one batch, so failing