Files

365 lines
8.6 KiB
Go

package main
import (
"bufio"
"context"
"io"
"log"
"net"
"strconv"
"strings"
"time"
"unicode/utf8"
)
var severityNames = [8]string{"emerg", "alert", "crit", "err", "warning", "notice", "info", "debug"}
var facilityNames = [24]string{
"kern", "user", "mail", "daemon", "auth", "syslog", "lpr", "news",
"uucp", "cron", "authpriv", "ftp", "ntp", "security", "console", "solaris-cron",
"local0", "local1", "local2", "local3", "local4", "local5", "local6", "local7",
}
// Entry is a decoded syslog message.
type Entry struct {
Time time.Time // timestamp from the message header (or reception time if absent)
Received time.Time // reception time, from the server clock
Host string
App string
ProcID string
Facility string
Severity string
SevNum int
SD string // raw RFC 5424 structured data
Message string
Source string // sender IP address
Proto string // udp or tcp
HostIP string // original host value when it was an IP resolved through DNS
SourceType string // "syslog" or "docker"
Extra map[string]string // additional fields (docker: container, image, …)
// Not stored. Wait: the producer can be slowed down when the queue is full
// (Docker, host logs) instead of losing the message. Done: called once the
// message is stored in VictoriaLogs or in the disk spool.
Wait bool
Done func()
}
// Record returns the entry in the shape stored in VictoriaLogs and returned by the API.
// The time axis (_time) is the reception time: device clocks can be wrong, and
// logs indexed in the past would escape time ranges. The device timestamp is
// kept in msg_time.
func (e *Entry) Record() map[string]string {
rcv := e.Received.UTC().Format(time.RFC3339Nano)
r := map[string]string{
"_time": rcv,
"received": rcv,
"msg_time": e.Time.UTC().Format(time.RFC3339Nano),
"_msg": e.Message,
"host": e.Host,
"app": e.App,
"facility": e.Facility,
"severity": e.Severity,
"sevnum": strconv.Itoa(e.SevNum),
"source": e.Source,
"proto": e.Proto,
}
if e.ProcID != "" {
r["procid"] = e.ProcID
}
if e.SD != "" {
r["sd"] = e.SD
}
if e.HostIP != "" {
r["host_ip"] = e.HostIP
}
if e.SourceType != "" {
r["source_type"] = e.SourceType
}
for k, v := range e.Extra {
if v != "" {
r[k] = v
}
}
return r
}
func (e *Entry) setPri(pri int) {
e.SevNum = pri % 8
e.Severity = severityNames[pri%8]
e.Facility = facilityNames[pri/8]
}
// ParseSyslog decodes an RFC 5424 or RFC 3164 (BSD) message. It never fails:
// a non-conforming message is kept as is.
func ParseSyslog(raw []byte, source, proto string, now time.Time) *Entry {
s := strings.TrimRight(string(raw), "\r\n\x00 ")
if !utf8.ValidString(s) {
s = strings.ToValidUTF8(s, "\uFFFD")
}
e := &Entry{Time: now, Received: now, Host: source, Source: source, Proto: proto, SourceType: "syslog"}
e.setPri(13) // user.notice, the RFC 3164 default
if len(s) > 2 && s[0] == '<' {
if end := strings.IndexByte(s, '>'); end > 1 && end <= 4 {
if pri, err := strconv.Atoi(s[1:end]); err == nil && pri >= 0 && pri <= 191 {
e.setPri(pri)
s = s[end+1:]
}
}
}
if strings.HasPrefix(s, "1 ") {
parse5424(e, s[2:])
} else {
parse3164(e, s)
}
if e.Message == "" {
e.Message = "(empty)"
}
return e
}
// parse5424: TIMESTAMP HOSTNAME APP-NAME PROCID MSGID STRUCTURED-DATA MSG
func parse5424(e *Entry, s string) {
f := strings.SplitN(s, " ", 6)
if len(f) < 6 {
e.Message = s
return
}
if t, err := time.Parse(time.RFC3339Nano, f[0]); err == nil {
e.Time = t
}
if f[1] != "-" {
e.Host = f[1]
}
if f[2] != "-" {
e.App = f[2]
}
if f[3] != "-" {
e.ProcID = f[3]
}
rest := f[5]
switch {
case rest == "-" || strings.HasPrefix(rest, "- "):
rest = strings.TrimPrefix(rest[1:], " ")
case strings.HasPrefix(rest, "["):
n := sdLen(rest)
e.SD = rest[:n]
rest = strings.TrimPrefix(rest[n:], " ")
}
e.Message = strings.TrimPrefix(rest, "\uFEFF") // optional UTF-8 BOM
}
// sdLen returns the length of the structured data block "[..][..]".
func sdLen(s string) int {
i := 0
for i < len(s) && s[i] == '[' {
inQuote := false
for i++; i < len(s); i++ {
c := s[i]
if inQuote && c == '\\' {
i++
continue
}
if c == '"' {
inQuote = !inQuote
continue
}
if c == ']' && !inQuote {
i++
break
}
}
}
if i > len(s) {
i = len(s)
}
return i
}
// parse3164: "Mmm dd hh:mm:ss HOST TAG[PID]: MSG", or an ISO timestamp (rsyslog).
func parse3164(e *Entry, s string) {
if len(s) >= 16 && s[15] == ' ' {
if t, err := time.ParseInLocation(time.Stamp, s[:15], time.Local); err == nil {
// No year in the timestamp: use the current one, unless that lands in the future.
now := e.Time
t = time.Date(now.Year(), t.Month(), t.Day(), t.Hour(), t.Minute(), t.Second(), 0, time.Local)
if t.After(now.Add(24 * time.Hour)) {
t = t.AddDate(-1, 0, 0)
}
e.Time = t
parseHostTag(e, s[16:], true)
return
}
}
if sp := strings.IndexByte(s, ' '); sp >= 19 {
if t, err := time.Parse(time.RFC3339Nano, s[:sp]); err == nil {
e.Time = t
parseHostTag(e, s[sp+1:], true)
return
}
}
parseHostTag(e, s, false)
}
func parseHostTag(e *Entry, s string, withHost bool) {
if withHost {
if sp := strings.IndexByte(s, ' '); sp > 0 {
tok := s[:sp]
if !strings.HasSuffix(tok, ":") && !strings.Contains(tok, "[") {
e.Host = tok
s = s[sp+1:]
}
}
}
if i := strings.IndexAny(s, ":[ "); i > 0 && i <= 48 {
switch s[i] {
case ':':
if i+1 == len(s) || s[i+1] == ' ' {
e.App = s[:i]
s = strings.TrimPrefix(s[i+1:], " ")
}
case '[':
if j := strings.IndexByte(s[i:], ']'); j > 1 {
e.App = s[:i]
e.ProcID = s[i+1 : i+j]
s = strings.TrimPrefix(strings.TrimPrefix(s[i+j+1:], ":"), " ")
}
}
}
e.Message = s
}
func serveUDP(ctx context.Context, pc net.PacketConn, sink func(*Entry)) {
buf := make([]byte, 65536)
for {
n, addr, err := pc.ReadFrom(buf)
if err != nil {
if ctx.Err() != nil {
return
}
log.Printf("syslog udp: %v", err)
time.Sleep(100 * time.Millisecond)
continue
}
sink(ParseSyslog(buf[:n], hostOf(addr), "udp", time.Now()))
}
}
// TCP limits: connections open at once, and how long a connection may stay
// silent before it is closed (senders reconnect on their own).
var (
tcpMaxConns = 512
tcpIdle = 30 * time.Minute
)
func serveTCP(ctx context.Context, ln net.Listener, sink func(*Entry)) {
slots := make(chan struct{}, tcpMaxConns)
for {
conn, err := ln.Accept()
if err != nil {
if ctx.Err() != nil {
return
}
log.Printf("syslog tcp: %v", err)
time.Sleep(100 * time.Millisecond)
continue
}
select {
case slots <- struct{}{}:
default:
log.Printf("syslog tcp: %d connections already open, refusing %s", tcpMaxConns, conn.RemoteAddr())
conn.Close()
continue
}
go func() {
defer func() { <-slots }()
handleTCP(ctx, conn, sink)
}()
}
}
// idleConn pushes the read deadline back before each read, so only a
// connection that stays silent for tcpIdle is closed.
type idleConn struct {
net.Conn
idle time.Duration
}
func (c idleConn) Read(p []byte) (int, error) {
_ = c.Conn.SetReadDeadline(time.Now().Add(c.idle))
return c.Conn.Read(p)
}
const maxFrame = 1 << 20
// handleTCP supports both RFC 6587 framings: octet counting
// ("123 <34>1 ...") and one message per line.
func handleTCP(ctx context.Context, conn net.Conn, sink func(*Entry)) {
defer conn.Close()
stop := context.AfterFunc(ctx, func() { conn.Close() })
defer stop()
src := hostOf(conn.RemoteAddr())
r := bufio.NewReaderSize(idleConn{conn, tcpIdle}, 64*1024)
for {
c, err := r.ReadByte()
if err != nil {
return
}
switch {
case c == '\n' || c == '\r' || c == ' ' || c == 0:
continue
case c >= '1' && c <= '9':
n := int(c - '0')
for {
c, err = r.ReadByte()
if err != nil {
return
}
if c == ' ' {
break
}
if c < '0' || c > '9' || n > maxFrame {
log.Printf("syslog tcp %s: invalid frame, closing connection", src)
return
}
n = n*10 + int(c-'0')
}
if n > maxFrame {
return
}
msg := make([]byte, n)
if _, err := io.ReadFull(r, msg); err != nil {
return
}
sink(ParseSyslog(msg, src, "tcp", time.Now()))
default:
_ = r.UnreadByte()
line, err := r.ReadSlice('\n')
if len(line) > 0 {
sink(ParseSyslog(line, src, "tcp", time.Now()))
}
// Line longer than the buffer: skip the rest of it.
for err == bufio.ErrBufferFull {
_, err = r.ReadSlice('\n')
}
if err != nil {
return
}
}
}
}
func hostOf(a net.Addr) string {
if a == nil {
return ""
}
h, _, err := net.SplitHostPort(a.String())
if err != nil {
return a.String()
}
return h
}