mirror of
https://github.com/netbirdio/netbird.git
synced 2026-08-04 19:55:09 -04:00
340 lines
11 KiB
Go
340 lines
11 KiB
Go
// Package appsec implements the CrowdSec AppSec (WAF) side of the remediation
|
|
// component protocol: each inspected HTTP request is mirrored to the Security
|
|
// Engine's AppSec endpoint, which replies with an allow / ban / captcha verdict
|
|
// for that request.
|
|
//
|
|
// This is a separate endpoint from the LAPI decision stream used by the
|
|
// crowdsec package: LAPI answers "is this IP known bad", AppSec answers "is
|
|
// this request an attack". The two are configured and enabled independently.
|
|
package appsec
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"net/netip"
|
|
"net/url"
|
|
"strings"
|
|
"time"
|
|
|
|
log "github.com/sirupsen/logrus"
|
|
|
|
"github.com/netbirdio/netbird/proxy/internal/restrict"
|
|
)
|
|
|
|
// Header names the AppSec component reads off the mirrored request. IP, URI and
|
|
// Verb are mandatory: the engine answers 500 when any of them is missing.
|
|
const (
|
|
headerAPIKey = "X-Crowdsec-Appsec-Api-Key" //nolint:gosec // G101: a header name, not a credential
|
|
headerIP = "X-Crowdsec-Appsec-Ip"
|
|
headerURI = "X-Crowdsec-Appsec-Uri"
|
|
headerVerb = "X-Crowdsec-Appsec-Verb"
|
|
headerHost = "X-Crowdsec-Appsec-Host"
|
|
headerUserAgent = "X-Crowdsec-Appsec-User-Agent"
|
|
headerHTTPVersion = "X-Crowdsec-Appsec-Http-Version"
|
|
headerTransactionID = "X-Crowdsec-Appsec-Transaction-Id"
|
|
)
|
|
|
|
// headerPrefix covers every protocol header. Any client-supplied header in this
|
|
// namespace is dropped before forwarding so a caller cannot influence the
|
|
// engine's view of its own address, or replay an API key.
|
|
const headerPrefix = "X-Crowdsec-Appsec-"
|
|
|
|
// Remediation actions the engine can return.
|
|
const (
|
|
actionAllow = "allow"
|
|
actionBan = "ban"
|
|
actionCaptcha = "captcha"
|
|
)
|
|
|
|
const (
|
|
// DefaultTimeout matches the 200ms budget CrowdSec's remediation component
|
|
// spec sets for the blocking AppSec call.
|
|
DefaultTimeout = 200 * time.Millisecond
|
|
// DefaultMaxBodyBytes caps the request body mirrored to the engine.
|
|
// Requests with a larger body are inspected on headers and URI only.
|
|
DefaultMaxBodyBytes int64 = 64 << 10
|
|
)
|
|
|
|
// ErrUnavailable reports that the engine could not produce a verdict: the call
|
|
// failed, timed out, or the engine rejected it (401 bad key, 500 malformed).
|
|
// Distinguished from a block verdict so the caller can apply the per-service
|
|
// mode: enforce fails closed, observe allows.
|
|
var ErrUnavailable = errors.New("appsec engine unavailable")
|
|
|
|
// Config configures a Client.
|
|
type Config struct {
|
|
// URL is the AppSec endpoint, e.g. http://127.0.0.1:7422/.
|
|
URL string
|
|
// APIKey is the CrowdSec bouncer API key. The AppSec component validates it
|
|
// against LAPI, so the same key used for the decision stream works here.
|
|
APIKey string
|
|
// Timeout bounds a single inspection call. Zero means DefaultTimeout.
|
|
Timeout time.Duration
|
|
// MaxBodyBytes caps the mirrored request body. Zero means
|
|
// DefaultMaxBodyBytes; negative disables body forwarding entirely.
|
|
MaxBodyBytes int64
|
|
Logger *log.Entry
|
|
}
|
|
|
|
// Client mirrors HTTP requests to a CrowdSec AppSec endpoint. It holds no
|
|
// per-service state and is safe for concurrent use.
|
|
type Client struct {
|
|
url string
|
|
apiKey string
|
|
maxBodyBytes int64
|
|
http *http.Client
|
|
logger *log.Entry
|
|
}
|
|
|
|
// New validates the config and returns a Client. The endpoint is not contacted
|
|
// here: the engine may come up after the proxy.
|
|
func New(cfg Config) (*Client, error) {
|
|
if cfg.URL == "" {
|
|
return nil, errors.New("appsec url is empty")
|
|
}
|
|
if cfg.APIKey == "" {
|
|
return nil, errors.New("appsec api key is empty")
|
|
}
|
|
parsed, err := url.Parse(cfg.URL)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("parse appsec url: %w", err)
|
|
}
|
|
if parsed.Scheme != "http" && parsed.Scheme != "https" {
|
|
return nil, fmt.Errorf("appsec url scheme %q is not http(s)", parsed.Scheme)
|
|
}
|
|
if parsed.Host == "" {
|
|
return nil, errors.New("appsec url has no host")
|
|
}
|
|
|
|
timeout := cfg.Timeout
|
|
if timeout <= 0 {
|
|
timeout = DefaultTimeout
|
|
}
|
|
maxBody := cfg.MaxBodyBytes
|
|
if maxBody == 0 {
|
|
maxBody = DefaultMaxBodyBytes
|
|
}
|
|
logger := cfg.Logger
|
|
if logger == nil {
|
|
logger = log.NewEntry(log.StandardLogger())
|
|
}
|
|
|
|
return &Client{
|
|
url: cfg.URL,
|
|
apiKey: cfg.APIKey,
|
|
maxBodyBytes: maxBody,
|
|
logger: logger,
|
|
http: &http.Client{
|
|
Timeout: timeout,
|
|
Transport: &http.Transport{
|
|
MaxIdleConns: 100,
|
|
MaxIdleConnsPerHost: 32,
|
|
IdleConnTimeout: 90 * time.Second,
|
|
},
|
|
},
|
|
}, nil
|
|
}
|
|
|
|
// Request is one inspection request.
|
|
type Request struct {
|
|
// HTTP is the in-flight client request. Inspect buffers and restores its
|
|
// body, so the request stays forwardable afterwards.
|
|
HTTP *http.Request
|
|
// ClientIP is the resolved client address (after trusted-proxy handling).
|
|
ClientIP netip.Addr
|
|
// TransactionID correlates the engine's alert with the proxy's access log
|
|
// entry. Empty lets the engine generate its own UUID.
|
|
TransactionID string
|
|
// OmitBodyFields suppresses body forwarding when the body is a form
|
|
// containing any of these field names. Used to keep credentials submitted
|
|
// to the proxy's own login form out of the engine.
|
|
OmitBodyFields []string
|
|
}
|
|
|
|
// Inspect mirrors r to the AppSec engine and returns its verdict. A nil error
|
|
// with restrict.Allow means the request passed. On failure it returns
|
|
// DenyAppSecUnavailable wrapped with ErrUnavailable; the caller decides whether
|
|
// that blocks, based on the per-service mode.
|
|
func (c *Client) Inspect(ctx context.Context, req Request) (restrict.Verdict, error) {
|
|
if c == nil {
|
|
return restrict.DenyAppSecUnavailable, ErrUnavailable
|
|
}
|
|
if req.HTTP == nil {
|
|
return restrict.DenyAppSecUnavailable, fmt.Errorf("%w: nil request", ErrUnavailable)
|
|
}
|
|
|
|
body, err := c.readBody(req)
|
|
if err != nil {
|
|
return restrict.DenyAppSecUnavailable, fmt.Errorf("%w: read body: %w", ErrUnavailable, err)
|
|
}
|
|
|
|
outbound, err := c.buildRequest(ctx, req, body)
|
|
if err != nil {
|
|
return restrict.DenyAppSecUnavailable, fmt.Errorf("%w: %w", ErrUnavailable, err)
|
|
}
|
|
|
|
resp, err := c.http.Do(outbound)
|
|
if err != nil {
|
|
return restrict.DenyAppSecUnavailable, fmt.Errorf("%w: %w", ErrUnavailable, err)
|
|
}
|
|
defer func() {
|
|
if err := resp.Body.Close(); err != nil {
|
|
c.logger.Tracef("close appsec response body: %v", err)
|
|
}
|
|
}()
|
|
|
|
return c.verdict(resp)
|
|
}
|
|
|
|
// readBody buffers the body so it can be mirrored, always restoring it on the
|
|
// original request. Returns nil when there is no body to forward: no body at
|
|
// all, an upgrade request, a body over the cap, or a credential form.
|
|
func (c *Client) readBody(req Request) ([]byte, error) {
|
|
r := req.HTTP
|
|
if c.maxBodyBytes < 0 || r.Body == nil || r.Body == http.NoBody {
|
|
return nil, nil
|
|
}
|
|
// Upgrade requests (websockets) have no meaningful request body and their
|
|
// stream must not be consumed here.
|
|
if r.Header.Get("Upgrade") != "" || strings.EqualFold(r.Header.Get("Connection"), "upgrade") {
|
|
return nil, nil
|
|
}
|
|
// A Content-Length over the cap is known to be too large before reading.
|
|
if r.ContentLength > c.maxBodyBytes {
|
|
return nil, nil
|
|
}
|
|
|
|
body, oversize, err := bufferBody(r, c.maxBodyBytes)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
// An oversize body was only partially read: a truncated prefix changes the
|
|
// engine's verdict in both directions, so inspect headers and URI only.
|
|
if oversize {
|
|
return nil, nil
|
|
}
|
|
if formHasField(r.Header.Get("Content-Type"), body, req.OmitBodyFields) {
|
|
return nil, nil
|
|
}
|
|
return body, nil
|
|
}
|
|
|
|
// buildRequest assembles the mirrored request. Per the protocol it is a GET
|
|
// when there is no body and a POST otherwise; bytes.Reader gives the outbound
|
|
// request an accurate Content-Length, which the engine relies on to read the
|
|
// body at all.
|
|
func (c *Client) buildRequest(ctx context.Context, req Request, body []byte) (*http.Request, error) {
|
|
method := http.MethodGet
|
|
var payload io.Reader
|
|
if len(body) > 0 {
|
|
method = http.MethodPost
|
|
payload = bytes.NewReader(body)
|
|
}
|
|
|
|
outbound, err := http.NewRequestWithContext(ctx, method, c.url, payload)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("build appsec request: %w", err)
|
|
}
|
|
|
|
r := req.HTTP
|
|
copyInspectableHeaders(outbound.Header, r.Header)
|
|
|
|
outbound.Header.Set(headerAPIKey, c.apiKey)
|
|
outbound.Header.Set(headerIP, req.ClientIP.Unmap().String())
|
|
outbound.Header.Set(headerURI, r.URL.RequestURI())
|
|
outbound.Header.Set(headerVerb, r.Method)
|
|
outbound.Header.Set(headerHost, r.Host)
|
|
if ua := r.UserAgent(); ua != "" {
|
|
outbound.Header.Set(headerUserAgent, ua)
|
|
}
|
|
outbound.Header.Set(headerHTTPVersion, httpVersion(r))
|
|
if req.TransactionID != "" {
|
|
outbound.Header.Set(headerTransactionID, req.TransactionID)
|
|
}
|
|
return outbound, nil
|
|
}
|
|
|
|
// verdict maps the engine's response to a restrict.Verdict. 200 is a pass and
|
|
// 401/500 are engine-side failures; every other status carries a remediation in
|
|
// the body. The blocked status code is operator-configurable
|
|
// (blocked_http_code), so the action field decides, not the status.
|
|
func (c *Client) verdict(resp *http.Response) (restrict.Verdict, error) {
|
|
switch resp.StatusCode {
|
|
case http.StatusOK:
|
|
return restrict.Allow, nil
|
|
case http.StatusUnauthorized:
|
|
return restrict.DenyAppSecUnavailable, fmt.Errorf("%w: rejected api key", ErrUnavailable)
|
|
case http.StatusInternalServerError:
|
|
return restrict.DenyAppSecUnavailable, fmt.Errorf("%w: engine error", ErrUnavailable)
|
|
}
|
|
|
|
var decoded struct {
|
|
Action string `json:"action"`
|
|
}
|
|
if err := json.NewDecoder(io.LimitReader(resp.Body, 4<<10)).Decode(&decoded); err != nil {
|
|
// A non-2xx status with an unreadable body is still a block: the engine
|
|
// answered, we just cannot tell which remediation it chose.
|
|
c.logger.Debugf("failed to decode appsec response (status %d): %v", resp.StatusCode, err)
|
|
return restrict.DenyAppSecBan, nil
|
|
}
|
|
|
|
switch decoded.Action {
|
|
case actionAllow:
|
|
return restrict.Allow, nil
|
|
case actionCaptcha:
|
|
return restrict.DenyAppSecCaptcha, nil
|
|
case actionBan:
|
|
return restrict.DenyAppSecBan, nil
|
|
default:
|
|
// Unknown remediation: the engine flagged the request, so deny.
|
|
c.logger.Debugf("unknown appsec action %q (status %d), treating as ban", decoded.Action, resp.StatusCode)
|
|
return restrict.DenyAppSecBan, nil
|
|
}
|
|
}
|
|
|
|
// copyInspectableHeaders copies the client's headers, which are what the WAF
|
|
// rules actually match on, dropping hop-by-hop headers that describe the
|
|
// proxy-to-engine connection rather than the client request, and any header in
|
|
// the AppSec protocol namespace.
|
|
func copyInspectableHeaders(dst, src http.Header) {
|
|
for name, values := range src {
|
|
if hopByHopHeaders[http.CanonicalHeaderKey(name)] {
|
|
continue
|
|
}
|
|
if strings.HasPrefix(http.CanonicalHeaderKey(name), headerPrefix) {
|
|
continue
|
|
}
|
|
dst[http.CanonicalHeaderKey(name)] = append([]string(nil), values...)
|
|
}
|
|
// Content-Length describes the mirrored payload, not the client's: net/http
|
|
// sets it from the body we actually attach. Content-Type is kept either way
|
|
// so rules matching on it still fire when the body was not forwarded.
|
|
dst.Del("Content-Length")
|
|
}
|
|
|
|
var hopByHopHeaders = map[string]bool{
|
|
"Connection": true,
|
|
"Keep-Alive": true,
|
|
"Proxy-Authenticate": true,
|
|
"Proxy-Authorization": true,
|
|
"Proxy-Connection": true,
|
|
"Te": true,
|
|
"Trailer": true,
|
|
"Transfer-Encoding": true,
|
|
"Upgrade": true,
|
|
}
|
|
|
|
// httpVersion renders the two-digit form the engine parses ("11", "20").
|
|
func httpVersion(r *http.Request) string {
|
|
major, minor := r.ProtoMajor, r.ProtoMinor
|
|
if major < 0 || major > 9 || minor < 0 || minor > 9 {
|
|
return ""
|
|
}
|
|
return fmt.Sprintf("%d%d", major, minor)
|
|
}
|