Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 21 additions & 7 deletions internal/generator/vector/output/syslog/syslog.go
Original file line number Diff line number Diff line change
Expand Up @@ -160,16 +160,30 @@ if err == null && value != null {
# KubernetesMinimal
# Adds namespace_name, pod_name, and container_name to the beginning of the message body (e.g. namespace_name=myproject, container_name=server, pod_name=pod-123, message={"foo":"bar"}).
# This may result in the message body being an invalid JSON structure.
# Fields are only added when present and non-empty (e.g. audit/journald logs have no Kubernetes metadata).

parts = []

if exists(.kubernetes.namespace_name) && ((to_string(.kubernetes.namespace_name) ?? "") != "") {
parts = push(parts, "namespace_name=" + to_string!(.kubernetes.namespace_name))
}

if exists(.kubernetes.container_name) && ((to_string(.kubernetes.container_name) ?? "") != "") {
parts = push(parts, "container_name=" + to_string!(.kubernetes.container_name))
}

if exists(.kubernetes.pod_name) && ((to_string(.kubernetes.pod_name) ?? "") != "") {
parts = push(parts, "pod_name=" + to_string!(.kubernetes.pod_name))
}

namespace = to_string(.kubernetes.namespace_name) ?? ""
container = to_string(.kubernetes.container_name) ?? ""
pod = to_string(.kubernetes.pod_name) ?? ""
msg_value = to_string(.message) ?? encode_json(.message)

.message = "namespace_name=" + namespace +
", container_name=" + container +
", pod_name=" + pod +
", message=" + msg_value`
if length(parts) > 0 {
parts = push(parts, "message=" + msg_value)
.message = join!(parts, ", ")
} else {
.message = msg_value
}`
)

