Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
88aaf89a93 |
@@ -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
|
||||
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 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.
|
||||
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.
|
||||
|
||||
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
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net"
|
||||
@@ -17,6 +18,7 @@ import (
|
||||
"os"
|
||||
"strconv"
|
||||
"sync/atomic"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"sneak.berlin/go/smallwebwaf/internal/requestlog"
|
||||
@@ -41,12 +43,15 @@ const (
|
||||
// dialTimeout bounds connecting to the endpoint, the TLS handshake
|
||||
// included.
|
||||
dialTimeout = 10 * time.Second
|
||||
// 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.
|
||||
// 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
|
||||
)
|
||||
|
||||
// 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
|
||||
// that is for the caller to bound.
|
||||
//
|
||||
// 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.
|
||||
// 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.
|
||||
func (s *Sender) Run(ctx context.Context, processLog *slog.Logger) {
|
||||
conn := s.send(ctx, processLog)
|
||||
if conn == nil && len(s.records) > 0 {
|
||||
@@ -218,12 +224,26 @@ func (s *Sender) send(ctx context.Context, processLog *slog.Logger) net.Conn {
|
||||
|
||||
for {
|
||||
conn, err := s.dial(ctx)
|
||||
|
||||
switch {
|
||||
case ctx.Err() != nil:
|
||||
if ctx.Err() != nil {
|
||||
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())
|
||||
|
||||
select {
|
||||
@@ -233,18 +253,6 @@ func (s *Sender) send(ctx context.Context, processLog *slog.Logger) net.Conn {
|
||||
}
|
||||
|
||||
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,
|
||||
// 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 {
|
||||
_, 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)
|
||||
}
|
||||
|
||||
|
||||
@@ -178,8 +178,7 @@ 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 at once, then a second later, then two seconds
|
||||
// after that.
|
||||
// connect again a second later, then two seconds after that.
|
||||
endpoint.refusing.Store(true)
|
||||
|
||||
_ = conn.Close()
|
||||
@@ -206,18 +205,104 @@ func TestReconnectsWithBackoffAfterTheEndpointGoesAway(t *testing.T) {
|
||||
conn = endpoint.next(t)
|
||||
wantFrame(t, bufio.NewReader(conn), record(t, local0Info, appName, "two"))
|
||||
wantRetries(t, logged, "1s", "2s")
|
||||
})
|
||||
}
|
||||
|
||||
// Having connected, the sender waits a second again after the
|
||||
// next failure.
|
||||
endpoint.refusing.Store(true)
|
||||
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)
|
||||
|
||||
_ = 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()
|
||||
|
||||
@@ -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
|
||||
// for each of delays, the time until the next attempt, in order, and
|
||||
// logged nothing else.
|
||||
// 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.
|
||||
func wantRetries(t *testing.T, logged *output, delays ...string) {
|
||||
t.Helper()
|
||||
|
||||
@@ -476,7 +586,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"] != "connecting to SWWAF_LOG_REMOTE_URL failed" {
|
||||
if err != nil || fields["msg"] != "sending to SWWAF_LOG_REMOTE_URL failed" {
|
||||
t.Fatalf("logged %q", line)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user