Show a target paused by its circuit breaker (closes #385)
check / check (push) Successful in 3m21s

While an http or slack target's circuit breaker is open, the target's
row on the webhook page says its deliveries are paused until the
cooldown ends, in UTC and from now. Each of its retrying deliveries
shows as waiting in the event log and on the event's page, until the
later of the cooldown's end and the end of its own backoff. While the
breaker is half-open, the row says deliveries are held while one
delivery tests the target, with no time, and deliveries keep their
plain status.

The engine gains one read, StateAndCooldown(targetID), taking a
breaker's state and remaining cooldown under one lock; the handlers
reach it through a one-method interface wired like Archives.

Model: opus-5-5
This commit is contained in:
2026-10-02 23:16:44 +00:00
parent 3489d6909a
commit d66292df43
19 changed files with 520 additions and 28 deletions
+14
View File
@@ -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() {
+38 -3
View File
@@ -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,30 @@ func (e *Engine) Rename(
return e.dbTarget.rename(targetID, webhookName, targetName)
}
// StateAndCooldown implements CircuitBreakers. It is
// CircuitBreaker.StateAndCooldown for the target's breaker. While the
// breaker is open, the pages show the target's deliveries as paused
// until its cooldown ends; while it is half-open, they show them as
// held, with no time, while one delivery tests whether the target has
// recovered. A target with no breaker reads as closed, 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.
+56
View File
@@ -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()
+8
View File
@@ -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,
+1
View File
@@ -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{
+6 -4
View File
@@ -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)