// Package remotelog sends the lines smallwebwaf writes on stdout to the // remote log endpoint, SWWAF_LOG_REMOTE_URL, as the "Request log" section // of SPEC.md describes: each line as the message of an RFC 5424 syslog // record, over UDP, TCP or TLS. Lines wait in a bounded buffer, so a slow // or unreachable endpoint never holds up a request or stdout. package remotelog import ( "bytes" "context" "crypto/tls" "crypto/x509" "errors" "fmt" "log/slog" "net" "net/url" "os" "strconv" "sync/atomic" "syscall" "time" "sneak.berlin/go/smallwebwaf/internal/requestlog" ) // The forms of SWWAF_LOG_REMOTE_URL, by its scheme. const ( SchemeUDP = "syslog+udp" SchemeTCP = "syslog+tcp" SchemeTLS = "syslog+tls" ) // A record's priority is the number of its facility times the number of // severities there are, plus the number of its severity. Every record's // severity is informational. const ( severities = 8 informational = 6 ) 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 ) // Params are what New needs. type Params struct { // URL is the endpoint (SWWAF_LOG_REMOTE_URL): SchemeUDP, SchemeTCP or // SchemeTLS, a host and a port. URL *url.URL // RootCAs are the certificates a SchemeTLS endpoint's certificate // must chain to (SWWAF_LOG_REMOTE_TLS_CA_FILE), nil for the host's. RootCAs *x509.CertPool // Buffer is the most lines held while they wait to be sent // (SWWAF_LOG_REMOTE_BUFFER). Buffer int // Facility is the number of the records' syslog facility // (SWWAF_LOG_REMOTE_FACILITY), and AppName their APP-NAME // (SWWAF_LOG_REMOTE_APP_NAME). Facility int AppName string } // Sender sends lines to the endpoint. Write puts them in its buffer, and // Run sends them from there. type Sender struct { url *url.URL tlsConfig *tls.Config // beforeTime and afterTime are the parts of every record's header // before and after its time, as RFC 5424 lays the header out. beforeTime string afterTime string // records is the buffer: each line's record, framed to be sent. records chan []byte sent atomic.Int64 dropped atomic.Int64 } // New returns a Sender for the endpoint params.URL. func New(params Params) *Sender { hostname, err := os.Hostname() if err != nil || hostname == "" { hostname = "-" // RFC 5424's value for a field that has none } priority := params.Facility*severities + informational return &Sender{ url: params.URL, tlsConfig: &tls.Config{ RootCAs: params.RootCAs, MinVersion: tls.VersionTLS12, }, // The 1 is the version of the format. The process id, the message // id and the structured data have no value. beforeTime: "<" + strconv.Itoa(priority) + ">1 ", afterTime: " " + hostname + " " + params.AppName + " - - - ", records: make(chan []byte, params.Buffer), } } // Write puts each line in p in the buffer, as the message of a record of // its own, and never waits: when the buffer is full, the oldest record in // it is dropped to make room. It is safe for concurrent use. func (s *Sender) Write(p []byte) (int, error) { at := requestlog.FormatTime(time.Now()) for line := range bytes.Lines(p) { line = bytes.TrimSuffix(line, []byte("\n")) if len(line) > 0 { s.put(s.record(at, line)) } } return len(p), nil } // Sent is how many records have been sent. func (s *Sender) Sent() int64 { return s.sent.Load() } // Dropped is how many records were dropped: the oldest in a full buffer, // and those whose sending failed. func (s *Sender) Dropped() int64 { return s.dropped.Load() } // Depth is how many records are in the buffer. func (s *Sender) Depth() int { return len(s.records) } // Run connects to the endpoint and sends each record as it comes into the // buffer, until ctx is done. Then it sends the records still in the buffer, // on the connection open at that time or, if there is none, on a new one, // 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. func (s *Sender) Run(ctx context.Context, processLog *slog.Logger) { conn := s.send(ctx, processLog) if conn == nil && len(s.records) > 0 { conn, _ = s.dial(context.WithoutCancel(ctx)) } if conn == nil { return } defer func() { _ = conn.Close() }() for { select { case record := <-s.records: if s.write(conn, record) != nil { return } default: return } } } // record returns line as an RFC 5424 record made at the time at, framed // for the endpoint: on its own over UDP, since each datagram holds one, // and over TCP and TLS after its length in bytes and a space, the // octet-counted framing of RFC 6587 and RFC 5425. func (s *Sender) record(at string, line []byte) []byte { record := make([]byte, 0, len(s.beforeTime)+len(at)+len(s.afterTime)+len(line)) record = append(record, s.beforeTime...) record = append(record, at...) record = append(record, s.afterTime...) record = append(record, line...) if s.url.Scheme == SchemeUDP { return record } return append([]byte(strconv.Itoa(len(record))+" "), record...) } // put adds record to the buffer, first dropping the oldest record in it // while it is full. func (s *Sender) put(record []byte) { for { select { case s.records <- record: return default: } select { case <-s.records: s.dropped.Add(1) default: } } } // send connects to the endpoint and sends each record as it comes into // the buffer, until ctx is done, and returns the connection then open, or // nil. func (s *Sender) send(ctx context.Context, processLog *slog.Logger) net.Conn { delay := firstRetryDelay for { conn, err := s.dial(ctx) if ctx.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) } } // sendOn sends each record on conn as it comes into the buffer, until one // fails, whose error it returns, or ctx is done. func (s *Sender) sendOn(ctx context.Context, conn net.Conn) error { for { select { case record := <-s.records: err := s.write(conn, record) if err != nil { return err } case <-ctx.Done(): return nil } } } // 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. 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) } s.sent.Add(1) return nil } // dial connects to the endpoint. func (s *Sender) dial(ctx context.Context) (net.Conn, error) { dialer := &net.Dialer{Timeout: dialTimeout} switch s.url.Scheme { case SchemeUDP: return dialer.DialContext(ctx, "udp", s.url.Host) case SchemeTLS: tlsDialer := &tls.Dialer{NetDialer: dialer, Config: s.tlsConfig} return tlsDialer.DialContext(ctx, "tcp", s.url.Host) default: return dialer.DialContext(ctx, "tcp", s.url.Host) } }