CRS 900xxx/901xxx (init, body inspection, paranoia setup) feuern auf JEDEM Request — keine Security-Events. Filter: nur rule_id >= 910000 wird als Alert in DB geschrieben. Buffer 512 → 2048. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
187 lines
4.8 KiB
Go
187 lines
4.8 KiB
Go
package waf
|
||
|
||
import (
|
||
"context"
|
||
"log/slog"
|
||
"net/http"
|
||
"strings"
|
||
|
||
"github.com/corazawaf/coraza/v3/types"
|
||
"github.com/dropmorepackets/haproxy-go/pkg/encoding"
|
||
"github.com/dropmorepackets/haproxy-go/spop"
|
||
)
|
||
|
||
// SPOEAgent wraps the haproxy-go SPOE server and dispatches each
|
||
// inspected request to the appropriate per-domain Coraza engine.
|
||
type SPOEAgent struct {
|
||
Manager *Manager
|
||
AlertWriter *AlertWriter
|
||
Addr string
|
||
}
|
||
|
||
// ListenAndServe starts the SPOE agent. Blocks until ctx is cancelled.
|
||
func (a *SPOEAgent) ListenAndServe(ctx context.Context) error {
|
||
agent := spop.Agent{
|
||
Addr: a.Addr,
|
||
Handler: spop.HandlerFunc(a.handle),
|
||
BaseContext: ctx,
|
||
}
|
||
return agent.ListenAndServe()
|
||
}
|
||
|
||
// handle is called by the haproxy-go SPOE library for every NOTIFY
|
||
// frame HAProxy sends. It extracts the request data, runs Coraza,
|
||
// and optionally sets a txn.waf.status variable to trigger a deny ACL.
|
||
func (a *SPOEAgent) handle(ctx context.Context, w *encoding.ActionWriter, m *encoding.Message) {
|
||
var (
|
||
clientIP string
|
||
method string
|
||
uri string // full request URI (path + optional ?query)
|
||
httpVer string
|
||
host string
|
||
rawHdrs string
|
||
)
|
||
|
||
// Iterate over the key-value pairs HAProxy sent with this message.
|
||
entry := encoding.AcquireKVEntry()
|
||
defer encoding.ReleaseKVEntry(entry)
|
||
for m.KV.Next(entry) {
|
||
switch {
|
||
case entry.NameEquals("src"):
|
||
addr := entry.ValueAddr()
|
||
if addr.IsValid() {
|
||
clientIP = addr.String()
|
||
}
|
||
case entry.NameEquals("method"):
|
||
method = string(entry.ValueBytes())
|
||
case entry.NameEquals("uri"):
|
||
uri = string(entry.ValueBytes())
|
||
case entry.NameEquals("ver"):
|
||
httpVer = string(entry.ValueBytes())
|
||
case entry.NameEquals("host"):
|
||
host = string(entry.ValueBytes())
|
||
case entry.NameEquals("headers"):
|
||
rawHdrs = string(entry.ValueBytes())
|
||
}
|
||
entry.Reset()
|
||
}
|
||
|
||
if host == "" {
|
||
return
|
||
}
|
||
|
||
de, ok := a.Manager.GetForHost(host)
|
||
if !ok {
|
||
return // WAF not configured or disabled for this domain
|
||
}
|
||
|
||
tx := de.WAF.NewTransaction()
|
||
defer func() {
|
||
tx.ProcessLogging()
|
||
if err := tx.Close(); err != nil {
|
||
slog.Warn("waf: tx.Close", "error", err)
|
||
}
|
||
}()
|
||
|
||
// Feed connection metadata.
|
||
if clientIP != "" {
|
||
tx.ProcessConnection(clientIP, 0, "", 0)
|
||
}
|
||
|
||
if uri == "" {
|
||
uri = "/"
|
||
}
|
||
if httpVer == "" {
|
||
httpVer = "HTTP/1.1"
|
||
}
|
||
tx.ProcessURI(uri, method, httpVer)
|
||
|
||
// Feed Host header first (required by many CRS rules).
|
||
tx.AddRequestHeader("Host", host)
|
||
|
||
// Parse and feed all raw headers.
|
||
parseHeaders(rawHdrs, func(name, val string) {
|
||
if !strings.EqualFold(name, "host") { // already added above
|
||
tx.AddRequestHeader(name, val)
|
||
}
|
||
})
|
||
|
||
// Evaluate request headers.
|
||
interruption := tx.ProcessRequestHeaders()
|
||
|
||
// Log all matched rules (detection + blocking).
|
||
for _, mr := range tx.MatchedRules() {
|
||
a.sendAlert(host, clientIP, method, uri, mr, interruption != nil)
|
||
}
|
||
|
||
if interruption != nil {
|
||
status := interruption.Status
|
||
if status == 0 {
|
||
status = http.StatusForbidden
|
||
}
|
||
slog.Info("waf: request blocked",
|
||
"host", host, "method", method, "uri", uri,
|
||
"client", clientIP, "status", status, "rule", interruption.RuleID,
|
||
)
|
||
if de.Mode == "blocking" {
|
||
if err := w.SetInt64(encoding.VarScopeTransaction, "status", int64(status)); err != nil {
|
||
slog.Warn("waf: SetInt64 status", "error", err)
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
// sendAlert enqueues a WAF alert for async DB write.
|
||
// Control-flow rules (pass+nolog with empty message) are skipped —
|
||
// they are CRS paranoia-level skip-markers, not real detections.
|
||
func (a *SPOEAgent) sendAlert(host, clientIP, method, uri string, mr types.MatchedRule, blocked bool) {
|
||
if a.AlertWriter == nil {
|
||
return
|
||
}
|
||
ruleID := mr.Rule().ID()
|
||
// Skip CRS setup/initialization rules (900xxx–909xxx) — they fire on
|
||
// every request as part of CRS init and are not security events.
|
||
// Real detection rules start at 910xxx (IP reputation) and above.
|
||
if ruleID > 0 && ruleID < 910000 {
|
||
return
|
||
}
|
||
// Skip control-flow rules with no message (PL-skip markers).
|
||
if mr.Message() == "" {
|
||
return
|
||
}
|
||
action := "detected"
|
||
if blocked && mr.Disruptive() {
|
||
action = "blocked"
|
||
}
|
||
a.AlertWriter.Send(Alert{
|
||
Hostname: host,
|
||
ClientIP: clientIP,
|
||
Method: method,
|
||
URI: uri,
|
||
RuleID: mr.Rule().ID(),
|
||
RuleMsg: mr.Message(),
|
||
Severity: mr.Rule().Severity().String(),
|
||
Action: action,
|
||
})
|
||
}
|
||
|
||
// parseHeaders splits HAProxy raw headers ("Name: value\r\n…") and
|
||
// calls fn for each valid header line.
|
||
func parseHeaders(raw string, fn func(name, val string)) {
|
||
for _, line := range strings.Split(raw, "\n") {
|
||
line = strings.TrimRight(line, "\r")
|
||
if line == "" {
|
||
continue
|
||
}
|
||
idx := strings.IndexByte(line, ':')
|
||
if idx <= 0 {
|
||
continue
|
||
}
|
||
name := strings.TrimSpace(line[:idx])
|
||
val := strings.TrimSpace(line[idx+1:])
|
||
if name != "" {
|
||
fn(name, val)
|
||
}
|
||
}
|
||
}
|