Compare commits

..
1 Commits
Author SHA1 Message Date
clawbot 51555acab0 Send every log line to a syslog server as well (closes #28)
check / check (push) Successful in 3m34s
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:09:05 +00:00
3 changed files with 51 additions and 177 deletions
+6 -8
View File
@@ -438,14 +438,12 @@ 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
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
whose sending fails is dropped too, and the connection closed. That failure,
like a failed attempt to connect, is logged and followed by the next attempt to
connect a second later, twice as long after each further failure up to a minute,
and a second again after a connection that stayed up for a minute before it
failed. A line too long for one UDP datagram is dropped alone, with no wait and
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.
whose sending fails is dropped too, and the connection is made again at once. A
failed attempt to connect is logged and followed by the next a second later,
twice as long after each further failure up to a minute, and a second again once
a connection is made. 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
or a new one, for at most two seconds, and gives up the rest; stdout has carried
+36 -50
View File
@@ -10,7 +10,6 @@ import (
"context"
"crypto/tls"
"crypto/x509"
"errors"
"fmt"
"log/slog"
"net"
@@ -18,7 +17,6 @@ import (
"os"
"strconv"
"sync/atomic"
"syscall"
"time"
"sneak.berlin/go/smallwebwaf/internal/requestlog"
@@ -43,15 +41,12 @@ const (
// dialTimeout bounds connecting to the endpoint, the TLS handshake
// included.
dialTimeout = 10 * time.Second
// After a failed attempt to connect, or a connection on which a record
// fails, the next attempt to connect is made a second later, and
// 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
retryDelayFactor = 2
maxRetryDelay = time.Minute
resetRetryDelayAfter = time.Minute
// After a failed attempt to connect, the next is made a second later,
// and retryDelayFactor times as long after each further failure in a
// row, up to a minute.
firstRetryDelay = time.Second
retryDelayFactor = 2
maxRetryDelay = time.Minute
)
// Params are what New needs.
@@ -148,12 +143,11 @@ func (s *Sender) Depth() int {
// until none is left or one fails, and returns. How long it may take over
// that is for the caller to bound.
//
// A connection on which a record fails is closed and the record dropped.
// That failure, like a failed attempt to connect, is logged to processLog
// and followed by the next attempt after firstRetryDelay, retryDelayFactor
// times as long after each further failure in a row up to maxRetryDelay,
// and firstRetryDelay again after a connection that stayed up for
// resetRetryDelayAfter. Meanwhile the records wait in the buffer.
// A connection on which a record fails is closed, the record dropped, and
// a new one made at once. A failed attempt to connect is logged to
// processLog and followed by the next after firstRetryDelay,
// retryDelayFactor times as long after each further failure in a row up
// to maxRetryDelay. Meanwhile the records wait in the buffer.
func (s *Sender) Run(ctx context.Context, processLog *slog.Logger) {
conn := s.send(ctx, processLog)
if conn == nil && len(s.records) > 0 {
@@ -224,35 +218,33 @@ func (s *Sender) send(ctx context.Context, processLog *slog.Logger) net.Conn {
for {
conn, err := s.dial(ctx)
if ctx.Err() != nil {
switch {
case ctx.Err() != nil:
return conn
case err != nil:
processLog.Warn("connecting to SWWAF_LOG_REMOTE_URL failed",
"error", err.Error(), "connecting_again_in", delay.String())
select {
case <-time.After(delay):
case <-ctx.Done():
return nil
}
delay = min(retryDelayFactor*delay, maxRetryDelay)
continue
}
delay = firstRetryDelay
err = s.sendOn(ctx, conn)
if err == nil {
return conn
}
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())
select {
case <-time.After(delay):
case <-ctx.Done():
return nil
}
delay = min(retryDelayFactor*delay, maxRetryDelay)
_ = conn.Close()
}
}
@@ -273,18 +265,12 @@ 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,
// 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.
// as dropped.
func (s *Sender) write(conn net.Conn, record []byte) error {
_, err := conn.Write(record)
if err != nil {
s.dropped.Add(1)
if errors.Is(err, syscall.EMSGSIZE) {
return nil
}
return fmt.Errorf("send a record: %w", err)
}
+9 -119
View File
@@ -178,7 +178,8 @@ func TestReconnectsWithBackoffAfterTheEndpointGoesAway(t *testing.T) {
// The endpoint goes away: it closes the connection, and refuses the
// next ones. The sender notices when a record fails, and tries to
// connect again a second later, then two seconds after that.
// connect again at once, then a second later, then two seconds
// after that.
endpoint.refusing.Store(true)
_ = conn.Close()
@@ -205,104 +206,18 @@ func TestReconnectsWithBackoffAfterTheEndpointGoesAway(t *testing.T) {
conn = endpoint.next(t)
wantFrame(t, bufio.NewReader(conn), record(t, local0Info, appName, "two"))
wantRetries(t, logged, "1s", "2s")
})
}
func TestAConnectionClosedAtOnceIsMadeAgainAfterAGrowingDelay(t *testing.T) {
t.Parallel()
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)
// Having connected, the sender waits a second again after the
// next failure.
endpoint.refusing.Store(true)
_ = conn.Close()
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")
})
}
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) {
t.Parallel()
@@ -549,34 +464,9 @@ func writeUntilDropped(t *testing.T, sender *remotelog.Sender, dropped int64) {
}
}
// wantConnectedAgainAfter writes a line while the sender waits to connect
// again, and checks that it connects, and takes the line from the buffer,
// 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.
// wantRetries checks that the sender logged a failed attempt to connect
// 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) {
t.Helper()
@@ -586,7 +476,7 @@ func wantRetries(t *testing.T, logged *output, delays ...string) {
var fields map[string]any
err := json.Unmarshal([]byte(line), &fields)
if err != nil || fields["msg"] != "sending to SWWAF_LOG_REMOTE_URL failed" {
if err != nil || fields["msg"] != "connecting to SWWAF_LOG_REMOTE_URL failed" {
t.Fatalf("logged %q", line)
}