Compare commits
4
Commits
| 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. |
|
| **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` |
|
||||||
|
|||||||
@@ -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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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",
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
Reference in New Issue
Block a user