check / check (push) Waiting to run
Zones in SWWAF_DNSBL_ZONES are asked about each client (RFC 5782 names) in the background, through the host's resolver or SWWAF_DNSBL_RESOLVER; no request waits. Verdicts last SWWAF_REPUTATION_CACHE_TTL and are kept in reputation.json, at most 100,000. After the blocklists, SWWAF_REPUTATION_ACTION (limit:25) denies, limits or logs a listed client; the log line names the zones, each raises reputation_hit, with metrics by zone. A failed, timed-out or refused query gives no verdict, raises source_failure, and pauses the zone a minute. Judgement call: answers in 127.255.255.0/24 or outside 127.0.0.0/8 are failures. Judgement call: the minute's pause after a failure; at most 1,000 queries at once. Rule suppressed: paralleltest on the DNSBL tests (Go's resolver shares state across synctest bubbles), funlen on the test of every logged setting. Model: opus-5-5
370 lines
10 KiB
Go
370 lines
10 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/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 {
|
|
processLog := requestlog.NewProcessLogger(params.Stdout,
|
|
config.InstanceName(params.LookupEnv))
|
|
|
|
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)
|
|
|
|
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,
|
|
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,
|
|
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
|
|
}
|