diff --git a/src/pkg/binding/poller.go b/src/pkg/binding/poller.go index b9f13ec13..19d8672fc 100644 --- a/src/pkg/binding/poller.go +++ b/src/pkg/binding/poller.go @@ -9,9 +9,11 @@ import ( "net" "net/http" "net/url" + "strings" "time" metrics "code.cloudfoundry.org/go-metric-registry" + "code.cloudfoundry.org/loggregator-agent-release/src/pkg/egress/syslog" v2 "code.cloudfoundry.org/loggregator-agent-release/src/pkg/ingress/v2" "code.cloudfoundry.org/loggregator-agent-release/src/pkg/simplecache" ) @@ -236,6 +238,17 @@ func (bc *bindingChecker) checkBindings(bindings []Binding) []Binding { continue } + if invalidLogFilter(u) { + bc.rejectBinding(b.Credentials, fmt.Sprintf("include-log-source-types and exclude-log-source-types cannot be used at the same time in syslog drain url %s", anonymousUrl.String()), true) + continue + } + + sourceTypes := getUnknownSourceTypes(u.Query()) + if sourceTypes != nil { + bc.rejectBinding(b.Credentials, fmt.Sprintf("Unknown source types '%s' in source type filter in syslog drain url %s", strings.Join(sourceTypes, ", "), anonymousUrl.String()), true) + continue + } + _, exists := bc.failedHostsCache.Get(u.Host) if exists { bc.rejectBinding(b.Credentials, fmt.Sprintf("Skipped resolve ip address for syslog drain with url %s due to prior failure", anonymousUrl.String()), false) @@ -308,6 +321,34 @@ func invalidScheme(scheme string) bool { return true } +// invalidLogFilter checks if both include-log-source-types and exclude-log-source-types +func invalidLogFilter(u *url.URL) bool { + includeSourceTypes := u.Query().Get("include-log-source-types") + excludeSourceTypes := u.Query().Get("exclude-log-source-types") + if excludeSourceTypes != "" && includeSourceTypes != "" { + return true + } + return false +} + +// assumes only one of include-log-source-types or exclude-log-source-types is set +func getUnknownSourceTypes(u url.Values) []string { + var sourceTypeList string + includeSourceTypes := u.Get("include-log-source-types") + excludeSourceTypes := u.Get("exclude-log-source-types") + + if includeSourceTypes != "" { + sourceTypeList = includeSourceTypes + } else if excludeSourceTypes != "" { + sourceTypeList = excludeSourceTypes + } else { + return nil + } + + _, unknownTypes := syslog.ParseSourceTypeList(sourceTypeList) + return unknownTypes +} + func CalculateBindingCount(bindings []Binding) int { apps := make(map[string]bool) for _, b := range bindings { diff --git a/src/pkg/binding/poller_test.go b/src/pkg/binding/poller_test.go index d60d1c654..fff7ce4cb 100644 --- a/src/pkg/binding/poller_test.go +++ b/src/pkg/binding/poller_test.go @@ -563,6 +563,125 @@ var _ = Describe("Poller", func() { Expect(bndChecker.invalidDrains).To(Equal(float64(0))) Expect(bndChecker.blacklistedDrains).To(Equal(float64(0))) }) + + Context("when both include-log-source-types and exclude-log-source-types are specified", func() { + It("ignores the drain and counts as invalid", func() { + bindings := []Binding{ + { + Url: "https://test.org/drain?include-log-source-types=app&exclude-log-source-types=rtr", + Credentials: []Credentials{ + { + Apps: []App{{Hostname: "app-hostname0", AppID: "app-id-0"}}, + }, + }, + }, + } + + filteredBindings := bndChecker.checkBindings(bindings) + + Expect(filteredBindings).To(BeEmpty()) + Expect(logClient.Message()).To(ContainElement(MatchRegexp("include-log-source-types and exclude-log-source-types cannot be used at the same time"))) + Expect(bndChecker.invalidDrains).To(BeNumerically("==", 1)) + Expect(bndChecker.blacklistedDrains).To(BeNumerically("==", 0)) + }) + + It("doesn't log the conflicting filters warning when warn is false", func() { + bndChecker.warn = false + bindings := []Binding{ + { + Url: "https://test.org/drain?include-log-source-types=app&exclude-log-source-types=rtr", + Credentials: []Credentials{ + { + Apps: []App{{Hostname: "app-hostname0", AppID: "app-id-0"}}, + }, + }, + }, + } + bndChecker.checkBindings(bindings) + + for _, msg := range logClient.Message() { + Expect(msg).ToNot(MatchRegexp("include-log-source-types and exclude-log-source-types cannot be used at the same time")) + } + }) + }) + + Context("when unknown source types are provided", func() { + It("logs a warning and ignores the drain in include mode", func() { + bindings := []Binding{ + { + Url: "https://test.org/drain?include-log-source-types=app,unknown,invalid,rtr", + Credentials: []Credentials{ + { + Apps: []App{{Hostname: "app-hostname0", AppID: "app-id-0"}}, + }, + }, + }, + } + filteredBindings := bndChecker.checkBindings(bindings) + + Expect(filteredBindings).To(BeEmpty()) + Expect(logClient.Message()).To(ContainElement(MatchRegexp("Unknown source types"))) + Expect(logClient.Message()).To(ContainElement(MatchRegexp("unknown"))) + Expect(logClient.Message()).To(ContainElement(MatchRegexp("invalid"))) + Expect(bndChecker.invalidDrains).To(BeNumerically("==", 1)) + }) + + It("logs a warning and ignores the drain in exclude mode", func() { + bindings := []Binding{ + { + Url: "https://test.org/drain?exclude-log-source-types=rtr,unknown", + Credentials: []Credentials{ + { + Apps: []App{{Hostname: "app-hostname0", AppID: "app-id-0"}}, + }, + }, + }, + } + filteredBindings := bndChecker.checkBindings(bindings) + + Expect(filteredBindings).To(BeEmpty()) + Expect(logClient.Message()).To(ContainElement(MatchRegexp("Unknown source types"))) + Expect(logClient.Message()).To(ContainElement(MatchRegexp("unknown"))) + Expect(bndChecker.invalidDrains).To(BeNumerically("==", 1)) + }) + + It("logs a warning and ignores the drain when source types have spaces", func() { + bindings := []Binding{ + { + Url: "https://test.org/drain?include-log-source-types=app, rtr", + Credentials: []Credentials{ + { + Apps: []App{{Hostname: "app-hostname0", AppID: "app-id-0"}}, + }, + }, + }, + } + filteredBindings := bndChecker.checkBindings(bindings) + + Expect(filteredBindings).To(BeEmpty()) + Expect(logClient.Message()).To(ContainElement(MatchRegexp("Unknown source types"))) + Expect(bndChecker.invalidDrains).To(BeNumerically("==", 1)) + }) + + It("doesn't log the warning when warn is false", func() { + bndChecker.warn = false + bindings := []Binding{ + { + Url: "https://test.org/drain?include-log-source-types=app,unknown,rtr", + Credentials: []Credentials{ + { + Apps: []App{{Hostname: "app-hostname0", AppID: "app-id-0"}}, + }, + }, + }, + } + bndChecker.checkBindings(bindings) + + for _, msg := range logClient.Message() { + Expect(msg).ToNot(MatchRegexp("Unknown source types")) + } + }) + }) }) }) diff --git a/src/pkg/egress/syslog/filtering_drain_writer.go b/src/pkg/egress/syslog/filtering_drain_writer.go index db6d56d90..236d3a3d7 100644 --- a/src/pkg/egress/syslog/filtering_drain_writer.go +++ b/src/pkg/egress/syslog/filtering_drain_writer.go @@ -35,22 +35,18 @@ func NewFilteringDrainWriter(binding Binding, writer egress.Writer) (*FilteringD } func (w *FilteringDrainWriter) Write(env *loggregator_v2.Envelope) error { - if w.binding.DrainData == ALL { - return w.writer.Write(env) - } - if env.GetTimer() != nil { - if w.binding.DrainData == TRACES { + if w.binding.DrainData == TRACES || w.binding.DrainData == ALL { return w.writer.Write(env) } } if env.GetEvent() != nil { - if w.binding.DrainData == LOGS { + if sendsEvents(w.binding.DrainData, w.binding.LogFilter, env.GetTags()["source_type"]) { return w.writer.Write(env) } } if env.GetLog() != nil { - if sendsLogs(w.binding.DrainData) { + if sendsLogs(w.binding.DrainData, w.binding.LogFilter, env.GetTags()["source_type"]) { return w.writer.Write(env) } } @@ -63,26 +59,25 @@ func (w *FilteringDrainWriter) Write(env *loggregator_v2.Envelope) error { return nil } -func sendsLogs(drainData DrainData) bool { - switch drainData { - case LOGS: - return true - case LOGS_AND_METRICS: - return true - case LOGS_NO_EVENTS: - return true - default: +func sendsEvents(drainData DrainData, logFilter *LogFilter, sourceTypeTag string) bool { + if drainData != LOGS && drainData != ALL { + return false + } + + return logFilter.ShouldInclude(sourceTypeTag) +} + +func sendsLogs(drainData DrainData, logFilter *LogFilter, sourceTypeTag string) bool { + if drainData == TRACES || drainData == METRICS { return false } + + return logFilter.ShouldInclude(sourceTypeTag) } func sendsMetrics(drainData DrainData) bool { - switch drainData { - case LOGS_AND_METRICS: - return true - case METRICS: - return true - default: + if drainData != LOGS_AND_METRICS && drainData != METRICS && drainData != ALL { return false } + return true } diff --git a/src/pkg/egress/syslog/filtering_drain_writer_test.go b/src/pkg/egress/syslog/filtering_drain_writer_test.go index 020e0b157..fc9d89011 100644 --- a/src/pkg/egress/syslog/filtering_drain_writer_test.go +++ b/src/pkg/egress/syslog/filtering_drain_writer_test.go @@ -41,7 +41,7 @@ var _ = Describe("Filtering Drain Writer", func() { Entry("traces", syslog.TRACES, false, false, false, true), Entry("all", syslog.ALL, true, true, true, true), Entry("log without events", syslog.LOGS_NO_EVENTS, true, false, false, false), - Entry("metrics and logs", syslog.LOGS_AND_METRICS, true, true, false, false), + Entry("metrics and logs without events", syslog.LOGS_AND_METRICS, true, true, false, false), ) It("errors on invalid binding type", func() { @@ -53,6 +53,299 @@ var _ = Describe("Filtering Drain Writer", func() { _, err := syslog.NewFilteringDrainWriter(binding, &fakeWriter{}) Expect(err).To(HaveOccurred()) }) + + Context("when source_type tag is missing", func() { + var envelope *loggregator_v2.Envelope + + BeforeEach(func() { + envelope = &loggregator_v2.Envelope{ + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{ + Payload: []byte("test log"), + }, + }, + Tags: map[string]string{ + // source_type tag is intentionally missing + }, + } + }) + + It("filters out the log with missing source_type when include filter is configured", func() { + binding := syslog.Binding{ + DrainData: syslog.LOGS, + LogFilter: syslog.NewLogFilter(syslog.LogSourceTypeSet{syslog.LOG_SOURCE_APP: struct{}{}}, syslog.LogFilterModeInclude), + } + fakeWriter := &fakeWriter{} + drainWriter, err := syslog.NewFilteringDrainWriter(binding, fakeWriter) + Expect(err).NotTo(HaveOccurred()) + + err = drainWriter.Write(envelope) + Expect(err).NotTo(HaveOccurred()) + Expect(fakeWriter.received).To(Equal(0)) + + appEnvelope := &loggregator_v2.Envelope{ + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("app log")}, + }, + Tags: map[string]string{"source_type": "APP/PROC/WEB/0"}, + } + err = drainWriter.Write(appEnvelope) + Expect(err).NotTo(HaveOccurred()) + Expect(fakeWriter.received).To(Equal(1)) + }) + + It("sends the log with missing source_type when exclude filter is configured", func() { + binding := syslog.Binding{ + DrainData: syslog.LOGS, + LogFilter: syslog.NewLogFilter(syslog.LogSourceTypeSet{syslog.LOG_SOURCE_RTR: struct{}{}}, syslog.LogFilterModeExclude), + } + fakeWriter := &fakeWriter{} + drainWriter, err := syslog.NewFilteringDrainWriter(binding, fakeWriter) + Expect(err).NotTo(HaveOccurred()) + + err = drainWriter.Write(envelope) + Expect(err).NotTo(HaveOccurred()) + Expect(fakeWriter.received).To(Equal(1)) + + rtrEnvelope := &loggregator_v2.Envelope{ + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("rtr log")}, + }, + Tags: map[string]string{"source_type": "RTR/1"}, + } + err = drainWriter.Write(rtrEnvelope) + Expect(err).NotTo(HaveOccurred()) + Expect(fakeWriter.received).To(Equal(1)) + }) + }) + + Context("when source_type tag is present", func() { + It("applies the include filter even when drain-data is ALL", func() { + binding := syslog.Binding{ + DrainData: syslog.ALL, + LogFilter: syslog.NewLogFilter(syslog.LogSourceTypeSet{syslog.LOG_SOURCE_APP: struct{}{}}, syslog.LogFilterModeInclude), + } + fakeWriter := &fakeWriter{} + drainWriter, err := syslog.NewFilteringDrainWriter(binding, fakeWriter) + Expect(err).NotTo(HaveOccurred()) + + envelopes := []*loggregator_v2.Envelope{ + { + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("app log")}, + }, + Tags: map[string]string{"source_type": "APP/PROC/WEB/0"}, + }, + { + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("rtr log")}, + }, + Tags: map[string]string{"source_type": "RTR/1"}, + }, + } + + for _, envelope := range envelopes { + err = drainWriter.Write(envelope) + Expect(err).NotTo(HaveOccurred()) + } + + // Only APP log should be sent; RTR must be dropped even with drain-data=all + Expect(fakeWriter.received).To(Equal(1)) + }) + + It("applies the exclude filter even when drain-data is ALL", func() { + binding := syslog.Binding{ + DrainData: syslog.ALL, + LogFilter: syslog.NewLogFilter(syslog.LogSourceTypeSet{syslog.LOG_SOURCE_RTR: struct{}{}}, syslog.LogFilterModeExclude), + } + fakeWriter := &fakeWriter{} + drainWriter, err := syslog.NewFilteringDrainWriter(binding, fakeWriter) + Expect(err).NotTo(HaveOccurred()) + + envelopes := []*loggregator_v2.Envelope{ + { + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("app log")}, + }, + Tags: map[string]string{"source_type": "APP/PROC/WEB/0"}, + }, + { + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("rtr log")}, + }, + Tags: map[string]string{"source_type": "RTR/1"}, + }, + } + + for _, envelope := range envelopes { + err = drainWriter.Write(envelope) + Expect(err).NotTo(HaveOccurred()) + } + + // APP log should be sent, RTR should be filtered out even with drain-data=all + Expect(fakeWriter.received).To(Equal(1)) + }) + + It("still forwards metrics and traces when drain-data is ALL with a log filter", func() { + binding := syslog.Binding{ + DrainData: syslog.ALL, + LogFilter: syslog.NewLogFilter(syslog.LogSourceTypeSet{syslog.LOG_SOURCE_APP: struct{}{}}, syslog.LogFilterModeInclude), + } + fakeWriter := &fakeWriter{} + drainWriter, err := syslog.NewFilteringDrainWriter(binding, fakeWriter) + Expect(err).NotTo(HaveOccurred()) + + envelopes := []*loggregator_v2.Envelope{ + {Message: &loggregator_v2.Envelope_Counter{Counter: &loggregator_v2.Counter{}}}, + {Message: &loggregator_v2.Envelope_Gauge{Gauge: &loggregator_v2.Gauge{}}}, + {Message: &loggregator_v2.Envelope_Timer{Timer: &loggregator_v2.Timer{}}}, + } + + for _, envelope := range envelopes { + err = drainWriter.Write(envelope) + Expect(err).NotTo(HaveOccurred()) + } + + Expect(fakeWriter.received).To(Equal(3)) + }) + + It("filters event envelopes by source type just like logs", func() { + binding := syslog.Binding{ + DrainData: syslog.ALL, + LogFilter: syslog.NewLogFilter(syslog.LogSourceTypeSet{syslog.LOG_SOURCE_APP: struct{}{}}, syslog.LogFilterModeInclude), + } + fakeWriter := &fakeWriter{} + drainWriter, err := syslog.NewFilteringDrainWriter(binding, fakeWriter) + Expect(err).NotTo(HaveOccurred()) + + appEvent := &loggregator_v2.Envelope{ + Message: &loggregator_v2.Envelope_Event{ + Event: &loggregator_v2.Event{Title: "app crash", Body: "exited"}, + }, + Tags: map[string]string{"source_type": "APP/PROC/WEB/0"}, + } + rtrEvent := &loggregator_v2.Envelope{ + Message: &loggregator_v2.Envelope_Event{ + Event: &loggregator_v2.Event{Title: "route", Body: "request"}, + }, + Tags: map[string]string{"source_type": "RTR/1"}, + } + + for _, envelope := range []*loggregator_v2.Envelope{appEvent, rtrEvent} { + err = drainWriter.Write(envelope) + Expect(err).NotTo(HaveOccurred()) + } + + // The include filter applies to events too: only the APP event passes, + // the RTR event is dropped. + Expect(fakeWriter.received).To(Equal(1)) + }) + + It("filters logs based on include filter - includes only APP logs", func() { + binding := syslog.Binding{ + DrainData: syslog.LOGS, + LogFilter: syslog.NewLogFilter(syslog.LogSourceTypeSet{syslog.LOG_SOURCE_APP: struct{}{}}, syslog.LogFilterModeInclude), + } + fakeWriter := &fakeWriter{} + drainWriter, err := syslog.NewFilteringDrainWriter(binding, fakeWriter) + Expect(err).NotTo(HaveOccurred()) + + envelopes := []*loggregator_v2.Envelope{ + { + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("app log")}, + }, + Tags: map[string]string{"source_type": "APP/PROC/WEB/0"}, + }, + { + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("rtr log")}, + }, + Tags: map[string]string{"source_type": "RTR/1"}, + }, + { + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("stg log")}, + }, + Tags: map[string]string{"source_type": "STG/0"}, + }, + } + + for _, envelope := range envelopes { + err = drainWriter.Write(envelope) + Expect(err).NotTo(HaveOccurred()) + } + + // Only APP log should be sent + Expect(fakeWriter.received).To(Equal(1)) + }) + + It("filters logs based on exclude filter - excludes RTR logs", func() { + binding := syslog.Binding{ + DrainData: syslog.LOGS, + LogFilter: syslog.NewLogFilter(syslog.LogSourceTypeSet{syslog.LOG_SOURCE_RTR: struct{}{}}, syslog.LogFilterModeExclude), + } + fakeWriter := &fakeWriter{} + drainWriter, err := syslog.NewFilteringDrainWriter(binding, fakeWriter) + Expect(err).NotTo(HaveOccurred()) + + envelopes := []*loggregator_v2.Envelope{ + { + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("app log")}, + }, + Tags: map[string]string{"source_type": "APP/PROC/WEB/0"}, + }, + { + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("rtr log")}, + }, + Tags: map[string]string{"source_type": "RTR/1"}, + }, + { + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{Payload: []byte("stg log")}, + }, + Tags: map[string]string{"source_type": "STG/0"}, + }, + } + + for _, envelope := range envelopes { + err = drainWriter.Write(envelope) + Expect(err).NotTo(HaveOccurred()) + } + + // APP and STG logs should be sent, RTR should be filtered out + Expect(fakeWriter.received).To(Equal(2)) + }) + + It("sends logs with unknown source_type prefix when filter is set", func() { + binding := syslog.Binding{ + DrainData: syslog.LOGS, + LogFilter: syslog.NewLogFilter(syslog.LogSourceTypeSet{syslog.LOG_SOURCE_APP: struct{}{}}, syslog.LogFilterModeExclude), + } + fakeWriter := &fakeWriter{} + drainWriter, err := syslog.NewFilteringDrainWriter(binding, fakeWriter) + Expect(err).NotTo(HaveOccurred()) + + envelope := &loggregator_v2.Envelope{ + Message: &loggregator_v2.Envelope_Log{ + Log: &loggregator_v2.Log{ + Payload: []byte("test log"), + }, + }, + Tags: map[string]string{ + "source_type": "UNKNOWN/some/path", + }, + } + + err = drainWriter.Write(envelope) + + // Should send the log because unknown types default to being included for exclude filter + Expect(err).NotTo(HaveOccurred()) + Expect(fakeWriter.received).To(Equal(1)) + }) + }) }) type fakeWriter struct { diff --git a/src/pkg/egress/syslog/log_source.go b/src/pkg/egress/syslog/log_source.go new file mode 100644 index 000000000..5e68cf549 --- /dev/null +++ b/src/pkg/egress/syslog/log_source.go @@ -0,0 +1,133 @@ +package syslog + +import "strings" + +// LogSourceType defines the source types used within Cloud Foundry +// Their order in the code is as documented in https://docs.cloudfoundry.org/devguide/deploy-apps/streaming-logs.html#format +type LogSourceType string + +const ( + LOG_SOURCE_API LogSourceType = "API" + LOG_SOURCE_STG LogSourceType = "STG" + LOG_SOURCE_RTR LogSourceType = "RTR" + LOG_SOURCE_LGR LogSourceType = "LGR" + LOG_SOURCE_APP LogSourceType = "APP" + LOG_SOURCE_SSH LogSourceType = "SSH" + LOG_SOURCE_CELL LogSourceType = "CELL" + LOG_SOURCE_PROXY LogSourceType = "PROXY" + LOG_SOURCE_HEALTH LogSourceType = "HEALTH" + LOG_SOURCE_SYS LogSourceType = "SYS" + LOG_SOURCE_STATS LogSourceType = "STATS" +) + +// validSourceTypes contains SourceType prefixes for efficient lookup +var validSourceTypes = map[LogSourceType]struct{}{ + LOG_SOURCE_API: {}, + LOG_SOURCE_STG: {}, + LOG_SOURCE_RTR: {}, + LOG_SOURCE_LGR: {}, + LOG_SOURCE_APP: {}, + LOG_SOURCE_SSH: {}, + LOG_SOURCE_CELL: {}, +} + +// IsValid checks if the provided SourceType is valid +func (lt LogSourceType) IsValid() bool { + _, ok := validSourceTypes[lt] + return ok +} + +// ParseSourceType parses a string into a SourceType value +func ParseSourceType(s string) (LogSourceType, bool) { + lt := LogSourceType(strings.ToUpper(s)) + return lt, lt.IsValid() +} + +// ParseSourceTypeList parses a comma-separated list of source types and returns +// the valid types as a set and any unknown types as a slice. +func ParseSourceTypeList(sourceTypeList string) (LogSourceTypeSet, []string) { + sourceTypes := strings.Split(sourceTypeList, ",") + set := make(LogSourceTypeSet, len(sourceTypes)) + var unknownTypes []string + + for _, sourceType := range sourceTypes { + t, ok := ParseSourceType(sourceType) + if !ok { + unknownTypes = append(unknownTypes, sourceType) + continue + } + set.Add(t) + } + + return set, unknownTypes +} + +// ExtractPrefix extracts the prefix from a source_type tag (e.g., "APP/PROC/WEB/0" -> "APP") +func ExtractPrefix(sourceTypeTag string) string { + if idx := strings.IndexByte(sourceTypeTag, '/'); idx != -1 { + return sourceTypeTag[:idx] + } + return sourceTypeTag +} + +// LogSourceTypeSet is a set of SourceTypes for efficient membership checking +type LogSourceTypeSet map[LogSourceType]struct{} + +// Add adds a SourceType to the set +func (s LogSourceTypeSet) Add(lt LogSourceType) { + s[lt] = struct{}{} +} + +// Contains checks if the set contains a SourceType +func (s LogSourceTypeSet) Contains(lt LogSourceType) bool { + _, exists := s[lt] + return exists +} + +// LogFilterMode determines how the log filter should be applied +type LogFilterMode int + +const ( + // LogFilterModeInclude only includes logs matching the specified types (strict) + LogFilterModeInclude LogFilterMode = iota + // LogFilterModeExclude excludes logs matching the specified types (permissive) + LogFilterModeExclude +) + +// LogFilter encapsulates source type filtering configuration +type LogFilter struct { + Types LogSourceTypeSet + Mode LogFilterMode +} + +// NewLogFilter creates a new LogFilter with the given types and mode +func NewLogFilter(types LogSourceTypeSet, mode LogFilterMode) *LogFilter { + return &LogFilter{ + Types: types, + Mode: mode, + } +} + +// ShouldInclude determines if a log with the given sourceTypeTag should be forwarded +// Include mode omits missing/unknown source types, exclude mode forwards them +func (f *LogFilter) ShouldInclude(sourceTypeTag string) bool { + if f == nil { + return true + } + + if sourceTypeTag == "" { + return f.Mode == LogFilterModeExclude + } + + prefix := ExtractPrefix(sourceTypeTag) + sourceType := LogSourceType(prefix) + if !sourceType.IsValid() { + return f.Mode == LogFilterModeExclude + } + + inSet := f.Types.Contains(sourceType) + if f.Mode == LogFilterModeInclude { + return inSet + } + return !inSet +} diff --git a/src/pkg/egress/syslog/syslog_connector.go b/src/pkg/egress/syslog/syslog_connector.go index 66134f24c..9f948bc00 100644 --- a/src/pkg/egress/syslog/syslog_connector.go +++ b/src/pkg/egress/syslog/syslog_connector.go @@ -19,6 +19,7 @@ type Binding struct { DrainData DrainData `json:"type,omitempty"` OmitMetadata bool InternalTls bool + LogFilter *LogFilter } type Drain struct { diff --git a/src/pkg/ingress/bindings/binding_config.go b/src/pkg/ingress/bindings/binding_config.go index 9a1931916..5a5b32328 100644 --- a/src/pkg/ingress/bindings/binding_config.go +++ b/src/pkg/ingress/bindings/binding_config.go @@ -35,6 +35,7 @@ func (d *DrainParamParser) FetchBindings() ([]syslog.Binding, error) { b.OmitMetadata = getOmitMetadata(urlParsed, d.defaultDrainMetadata) b.InternalTls = getInternalTLS(urlParsed) b.DrainData = getBindingType(urlParsed) + b.LogFilter = d.getLogFilter(urlParsed) processed = append(processed, b) } @@ -57,6 +58,7 @@ func getOmitMetadata(url *url.URL, defaultDrainMetadata bool) bool { } func getBindingType(u *url.URL) syslog.DrainData { + // Legacy drain-type query param does not forward events drainData := syslog.LOGS switch u.Query().Get("drain-type") { case "logs": @@ -84,6 +86,29 @@ func getBindingType(u *url.URL) syslog.DrainData { return drainData } +func (d *DrainParamParser) getLogFilter(u *url.URL) *syslog.LogFilter { + includeSourceTypes := u.Query().Get("include-log-source-types") + excludeSourceTypes := u.Query().Get("exclude-log-source-types") + + if excludeSourceTypes != "" { + return d.newLogFilter(excludeSourceTypes, syslog.LogFilterModeExclude) + } else if includeSourceTypes != "" { + return d.newLogFilter(includeSourceTypes, syslog.LogFilterModeInclude) + } + return nil +} + +// newLogFilter parses a URL query parameter into a LogFilter. +// sourceTypeList is assumed to be a comma-separated list of valid source types. +func (d *DrainParamParser) newLogFilter(sourceTypeList string, mode syslog.LogFilterMode) *syslog.LogFilter { + if sourceTypeList == "" { + return nil + } + + set, _ := syslog.ParseSourceTypeList(sourceTypeList) + return syslog.NewLogFilter(set, mode) +} + func getRemoveMetadataQuery(u *url.URL) string { q := u.Query().Get("disable-metadata") if q == "" { diff --git a/src/pkg/ingress/bindings/binding_config_test.go b/src/pkg/ingress/bindings/binding_config_test.go index ef9146e74..57b0b87d8 100644 --- a/src/pkg/ingress/bindings/binding_config_test.go +++ b/src/pkg/ingress/bindings/binding_config_test.go @@ -15,9 +15,9 @@ var _ = Describe("Drain Param Config", func() { {Drain: syslog.Drain{Url: "https://test.org/drain"}}, } f := newStubFetcher(bs, nil) - wf := bindings.NewDrainParamParser(f, true) + dp := bindings.NewDrainParamParser(f, true) - configedBindings, _ := wf.FetchBindings() + configedBindings, _ := dp.FetchBindings() Expect(configedBindings[0].OmitMetadata).To(BeFalse()) }) @@ -27,9 +27,9 @@ var _ = Describe("Drain Param Config", func() { {Drain: syslog.Drain{Url: "https://test.org/drain?omit-metadata=true"}}, } f := newStubFetcher(bs, nil) - wf := bindings.NewDrainParamParser(f, true) + dp := bindings.NewDrainParamParser(f, true) - configedBindings, _ := wf.FetchBindings() + configedBindings, _ := dp.FetchBindings() Expect(configedBindings[0].OmitMetadata).To(BeTrue()) Expect(configedBindings[1].OmitMetadata).To(BeTrue()) }) @@ -39,9 +39,9 @@ var _ = Describe("Drain Param Config", func() { {Drain: syslog.Drain{Url: "https://test.org/drain"}}, } f := newStubFetcher(bs, nil) - wf := bindings.NewDrainParamParser(f, false) + dp := bindings.NewDrainParamParser(f, false) - configedBindings, _ := wf.FetchBindings() + configedBindings, _ := dp.FetchBindings() Expect(configedBindings[0].OmitMetadata).To(BeTrue()) }) @@ -51,9 +51,9 @@ var _ = Describe("Drain Param Config", func() { {Drain: syslog.Drain{Url: "https://test.org/drain?omit-metadata=false"}}, } f := newStubFetcher(bs, nil) - wf := bindings.NewDrainParamParser(f, false) + dp := bindings.NewDrainParamParser(f, false) - configedBindings, _ := wf.FetchBindings() + configedBindings, _ := dp.FetchBindings() Expect(configedBindings[0].OmitMetadata).To(BeFalse()) Expect(configedBindings[1].OmitMetadata).To(BeFalse()) }) @@ -63,48 +63,158 @@ var _ = Describe("Drain Param Config", func() { {Drain: syslog.Drain{Url: "https://test.org/drain?ssl-strict-internal=true"}}, } f := newStubFetcher(bs, nil) - wf := bindings.NewDrainParamParser(f, true) + dp := bindings.NewDrainParamParser(f, true) - configedBindings, _ := wf.FetchBindings() + configedBindings, _ := dp.FetchBindings() Expect(configedBindings[0].InternalTls).To(BeTrue()) }) It("sets drain data appropriately'", func() { - bs := []syslog.Binding{ - {Drain: syslog.Drain{Url: "https://test.org/drain"}}, - {Drain: syslog.Drain{Url: "https://test.org/drain?drain-data=logs"}}, - {Drain: syslog.Drain{Url: "https://test.org/drain?drain-data=metrics"}}, - {Drain: syslog.Drain{Url: "https://test.org/drain?drain-data=traces"}}, - {Drain: syslog.Drain{Url: "https://test.org/drain?drain-data=all"}}, + testCases := []struct { + name string + url string + expected syslog.DrainData + }{ + { + name: "no drain-data parameter defaults to logs", + url: "https://test.org/drain", + expected: syslog.LOGS, + }, + { + name: "drain-data=logs", + url: "https://test.org/drain?drain-data=logs", + expected: syslog.LOGS, + }, + { + name: "drain-data=metrics", + url: "https://test.org/drain?drain-data=metrics", + expected: syslog.METRICS, + }, + { + name: "drain-data=traces", + url: "https://test.org/drain?drain-data=traces", + expected: syslog.TRACES, + }, + { + name: "drain-data=all", + url: "https://test.org/drain?drain-data=all", + expected: syslog.ALL, + }, + } + + for _, tc := range testCases { + By(tc.name) + bs := []syslog.Binding{ + {Drain: syslog.Drain{Url: tc.url}}, + } + f := newStubFetcher(bs, nil) + dp := bindings.NewDrainParamParser(f, true) + + configedBindings, _ := dp.FetchBindings() + Expect(configedBindings[0].DrainData).To(Equal(tc.expected)) + } + }) + + It("sets drain filter appropriately", func() { + testCases := []struct { + name string + url string + expected *syslog.LogFilter + }{ + { + name: "empty drain URL defaults to no filtering", + url: "https://test.org/drain", + expected: nil, + }, + { + name: "include-log-source-types= defaults to no filtering", + url: "https://test.org/drain?include-log-source-types=", + expected: nil, + }, + { + name: "exclude-log-source-types= defaults to no filtering", + url: "https://test.org/drain?exclude-log-source-types=", + expected: nil, + }, + { + name: "include-log-source-types=app", + url: "https://test.org/drain?include-log-source-types=app", + expected: NewLogFilter(syslog.LogFilterModeInclude, syslog.LOG_SOURCE_APP), + }, + { + name: "include-log-source-types=app,stg,cell", + url: "https://test.org/drain?include-log-source-types=app,stg,cell", + expected: NewLogFilter(syslog.LogFilterModeInclude, syslog.LOG_SOURCE_APP, syslog.LOG_SOURCE_STG, syslog.LOG_SOURCE_CELL), + }, + { + name: "exclude-log-source-types=rtr,cell,stg", + url: "https://test.org/drain?exclude-log-source-types=rtr,cell,stg", + expected: NewLogFilter(syslog.LogFilterModeExclude, syslog.LOG_SOURCE_RTR, syslog.LOG_SOURCE_CELL, syslog.LOG_SOURCE_STG), + }, + { + name: "exclude-log-source-types=rtr", + url: "https://test.org/drain?exclude-log-source-types=rtr", + expected: NewLogFilter(syslog.LogFilterModeExclude, syslog.LOG_SOURCE_RTR), + }, + } + + for _, tc := range testCases { + By(tc.name) + bs := []syslog.Binding{ + {Drain: syslog.Drain{Url: tc.url}}, + } + f := newStubFetcher(bs, nil) + dp := bindings.NewDrainParamParser(f, true) + + configedBindings, _ := dp.FetchBindings() + Expect(configedBindings[0].LogFilter).To(Equal(tc.expected), "failed for case: %s", tc.name) } - f := newStubFetcher(bs, nil) - wf := bindings.NewDrainParamParser(f, true) - - configedBindings, _ := wf.FetchBindings() - Expect(configedBindings[0].DrainData).To(Equal(syslog.LOGS)) - Expect(configedBindings[1].DrainData).To(Equal(syslog.LOGS)) - Expect(configedBindings[2].DrainData).To(Equal(syslog.METRICS)) - Expect(configedBindings[3].DrainData).To(Equal(syslog.TRACES)) - Expect(configedBindings[4].DrainData).To(Equal(syslog.ALL)) }) It("sets drain data for old parameter appropriately'", func() { - bs := []syslog.Binding{ - {Drain: syslog.Drain{Url: "https://test.org/drain?drain-type=metrics"}}, - {Drain: syslog.Drain{Url: "https://test.org/drain?drain-type=logs"}}, - {Drain: syslog.Drain{Url: "https://test.org/drain"}}, - {Drain: syslog.Drain{Url: "https://test.org/drain?drain-type=all"}}, - {Drain: syslog.Drain{Url: "https://test.org/drain?include-metrics-deprecated=true"}}, + testCases := []struct { + name string + url string + expected syslog.DrainData + }{ + { + name: "drain-type=metrics", + url: "https://test.org/drain?drain-type=metrics", + expected: syslog.METRICS, + }, + { + name: "drain-type=logs", + url: "https://test.org/drain?drain-type=logs", + expected: syslog.LOGS_NO_EVENTS, + }, + { + name: "no drain-type parameter", + url: "https://test.org/drain", + expected: syslog.LOGS, + }, + { + name: "drain-type=all", + url: "https://test.org/drain?drain-type=all", + expected: syslog.LOGS_AND_METRICS, + }, + { + name: "include-metrics-deprecated=true", + url: "https://test.org/drain?include-metrics-deprecated=true", + expected: syslog.ALL, + }, + } + + for _, tc := range testCases { + By(tc.name) + bs := []syslog.Binding{ + {Drain: syslog.Drain{Url: tc.url}}, + } + f := newStubFetcher(bs, nil) + dp := bindings.NewDrainParamParser(f, true) + + configedBindings, _ := dp.FetchBindings() + Expect(configedBindings[0].DrainData).To(Equal(tc.expected)) } - f := newStubFetcher(bs, nil) - wf := bindings.NewDrainParamParser(f, true) - - configedBindings, _ := wf.FetchBindings() - Expect(configedBindings[0].DrainData).To(Equal(syslog.METRICS)) - Expect(configedBindings[1].DrainData).To(Equal(syslog.LOGS_NO_EVENTS)) - Expect(configedBindings[2].DrainData).To(Equal(syslog.LOGS)) - Expect(configedBindings[3].DrainData).To(Equal(syslog.LOGS_AND_METRICS)) - Expect(configedBindings[4].DrainData).To(Equal(syslog.ALL)) }) It("omits bindings with bad Drain URLs", func() { @@ -114,9 +224,9 @@ var _ = Describe("Drain Param Config", func() { {Drain: syslog.Drain{Url: "https://test.org/drain?omit-metadata=true"}}, } f := newStubFetcher(bs, nil) - wf := bindings.NewDrainParamParser(f, true) + dp := bindings.NewDrainParamParser(f, true) - configedBindings, err := wf.FetchBindings() + configedBindings, err := dp.FetchBindings() Expect(err).ToNot(HaveOccurred()) Expect(configedBindings).To(HaveLen(2)) Expect(configedBindings[0].Drain).To(Equal(syslog.Drain{Url: "https://test.org/drain?disable-metadata=true"})) @@ -125,9 +235,9 @@ var _ = Describe("Drain Param Config", func() { It("Returns a error when fetching fails", func() { f := newStubFetcher(nil, errors.New("Ahhh an error")) - wf := bindings.NewDrainParamParser(f, true) + dp := bindings.NewDrainParamParser(f, true) - _, err := wf.FetchBindings() + _, err := dp.FetchBindings() Expect(err).To(MatchError("Ahhh an error")) }) }) @@ -151,3 +261,11 @@ func (f *stubFetcher) FetchBindings() ([]syslog.Binding, error) { func (f *stubFetcher) DrainLimit() int { return -1 } + +func NewLogFilter(mode syslog.LogFilterMode, sourceTypes ...syslog.LogSourceType) *syslog.LogFilter { + set := make(syslog.LogSourceTypeSet, len(sourceTypes)) + for _, t := range sourceTypes { + set[t] = struct{}{} + } + return syslog.NewLogFilter(set, mode) +} diff --git a/src/pkg/ingress/bindings/filtered_binding_fetcher.go b/src/pkg/ingress/bindings/filtered_binding_fetcher.go index c9a08f0a3..4ae176042 100644 --- a/src/pkg/ingress/bindings/filtered_binding_fetcher.go +++ b/src/pkg/ingress/bindings/filtered_binding_fetcher.go @@ -3,6 +3,7 @@ package bindings import ( "log" "net/url" + "strings" "time" "code.cloudfoundry.org/loggregator-agent-release/src/pkg/binding" @@ -55,7 +56,7 @@ func (f *FilteredBindingFetcher) FetchBindings() ([]syslog.Binding, error) { continue } - anonymousUrl := u + anonymousUrl := *u anonymousUrl.User = nil anonymousUrl.RawQuery = "" @@ -69,6 +70,19 @@ func (f *FilteredBindingFetcher) FetchBindings() ([]syslog.Binding, error) { continue } + if invalidLogFilter(u) { + invalidDrains += 1 + f.printWarning("include-log-source-types and exclude-log-source-types cannot be used at the same time in syslog drain url %s for application %s", anonymousUrl.String(), b.AppId) + continue + } + + sourceTypes := getUnknownSourceTypes(u.Query()) + if sourceTypes != nil { + invalidDrains += 1 + f.printWarning("Unknown source types '%s' in source type filter in syslog drain url %s for application %s", strings.Join(sourceTypes, ", "), anonymousUrl.String(), b.AppId) + continue + } + _, exists := f.failedHostsCache.Get(u.Host) if exists { invalidDrains += 1 @@ -99,6 +113,34 @@ func (f *FilteredBindingFetcher) FetchBindings() ([]syslog.Binding, error) { } +// invalidLogFilter checks if both include-log-source-types and exclude-log-source-types +func invalidLogFilter(u *url.URL) bool { + includeSourceTypes := u.Query().Get("include-log-source-types") + excludeSourceTypes := u.Query().Get("exclude-log-source-types") + if excludeSourceTypes != "" && includeSourceTypes != "" { + return true + } + return false +} + +// assumes only one of include-log-source-types or exclude-log-source-types is set +func getUnknownSourceTypes(u url.Values) []string { + var sourceTypeList string + includeSourceTypes := u.Get("include-log-source-types") + excludeSourceTypes := u.Get("exclude-log-source-types") + + if includeSourceTypes != "" { + sourceTypeList = includeSourceTypes + } else if excludeSourceTypes != "" { + sourceTypeList = excludeSourceTypes + } else { + return nil + } + + _, unknownTypes := syslog.ParseSourceTypeList(sourceTypeList) + return unknownTypes +} + func (f FilteredBindingFetcher) printWarning(format string, v ...any) { if f.warn { f.logger.Printf(format, v...) diff --git a/src/pkg/ingress/bindings/filtered_binding_fetcher_test.go b/src/pkg/ingress/bindings/filtered_binding_fetcher_test.go index 767b9fdc9..703a3183f 100644 --- a/src/pkg/ingress/bindings/filtered_binding_fetcher_test.go +++ b/src/pkg/ingress/bindings/filtered_binding_fetcher_test.go @@ -239,6 +239,146 @@ var _ = Describe("FilteredBindingFetcher", func() { }) }) + Context("when both include-log-source-types and exclude-log-source-types are specified", func() { + var logBuffer bytes.Buffer + var warn bool + var mockic *bindingfakes.FakeIPChecker + + BeforeEach(func() { + logBuffer = bytes.Buffer{} + log.SetOutput(&logBuffer) + warn = true + mockic = &bindingfakes.FakeIPChecker{} + mockic.ResolveAddrReturns(net.ParseIP("10.10.10.10"), nil) + mockic.CheckBlacklistReturns(nil) + }) + + JustBeforeEach(func() { + input := []syslog.Binding{ + {AppId: "app-id", Hostname: "we.dont.care", Drain: syslog.Drain{Url: "https://test.org/drain?include-log-source-types=app&exclude-log-source-types=rtr"}}, + } + filter = bindings.NewFilteredBindingFetcher( + mockic, + &SpyBindingReader{bindings: input}, + warn, + log, + ) + }) + + It("ignores the drain", func() { + actual, err := filter.FetchBindings() + + Expect(err).ToNot(HaveOccurred()) + Expect(actual).To(HaveLen(0)) + Expect(logBuffer.String()).Should(MatchRegexp("include-log-source-types and exclude-log-source-types cannot be used at the same time")) + + }) + + Context("when configured not to warn", func() { + BeforeEach(func() { + warn = false + }) + It("doesn't log the conflicting filters warning", func() { + _, err := filter.FetchBindings() + Expect(err).ToNot(HaveOccurred()) + Expect(logBuffer.String()).ToNot(MatchRegexp("include-log-source-types and exclude-log-source-types cannot be used at the same time")) + }) + }) + }) + + Context("when unknown source types are provided", func() { + var logBuffer bytes.Buffer + var warn bool + var mockic *bindingfakes.FakeIPChecker + + BeforeEach(func() { + logBuffer = bytes.Buffer{} + log.SetOutput(&logBuffer) + warn = true + mockic = &bindingfakes.FakeIPChecker{} + mockic.ResolveAddrReturns(net.ParseIP("10.10.10.10"), nil) + mockic.CheckBlacklistReturns(nil) + }) + + It("logs a warning and ignores the drain in include mode", func() { + input := []syslog.Binding{ + {AppId: "app-id", Hostname: "we.dont.care", Drain: syslog.Drain{Url: "https://test.org/drain?include-log-source-types=app,unknown,invalid,rtr"}}, + } + filter = bindings.NewFilteredBindingFetcher( + mockic, + &SpyBindingReader{bindings: input}, + warn, + log, + ) + + actual, err := filter.FetchBindings() + + Expect(err).ToNot(HaveOccurred()) + Expect(actual).To(HaveLen(0)) + Expect(logBuffer.String()).Should(MatchRegexp("Unknown source types")) + Expect(logBuffer.String()).Should(MatchRegexp("unknown")) + Expect(logBuffer.String()).Should(MatchRegexp("invalid")) + }) + + It("logs a warning and ignores the drain in exclude mode", func() { + input := []syslog.Binding{ + {AppId: "app-id", Hostname: "we.dont.care", Drain: syslog.Drain{Url: "https://test.org/drain?exclude-log-source-types=rtr,unknown"}}, + } + filter = bindings.NewFilteredBindingFetcher( + mockic, + &SpyBindingReader{bindings: input}, + warn, + log, + ) + + actual, err := filter.FetchBindings() + + Expect(err).ToNot(HaveOccurred()) + Expect(actual).To(HaveLen(0)) + Expect(logBuffer.String()).Should(MatchRegexp("Unknown source types")) + Expect(logBuffer.String()).Should(MatchRegexp("unknown")) + }) + + It("logs a warning and ignores the drain when source types have spaces", func() { + input := []syslog.Binding{ + {AppId: "app-id", Hostname: "we.dont.care", Drain: syslog.Drain{Url: "https://test.org/drain?include-log-source-types=app, rtr"}}, + } + filter = bindings.NewFilteredBindingFetcher( + mockic, + &SpyBindingReader{bindings: input}, + warn, + log, + ) + + actual, err := filter.FetchBindings() + + Expect(err).ToNot(HaveOccurred()) + Expect(actual).To(HaveLen(0)) + Expect(logBuffer.String()).Should(MatchRegexp("Unknown source types")) + }) + + Context("when configured not to warn", func() { + BeforeEach(func() { + warn = false + }) + It("doesn't log the warning", func() { + input := []syslog.Binding{ + {AppId: "app-id", Hostname: "we.dont.care", Drain: syslog.Drain{Url: "https://test.org/drain?include-log-source-types=app,unknown,rtr"}}, + } + filter = bindings.NewFilteredBindingFetcher( + mockic, + &SpyBindingReader{bindings: input}, + warn, + log, + ) + + _, err := filter.FetchBindings() + Expect(err).ToNot(HaveOccurred()) + Expect(logBuffer.String()).ToNot(MatchRegexp("Unknown source types")) + }) + }) + }) + Context("when the syslog drain has been blacklisted", func() { var logBuffer bytes.Buffer var warn bool