Show a target paused by its circuit breaker (closes #385)
check / check (push) Successful in 3m16s
check / check (push) Successful in 3m16s
A target whose circuit breaker had tripped still showed as Active, and its deliveries sat at a bare "retrying" with no attempts. Its row on the webhook page now says deliveries are paused after repeated failures and until when the cooldown ends, adding that one waiting delivery is then sent to test the target; while half-open it says deliveries are held while one tests it, with no time. Waiting deliveries show "next try no earlier than" the later of the cooldown and their own backoff, with the date when not today. The engine gains one read of a breaker's state and remaining cooldown under one lock, and shares the backoff formula. Model: opus-5-5
This commit was merged in pull request #478.
This commit is contained in:
@@ -102,6 +102,20 @@ func (cb *CircuitBreaker) CooldownRemaining() time.Duration {
|
||||
return remaining
|
||||
}
|
||||
|
||||
// StateAndCooldown returns the circuit state and, while the circuit is
|
||||
// open, what is left of the cooldown, or zero once that has passed.
|
||||
// Both are read under one lock, so they always agree.
|
||||
func (cb *CircuitBreaker) StateAndCooldown() (CircuitState, time.Duration) {
|
||||
cb.mu.Lock()
|
||||
defer cb.mu.Unlock()
|
||||
|
||||
if cb.state != CircuitOpen {
|
||||
return cb.state, 0
|
||||
}
|
||||
|
||||
return cb.state, max(cb.cooldown-time.Since(cb.lastFailure), 0)
|
||||
}
|
||||
|
||||
// RecordSuccess records a successful delivery and resets
|
||||
// the circuit breaker to closed state.
|
||||
func (cb *CircuitBreaker) RecordSuccess() {
|
||||
|
||||
@@ -143,6 +143,15 @@ type Archives interface {
|
||||
Rename(targetID, webhookName, targetName string) error
|
||||
}
|
||||
|
||||
// CircuitBreakers is how the handlers read a target's circuit
|
||||
// breaker, so the webhook page and the event log can say that
|
||||
// deliveries to the target are paused and until when. Like Archives,
|
||||
// it keeps the handlers free of the engine's internals and is
|
||||
// trivially faked in tests.
|
||||
type CircuitBreakers interface {
|
||||
StateAndCooldown(targetID string) (CircuitState, time.Duration)
|
||||
}
|
||||
|
||||
// EngineParams are the fx dependencies for the delivery
|
||||
// engine.
|
||||
type EngineParams struct {
|
||||
@@ -186,9 +195,11 @@ type Engine struct {
|
||||
// targets maps each target type to its implementation.
|
||||
targets map[database.TargetType]Target
|
||||
|
||||
// httpTarget is retained so tests can reach the HTTP
|
||||
// target's shared client and circuit breakers.
|
||||
httpTarget *httpTarget
|
||||
// httpTarget and slackTarget are retained so StateAndCooldown
|
||||
// can read their circuit breakers, and so tests can reach the
|
||||
// HTTP target's shared client.
|
||||
httpTarget *httpTarget
|
||||
slackTarget *slackTarget
|
||||
|
||||
// dbTarget is retained so the engine can reach the archive
|
||||
// writer registry for eviction, renames and the idle sweep.
|
||||
@@ -300,6 +311,28 @@ func (e *Engine) Rename(
|
||||
return e.dbTarget.rename(targetID, webhookName, targetName)
|
||||
}
|
||||
|
||||
// StateAndCooldown implements CircuitBreakers. It returns the state of
|
||||
// the target's circuit breaker and, while the breaker is open, what is
|
||||
// left of its cooldown; the cooldown is zero once that has passed and
|
||||
// in any other state. A target with no breaker reads as closed with no
|
||||
// cooldown, and reading never creates one.
|
||||
func (e *Engine) StateAndCooldown(
|
||||
targetID string,
|
||||
) (CircuitState, time.Duration) {
|
||||
for _, core := range []*httpCore{
|
||||
e.httpTarget.httpCore, e.slackTarget.httpCore,
|
||||
} {
|
||||
val, ok := core.circuitBreakers.Load(targetID)
|
||||
if ok {
|
||||
cb, _ := val.(*CircuitBreaker)
|
||||
|
||||
return cb.StateAndCooldown()
|
||||
}
|
||||
}
|
||||
|
||||
return CircuitClosed, 0
|
||||
}
|
||||
|
||||
// ScheduleRetry schedules a task to be re-enqueued onto the
|
||||
// retry channel after delay. It implements the Scheduler
|
||||
// interface the targets use to own their durable retries.
|
||||
|
||||
@@ -1018,6 +1018,62 @@ func TestGetCircuitBreaker_CreatesOnDemand(t *testing.T) {
|
||||
)
|
||||
}
|
||||
|
||||
// TestStateAndCooldown_ReadsHTTPAndSlackBreakers proves the engine
|
||||
// reads the state of an http or a slack target's circuit breaker, with
|
||||
// what is left of its cooldown while it is open, and no cooldown while
|
||||
// it is half-open, once it closes, or for a target with no breaker.
|
||||
func TestStateAndCooldown_ReadsHTTPAndSlackBreakers(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
e := testEngine(t, 1)
|
||||
|
||||
httpID := uuid.New().String()
|
||||
slackID := uuid.New().String()
|
||||
|
||||
state, cooldown := e.StateAndCooldown(httpID)
|
||||
assert.Equal(t, delivery.CircuitClosed, state, "no breaker")
|
||||
assert.Zero(t, cooldown, "no breaker")
|
||||
|
||||
httpCB := delivery.NewTestCircuitBreaker(1, time.Hour)
|
||||
e.ExportSetCircuitBreaker(httpID, httpCB)
|
||||
|
||||
slackCB := delivery.NewTestCircuitBreaker(1, time.Hour)
|
||||
e.ExportSetSlackCircuitBreaker(slackID, slackCB)
|
||||
|
||||
httpCB.RecordFailure()
|
||||
slackCB.RecordFailure()
|
||||
|
||||
for _, id := range []string{httpID, slackID} {
|
||||
state, cooldown := e.StateAndCooldown(id)
|
||||
assert.Equal(t, delivery.CircuitOpen, state)
|
||||
assert.Greater(t, cooldown, 59*time.Minute)
|
||||
assert.LessOrEqual(t, cooldown, time.Hour)
|
||||
}
|
||||
|
||||
httpCB.RecordSuccess()
|
||||
slackCB.RecordSuccess()
|
||||
|
||||
for _, id := range []string{httpID, slackID} {
|
||||
state, cooldown := e.StateAndCooldown(id)
|
||||
assert.Equal(t, delivery.CircuitClosed, state, "closed")
|
||||
assert.Zero(t, cooldown, "closed")
|
||||
}
|
||||
|
||||
// A breaker with no cooldown goes half-open on the first Allow
|
||||
// after it trips, letting that one delivery through to test the
|
||||
// target.
|
||||
halfOpenID := uuid.New().String()
|
||||
halfOpenCB := delivery.NewTestCircuitBreaker(1, 0)
|
||||
e.ExportSetCircuitBreaker(halfOpenID, halfOpenCB)
|
||||
|
||||
halfOpenCB.RecordFailure()
|
||||
require.True(t, halfOpenCB.Allow())
|
||||
|
||||
state, cooldown = e.StateAndCooldown(halfOpenID)
|
||||
assert.Equal(t, delivery.CircuitHalfOpen, state)
|
||||
assert.Zero(t, cooldown, "half-open")
|
||||
}
|
||||
|
||||
func TestParseHTTPConfig_Valid(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
@@ -212,6 +212,14 @@ func (e *Engine) ExportSetCircuitBreaker(
|
||||
e.httpTarget.circuitBreakers.Store(targetID, cb)
|
||||
}
|
||||
|
||||
// ExportSetSlackCircuitBreaker is ExportSetCircuitBreaker for the
|
||||
// slack target.
|
||||
func (e *Engine) ExportSetSlackCircuitBreaker(
|
||||
targetID string, cb *CircuitBreaker,
|
||||
) {
|
||||
e.slackTarget.circuitBreakers.Store(targetID, cb)
|
||||
}
|
||||
|
||||
// ExportParseHTTPConfig exposes parseHTTPConfig.
|
||||
func (e *Engine) ExportParseHTTPConfig(
|
||||
configJSON string,
|
||||
|
||||
@@ -105,6 +105,7 @@ func (e *Engine) initTargets(client *http.Client) {
|
||||
dbT := &databaseTarget{eng: e}
|
||||
|
||||
e.httpTarget = httpT
|
||||
e.slackTarget = slackT
|
||||
e.dbTarget = dbT
|
||||
|
||||
e.targets = map[database.TargetType]Target{
|
||||
|
||||
@@ -234,7 +234,7 @@ func (c *httpCore) handleRetry(
|
||||
database.DeliveryStatusRetrying,
|
||||
)
|
||||
|
||||
backoff := calcBackoff(attemptNum)
|
||||
backoff := Backoff(attemptNum)
|
||||
|
||||
retryTask := *task
|
||||
retryTask.AttemptNum = attemptNum + 1
|
||||
@@ -301,7 +301,7 @@ func (c *httpCore) remainingBackoff(
|
||||
return 0
|
||||
}
|
||||
|
||||
backoff := calcBackoff(attemptNum)
|
||||
backoff := Backoff(attemptNum)
|
||||
elapsed := time.Since(lastResult.CreatedAt)
|
||||
remaining := backoff - elapsed
|
||||
|
||||
@@ -326,12 +326,14 @@ func (c *httpCore) backoffElapsed(
|
||||
return true
|
||||
}
|
||||
|
||||
backoff := calcBackoff(attemptNum)
|
||||
backoff := Backoff(attemptNum)
|
||||
|
||||
return time.Since(lastResult.CreatedAt) >= backoff
|
||||
}
|
||||
|
||||
func calcBackoff(attemptNum int) time.Duration {
|
||||
// Backoff is how long an http or slack target with retries waits after
|
||||
// a delivery's failed attempt attemptNum before trying it again.
|
||||
func Backoff(attemptNum int) time.Duration {
|
||||
shift := max(attemptNum-1, 0)
|
||||
shift = min(shift, maxBackoffShift)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user