check / check (push) Canceled after 0s
The IPv6 group that is one client, the size of the table of clients and the level of the process's own lines become settings. clientGroup reads the group length from them, so limits, bans, history, lookups, AbuseIPDB scores and per-client anomaly counters all follow it; ratelimit.New takes the table size; the process logger takes the level once the settings are read, and request lines, written apart from it, are never held back. Judgement call: SWWAF_IPV6_GROUP_PREFIX accepts 32 to 128, the issue's example range. Model: opus-5-5
378 lines
11 KiB
Go
378 lines
11 KiB
Go
// Package smallwebwaf runs the smallwebwaf process: it reads the settings,
|
|
// the rule files, the lookup database and the state files, serves requests
|
|
// until it is told to stop, and then stops in an orderly way, writing the
|
|
// state files.
|
|
package smallwebwaf
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"io"
|
|
"log/slog"
|
|
"net"
|
|
"net/http"
|
|
"os"
|
|
"os/signal"
|
|
"syscall"
|
|
"time"
|
|
|
|
"sneak.berlin/go/smallwebwaf/internal/alerts"
|
|
"sneak.berlin/go/smallwebwaf/internal/config"
|
|
"sneak.berlin/go/smallwebwaf/internal/lookup"
|
|
"sneak.berlin/go/smallwebwaf/internal/proxy"
|
|
"sneak.berlin/go/smallwebwaf/internal/remotelog"
|
|
"sneak.berlin/go/smallwebwaf/internal/reputation"
|
|
"sneak.berlin/go/smallwebwaf/internal/requestlog"
|
|
"sneak.berlin/go/smallwebwaf/internal/rules"
|
|
"sneak.berlin/go/smallwebwaf/internal/state"
|
|
)
|
|
|
|
// shutdownTimeout is how long requests in progress may take to finish
|
|
// once smallwebwaf is told to stop, before their connections are closed.
|
|
// runit and docker wait a little longer before they kill the process.
|
|
const shutdownTimeout = 5 * time.Second
|
|
|
|
// remoteLogStopTimeout is how long, as smallwebwaf stops, the log lines
|
|
// still waiting are sent to SWWAF_LOG_REMOTE_URL before they are given
|
|
// up. stdout has carried them.
|
|
const remoteLogStopTimeout = 2 * time.Second
|
|
|
|
// Params are what Run needs from the process.
|
|
type Params struct {
|
|
// Version is the version of the binary, set when it is built.
|
|
Version string
|
|
// LookupEnv reads an environment variable, normally os.LookupEnv.
|
|
LookupEnv func(string) (string, bool)
|
|
// Stdout receives the request log and the process's own messages.
|
|
Stdout io.Writer
|
|
}
|
|
|
|
// Main runs smallwebwaf until SIGTERM or SIGINT, and returns the
|
|
// process's exit status. Run as `smallwebwaf healthcheck`, it is the
|
|
// container's health check instead.
|
|
func Main(version string) int {
|
|
if len(os.Args) > 1 && os.Args[1] == "healthcheck" {
|
|
return HealthCheck(context.Background(), os.Args[2:], os.LookupEnv, os.Stderr)
|
|
}
|
|
|
|
ctx, stop := signal.NotifyContext(context.Background(),
|
|
syscall.SIGTERM, os.Interrupt)
|
|
defer stop()
|
|
|
|
return Run(ctx, Params{
|
|
Version: version,
|
|
LookupEnv: os.LookupEnv,
|
|
Stdout: os.Stdout,
|
|
})
|
|
}
|
|
|
|
// Run reads the settings, the rule files, the lookup database and the
|
|
// state files, then serves requests until ctx is done. It returns the
|
|
// process's exit status, 1 when smallwebwaf cannot start.
|
|
func Run(ctx context.Context, params Params) int {
|
|
// Until the settings are read, the one message is an invalid setting's
|
|
// error, which every SWWAF_LOG_LEVEL lets through.
|
|
processLog := requestlog.NewProcessLogger(params.Stdout,
|
|
config.InstanceName(params.LookupEnv), slog.LevelError)
|
|
|
|
cfg, err := config.FromEnvironment(params.LookupEnv)
|
|
if err != nil {
|
|
processLog.Error("invalid setting", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
// While SWWAF_LOG_REMOTE_URL is set, every line on stdout from here on
|
|
// is sent there too.
|
|
stdout := params.Stdout
|
|
|
|
var remote *remotelog.Sender
|
|
|
|
if cfg.LogRemoteURL != nil {
|
|
remote = newRemoteLogSender(cfg)
|
|
stdout = io.MultiWriter(params.Stdout, remote)
|
|
}
|
|
|
|
processLog = requestlog.NewProcessLogger(stdout, cfg.InstanceName, cfg.LogLevel)
|
|
|
|
if remote != nil {
|
|
stopSending := startSending(ctx, remote, processLog)
|
|
defer stopSending()
|
|
}
|
|
|
|
// The state files and the alerts give times in UTC.
|
|
now := func() time.Time { return time.Now().UTC() }
|
|
|
|
alertQueue := newAlertQueue(cfg, now, processLog)
|
|
|
|
ruleFiles, err := rules.Load(rules.Params{
|
|
Dir: cfg.RulesDir,
|
|
Enabled: cfg.RulesEnabled,
|
|
ProcessLog: processLog,
|
|
Alerts: alertQueue,
|
|
})
|
|
if err != nil {
|
|
processLog.Error("cannot use the rule files", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
server, err := newServer(cfg, stdout, processLog, now, ruleFiles, alertQueue)
|
|
if err != nil {
|
|
processLog.Error("cannot use the lookup database", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
if remote != nil {
|
|
server.Metrics.AddRemoteLog(remote)
|
|
}
|
|
|
|
files, err := loadStateFiles(cfg, server, alertQueue, now, processLog)
|
|
if err != nil {
|
|
processLog.Error("cannot use the state files", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
listener, err := (&net.ListenConfig{}).Listen(ctx, "tcp", cfg.ListenAddr)
|
|
if err != nil {
|
|
processLog.Error("cannot listen on SWWAF_LISTEN_ADDR",
|
|
"error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
processLog.Info("starting",
|
|
"version", params.Version,
|
|
"address", listener.Addr().String(),
|
|
"settings", cfg)
|
|
|
|
return serve(ctx, server, listener, files, ruleFiles, alertQueue, processLog)
|
|
}
|
|
|
|
// newServer returns the server smallwebwaf runs, with the metrics of the
|
|
// alerts, after reading the lookup database while SWWAF_LOOKUP_SOURCE is
|
|
// file. A lookup database that cannot be read is an error.
|
|
func newServer(
|
|
cfg *config.Config, stdout io.Writer, processLog *slog.Logger,
|
|
now func() time.Time, ruleFiles *rules.Files, alertQueue *alerts.Queue,
|
|
) (*proxy.Server, error) {
|
|
var lookupFile *lookup.File
|
|
|
|
if cfg.LookupSource == "file" {
|
|
var err error
|
|
|
|
lookupFile, err = lookup.OpenFile(lookup.FileParams{
|
|
Path: cfg.LookupDBPath,
|
|
Now: now,
|
|
ProcessLog: processLog,
|
|
Alerts: alertQueue,
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
|
|
server := proxy.New(proxy.Params{
|
|
Config: cfg,
|
|
RequestLog: stdout,
|
|
ProcessLog: processLog,
|
|
GeoJSURL: lookup.URL,
|
|
AbuseIPDBURL: reputation.AbuseIPDBURL,
|
|
LookupFile: lookupFile,
|
|
Now: now,
|
|
Rules: ruleFiles,
|
|
Alerts: alertQueue,
|
|
})
|
|
server.Metrics.AddAlerts(alertQueue)
|
|
|
|
if lookupFile != nil {
|
|
server.Metrics.AddLookupFile(lookupFile.LastRead, lookupFile.ReadFailures)
|
|
}
|
|
|
|
return server, nil
|
|
}
|
|
|
|
// loadStateFiles reads the state files into the parts of server and into
|
|
// alertQueue, as state.Load does.
|
|
func loadStateFiles(
|
|
cfg *config.Config, server *proxy.Server, alertQueue *alerts.Queue,
|
|
now func() time.Time, processLog *slog.Logger,
|
|
) (*state.Files, error) {
|
|
return state.Load(state.Params{
|
|
Dir: cfg.StateDir,
|
|
WriteDelay: cfg.StateWriteDelay,
|
|
CounterInterval: cfg.StateCounterInterval,
|
|
Ledger: server.Ledger,
|
|
Limiter: server.Limiter,
|
|
GeoJS: server.GeoJS,
|
|
Lists: server.Lists,
|
|
DNSBL: server.DNSBL,
|
|
AbuseIPDB: server.AbuseIPDB,
|
|
Alerts: alertQueue,
|
|
Anomalies: server.Anomalies,
|
|
Now: now,
|
|
ProcessLog: processLog,
|
|
Metrics: server.Metrics,
|
|
})
|
|
}
|
|
|
|
// newAlertQueue returns the queue of the alerts to the webhook, Slack and
|
|
// ntfy, with the settings for them.
|
|
func newAlertQueue(
|
|
cfg *config.Config, now func() time.Time, processLog *slog.Logger,
|
|
) *alerts.Queue {
|
|
return alerts.New(alerts.Params{
|
|
WebhookURL: cfg.AlertWebhookURL,
|
|
WebhookHeaders: cfg.AlertWebhookHeaders,
|
|
SlackURL: cfg.AlertSlackWebhookURL,
|
|
NtfyURL: cfg.AlertNtfyURL,
|
|
NtfyToken: cfg.AlertNtfyToken,
|
|
Events: cfg.AlertEvents,
|
|
Cooldown: cfg.AlertCooldown,
|
|
MaxPerHour: cfg.AlertMaxPerHour,
|
|
Instance: cfg.InstanceName,
|
|
Now: now,
|
|
ProcessLog: processLog,
|
|
})
|
|
}
|
|
|
|
// newRemoteLogSender returns a sender of the log lines to
|
|
// SWWAF_LOG_REMOTE_URL, with the settings for it.
|
|
func newRemoteLogSender(cfg *config.Config) *remotelog.Sender {
|
|
return remotelog.New(remotelog.Params{
|
|
URL: cfg.LogRemoteURL,
|
|
RootCAs: cfg.LogRemoteTLSCAs,
|
|
Buffer: cfg.LogRemoteBuffer,
|
|
Facility: cfg.LogRemoteFacility,
|
|
AppName: cfg.LogRemoteAppName,
|
|
})
|
|
}
|
|
|
|
// startSending runs remote until the function it returns is called, which
|
|
// then waits at most remoteLogStopTimeout for the lines still waiting to
|
|
// be sent. Sending goes on after ctx is done, so that the lines written
|
|
// while smallwebwaf stops are sent too.
|
|
func startSending(
|
|
ctx context.Context, remote *remotelog.Sender, processLog *slog.Logger,
|
|
) func() {
|
|
sending, stop := context.WithCancel(context.WithoutCancel(ctx))
|
|
sent := make(chan struct{})
|
|
|
|
go func() {
|
|
remote.Run(sending, processLog)
|
|
close(sent)
|
|
}()
|
|
|
|
return func() {
|
|
stop()
|
|
|
|
select {
|
|
case <-sent:
|
|
case <-time.After(remoteLogStopTimeout):
|
|
}
|
|
}
|
|
}
|
|
|
|
// serve serves requests on listener, writes the state files as they are
|
|
// due, takes in an admin's edits of them, reads the rule files again as
|
|
// they change, and the lookup database when it is replaced, fetches the
|
|
// lists the settings name by URL as they are due, and sends the alerts,
|
|
// until ctx is done. Then it gives the requests in progress
|
|
// shutdownTimeout to finish, and writes every state file, alerts.json with
|
|
// the alerts still waiting.
|
|
func serve(
|
|
ctx context.Context, server *proxy.Server, listener net.Listener,
|
|
files *state.Files, ruleFiles *rules.Files, alertQueue *alerts.Queue,
|
|
processLog *slog.Logger,
|
|
) int {
|
|
served := make(chan error, 1)
|
|
|
|
go func() {
|
|
served <- server.Serve(listener)
|
|
}()
|
|
|
|
writing, stopWriting := context.WithCancel(ctx)
|
|
defer stopWriting()
|
|
|
|
written := inBackground(func() { files.Run(writing) })
|
|
watched := inBackground(func() { files.Watch(writing) })
|
|
rulesWatched := inBackground(func() { ruleFiles.Watch(writing) })
|
|
lookupFileWatched := inBackground(func() {
|
|
if server.LookupFile != nil {
|
|
server.LookupFile.Watch(writing)
|
|
}
|
|
})
|
|
listsFetched := inBackground(func() { server.Lists.Run(writing) })
|
|
alertsSent := inBackground(func() { alertQueue.Run(writing) })
|
|
|
|
select {
|
|
case err := <-served:
|
|
processLog.Error("serving failed", "error", err.Error())
|
|
|
|
return 1
|
|
case <-ctx.Done():
|
|
}
|
|
|
|
processLog.Info("stopping")
|
|
|
|
shutdownCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx),
|
|
shutdownTimeout)
|
|
defer cancel()
|
|
|
|
err := server.Shutdown(shutdownCtx)
|
|
if err != nil {
|
|
processLog.Warn("requests still in progress were cut off",
|
|
"error", err.Error())
|
|
|
|
_ = server.Close()
|
|
}
|
|
|
|
err = <-served
|
|
if !errors.Is(err, http.ErrServerClosed) {
|
|
processLog.Error("serving failed", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
// Run and Watch have ended, so nothing else reads or writes the
|
|
// files, and no alert is being sent, so that alerts.json keeps every
|
|
// alert not yet sent. Every request has ended too, but for two kinds
|
|
// that Go's server does not wait for: one cut off because Shutdown
|
|
// timed out, and one whose connection switched protocols, such as a
|
|
// WebSocket. Such a request adds to its client's history only as it
|
|
// ends, which can be after this write, and then that request is
|
|
// missing from clients.json.
|
|
<-written
|
|
<-watched
|
|
<-rulesWatched
|
|
<-lookupFileWatched
|
|
<-listsFetched
|
|
<-alertsSent
|
|
|
|
err = files.WriteAll()
|
|
if err != nil {
|
|
processLog.Error("writing the state files failed", "error", err.Error())
|
|
|
|
return 1
|
|
}
|
|
|
|
processLog.Info("stopped")
|
|
|
|
return 0
|
|
}
|
|
|
|
// inBackground runs task on a goroutine of its own, and returns a channel
|
|
// that is closed once task has returned.
|
|
func inBackground(task func()) <-chan struct{} {
|
|
done := make(chan struct{})
|
|
|
|
go func() {
|
|
task()
|
|
close(done)
|
|
}()
|
|
|
|
return done
|
|
}
|