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
|
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
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user