check / check (push) Successful in 3m20s
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, then one is sent to test the target. Each of its retrying deliveries shows as waiting in the event log and on the event's page, with the earliest it can be tried next: the later of the cooldown's end and the end of its own backoff, dated when not today in UTC. 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
175 lines
3.7 KiB
Go
175 lines
3.7 KiB
Go
package delivery
|
|
|
|
import (
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// CircuitState represents the current state of a circuit
|
|
// breaker.
|
|
type CircuitState int
|
|
|
|
const (
|
|
// CircuitClosed is the normal operating state.
|
|
CircuitClosed CircuitState = iota
|
|
// CircuitOpen means the circuit has tripped.
|
|
CircuitOpen
|
|
// CircuitHalfOpen allows a single probe delivery to
|
|
// test whether the target has recovered.
|
|
CircuitHalfOpen
|
|
)
|
|
|
|
const (
|
|
// defaultFailureThreshold is the number of consecutive
|
|
// failures before a circuit breaker trips open.
|
|
defaultFailureThreshold = 5
|
|
|
|
// defaultCooldown is how long a circuit stays open
|
|
// before transitioning to half-open.
|
|
defaultCooldown = 30 * time.Second
|
|
)
|
|
|
|
// CircuitBreaker implements the circuit breaker pattern
|
|
// for a single delivery target.
|
|
type CircuitBreaker struct {
|
|
mu sync.Mutex
|
|
state CircuitState
|
|
failures int
|
|
threshold int
|
|
cooldown time.Duration
|
|
lastFailure time.Time
|
|
}
|
|
|
|
// NewCircuitBreaker creates a circuit breaker with default
|
|
// settings.
|
|
func NewCircuitBreaker() *CircuitBreaker {
|
|
return &CircuitBreaker{
|
|
state: CircuitClosed,
|
|
threshold: defaultFailureThreshold,
|
|
cooldown: defaultCooldown,
|
|
}
|
|
}
|
|
|
|
// Allow checks whether a delivery attempt should proceed.
|
|
func (cb *CircuitBreaker) Allow() bool {
|
|
cb.mu.Lock()
|
|
defer cb.mu.Unlock()
|
|
|
|
switch cb.state {
|
|
case CircuitClosed:
|
|
return true
|
|
|
|
case CircuitOpen:
|
|
if time.Since(cb.lastFailure) >= cb.cooldown {
|
|
cb.state = CircuitHalfOpen
|
|
|
|
return true
|
|
}
|
|
|
|
return false
|
|
|
|
case CircuitHalfOpen:
|
|
return false
|
|
|
|
default:
|
|
return true
|
|
}
|
|
}
|
|
|
|
// CooldownRemaining returns how long a delivery that Allow refused
|
|
// 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 {
|
|
cb.mu.Lock()
|
|
defer cb.mu.Unlock()
|
|
|
|
if cb.state == CircuitHalfOpen {
|
|
return cb.cooldown
|
|
}
|
|
|
|
if cb.state != CircuitOpen {
|
|
return 0
|
|
}
|
|
|
|
remaining := cb.cooldown - time.Since(cb.lastFailure)
|
|
if remaining < 0 {
|
|
return 0
|
|
}
|
|
|
|
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() {
|
|
cb.mu.Lock()
|
|
defer cb.mu.Unlock()
|
|
|
|
cb.failures = 0
|
|
cb.state = CircuitClosed
|
|
}
|
|
|
|
// RecordFailure records a failed delivery. If the failure
|
|
// count reaches the threshold, the circuit trips open.
|
|
func (cb *CircuitBreaker) RecordFailure() {
|
|
cb.mu.Lock()
|
|
defer cb.mu.Unlock()
|
|
|
|
cb.failures++
|
|
cb.lastFailure = time.Now()
|
|
|
|
switch cb.state {
|
|
case CircuitClosed:
|
|
if cb.failures >= cb.threshold {
|
|
cb.state = CircuitOpen
|
|
}
|
|
|
|
case CircuitOpen:
|
|
// Already open; no state change needed.
|
|
|
|
case CircuitHalfOpen:
|
|
// Probe failed -- reopen immediately.
|
|
cb.state = CircuitOpen
|
|
}
|
|
}
|
|
|
|
// State returns the current circuit state.
|
|
func (cb *CircuitBreaker) State() CircuitState {
|
|
cb.mu.Lock()
|
|
defer cb.mu.Unlock()
|
|
|
|
return cb.state
|
|
}
|
|
|
|
// String returns the human-readable name of a circuit
|
|
// state.
|
|
func (s CircuitState) String() string {
|
|
switch s {
|
|
case CircuitClosed:
|
|
return "closed"
|
|
case CircuitOpen:
|
|
return "open"
|
|
case CircuitHalfOpen:
|
|
return "half-open"
|
|
default:
|
|
return "unknown"
|
|
}
|
|
}
|