Deploy: main into prod #343
@@ -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,
|
||||||
|
|||||||
@@ -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,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// 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(
|
c.eng.settleStatus(
|
||||||
webhookDB, d, d.Target.Type,
|
webhookDB, d, d.Target.Type,
|
||||||
database.DeliveryStatusRetrying,
|
database.DeliveryStatusRetrying,
|
||||||
)
|
)
|
||||||
|
}
|
||||||
|
|
||||||
retryTask := *task
|
retryTask := *task
|
||||||
sched.ScheduleRetry(retryTask, remaining)
|
sched.ScheduleRetry(retryTask, remaining)
|
||||||
|
|||||||
Reference in New Issue
Block a user