func New(id string, o *adapters.Output, inputs []string, secrets observability.Secrets, op utils.Options) (_ string, sink types.Sink, tfs api.Transforms) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -81,16 +81,30 @@ for_each(excluded_fields) -> |_index, field| {
# KubernetesMinimal
# Adds namespace_name, pod_name, and container_name to the beginning of the message body (e.g. namespace_name=myproject, container_name=server, pod_name=pod-123, message={"foo":"bar"}).
# This may result in the message body being an invalid JSON structure.
# Fields are only added when present and non-empty (e.g. audit/journald logs have no Kubernetes metadata).

parts = []

if exists(.kubernetes.namespace_name) && ((to_string(.kubernetes.namespace_name) ?? "") != "") {
parts = push(parts, "namespace_name=" + to_string!(.kubernetes.namespace_name))
}

if exists(.kubernetes.container_name) && ((to_string(.kubernetes.container_name) ?? "") != "") {
parts = push(parts, "container_name=" + to_string!(.kubernetes.container_name))
}

if exists(.kubernetes.pod_name) && ((to_string(.kubernetes.pod_name) ?? "") != "") {
parts = push(parts, "pod_name=" + to_string!(.kubernetes.pod_name))
}

namespace = to_string(.kubernetes.namespace_name) ?? ""
container = to_string(.kubernetes.container_name) ?? ""
pod = to_string(.kubernetes.pod_name) ?? ""
msg_value = to_string(.message) ?? encode_json(.message)

.message = "namespace_name=" + namespace +
", container_name=" + container +
", pod_name=" + pod +
", message=" + msg_value
if length(parts) > 0 {
parts = push(parts, "message=" + msg_value)
.message = join!(parts, ", ")
} else {
.message = msg_value
}
'''

[sinks.example]
Expand Down
58 changes: 58 additions & 0 deletions test/functional/outputs/syslog/rfc3164_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -129,4 +129,62 @@ var _ = Describe("[Functional][Outputs][Syslog] RFC3164 tests", func() {
Entry("should enrich application logs with container source info", obs.InputTypeApplication, obs.EnrichmentTypeKubernetesMinimal),
Entry("should do nothing additional to the application logs", obs.InputTypeApplication, obs.EnrichmentTypeNone),
)

Context("KubernetesMinimal enrichment should not add empty prefixes to non-container logs", func() {

It("should not add empty kubernetes prefixes to audit logs", func() {
obstestruntime.NewClusterLogForwarderBuilder(framework.Forwarder).
FromInput(obs.InputTypeAudit).
ToSyslogOutput(obs.SyslogRFC3164, func(spec *obs.OutputSpec) {
spec.Syslog.Facility = "user"
spec.Syslog.Severity = "debug"
spec.Syslog.PayloadKey = "{.message}"
spec.Syslog.Enrichment = obs.EnrichmentTypeKubernetesMinimal
})
Expect(framework.Deploy()).To(BeNil())

Expect(framework.WriteK8sAuditLog(1)).To(BeNil())

logs, err := framework.ReadAuditLogsFrom(string(obs.OutputTypeSyslog))
Expect(err).To(BeNil(), "Expected no errors reading the logs")
Expect(logs).ToNot(BeEmpty(), "Expected the receiver to receive the message")

logLine := strings.TrimSpace(logs[0])
Expect(logLine).ToNot(ContainSubstring("namespace_name=,"), "Audit logs should not have empty namespace_name prefix")
Expect(logLine).ToNot(ContainSubstring("container_name=,"), "Audit logs should not have empty container_name prefix")
Expect(logLine).ToNot(ContainSubstring("pod_name=,"), "Audit logs should not have empty pod_name prefix")

collectorLogs, err := framework.ReadCollectorLogs()
Expect(err).To(BeNil())
Expect(collectorLogs).ToNot(ContainSubstring("VRL compilation warning"))
})

It("should not add empty kubernetes prefixes to infrastructure journal logs", func() {
obstestruntime.NewClusterLogForwarderBuilder(framework.Forwarder).
FromInput(obs.InputTypeInfrastructure).
ToSyslogOutput(obs.SyslogRFC3164, func(spec *obs.OutputSpec) {
spec.Syslog.Facility = "user"
spec.Syslog.Severity = "debug"
spec.Syslog.PayloadKey = "{.message}"
spec.Syslog.Enrichment = obs.EnrichmentTypeKubernetesMinimal
})
Expect(framework.Deploy()).To(BeNil())

logline := functional.NewJournalLog(3, "*", "*")
Expect(framework.WriteMessagesToInfraJournalLog(logline, 1)).To(BeNil())

logs, err := framework.ReadInfrastructureLogsFrom(string(obs.OutputTypeSyslog))
Expect(err).To(BeNil(), "Expected no errors reading the logs")
Expect(logs).ToNot(BeEmpty(), "Expected the receiver to receive the message")

logLine := strings.TrimSpace(logs[0])
Expect(logLine).ToNot(ContainSubstring("namespace_name=,"), "Journal logs should not have empty namespace_name prefix")
Expect(logLine).ToNot(ContainSubstring("container_name=,"), "Journal logs should not have empty container_name prefix")
Expect(logLine).ToNot(ContainSubstring("pod_name=,"), "Journal logs should not have empty pod_name prefix")

collectorLogs, err := framework.ReadCollectorLogs()
Expect(err).To(BeNil())
Expect(collectorLogs).ToNot(ContainSubstring("VRL compilation warning"))
})
})
})
58 changes: 58 additions & 0 deletions test/functional/outputs/syslog/rfc5424_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -242,6 +242,64 @@ var _ = Describe("[Functional][Outputs][Syslog] RFC5424 tests", func() {
Entry("capitalized: Informational", "Informational", "<14>1 "), // 8 + 6
)

Context("KubernetesMinimal enrichment should not add empty prefixes to non-container logs", func() {

It("should not add empty kubernetes prefixes to audit logs", func() {
obstestruntime.NewClusterLogForwarderBuilder(framework.Forwarder).
FromInput(obs.InputTypeAudit).
ToSyslogOutput(obs.SyslogRFC5424, func(spec *obs.OutputSpec) {
spec.Syslog.Facility = "user"
spec.Syslog.Severity = "debug"
spec.Syslog.PayloadKey = "{.message}"
spec.Syslog.Enrichment = obs.EnrichmentTypeKubernetesMinimal
})
Expect(framework.Deploy()).To(BeNil())

Expect(framework.WriteK8sAuditLog(1)).To(BeNil())

logs, err := framework.ReadAuditLogsFrom(string(obs.OutputTypeSyslog))
Expect(err).To(BeNil(), "Expected no errors reading the logs")
Expect(logs).ToNot(BeEmpty(), "Expected the receiver to receive the message")

logLine := strings.TrimSpace(logs[0])
Expect(logLine).ToNot(ContainSubstring("namespace_name=,"), "Audit logs should not have empty namespace_name prefix")
Expect(logLine).ToNot(ContainSubstring("container_name=,"), "Audit logs should not have empty container_name prefix")
Expect(logLine).ToNot(ContainSubstring("pod_name=,"), "Audit logs should not have empty pod_name prefix")

collectorLogs, err := framework.ReadCollectorLogs()
Expect(err).To(BeNil())
Expect(collectorLogs).ToNot(ContainSubstring("VRL compilation warning"))
})

It("should not add empty kubernetes prefixes to infrastructure journal logs", func() {
obstestruntime.NewClusterLogForwarderBuilder(framework.Forwarder).
FromInput(obs.InputTypeInfrastructure).
ToSyslogOutput(obs.SyslogRFC5424, func(spec *obs.OutputSpec) {
spec.Syslog.Facility = "user"
spec.Syslog.Severity = "debug"
spec.Syslog.PayloadKey = "{.message}"
spec.Syslog.Enrichment = obs.EnrichmentTypeKubernetesMinimal
})
Expect(framework.Deploy()).To(BeNil())

logline := functional.NewJournalLog(3, "*", "*")
Expect(framework.WriteMessagesToInfraJournalLog(logline, 1)).To(BeNil())

logs, err := framework.ReadInfrastructureLogsFrom(string(obs.OutputTypeSyslog))
Expect(err).To(BeNil(), "Expected no errors reading the logs")
Expect(logs).ToNot(BeEmpty(), "Expected the receiver to receive the message")

logLine := strings.TrimSpace(logs[0])
Expect(logLine).ToNot(ContainSubstring("namespace_name=,"), "Journal logs should not have empty namespace_name prefix")
Expect(logLine).ToNot(ContainSubstring("container_name=,"), "Journal logs should not have empty container_name prefix")
Expect(logLine).ToNot(ContainSubstring("pod_name=,"), "Journal logs should not have empty pod_name prefix")

collectorLogs, err := framework.ReadCollectorLogs()
Expect(err).To(BeNil())
Expect(collectorLogs).ToNot(ContainSubstring("VRL compilation warning"))
})
})

It("should be able to send a large payload", func() {
obstestruntime.NewClusterLogForwarderBuilder(framework.Forwarder).
FromInput(obs.InputTypeApplication).
Expand Down
Loading