2018-11-14 09:18:03 +00:00
|
|
|
package circuitbreaker
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"net/http"
|
2022-04-05 10:30:08 +00:00
|
|
|
"time"
|
2018-11-14 09:18:03 +00:00
|
|
|
|
2022-11-21 17:36:05 +00:00
|
|
|
"github.com/rs/zerolog"
|
|
|
|
"github.com/rs/zerolog/log"
|
2023-02-03 14:24:05 +00:00
|
|
|
"github.com/traefik/traefik/v3/pkg/config/dynamic"
|
|
|
|
"github.com/traefik/traefik/v3/pkg/logs"
|
|
|
|
"github.com/traefik/traefik/v3/pkg/middlewares"
|
2024-03-12 08:48:04 +00:00
|
|
|
"github.com/traefik/traefik/v3/pkg/middlewares/observability"
|
2022-11-21 17:36:05 +00:00
|
|
|
"github.com/vulcand/oxy/v2/cbreaker"
|
2024-01-08 08:10:06 +00:00
|
|
|
"go.opentelemetry.io/otel/trace"
|
2018-11-14 09:18:03 +00:00
|
|
|
)
|
|
|
|
|
2022-04-05 10:30:08 +00:00
|
|
|
const typeName = "CircuitBreaker"
|
2018-11-14 09:18:03 +00:00
|
|
|
|
|
|
|
type circuitBreaker struct {
|
|
|
|
circuitBreaker *cbreaker.CircuitBreaker
|
|
|
|
name string
|
|
|
|
}
|
|
|
|
|
|
|
|
// New creates a new circuit breaker middleware.
|
2019-07-10 07:26:04 +00:00
|
|
|
func New(ctx context.Context, next http.Handler, confCircuitBreaker dynamic.CircuitBreaker, name string) (http.Handler, error) {
|
2018-11-14 09:18:03 +00:00
|
|
|
expression := confCircuitBreaker.Expression
|
|
|
|
|
2022-11-21 17:36:05 +00:00
|
|
|
logger := middlewares.GetLogger(ctx, name, typeName)
|
|
|
|
logger.Debug().Msg("Creating middleware")
|
|
|
|
logger.Debug().Msgf("Setting up with expression: %s", expression)
|
2022-04-05 10:30:08 +00:00
|
|
|
|
2024-01-29 09:58:05 +00:00
|
|
|
responseCode := confCircuitBreaker.ResponseCode
|
|
|
|
|
2022-11-21 17:36:05 +00:00
|
|
|
cbOpts := []cbreaker.Option{
|
2022-04-05 10:30:08 +00:00
|
|
|
cbreaker.Fallback(http.HandlerFunc(func(rw http.ResponseWriter, req *http.Request) {
|
2024-03-12 08:48:04 +00:00
|
|
|
observability.SetStatusErrorf(req.Context(), "blocked by circuit-breaker (%q)", expression)
|
2024-01-29 09:58:05 +00:00
|
|
|
rw.WriteHeader(responseCode)
|
2022-04-05 10:30:08 +00:00
|
|
|
|
2024-01-29 09:58:05 +00:00
|
|
|
if _, err := rw.Write([]byte(http.StatusText(responseCode))); err != nil {
|
2022-11-21 17:36:05 +00:00
|
|
|
log.Ctx(req.Context()).Error().Err(err).Send()
|
2022-04-05 10:30:08 +00:00
|
|
|
}
|
|
|
|
})),
|
2022-11-21 17:36:05 +00:00
|
|
|
cbreaker.Logger(logs.NewOxyWrapper(*logger)),
|
|
|
|
cbreaker.Verbose(logger.GetLevel() == zerolog.TraceLevel),
|
2022-04-05 10:30:08 +00:00
|
|
|
}
|
2018-11-14 09:18:03 +00:00
|
|
|
|
2022-04-05 10:30:08 +00:00
|
|
|
if confCircuitBreaker.CheckPeriod > 0 {
|
|
|
|
cbOpts = append(cbOpts, cbreaker.CheckPeriod(time.Duration(confCircuitBreaker.CheckPeriod)))
|
|
|
|
}
|
|
|
|
|
|
|
|
if confCircuitBreaker.FallbackDuration > 0 {
|
|
|
|
cbOpts = append(cbOpts, cbreaker.FallbackDuration(time.Duration(confCircuitBreaker.FallbackDuration)))
|
|
|
|
}
|
|
|
|
|
|
|
|
if confCircuitBreaker.RecoveryDuration > 0 {
|
|
|
|
cbOpts = append(cbOpts, cbreaker.RecoveryDuration(time.Duration(confCircuitBreaker.RecoveryDuration)))
|
|
|
|
}
|
|
|
|
|
|
|
|
oxyCircuitBreaker, err := cbreaker.New(next, expression, cbOpts...)
|
2018-11-14 09:18:03 +00:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2022-11-21 17:36:05 +00:00
|
|
|
|
2018-11-14 09:18:03 +00:00
|
|
|
return &circuitBreaker{
|
|
|
|
circuitBreaker: oxyCircuitBreaker,
|
|
|
|
name: name,
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2024-01-08 08:10:06 +00:00
|
|
|
func (c *circuitBreaker) GetTracingInformation() (string, string, trace.SpanKind) {
|
|
|
|
return c.name, typeName, trace.SpanKindInternal
|
2018-11-14 09:18:03 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func (c *circuitBreaker) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
|
|
|
|
c.circuitBreaker.ServeHTTP(rw, req)
|
|
|
|
}
|