Take in an admin's edits of the state files while running (closes #68)
check / check (push) Successful in 4m1s
check / check (push) Successful in 4m1s
smallwebwaf watches SWWAF_STATE_DIR with fsnotify and takes in a saved edit of a state file in place of what it held. It knows its own writes by the SHA-256 of what it last read or wrote; each write first takes in an edit made since. An edit that does not parse is renamed to <name>.bad at the next write. Each edit taken in or set aside is logged and counted. Every ban on a netblock is checked, and the next ban is worked out from the one that ended last. README.md says how to add and lift a ban. Judgement call: a broken edit is set aside at the next write, since an editor's file can be read half written. Model: opus-5-5
This commit was merged in pull request #75.
This commit is contained in:
+270
-76
@@ -1,14 +1,17 @@
|
||||
// Package state keeps smallwebwaf's state in JSON files in
|
||||
// SWWAF_STATE_DIR, as the "Persistent state" section of SPEC.md describes:
|
||||
// bans.json holds the bans, clients.json each client's counters and
|
||||
// history, and lookups.json GeoJS's answers. Load reads them at start, and
|
||||
// Run and WriteAll write them, each from a snapshot its part takes under
|
||||
// its own lock, so that no request waits on the disk.
|
||||
// history, and lookups.json GeoJS's answers. Load reads them at start,
|
||||
// Watch takes in an admin's edit of one while smallwebwaf runs, and Run
|
||||
// and WriteAll write them. The disk is read and written outside the
|
||||
// parts' locks, which are held only to take a snapshot or to put in what
|
||||
// a file holds, so that no request waits on the disk.
|
||||
package state
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
@@ -17,8 +20,10 @@ import (
|
||||
"net/netip"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/fsnotify/fsnotify"
|
||||
"sneak.berlin/go/smallwebwaf/internal/bans"
|
||||
"sneak.berlin/go/smallwebwaf/internal/lookup"
|
||||
"sneak.berlin/go/smallwebwaf/internal/metrics"
|
||||
@@ -61,15 +66,26 @@ type Params struct {
|
||||
// Now tells the time by which the counters' buckets run out, normally
|
||||
// time.Now in UTC.
|
||||
Now func() time.Time
|
||||
// ProcessLog receives what was read, and the writes that fail.
|
||||
// ProcessLog receives what was read and taken in, the edits set aside,
|
||||
// and the writes that fail.
|
||||
ProcessLog *slog.Logger
|
||||
// Metrics count each file's writes.
|
||||
// Metrics count each file's writes, and the edits taken in and set
|
||||
// aside.
|
||||
Metrics *metrics.Metrics
|
||||
}
|
||||
|
||||
// Files are the state files of a running smallwebwaf.
|
||||
type Files struct {
|
||||
params Params
|
||||
|
||||
// mu is held while a file is read for an edit, and while it is
|
||||
// written, so that Watch and the writes take turns. No request takes
|
||||
// it.
|
||||
mu sync.Mutex
|
||||
// sums are the SHA-256 sums of what each file held, by name, when
|
||||
// smallwebwaf last read or wrote it. A file that holds anything else
|
||||
// has been edited since.
|
||||
sums map[string][sha256.Size]byte
|
||||
}
|
||||
|
||||
// bansFile is bans.json, indented for an admin to read and edit.
|
||||
@@ -120,41 +136,28 @@ func Load(params Params) (*Files, error) {
|
||||
return nil, fmt.Errorf("SWWAF_STATE_DIR cannot be written: %w", err)
|
||||
}
|
||||
|
||||
var (
|
||||
bansIn bansFile
|
||||
clientsIn clientsFile
|
||||
lookupsIn lookupsFile
|
||||
)
|
||||
f := &Files{params: params, sums: map[string][sha256.Size]byte{}}
|
||||
|
||||
err = errors.Join(
|
||||
read(params.Dir, bansJSON, &bansIn),
|
||||
read(params.Dir, clientsJSON, &clientsIn),
|
||||
read(params.Dir, lookupsJSON, &lookupsIn),
|
||||
)
|
||||
bansRead, bansErr := f.read(bansJSON)
|
||||
clientsRead, clientsErr := f.read(clientsJSON)
|
||||
lookupsRead, lookupsErr := f.read(lookupsJSON)
|
||||
|
||||
err = errors.Join(bansErr, clientsErr, lookupsErr)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
held := make([]bans.Ban, 0, len(bansIn.Bans))
|
||||
for _, entry := range bansIn.Bans {
|
||||
held = append(held, entry.ban())
|
||||
}
|
||||
|
||||
params.Ledger.Load(held)
|
||||
params.Limiter.Load(clientsIn.Clients, params.Now())
|
||||
params.GeoJS.Load(lookupsIn.Lookups)
|
||||
|
||||
params.ProcessLog.Info("read the state files", "directory", params.Dir,
|
||||
"bans", len(bansIn.Bans), "clients", len(clientsIn.Clients),
|
||||
"lookups", len(lookupsIn.Lookups))
|
||||
"bans", bansRead, "clients", clientsRead, "lookups", lookupsRead)
|
||||
|
||||
return &Files{params: params}, nil
|
||||
return f, nil
|
||||
}
|
||||
|
||||
// Run writes bans.json WriteDelay after a ban is made, with every ban
|
||||
// made in between, and every file every CounterInterval, until ctx is
|
||||
// done. A write that fails is logged, and the file is written again at
|
||||
// its next write.
|
||||
// its next write. Each write takes in an admin's edit of its file first,
|
||||
// as writeFile describes.
|
||||
func (f *Files) Run(ctx context.Context) {
|
||||
interval := time.NewTicker(f.params.CounterInterval)
|
||||
defer interval.Stop()
|
||||
@@ -172,7 +175,7 @@ func (f *Files) Run(ctx context.Context) {
|
||||
case <-bansDue:
|
||||
bansDue = nil
|
||||
|
||||
f.logFailure(f.writeBans())
|
||||
f.logFailure(f.writeFile(bansJSON))
|
||||
case <-interval.C:
|
||||
f.logFailure(f.WriteAll())
|
||||
}
|
||||
@@ -182,7 +185,50 @@ func (f *Files) Run(ctx context.Context) {
|
||||
// WriteAll writes every state file, as smallwebwaf stops. A file that
|
||||
// fails does not keep the others from being written.
|
||||
func (f *Files) WriteAll() error {
|
||||
return errors.Join(f.writeBans(), f.writeClients(), f.writeLookups())
|
||||
return errors.Join(f.writeFile(bansJSON), f.writeFile(clientsJSON),
|
||||
f.writeFile(lookupsJSON))
|
||||
}
|
||||
|
||||
// Watch watches Dir until ctx is done, and takes in an admin's edit of a
|
||||
// state file as soon as it is saved: what the file holds replaces what
|
||||
// smallwebwaf held for it. An edit that does not parse is left for the
|
||||
// file's next write, which sets it aside, since a file can be read while
|
||||
// an editor is still writing it. If Dir cannot be watched, that is
|
||||
// logged, and an edit is taken in only before its file is written.
|
||||
func (f *Files) Watch(ctx context.Context) {
|
||||
watcher, err := fsnotify.NewWatcher()
|
||||
if err == nil {
|
||||
defer func() {
|
||||
_ = watcher.Close()
|
||||
}()
|
||||
|
||||
err = watcher.Add(f.params.Dir)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
f.params.ProcessLog.Error("cannot watch the state files for edits",
|
||||
"error", err.Error())
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
f.params.ProcessLog.Info("watching the state files for edits",
|
||||
"directory", f.params.Dir)
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case event := <-watcher.Events:
|
||||
switch name := filepath.Base(event.Name); name {
|
||||
case bansJSON, clientsJSON, lookupsJSON:
|
||||
f.fileChanged(name)
|
||||
}
|
||||
case err = <-watcher.Errors:
|
||||
f.params.ProcessLog.Warn("watching the state files failed",
|
||||
"error", err.Error())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// logFailure logs a write that failed.
|
||||
@@ -193,52 +239,209 @@ func (f *Files) logFailure(err error) {
|
||||
}
|
||||
}
|
||||
|
||||
// writeBans writes bans.json.
|
||||
func (f *Files) writeBans() error {
|
||||
held := f.params.Ledger.Snapshot()
|
||||
// fileChanged takes in what the state file name holds, as Watch sees it
|
||||
// change, if that is an edit made since smallwebwaf last read or wrote
|
||||
// the file. A file that cannot be read or does not parse is left for its
|
||||
// next write.
|
||||
func (f *Files) fileChanged(name string) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
file := bansFile{Version: version, Bans: make([]banEntry, 0, len(held))}
|
||||
for _, ban := range held {
|
||||
file.Bans = append(file.Bans, newBanEntry(ban))
|
||||
data, changed, err := f.readChanged(name)
|
||||
if err != nil || !changed {
|
||||
return
|
||||
}
|
||||
|
||||
data, err := json.MarshalIndent(file, "", " ")
|
||||
if err != nil {
|
||||
return fmt.Errorf("encode %s: %w", bansJSON, err)
|
||||
}
|
||||
|
||||
return f.writeCounted(bansJSON, append(data, '\n'))
|
||||
_ = f.takeInEdit(name, data)
|
||||
}
|
||||
|
||||
// writeClients writes clients.json.
|
||||
func (f *Files) writeClients() error {
|
||||
data, err := encodeOnePerLine("clients", f.params.Limiter.Snapshot())
|
||||
// takeInEdit takes in data, an edit of the state file name, as takeIn
|
||||
// does, and counts and logs it. Every edit taken in while smallwebwaf
|
||||
// runs, by Watch or by a write, is taken in here. An edit that does not
|
||||
// parse is neither counted nor logged, and takeIn's error returned.
|
||||
func (f *Files) takeInEdit(name string, data []byte) error {
|
||||
_, err := f.takeIn(name, data)
|
||||
if err != nil {
|
||||
return fmt.Errorf("encode %s: %w", clientsJSON, err)
|
||||
return err
|
||||
}
|
||||
|
||||
return f.writeCounted(clientsJSON, data)
|
||||
// Counted before it is logged, so that the count is there once the
|
||||
// log line is.
|
||||
f.params.Metrics.StateFileEditTakenIn(name)
|
||||
f.params.ProcessLog.Info("took in an edit of a state file",
|
||||
"file", filepath.Join(f.params.Dir, name))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// writeLookups writes lookups.json.
|
||||
func (f *Files) writeLookups() error {
|
||||
data, err := encodeOnePerLine("lookups", f.params.GeoJS.Snapshot())
|
||||
if err != nil {
|
||||
return fmt.Errorf("encode %s: %w", lookupsJSON, err)
|
||||
// read takes in the state file name at start, and returns how many
|
||||
// entries it holds. A missing file holds none.
|
||||
func (f *Files) read(name string) (int, error) {
|
||||
data, changed, err := f.readChanged(name)
|
||||
if err != nil || !changed {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
return f.writeCounted(lookupsJSON, data)
|
||||
return f.takeIn(name, data)
|
||||
}
|
||||
|
||||
// writeCounted writes data to the state file name, as write does, and
|
||||
// counts the write in the metrics.
|
||||
func (f *Files) writeCounted(name string, data []byte) error {
|
||||
err := write(f.params.Dir, name, data)
|
||||
// readChanged returns what the state file name holds, and whether that
|
||||
// has changed since smallwebwaf last read or wrote the file, as it has
|
||||
// for a file smallwebwaf never read or wrote. A missing file has not
|
||||
// changed: it is written again at its next write.
|
||||
func (f *Files) readChanged(name string) ([]byte, bool, error) {
|
||||
path := filepath.Join(f.params.Dir, name)
|
||||
|
||||
data, err := os.ReadFile(path) //nolint:gosec // a state file, in SWWAF_STATE_DIR
|
||||
if errors.Is(err, fs.ErrNotExist) {
|
||||
return nil, false, nil
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
|
||||
return data, sha256.Sum256(data) != f.sums[name], nil
|
||||
}
|
||||
|
||||
// takeIn parses data, what the state file name holds, puts it into the
|
||||
// part that keeps that state, in place of what the part held, and returns
|
||||
// how many entries the file holds. An error names the file and, where the
|
||||
// JSON decoder tells it, the line and column, or else the entry.
|
||||
func (f *Files) takeIn(name string, data []byte) (int, error) {
|
||||
path := filepath.Join(f.params.Dir, name)
|
||||
|
||||
var entries int
|
||||
|
||||
switch name {
|
||||
case bansJSON:
|
||||
var file bansFile
|
||||
|
||||
err := parse(path, data, &file)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
held := make([]bans.Ban, 0, len(file.Bans))
|
||||
for _, entry := range file.Bans {
|
||||
held = append(held, entry.ban())
|
||||
}
|
||||
|
||||
f.params.Ledger.Load(held)
|
||||
entries = len(held)
|
||||
case clientsJSON:
|
||||
var file clientsFile
|
||||
|
||||
err := parse(path, data, &file)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
f.params.Limiter.Load(file.Clients, f.params.Now())
|
||||
entries = len(file.Clients)
|
||||
case lookupsJSON:
|
||||
var file lookupsFile
|
||||
|
||||
err := parse(path, data, &file)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
f.params.GeoJS.Load(file.Lookups)
|
||||
entries = len(file.Lookups)
|
||||
}
|
||||
|
||||
f.sums[name] = sha256.Sum256(data)
|
||||
|
||||
return entries, nil
|
||||
}
|
||||
|
||||
// writeFile writes the state file name from what smallwebwaf holds. An
|
||||
// edit made since smallwebwaf last read or wrote the file is taken in
|
||||
// first, so that it is not overwritten, or set aside if it does not
|
||||
// parse. A file that cannot be read, or an edit that cannot be set
|
||||
// aside, is left as it is, and the write given up. Every write is counted
|
||||
// in the metrics, and one that fails or is given up as a failure.
|
||||
func (f *Files) writeFile(name string) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
data, changed, err := f.readChanged(name)
|
||||
if err == nil && changed {
|
||||
err = f.takeInEdit(name, data)
|
||||
if err != nil {
|
||||
err = f.setAside(name, err)
|
||||
}
|
||||
}
|
||||
|
||||
if err == nil {
|
||||
data, err = f.encode(name)
|
||||
if err != nil {
|
||||
err = fmt.Errorf("encode %s: %w", name, err)
|
||||
}
|
||||
}
|
||||
|
||||
if err == nil {
|
||||
err = write(f.params.Dir, name, data)
|
||||
}
|
||||
|
||||
if err == nil {
|
||||
// The file holds data from here on, even if the directory sync
|
||||
// fails, so that its next read does not take it for an admin's
|
||||
// edit.
|
||||
f.sums[name] = sha256.Sum256(data)
|
||||
err = syncDirectory(f.params.Dir)
|
||||
}
|
||||
|
||||
f.params.Metrics.StateFileWritten(name, len(data), err)
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
// setAside renames the state file name, an edit that does not parse with
|
||||
// parseErr, to name.bad, for the admin to mend, and logs it with where in
|
||||
// the file the error is. If the rename fails, the edit is left as it is,
|
||||
// and the error returned is parseErr joined with the rename's.
|
||||
func (f *Files) setAside(name string, parseErr error) error {
|
||||
path := filepath.Join(f.params.Dir, name)
|
||||
|
||||
err := os.Rename(path, path+".bad")
|
||||
if err != nil {
|
||||
return errors.Join(parseErr, err)
|
||||
}
|
||||
|
||||
f.params.ProcessLog.Error("set aside an edit of a state file that does not parse",
|
||||
"file", path+".bad", "error", parseErr.Error())
|
||||
f.params.Metrics.StateFileEditSetAside(name)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// encode returns the state file name as smallwebwaf writes it, from a
|
||||
// snapshot of the part that keeps that state.
|
||||
func (f *Files) encode(name string) ([]byte, error) {
|
||||
switch name {
|
||||
case bansJSON:
|
||||
held := f.params.Ledger.Snapshot()
|
||||
|
||||
file := bansFile{Version: version, Bans: make([]banEntry, 0, len(held))}
|
||||
for _, ban := range held {
|
||||
file.Bans = append(file.Bans, newBanEntry(ban))
|
||||
}
|
||||
|
||||
data, err := json.MarshalIndent(file, "", " ")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return append(data, '\n'), nil
|
||||
case clientsJSON:
|
||||
return encodeOnePerLine("clients", f.params.Limiter.Snapshot())
|
||||
default: // lookups.json
|
||||
return encodeOnePerLine("lookups", f.params.GeoJS.Snapshot())
|
||||
}
|
||||
}
|
||||
|
||||
// newBanEntry returns ban as bans.json holds it.
|
||||
func newBanEntry(ban bans.Ban) banEntry {
|
||||
entry := banEntry{Netblock: ban.Netblock, Start: ban.Start, Notes: ban.Notes}
|
||||
@@ -389,28 +592,16 @@ func checkWritable(dir string) error {
|
||||
return errors.Join(file.Close(), os.Remove(file.Name()))
|
||||
}
|
||||
|
||||
// read reads the state file name in dir into file, a pointer to that
|
||||
// file's struct, and checks its entries. A missing file leaves file as it
|
||||
// is.
|
||||
func read(dir, name string, file stateFile) error {
|
||||
path := filepath.Join(dir, name)
|
||||
|
||||
data, err := os.ReadFile(path) //nolint:gosec // a state file, in SWWAF_STATE_DIR
|
||||
if errors.Is(err, fs.ErrNotExist) {
|
||||
return nil
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// parse reads data, what the state file at path holds, into file, a
|
||||
// pointer to that file's struct, and checks its entries.
|
||||
func parse(path string, data []byte, file stateFile) error {
|
||||
// The version is read first, so that a file of another version is
|
||||
// refused for that, and not for an entry this version cannot read.
|
||||
var header struct {
|
||||
Version int `json:"version"`
|
||||
}
|
||||
|
||||
err = json.Unmarshal(data, &header)
|
||||
err := json.Unmarshal(data, &header)
|
||||
if err == nil && header.Version != version {
|
||||
err = fmt.Errorf("%w %d, where this smallwebwaf reads version %d",
|
||||
errVersion, header.Version, version)
|
||||
@@ -463,7 +654,7 @@ func position(data []byte, err error) string {
|
||||
// write writes data to the file name in dir so that a crash at any
|
||||
// moment leaves either the old file or the new one, whole: data goes to a
|
||||
// temporary file in the same directory, which is synced and renamed over
|
||||
// name, and then the directory is synced, so that the rename lasts.
|
||||
// name. syncDirectory must follow, so that the rename lasts.
|
||||
func write(dir, name string, data []byte) error {
|
||||
path := filepath.Join(dir, name)
|
||||
temporary := path + ".tmp"
|
||||
@@ -475,10 +666,13 @@ func write(dir, name string, data []byte) error {
|
||||
|
||||
if err != nil {
|
||||
_ = os.Remove(temporary)
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
// syncDirectory syncs dir to the disk, so that a rename in it lasts.
|
||||
func syncDirectory(dir string) error {
|
||||
directory, err := os.Open(dir) //nolint:gosec // SWWAF_STATE_DIR itself
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
Reference in New Issue
Block a user