2019-08-26 12:20:06 +02:00
|
|
|
package inflightreq
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"fmt"
|
|
|
|
"net/http"
|
|
|
|
|
|
|
|
"github.com/opentracing/opentracing-go/ext"
|
2022-11-21 18:36:05 +01:00
|
|
|
"github.com/rs/zerolog"
|
2020-09-16 15:46:04 +02:00
|
|
|
"github.com/traefik/traefik/v2/pkg/config/dynamic"
|
2022-11-21 18:36:05 +01:00
|
|
|
"github.com/traefik/traefik/v2/pkg/logs"
|
2020-09-16 15:46:04 +02:00
|
|
|
"github.com/traefik/traefik/v2/pkg/middlewares"
|
|
|
|
"github.com/traefik/traefik/v2/pkg/tracing"
|
2022-11-21 18:36:05 +01:00
|
|
|
"github.com/vulcand/oxy/v2/connlimit"
|
2019-08-26 12:20:06 +02:00
|
|
|
)
|
|
|
|
|
|
|
|
const (
|
|
|
|
typeName = "InFlightReq"
|
|
|
|
)
|
|
|
|
|
|
|
|
type inFlightReq struct {
|
|
|
|
handler http.Handler
|
|
|
|
name string
|
|
|
|
}
|
|
|
|
|
|
|
|
// New creates a max request middleware.
|
2020-04-29 18:32:05 +02:00
|
|
|
// If no source criterion is provided in the config, it defaults to RequestHost.
|
2019-08-26 12:20:06 +02:00
|
|
|
func New(ctx context.Context, next http.Handler, config dynamic.InFlightReq, name string) (http.Handler, error) {
|
2022-11-21 18:36:05 +01:00
|
|
|
logger := middlewares.GetLogger(ctx, name, typeName)
|
|
|
|
logger.Debug().Msg("Creating middleware")
|
|
|
|
|
|
|
|
ctxLog := logger.WithContext(ctx)
|
2019-08-26 12:20:06 +02:00
|
|
|
|
|
|
|
if config.SourceCriterion == nil ||
|
|
|
|
config.SourceCriterion.IPStrategy == nil &&
|
|
|
|
config.SourceCriterion.RequestHeaderName == "" && !config.SourceCriterion.RequestHost {
|
|
|
|
config.SourceCriterion = &dynamic.SourceCriterion{
|
|
|
|
RequestHost: true,
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
sourceMatcher, err := middlewares.GetSourceExtractor(ctxLog, config.SourceCriterion)
|
|
|
|
if err != nil {
|
2020-05-11 12:06:07 +02:00
|
|
|
return nil, fmt.Errorf("error creating requests limiter: %w", err)
|
2019-08-26 12:20:06 +02:00
|
|
|
}
|
|
|
|
|
2022-11-21 18:36:05 +01:00
|
|
|
handler, err := connlimit.New(next, sourceMatcher, config.Amount,
|
|
|
|
connlimit.Logger(logs.NewOxyWrapper(*logger)),
|
|
|
|
connlimit.Verbose(logger.GetLevel() == zerolog.TraceLevel))
|
2019-08-26 12:20:06 +02:00
|
|
|
if err != nil {
|
2020-05-11 12:06:07 +02:00
|
|
|
return nil, fmt.Errorf("error creating connection limit: %w", err)
|
2019-08-26 12:20:06 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
return &inFlightReq{handler: handler, name: name}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (i *inFlightReq) GetTracingInformation() (string, ext.SpanKindEnum) {
|
|
|
|
return i.name, tracing.SpanKindNoneEnum
|
|
|
|
}
|
|
|
|
|
|
|
|
func (i *inFlightReq) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
|
|
|
|
i.handler.ServeHTTP(rw, req)
|
|
|
|
}
|