Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
51555acab0 |
@@ -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
|
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 closed. That failure,
|
whose sending fails is dropped too, and the connection is made again at once. A
|
||||||
like a failed attempt to connect, is logged and followed by the next attempt to
|
failed attempt to connect is logged and followed by the next a second later,
|
||||||
connect a second later, twice as long after each further failure up to a minute,
|
twice as long after each further failure up to a minute, and a second again once
|
||||||
and a second again after a connection that stayed up for a minute before it
|
a connection is made. UDP gives no sign of what arrives, and over TCP and TLS a
|
||||||
failed. A line too long for one UDP datagram is dropped alone, with no wait and
|
line sent on a connection the server has just closed can be lost before a
|
||||||
nothing logged. UDP gives no sign of what arrives, and over TCP and TLS a line
|
failure shows; such a loss is not counted.
|
||||||
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,7 +10,6 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"crypto/tls"
|
"crypto/tls"
|
||||||
"crypto/x509"
|
"crypto/x509"
|
||||||
"errors"
|
|
||||||
"fmt"
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net"
|
"net"
|
||||||
@@ -18,7 +17,6 @@ import (
|
|||||||
"os"
|
"os"
|
||||||
"strconv"
|
"strconv"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"syscall"
|
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"sneak.berlin/go/smallwebwaf/internal/requestlog"
|
"sneak.berlin/go/smallwebwaf/internal/requestlog"
|
||||||
@@ -43,15 +41,12 @@ 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, or a connection on which a record
|
// After a failed attempt to connect, the next is made a second later,
|
||||||
// fails, the next attempt to connect is made a second later, and
|
// and retryDelayFactor times as long after each further failure in a
|
||||||
// retryDelayFactor times as long after each further failure in a row,
|
// row, up to a minute.
|
||||||
// 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.
|
||||||
@@ -148,12 +143,11 @@ 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 and the record dropped.
|
// A connection on which a record fails is closed, the record dropped, and
|
||||||
// That failure, like a failed attempt to connect, is logged to processLog
|
// a new one made at once. A failed attempt to connect is logged to
|
||||||
// and followed by the next attempt after firstRetryDelay, retryDelayFactor
|
// processLog and followed by the next after firstRetryDelay,
|
||||||
// times as long after each further failure in a row up to maxRetryDelay,
|
// retryDelayFactor times as long after each further failure in a row up
|
||||||
// and firstRetryDelay again after a connection that stayed up for
|
// to maxRetryDelay. Meanwhile the records wait in the buffer.
|
||||||
// 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 {
|
||||||
@@ -224,26 +218,12 @@ 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 {
|
||||||
@@ -253,6 +233,18 @@ 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()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -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,
|
// 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
|
// as dropped.
|
||||||
// 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,7 +178,8 @@ 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 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)
|
endpoint.refusing.Store(true)
|
||||||
|
|
||||||
_ = conn.Close()
|
_ = conn.Close()
|
||||||
@@ -205,104 +206,18 @@ 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")
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestAConnectionClosedAtOnceIsMadeAgainAfterAGrowingDelay(t *testing.T) {
|
// Having connected, the sender waits a second again after the
|
||||||
t.Parallel()
|
// next failure.
|
||||||
|
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()
|
||||||
|
|
||||||
@@ -549,34 +464,9 @@ func writeUntilDropped(t *testing.T, sender *remotelog.Sender, dropped int64) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// wantConnectedAgainAfter writes a line while the sender waits to connect
|
// wantRetries checks that the sender logged a failed attempt to connect
|
||||||
// again, and checks that it connects, and takes the line from the buffer,
|
// for each of delays, the time until the next attempt, in order, and
|
||||||
// only once delay is over.
|
// logged nothing else.
|
||||||
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()
|
||||||
|
|
||||||
@@ -586,7 +476,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"] != "sending to SWWAF_LOG_REMOTE_URL failed" {
|
if err != nil || fields["msg"] != "connecting to SWWAF_LOG_REMOTE_URL failed" {
|
||||||
t.Fatalf("logged %q", line)
|
t.Fatalf("logged %q", line)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user