Skip to content
Merged
2 changes: 1 addition & 1 deletion common/log_parser_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ import (
func TestLogParser(t *testing.T) {
const defaultPatternsPerLevel = 256
ch := make(chan logparser.LogEntry)
parser := logparser.NewParser(ch, nil, nil, 1*time.Second, defaultPatternsPerLevel, logparser.SensitiveConfig{
parser := logparser.NewParser(ch, nil, nil, 1*time.Second, defaultPatternsPerLevel, false, nil, logparser.SensitiveConfig{
Enabled: true,
MinConfidence: "high",
MaxDetections: 100,
Expand Down
42 changes: 34 additions & 8 deletions containers/container.go
Original file line number Diff line number Diff line change
Expand Up @@ -682,7 +682,8 @@ func (c *Container) onFileOpen(pid uint32, fd uint64, mnt uint64, log bool) {
return
}
}
mntId, logPath := resolveFd(pid, fd)
info := proc.GetFdInfo(pid, fd)
mntId, logPath := resolveFd(info)
func() {
if mntId == "" {
return
Expand All @@ -706,8 +707,17 @@ func (c *Container) onFileOpen(pid uint32, fd uint64, mnt uint64, log bool) {
c.lock.Unlock()
}
}()
if logPath != "" {
if *flags.EnableDynamicLogTailing {
if *flags.EnableDynamicLogTailing {
var logPaths []string
switch {
case logPath != "":
logPaths = []string{logPath}
case log && (info == nil || !strings.HasPrefix(info.Dest, "/var/log/")):
// The kernel saw a /var/log/ file opened, but the fd now points
// elsewhere (freopen): find the log files the process holds.
logPaths = findLogFiles(pid)
}
for _, logPath := range logPaths {
c.lock.Lock()
c.runLogParser(logPath)
c.lock.Unlock()
Expand Down Expand Up @@ -2112,7 +2122,7 @@ func (c *Container) runLogParser(logPath string) {
return
}
ch := make(chan logparser.LogEntry)
parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, sensitiveCfg)
parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter(), sensitiveCfg)
reader, err := logs.NewTailReader(proc.HostPath(logPath), ch)
if err != nil {
klog.Warningln(err)
Expand All @@ -2132,7 +2142,7 @@ func (c *Container) runLogParser(logPath string) {
klog.Warningln(err)
return
}
parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, sensitiveCfg)
parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter(), sensitiveCfg)
stop := func() {
JournaldUnsubscribe(c.metadata.systemd.Unit)
}
Expand All @@ -2149,7 +2159,7 @@ func (c *Container) runLogParser(logPath string) {
delete(c.logParsers, "stdout/stderr")
}
ch := make(chan logparser.LogEntry)
parser := logparser.NewParser(ch, c.metadata.logDecoder, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, sensitiveCfg)
parser := logparser.NewParser(ch, c.metadata.logDecoder, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter(), sensitiveCfg)
reader, err := logs.NewTailReader(proc.HostPath(c.metadata.logPath), ch)
if err != nil {
klog.Warningln(err)
Expand Down Expand Up @@ -2445,8 +2455,7 @@ func countTLSAttach(lib string, result ebpftracer.TLSAttachResult) {
}
}

func resolveFd(pid uint32, fd uint64) (mntId string, logPath string) {
info := proc.GetFdInfo(pid, fd)
func resolveFd(info *proc.FdInfo) (mntId string, logPath string) {
if info == nil {
return
}
Expand Down Expand Up @@ -2498,3 +2507,20 @@ func sampleString(v interface{}) string {
}
return ""
}

func findLogFiles(pid uint32) []string {
fds, err := proc.ReadFds(pid)
if err != nil {
return nil
}
var res []string
for _, fd := range fds {
if !strings.HasPrefix(fd.Dest, "/var/log/") {
continue
}
if _, logPath := resolveFd(proc.GetFdInfo(pid, fd.Fd)); logPath != "" {
res = append(res, logPath)
}
}
return res
}
Comment thread
mayankpande88 marked this conversation as resolved.
34 changes: 18 additions & 16 deletions flags/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,13 @@ import (
)

var (
ListenAddress = kingpin.Flag("listen", "Listen address - ip:port or :port (default 0.0.0.0:80, or 127.0.0.1:10300 when --metrics-endpoint is set)").Envar("LISTEN").String()
CgroupRoot = kingpin.Flag("cgroupfs-root", "The mount point of the host cgroupfs root").Default("/sys/fs/cgroup").Envar("CGROUPFS_ROOT").String()
DisableLogParsing = kingpin.Flag("disable-log-parsing", "Disable container log parsing").Default("false").Envar("DISABLE_LOG_PARSING").Bool()
DisablePinger = kingpin.Flag("disable-pinger", "Don't ping upstreams").Default("true").Envar("DISABLE_PINGER").Bool()
DisableL7Tracing = kingpin.Flag("disable-l7-tracing", "Disable L7 tracing").Default("false").Envar("DISABLE_L7_TRACING").Bool()
EnableDotNetTracing = kingpin.Flag("enable-dotnet-tracing", "Enable .NET CLR tracing").Default("false").Envar("ENABLE_DOTNET_TRACING").Bool()
ListenAddress = kingpin.Flag("listen", "Listen address - ip:port or :port (default 0.0.0.0:80, or 127.0.0.1:10300 when --metrics-endpoint is set)").Envar("LISTEN").String()
CgroupRoot = kingpin.Flag("cgroupfs-root", "The mount point of the host cgroupfs root").Default("/sys/fs/cgroup").Envar("CGROUPFS_ROOT").String()
DisableLogParsing = kingpin.Flag("disable-log-parsing", "Disable container log parsing").Default("false").Envar("DISABLE_LOG_PARSING").Bool()
DisableJsonLogParsing = kingpin.Flag("disable-json-log-parsing", "Disable extracting the message, severity, and attributes from JSON-formatted logs").Default("false").Envar("DISABLE_JSON_LOG_PARSING").Bool()
DisablePinger = kingpin.Flag("disable-pinger", "Don't ping upstreams").Default("true").Envar("DISABLE_PINGER").Bool()
DisableL7Tracing = kingpin.Flag("disable-l7-tracing", "Disable L7 tracing").Default("false").Envar("DISABLE_L7_TRACING").Bool()
EnableDotNetTracing = kingpin.Flag("enable-dotnet-tracing", "Enable .NET CLR tracing").Default("false").Envar("ENABLE_DOTNET_TRACING").Bool()
// Off by default: the only thing it produces is
// container_nodejs_event_loop_blocked_time_seconds_total, which nothing
// currently consumes, and attaching the probes reads the whole ELF symbol
Expand Down Expand Up @@ -44,16 +45,17 @@ var (
Strings()
EphemeralPortRange = kingpin.Flag("ephemeral-port-range", `Destination and Listen TCP ports from these ranges will be skipped, e.g. "32768-60999" or "1024-23768 30000-65535"`).Default("32768-60999").Envar("EPHEMERAL_PORT_RANGE").String()

Provider = kingpin.Flag("provider", "`provider` label for `node_cloud_info` metric").Envar("PROVIDER").String()
Region = kingpin.Flag("region", "`region` label for `node_cloud_info` metric").Envar("REGION").String()
AvailabilityZone = kingpin.Flag("availability-zone", "`availability_zone` label for `node_cloud_info` metric").Envar("AVAILABILITY_ZONE").String()
AccountId = kingpin.Flag("account-id", "`account_id` label for `node_cloud_info` metric").Envar("ACCOUNT_ID").String()
InstanceType = kingpin.Flag("instance-type", "`instance_type` label for `node_cloud_info` metric").Envar("INSTANCE_TYPE").String()
InstanceLifeCycle = kingpin.Flag("instance-life-cycle", "`instance_life_cycle` label for `node_cloud_info` metric").Envar("INSTANCE_LIFE_CYCLE").String()
LogPerSecond = kingpin.Flag("log-per-second", "The number of logs per second").Default("10.0").Envar("LOG_PER_SECOND").Float64()
LogBurst = kingpin.Flag("log-burst", "The maximum number of tokens that can be consumed in a single call to allow").Default("100").Envar("LOG_BURST").Int()
LogPatternsPerContainer = kingpin.Flag("log-patterns-per-container", "Max unique log patterns per container per level").Default("256").Envar("LOG_PATTERNS_PER_CONTAINER").Int()
MaxFQDNsPerContainer = kingpin.Flag("max-fqdns-per-container", "Max unique FQDN values per container, extras are bucketed under '~other'").Default("50").Envar("MAX_FQDNS_PER_CONTAINER").Int()
Provider = kingpin.Flag("provider", "`provider` label for `node_cloud_info` metric").Envar("PROVIDER").String()
Region = kingpin.Flag("region", "`region` label for `node_cloud_info` metric").Envar("REGION").String()
AvailabilityZone = kingpin.Flag("availability-zone", "`availability_zone` label for `node_cloud_info` metric").Envar("AVAILABILITY_ZONE").String()
AccountId = kingpin.Flag("account-id", "`account_id` label for `node_cloud_info` metric").Envar("ACCOUNT_ID").String()
InstanceType = kingpin.Flag("instance-type", "`instance_type` label for `node_cloud_info` metric").Envar("INSTANCE_TYPE").String()
InstanceLifeCycle = kingpin.Flag("instance-life-cycle", "`instance_life_cycle` label for `node_cloud_info` metric").Envar("INSTANCE_LIFE_CYCLE").String()
LogPerSecond = kingpin.Flag("log-per-second", "The number of logs per second").Default("10.0").Envar("LOG_PER_SECOND").Float64()
LogBurst = kingpin.Flag("log-burst", "The maximum number of tokens that can be consumed in a single call to allow").Default("100").Envar("LOG_BURST").Int()
LogPatternsPerContainer = kingpin.Flag("log-patterns-per-container", "Max unique log patterns per container per level").Default("256").Envar("LOG_PATTERNS_PER_CONTAINER").Int()
LogPatternExtractionLimit = kingpin.Flag("log-pattern-extraction-limit", "Max warning/error log messages per second per container for which patterns are extracted. Over-limit messages are still exported and counted, under a dedicated 'event was sampled' pattern (0 - unlimited)").Default("0").Envar("LOG_PATTERN_EXTRACTION_LIMIT").Float64()
MaxFQDNsPerContainer = kingpin.Flag("max-fqdns-per-container", "Max unique FQDN values per container, extras are bucketed under '~other'").Default("50").Envar("MAX_FQDNS_PER_CONTAINER").Int()

MaxLabelLength = kingpin.Flag("max-label-length", "Maximum length of a metric label value").Default("4096").Envar("MAX_LABEL_LENGTH").Int()

Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ require (
github.com/hashicorp/golang-lru/v2 v2.0.7
github.com/jpillora/backoff v1.0.0
github.com/mdlayher/taskstats v0.0.0-20230712191918-387b3d561d14
github.com/nudgebee/logparser v0.0.0-20260406040008-c408d0a5c69c
github.com/nudgebee/logparser v0.0.0-20261008200612-6da0f06ad2f3
github.com/opencontainers/runtime-spec v1.1.0
github.com/prometheus/client_golang v1.20.5
github.com/prometheus/client_model v0.6.2
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -295,8 +295,8 @@ github.com/morikuni/aec v1.0.0 h1:nP9CBfwrvYnBRgY6qfDQkygYDmYwOilePFkwzv4dU8A=
github.com/morikuni/aec v1.0.0/go.mod h1:BbKIizmSmc5MMPqRYbxO4ZU0S0+P200+tUnFx7PXmsc=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
github.com/nudgebee/logparser v0.0.0-20260406040008-c408d0a5c69c h1:drsGFCiBXp5OZm92HE0iVx8U4XhpzAxL5vVKubBGl+I=
github.com/nudgebee/logparser v0.0.0-20260406040008-c408d0a5c69c/go.mod h1:oFnM9D6YEjZzb1jy0kJ/NUkTbkZA+BgNYjpLO4n4szA=
github.com/nudgebee/logparser v0.0.0-20261008200612-6da0f06ad2f3 h1:LsztuTRobqh4hdv5d1hCRFs++NUET+AJiKSzh1za2KA=
github.com/nudgebee/logparser v0.0.0-20261008200612-6da0f06ad2f3/go.mod h1:mUWBZ/SsHHevKDuVZ7MzFwIT5oLprAFRlp8oYq/Q/0U=
github.com/oklog/ulid v1.3.1 h1:EGfNDEx6MqHz8B3uNV6QAib1UR2Lm97sHi3ocA6ESJ4=
github.com/oklog/ulid v1.3.1/go.mod h1:CirwcVhetQ6Lv90oh/F+FBtV6XMibvdAFo93nm5qn4U=
github.com/onsi/ginkgo v1.16.5 h1:8xi0RTUf59SOSfEtZMvwTvXYMzG4gV23XVHOZiXNtnE=
Expand Down
86 changes: 82 additions & 4 deletions logs/otel.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package logs

import (
"context"
"strings"
"time"

otel "github.com/agoda-com/opentelemetry-logs-go"
Expand All @@ -15,11 +16,25 @@ import (
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/sdk/resource"
semconv "go.opentelemetry.io/otel/semconv/v1.18.0"
"go.opentelemetry.io/otel/trace"
"golang.org/x/time/rate"
"k8s.io/klog/v2"
)

var otelLogger otelLogs.Logger

// PatternExtractionRateLimiter caps, per container, how many warning and
// error messages per second get a log pattern extracted (nil: unlimited).
func PatternExtractionRateLimiter() *rate.Limiter {
limit := *flags.LogPatternExtractionLimit
if limit <= 0 {
return nil
}
// At least 1: with a burst of 0 the limiter rejects every event, so a
// limit below 0.1/s would turn pattern extraction off entirely.
return rate.NewLimiter(rate.Limit(limit), max(1, int(limit*10)))
}
Comment thread
mayankpande88 marked this conversation as resolved.

func Init(machineId, hostname, version string) {
endpointUrl := *flags.LogsEndpoint
if endpointUrl == nil {
Expand Down Expand Up @@ -60,11 +75,72 @@ func Init(machineId, hostname, version string) {
otelLogger = loggerProvider.Logger("nudgebee-node-agent", otelLogs.WithInstrumentationVersion(version))
}

const (
traceIdKey = "traceid"
spanIdKey = "spanid"
)

var (
traceIdKeys = []string{traceIdKey, "trace_id", "trace-id", "trace.id"}
spanIdKeys = []string{spanIdKey, "span_id", "span-id", "span.id"}
)

func normalizeTraceContextKey(k string) string {
k = strings.ToLower(k)
switch k {
case "@tr":
return traceIdKey
case "@sp":
return spanIdKey
}
if strings.Contains(k, "parent") { // e.g. parent_span_id is not the record's span id
return ""
}
for _, s := range traceIdKeys {
if k == s || strings.HasSuffix(k, "."+s) {
return traceIdKey
}
}
for _, s := range spanIdKeys {
if k == s || strings.HasSuffix(k, "."+s) {
return spanIdKey
}
}
return ""
}

func logRecordAttrs(patternHash string, attributes map[string]string) ([]attribute.KeyValue, *trace.TraceID, *trace.SpanID) {
var traceId *trace.TraceID
var spanId *trace.SpanID
attrs := make([]attribute.KeyValue, 0, len(attributes)+1)
attrs = append(attrs, attribute.Key("pattern.hash").String(patternHash))
for k, v := range attributes {
switch normalizeTraceContextKey(k) {
case traceIdKey:
if traceId == nil {
if id, err := trace.TraceIDFromHex(strings.ToLower(v)); err == nil {
traceId = &id
continue
}
}
case spanIdKey:
if spanId == nil {
if id, err := trace.SpanIDFromHex(strings.ToLower(v)); err == nil {
spanId = &id
continue
}
}
}
attrs = append(attrs, attribute.Key(k).String(v))
}
return attrs, traceId, spanId
}

func OtelLogEmitter(containerId string) logparser.OnMsgCallbackF {
if otelLogger == nil {
return nil
}
return func(ts time.Time, level logparser.Level, patternHash string, msg string) {
return func(ts time.Time, level logparser.Level, patternHash string, msg string, attributes map[string]string) {
severityText := level.String()
severityNumber := otelLogs.UNSPECIFIED
switch level {
Expand All @@ -80,19 +156,21 @@ func OtelLogEmitter(containerId string) logparser.OnMsgCallbackF {
severityNumber = otelLogs.DEBUG
}

attrs, traceId, spanId := logRecordAttrs(patternHash, attributes)

otelLogger.Emit(
otelLogs.NewLogRecord(otelLogs.LogRecordConfig{
ObservedTimestamp: ts,
TraceId: traceId,
SpanId: spanId,
SeverityText: &severityText,
SeverityNumber: &severityNumber,
Body: &msg,
Resource: resource.NewSchemaless(
semconv.ServiceName(common.ContainerIdToOtelServiceName(containerId)),
semconv.ContainerID(containerId),
),
Attributes: &[]attribute.KeyValue{
attribute.Key("pattern.hash").String(patternHash),
},
Attributes: &attrs,
}),
)
}
Expand Down
24 changes: 15 additions & 9 deletions logs/tail_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,18 +37,20 @@ func NewTailReader(fileName string, ch chan<- logparser.LogEntry) (*TailReader,
stopped: make(chan struct{}),
}
var err error
if r.file, err = os.Open(fileName); err != nil {
if r.file, err = os.Open(fileName); err != nil && !os.IsNotExist(err) {
return nil, err
}
if r.info, err = r.file.Stat(); err != nil {
_ = r.file.Close()
return nil, err
}
if _, err = r.file.Seek(0, io.SeekEnd); err != nil {
_ = r.file.Close()
return nil, err
if r.file != nil { // the file may not exist yet, poll() waits for it to appear
if r.info, err = r.file.Stat(); err != nil {
_ = r.file.Close()
return nil, err
}
if _, err = r.file.Seek(0, io.SeekEnd); err != nil {
_ = r.file.Close()
return nil, err
}
r.reader = bufio.NewReader(r.file)
}
r.reader = bufio.NewReader(r.file)

go func() {
const maxPrefixLen = 64 * 1024 // 64KB cap on partial line buffer
Expand All @@ -59,6 +61,10 @@ func NewTailReader(fileName string, ch chan<- logparser.LogEntry) (*TailReader,
r.stopped <- struct{}{}
return
default:
if r.reader == nil {
r.poll(ctx)
continue
}
line, err := r.reader.ReadString('\n')
if err != nil {
if len(prefix)+len(line) > maxPrefixLen {
Expand Down
9 changes: 9 additions & 0 deletions logs/tail_reader_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,10 @@ func TestTailReader(t *testing.T) {
assert.NoError(t, err)
defer os.Remove(f.Name())

// the file doesn't exist yet
err = os.Remove(f.Name())
assert.NoError(t, err)

tailPollInterval = time.Millisecond * 100
ch := make(chan logparser.LogEntry, 10)
tr, err := NewTailReader(f.Name(), ch)
Expand All @@ -35,6 +39,11 @@ func TestTailReader(t *testing.T) {
assert.Equal(t, expected, entry.Content)
}

// the file is created
wait()
f, err = os.Create(f.Name())
assert.NoError(t, err)

write("foo 1\n")
get("foo 1")

Expand Down
Loading