Compare commits

1 Commits
Author SHA1 Message Date
clawbot 88aaf89a93 Send every log line to a syslog server as well (closes #28)
check / check (push) Successful in 3m5s
With SWWAF_LOG_REMOTE_URL set (syslog+udp, syslog+tcp or syslog+tls),
every line on stdout is also sent as the message of an RFC 5424 record,
octet-counted over TCP and TLS, from a bounded buffer that drops its
oldest line when full, so a slow or unreachable server holds up nothing.
Failed connections are retried with backoff; lines sent, dropped and
waiting are metrics. At a stop the lines still waiting get at most two
seconds. SWWAF_LOG_REMOTE_APP_NAME defaults to SWWAF_INSTANCE_NAME; while
sending, an app name RFC 5424 does not allow stops the start. Standard
library only: log/syslog writes only the older format.

Model: opus-5-5
2026-10-06 19:49:24 +00:00
3 changed files with 170 additions and 44 deletions
+8 -6
View File
@@ -438,12 +438,14 @@ the app name `SWWAF_LOG_REMOTE_APP_NAME` gives. stdout is unchanged.
The lines wait in a buffer of `SWWAF_LOG_REMOTE_BUFFER` lines and are sent from The lines wait in a buffer of `SWWAF_LOG_REMOTE_BUFFER` lines and are sent from
there, so a server that is slow or cannot be reached never holds up a request or there, so a server that is slow or cannot be reached never holds up a request or
stdout. When the buffer is full, its oldest line is dropped to make room. A line stdout. When the buffer is full, its oldest line is dropped to make room. A line
whose sending fails is dropped too, and the connection is made again at once. A whose sending fails is dropped too, and the connection closed. That failure,
failed attempt to connect is logged and followed by the next a second later, like a failed attempt to connect, is logged and followed by the next attempt to
twice as long after each further failure up to a minute, and a second again once connect a second later, twice as long after each further failure up to a minute,
a connection is made. UDP gives no sign of what arrives, and over TCP and TLS a and a second again after a connection that stayed up for a minute before it
line sent on a connection the server has just closed can be lost before a failed. A line too long for one UDP datagram is dropped alone, with no wait and
failure shows; such a loss is not counted. nothing logged. UDP gives no sign of what arrives, and over TCP and TLS a line
sent on a connection the server has just closed can be lost before a failure
shows; such a loss is not counted.
As `smallwebwaf` stops, it sends the lines still waiting, on the connection open As `smallwebwaf` stops, it sends the lines still waiting, on the connection open
or a new one, for at most two seconds, and gives up the rest; stdout has carried or a new one, for at most two seconds, and gives up the rest; stdout has carried
+40 -26
View File
@@ -10,6 +10,7 @@ import (
"context" "context"
"crypto/tls" "crypto/tls"
"crypto/x509" "crypto/x509"
"errors"
"fmt" "fmt"
"log/slog" "log/slog"
"net" "net"
@@ -17,6 +18,7 @@ import (
"os" "os"
"strconv" "strconv"
"sync/atomic" "sync/atomic"
"syscall"
"time" "time"
"sneak.berlin/go/smallwebwaf/internal/requestlog" "sneak.berlin/go/smallwebwaf/internal/requestlog"
@@ -41,12 +43,15 @@ const (
// dialTimeout bounds connecting to the endpoint, the TLS handshake // dialTimeout bounds connecting to the endpoint, the TLS handshake
// included. // included.
dialTimeout = 10 * time.Second dialTimeout = 10 * time.Second
// After a failed attempt to connect, the next is made a second later, // After a failed attempt to connect, or a connection on which a record
// and retryDelayFactor times as long after each further failure in a // fails, the next attempt to connect is made a second later, and
// row, up to a minute. // retryDelayFactor times as long after each further failure in a row,
// up to a minute. A connection that fails after it has stayed up for
// resetRetryDelayAfter ends the row.
firstRetryDelay = time.Second firstRetryDelay = time.Second
retryDelayFactor = 2 retryDelayFactor = 2
maxRetryDelay = time.Minute maxRetryDelay = time.Minute
resetRetryDelayAfter = time.Minute
) )
// Params are what New needs. // Params are what New needs.
@@ -143,11 +148,12 @@ func (s *Sender) Depth() int {
// until none is left or one fails, and returns. How long it may take over // until none is left or one fails, and returns. How long it may take over
// that is for the caller to bound. // that is for the caller to bound.
// //
// A connection on which a record fails is closed, the record dropped, and // A connection on which a record fails is closed and the record dropped.
// a new one made at once. A failed attempt to connect is logged to // That failure, like a failed attempt to connect, is logged to processLog
// processLog and followed by the next after firstRetryDelay, // and followed by the next attempt after firstRetryDelay, retryDelayFactor
// retryDelayFactor times as long after each further failure in a row up // times as long after each further failure in a row up to maxRetryDelay,
// to maxRetryDelay. Meanwhile the records wait in the buffer. // and firstRetryDelay again after a connection that stayed up for
// resetRetryDelayAfter. Meanwhile the records wait in the buffer.
func (s *Sender) Run(ctx context.Context, processLog *slog.Logger) { func (s *Sender) Run(ctx context.Context, processLog *slog.Logger) {
conn := s.send(ctx, processLog) conn := s.send(ctx, processLog)
if conn == nil && len(s.records) > 0 { if conn == nil && len(s.records) > 0 {
@@ -218,12 +224,26 @@ func (s *Sender) send(ctx context.Context, processLog *slog.Logger) net.Conn {
for { for {
conn, err := s.dial(ctx) conn, err := s.dial(ctx)
if ctx.Err() != nil {
switch {
case ctx.Err() != nil:
return conn return conn
case err != nil: }
processLog.Warn("connecting to SWWAF_LOG_REMOTE_URL failed",
if err == nil {
connected := time.Now()
err = s.sendOn(ctx, conn)
if err == nil {
return conn
}
_ = conn.Close()
if time.Since(connected) >= resetRetryDelayAfter {
delay = firstRetryDelay
}
}
processLog.Warn("sending to SWWAF_LOG_REMOTE_URL failed",
"error", err.Error(), "connecting_again_in", delay.String()) "error", err.Error(), "connecting_again_in", delay.String())
select { select {
@@ -233,18 +253,6 @@ func (s *Sender) send(ctx context.Context, processLog *slog.Logger) net.Conn {
} }
delay = min(retryDelayFactor*delay, maxRetryDelay) delay = min(retryDelayFactor*delay, maxRetryDelay)
continue
}
delay = firstRetryDelay
err = s.sendOn(ctx, conn)
if err == nil {
return conn
}
_ = conn.Close()
} }
} }
@@ -265,12 +273,18 @@ func (s *Sender) sendOn(ctx context.Context, conn net.Conn) error {
} }
// write sends record on conn, and counts it as sent or, if that fails, // write sends record on conn, and counts it as sent or, if that fails,
// as dropped. // as dropped. A record too long for one UDP datagram is dropped without
// an error, since the connection has not failed: a long request must not
// hold up the lines after it.
func (s *Sender) write(conn net.Conn, record []byte) error { func (s *Sender) write(conn net.Conn, record []byte) error {
_, err := conn.Write(record) _, err := conn.Write(record)
if err != nil { if err != nil {
s.dropped.Add(1) s.dropped.Add(1)
if errors.Is(err, syscall.EMSGSIZE) {
return nil
}
return fmt.Errorf("send a record: %w", err) return fmt.Errorf("send a record: %w", err)
} }
+119 -9
View File
@@ -178,8 +178,7 @@ func TestReconnectsWithBackoffAfterTheEndpointGoesAway(t *testing.T) {
// The endpoint goes away: it closes the connection, and refuses the // The endpoint goes away: it closes the connection, and refuses the
// next ones. The sender notices when a record fails, and tries to // next ones. The sender notices when a record fails, and tries to
// connect again at once, then a second later, then two seconds // connect again a second later, then two seconds after that.
// after that.
endpoint.refusing.Store(true) endpoint.refusing.Store(true)
_ = conn.Close() _ = conn.Close()
@@ -206,18 +205,104 @@ func TestReconnectsWithBackoffAfterTheEndpointGoesAway(t *testing.T) {
conn = endpoint.next(t) conn = endpoint.next(t)
wantFrame(t, bufio.NewReader(conn), record(t, local0Info, appName, "two")) wantFrame(t, bufio.NewReader(conn), record(t, local0Info, appName, "two"))
wantRetries(t, logged, "1s", "2s") wantRetries(t, logged, "1s", "2s")
})
}
// Having connected, the sender waits a second again after the func TestAConnectionClosedAtOnceIsMadeAgainAfterAGrowingDelay(t *testing.T) {
// next failure. t.Parallel()
endpoint.refusing.Store(true)
synctest.Test(t, func(t *testing.T) {
endpoint := listen(t)
sender, logged, _ := run(t, params(remotelog.SchemeTCP, endpoint.Addr()))
// The endpoint closes each connection as soon as it takes it. The
// sender notices when a record fails, and connects again a second
// later, then two seconds after that, then four.
delays := []time.Duration{time.Second, 2 * time.Second, 4 * time.Second}
for i, delay := range delays {
_ = accept(t, endpoint).Close()
writeUntilDropped(t, sender, int64(i+1))
wantConnectedAgainAfter(t, sender, delay)
}
wantRetries(t, logged, "1s", "2s", "4s")
})
}
func TestTheDelayStartsAgainAfterAConnectionThatStayedUpAMinute(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
endpoint := listen(t)
sender, logged, _ := run(t, params(remotelog.SchemeTCP, endpoint.Addr()))
_ = accept(t, endpoint).Close()
writeUntilDropped(t, sender, 1)
wantConnectedAgainAfter(t, sender, time.Second)
// A connection that fails just short of a minute after it was made
// leaves the delay growing.
conn := accept(t, endpoint)
time.Sleep(time.Minute - time.Nanosecond)
_ = conn.Close() _ = conn.Close()
writeUntilDropped(t, sender, 2) writeUntilDropped(t, sender, 2)
wantConnectedAgainAfter(t, sender, 2*time.Second)
// One that fails a minute after it was made starts it again from a
// second.
conn = accept(t, endpoint)
time.Sleep(time.Minute)
_ = conn.Close()
writeUntilDropped(t, sender, 3)
wantConnectedAgainAfter(t, sender, time.Second)
wantRetries(t, logged, "1s", "2s", "1s") wantRetries(t, logged, "1s", "2s", "1s")
}) })
} }
func TestALineTooLongForADatagramIsDroppedAlone(t *testing.T) {
t.Parallel()
synctest.Test(t, func(t *testing.T) {
endpoint, err := (&net.ListenConfig{}).ListenPacket(t.Context(), "udp", loopback)
if err != nil {
t.Fatalf("listen: %v", err)
}
t.Cleanup(func() { _ = endpoint.Close() })
sender, logged, _ := run(t, params(remotelog.SchemeUDP, endpoint.LocalAddr()))
// With its header, the first line's record is longer than the 65507
// bytes a UDP datagram over IPv4 holds. It is dropped, nothing is
// logged, and the next line is sent at once.
_, _ = sender.Write([]byte(strings.Repeat("x", 65507) + "\nnext\n"))
synctest.Wait()
wantCounts(t, sender, 1, 1, 0)
wantRetries(t, logged)
datagram := make([]byte, 1024)
n, _, err := endpoint.ReadFrom(datagram)
if err != nil {
t.Fatalf("read: %v", err)
}
want := record(t, local0Info, appName, "next")
if string(datagram[:n]) != want {
t.Errorf("datagram %q, want %q", datagram[:n], want)
}
})
}
func TestRecordsWaitingAtTheStopAreSent(t *testing.T) { func TestRecordsWaitingAtTheStopAreSent(t *testing.T) {
t.Parallel() t.Parallel()
@@ -464,9 +549,34 @@ func writeUntilDropped(t *testing.T, sender *remotelog.Sender, dropped int64) {
} }
} }
// wantRetries checks that the sender logged a failed attempt to connect // wantConnectedAgainAfter writes a line while the sender waits to connect
// for each of delays, the time until the next attempt, in order, and // again, and checks that it connects, and takes the line from the buffer,
// logged nothing else. // only once delay is over.
func wantConnectedAgainAfter(
t *testing.T, sender *remotelog.Sender, delay time.Duration,
) {
t.Helper()
_, _ = sender.Write([]byte("waiting\n"))
time.Sleep(delay - time.Nanosecond)
synctest.Wait()
if sender.Depth() != 1 {
t.Fatalf("connected again before %v", delay)
}
time.Sleep(time.Nanosecond)
synctest.Wait()
if sender.Depth() != 0 {
t.Fatalf("not connected again after %v", delay)
}
}
// wantRetries checks that the sender logged a failure, of an attempt to
// connect or of a connection, for each of delays, the time until the next
// attempt, in order, and logged nothing else.
func wantRetries(t *testing.T, logged *output, delays ...string) { func wantRetries(t *testing.T, logged *output, delays ...string) {
t.Helper() t.Helper()
@@ -476,7 +586,7 @@ func wantRetries(t *testing.T, logged *output, delays ...string) {
var fields map[string]any var fields map[string]any
err := json.Unmarshal([]byte(line), &fields) err := json.Unmarshal([]byte(line), &fields)
if err != nil || fields["msg"] != "connecting to SWWAF_LOG_REMOTE_URL failed" { if err != nil || fields["msg"] != "sending to SWWAF_LOG_REMOTE_URL failed" {
t.Fatalf("logged %q", line) t.Fatalf("logged %q", line)
} }