Skip to content

Commit f82af4c

Browse files
committed
SCFF-44 added the include-only and exclude-always filters
1 parent a368f38 commit f82af4c

5 files changed

Lines changed: 245 additions & 61 deletions

File tree

README.md

Lines changed: 10 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,6 @@
22

33
This Nozzle aggregates all the events from the _Firehose_ feature in Cloud Foundry towards Sumo Logic
44

5-
## Getting Started
6-
75
### Options of use
86

97
```
@@ -18,17 +16,22 @@ Flags:
1816
Cloud Foundry User
1917
--cloudfoundry-password=CLOUDFOUNDRY-PASSWORD
2018
Cloud Foundry Password
21-
--events="LogMessage" Comma separated list of events you would like. Valid options are ContainerMetric, CounterEvent, Error, HttpStart,HttpStartStop, HttpStop, LogMessage, ValueMetric
19+
--events="LogMessage" Comma separated list of events you would like. Valid options are ContainerMetric,
20+
CounterEvent, Error, HttpStart, HttpStartStop, HttpStop, LogMessage, ValueMetric
2221
--nozzle-polling-period=15s Nozzle Polling Period
23-
--log-events-batch-size=LOG-EVENTS-BATCH-SIZE
24-
Log Events Batch Size
25-
--sumo-post-minimum-delay=SUMO-POST-MINIMUM-DELAY
22+
--log-events-batch-size=200 Log Events Batch Size to send to Sumo
23+
--sumo-post-minimum-delay=200ms
2624
Sumo Logic HTTP Post Minimum Delay
2725
--sumo-category="" Sumo Logic Category
2826
--sumo-name="" Sumo Logic Name
2927
--sumo-host="" Sumo Logic Host
3028
--verbose-log-messages Allow Verbose Log Messages
31-
--custom-metadata="" Custom Metadata
29+
--custom-metadata="" Use this flag for addingCustom Metadata (key1:value1,key2:value2, etc...)
30+
--include-only-matching-filter=""
31+
Adds an 'Include only' filter to Events content (key1:value1,key2:value2, etc...)
32+
--exclude-always-matching-filter=""
33+
Adds an 'Exclude always' filter to Events content (key1:value1,key2:value2,
34+
etc...)
3235
--version Show application version.
3336
```
3437

main.go

Lines changed: 17 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -21,20 +21,22 @@ var (
2121
apiEndpoint = kingpin.Flag("api-endpoint", "CF API Endpoint").OverrideDefaultFromEnvar("API_ENDPOINT").String()
2222
sumoEndpoint = kingpin.Flag("sumo-endpoint", "Sumo Logic Endpoint").OverrideDefaultFromEnvar("SUMO_ENDPOINT").String()
2323
//dopplerEndpoint = kingpin.Flag("doppler-endpoint", "Overwrite default doppler endpoint return by /v2/info").OverrideDefaultFromEnvar("DOPPLER_ENDPOINT").String()
24-
subscriptionId = kingpin.Flag("subscription-id", "Cloud Foundry ID for the subscription.").Default("firehose").OverrideDefaultFromEnvar("FIREHOSE_SUBSCRIPTION_ID").String()
25-
user = kingpin.Flag("cloudfoundry-user", "Cloud Foundry User").OverrideDefaultFromEnvar("CLOUDFOUNDRY_USER").String() //user created in CF, authorized to connect the firehose
26-
password = kingpin.Flag("cloudfoundry-password", "Cloud Foundry Password").OverrideDefaultFromEnvar("CLOUDFOUNDRY_PASSWORD").String() // password created along with the firehose_user //kingpin.Flag("skip-ssl-validation", "Please don't").Default("false").OverrideDefaultFromEnvar("SKIP_SSL_VALIDATION").Bool()
27-
keepAlive, errK = time.ParseDuration("25s") //default Error, ContainerMetric, HttpStart, HttpStop, HttpStartStop, LogMessage, ValueMetric, CounterEvent
28-
wantedEvents = kingpin.Flag("events", fmt.Sprintf("Comma separated list of events you would like. Valid options are %s", eventRouting.GetListAuthorizedEventEvents())).Default("LogMessage").OverrideDefaultFromEnvar("EVENTS").String()
29-
boltDatabasePath = "event.db"
30-
tickerTime = kingpin.Flag("nozzle-polling-period", "Nozzle Polling Period").Default("15s").OverrideDefaultFromEnvar("NOZZLE_POLLING_PERIOD").Duration()
31-
eventsBatchSize = kingpin.Flag("log-events-batch-size", "Log Events Batch Size").OverrideDefaultFromEnvar("LOG_EVENTS_BATCH_SIZE").Int()
32-
sumoPostMinimumDelay = kingpin.Flag("sumo-post-minimum-delay", "Sumo Logic HTTP Post Minimum Delay").OverrideDefaultFromEnvar("SUMO_POST_MINIMUM_DELAY").Duration()
33-
sumoCategory = kingpin.Flag("sumo-category", "Sumo Logic Category").Default("").OverrideDefaultFromEnvar("SUMO_CATEGORY").String()
34-
sumoName = kingpin.Flag("sumo-name", "Sumo Logic Name").Default("").OverrideDefaultFromEnvar("SUMO_NAME").String()
35-
sumoHost = kingpin.Flag("sumo-host", "Sumo Logic Host").Default("").OverrideDefaultFromEnvar("SUMO_HOST").String()
36-
verboseLogMessages = kingpin.Flag("verbose-log-messages", "Allow Verbose Log Messages").Default("false").OverrideDefaultFromEnvar("VERBOSE_LOG_MESSAGES").Bool()
37-
customMetadata = kingpin.Flag("custom-metadata", "Custom Metadata").Default("").OverrideDefaultFromEnvar("CUSTOM_METADATA").String()
24+
subscriptionId = kingpin.Flag("subscription-id", "Cloud Foundry ID for the subscription.").Default("firehose").OverrideDefaultFromEnvar("FIREHOSE_SUBSCRIPTION_ID").String()
25+
user = kingpin.Flag("cloudfoundry-user", "Cloud Foundry User").OverrideDefaultFromEnvar("CLOUDFOUNDRY_USER").String() //user created in CF, authorized to connect the firehose
26+
password = kingpin.Flag("cloudfoundry-password", "Cloud Foundry Password").OverrideDefaultFromEnvar("CLOUDFOUNDRY_PASSWORD").String() // password created along with the firehose_user //kingpin.Flag("skip-ssl-validation", "Please don't").Default("false").OverrideDefaultFromEnvar("SKIP_SSL_VALIDATION").Bool()
27+
keepAlive, errK = time.ParseDuration("25s") //default Error, ContainerMetric, HttpStart, HttpStop, HttpStartStop, LogMessage, ValueMetric, CounterEvent
28+
wantedEvents = kingpin.Flag("events", fmt.Sprintf("Comma separated list of events you would like. Valid options are %s", eventRouting.GetListAuthorizedEventEvents())).Default("LogMessage").OverrideDefaultFromEnvar("EVENTS").String()
29+
boltDatabasePath = "event.db"
30+
tickerTime = kingpin.Flag("nozzle-polling-period", "Nozzle Polling Period").Default("15s").OverrideDefaultFromEnvar("NOZZLE_POLLING_PERIOD").Duration()
31+
eventsBatchSize = kingpin.Flag("log-events-batch-size", "Log Events Batch Size to send to Sumo").Default("200").OverrideDefaultFromEnvar("LOG_EVENTS_BATCH_SIZE").Int()
32+
sumoPostMinimumDelay = kingpin.Flag("sumo-post-minimum-delay", "Sumo Logic HTTP Post Minimum Delay").Default("200ms").OverrideDefaultFromEnvar("SUMO_POST_MINIMUM_DELAY").Duration()
33+
sumoCategory = kingpin.Flag("sumo-category", "Sumo Logic Category").Default("").OverrideDefaultFromEnvar("SUMO_CATEGORY").String()
34+
sumoName = kingpin.Flag("sumo-name", "Sumo Logic Name").Default("").OverrideDefaultFromEnvar("SUMO_NAME").String()
35+
sumoHost = kingpin.Flag("sumo-host", "Sumo Logic Host").Default("").OverrideDefaultFromEnvar("SUMO_HOST").String()
36+
verboseLogMessages = kingpin.Flag("verbose-log-messages", "Allow Verbose Log Messages").Default("false").OverrideDefaultFromEnvar("VERBOSE_LOG_MESSAGES").Bool()
37+
customMetadata = kingpin.Flag("custom-metadata", "Use this flag for addingCustom Metadata (key1:value1,key2:value2, etc...)").Default("").OverrideDefaultFromEnvar("CUSTOM_METADATA").String()
38+
includeOnlyMatchingFilter = kingpin.Flag("include-only-matching-filter", "Adds an 'Include only' filter to Events content (key1:value1,key2:value2, etc...)").Default("").OverrideDefaultFromEnvar("INCLUDE_ONLY_MATCHING_FILTER").String()
39+
excludeAlwaysMatchingFilter = kingpin.Flag("exclude-always-matching-filter", "Adds an 'Exclude always' filter to Events content (key1:value1,key2:value2, etc...)").Default("").OverrideDefaultFromEnvar("EXCLUDE_ALWAYS_MATCHING_FILTER").String()
3840
)
3941

4042
var (
@@ -88,7 +90,7 @@ func main() {
8890

8991
logging.Info.Println("Creating queue")
9092
queue := eventQueue.NewQueue(make([]*events.Event, 100))
91-
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 5000, &queue, *eventsBatchSize, *sumoPostMinimumDelay, *sumoCategory, *sumoName, *sumoHost, *verboseLogMessages, *customMetadata)
93+
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 5000, &queue, *eventsBatchSize, *sumoPostMinimumDelay, *sumoCategory, *sumoName, *sumoHost, *verboseLogMessages, *customMetadata, *includeOnlyMatchingFilter, *excludeAlwaysMatchingFilter)
9294
go loggingClientSumo.Start() //multi
9395

9496
logging.Info.Println("Creating Events")

sumoCFFirehose/sumoLogicAppender.go

Lines changed: 86 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -15,18 +15,20 @@ import (
1515
)
1616

1717
type SumoLogicAppender struct {
18-
url string
19-
connectionTimeout int //10000
20-
httpClient http.Client
21-
nozzleQueue *eventQueue.Queue
22-
eventsBatchSize int
23-
sumoPostMinimumDelay time.Duration
24-
timerBetweenPost time.Time
25-
sumoCategory string
26-
sumoName string
27-
sumoHost string
28-
verboseLogMessages bool
29-
customMetadata string
18+
url string
19+
connectionTimeout int //10000
20+
httpClient http.Client
21+
nozzleQueue *eventQueue.Queue
22+
eventsBatchSize int
23+
sumoPostMinimumDelay time.Duration
24+
timerBetweenPost time.Time
25+
sumoCategory string
26+
sumoName string
27+
sumoHost string
28+
verboseLogMessages bool
29+
customMetadata string
30+
includeOnlyMatchingFilter string
31+
excludeAlwaysMatchingFilter string
3032
}
3133

3234
type SumoBuffer struct {
@@ -35,19 +37,21 @@ type SumoBuffer struct {
3537
timerIdlebuffer time.Time
3638
}
3739

38-
func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQueue *eventQueue.Queue, eventsBatchSize int, sumoPostMinimumDelay time.Duration, sumoCategory string, sumoName string, sumoHost string, verboseLogMessages bool, customMetadata string) *SumoLogicAppender {
40+
func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQueue *eventQueue.Queue, eventsBatchSize int, sumoPostMinimumDelay time.Duration, sumoCategory string, sumoName string, sumoHost string, verboseLogMessages bool, customMetadata string, includeOnlyMatchingFilter string, excludeAlwaysMatchingFilter string) *SumoLogicAppender {
3941
return &SumoLogicAppender{
40-
url: urlValue,
41-
connectionTimeout: connectionTimeoutValue,
42-
httpClient: http.Client{Timeout: time.Duration(connectionTimeoutValue * int(time.Millisecond))},
43-
nozzleQueue: nozzleQueue,
44-
eventsBatchSize: eventsBatchSize,
45-
sumoPostMinimumDelay: sumoPostMinimumDelay,
46-
sumoCategory: sumoCategory,
47-
sumoName: sumoName,
48-
sumoHost: sumoHost,
49-
verboseLogMessages: verboseLogMessages,
50-
customMetadata: customMetadata,
42+
url: urlValue,
43+
connectionTimeout: connectionTimeoutValue,
44+
httpClient: http.Client{Timeout: time.Duration(connectionTimeoutValue * int(time.Millisecond))},
45+
nozzleQueue: nozzleQueue,
46+
eventsBatchSize: eventsBatchSize,
47+
sumoPostMinimumDelay: sumoPostMinimumDelay,
48+
sumoCategory: sumoCategory,
49+
sumoName: sumoName,
50+
sumoHost: sumoHost,
51+
verboseLogMessages: verboseLogMessages,
52+
customMetadata: customMetadata,
53+
includeOnlyMatchingFilter: includeOnlyMatchingFilter,
54+
excludeAlwaysMatchingFilter: excludeAlwaysMatchingFilter,
5155
}
5256
}
5357

@@ -104,9 +108,45 @@ func (s *SumoLogicAppender) Start() {
104108

105109
}
106110

111+
}
112+
func WantedEvent(event string, includeOnlyMatchingFilter string, excludeAlwaysMatchingFilter string) bool {
113+
if includeOnlyMatchingFilter != "" && excludeAlwaysMatchingFilter != "" {
114+
subsliceInclude := ParseCustomInput(includeOnlyMatchingFilter)
115+
subsliceExclude := ParseCustomInput(excludeAlwaysMatchingFilter)
116+
for key, value := range subsliceInclude {
117+
if strings.Contains(event, "\""+key+"\":\""+value+"\"") {
118+
return true
119+
}
120+
}
121+
for key, value := range subsliceExclude {
122+
if strings.Contains(event, "\""+key+"\":\""+value+"\"") {
123+
return false
124+
}
125+
}
126+
return false
127+
} else if includeOnlyMatchingFilter != "" {
128+
subslice := ParseCustomInput(includeOnlyMatchingFilter)
129+
for key, value := range subslice {
130+
if strings.Contains(event, "\""+key+"\":\""+value+"\"") {
131+
return true
132+
}
133+
}
134+
return false
135+
136+
} else if excludeAlwaysMatchingFilter != "" {
137+
subslice := ParseCustomInput(excludeAlwaysMatchingFilter)
138+
for key, value := range subslice {
139+
if strings.Contains(event, "\""+key+"\":\""+value+"\"") {
140+
return false
141+
}
142+
}
143+
return true
144+
}
145+
return true
146+
107147
}
108148

109-
func StringBuilder(event *events.Event, verboseLogMessages bool) string {
149+
func StringBuilder(event *events.Event, verboseLogMessages bool, includeOnlyMatchingFilter string, excludeAlwaysMatchingFilter string) string {
110150
eventType := event.Type
111151
var msg []byte
112152
switch eventType {
@@ -176,23 +216,29 @@ func StringBuilder(event *events.Event, verboseLogMessages bool) string {
176216
msg = message
177217
}
178218
}
219+
179220
buf := new(bytes.Buffer)
180221
buf.Write(msg)
181-
return buf.String() + "\n"
222+
if WantedEvent(buf.String(), includeOnlyMatchingFilter, excludeAlwaysMatchingFilter) {
223+
return buf.String() + "\n"
224+
} else {
225+
return ""
226+
}
227+
182228
}
183229

184230
func (s *SumoLogicAppender) AppendLogs(buffer *SumoBuffer) {
185-
buffer.logStringToSend.Write([]byte(StringBuilder(s.nozzleQueue.Pop(), s.verboseLogMessages)))
231+
buffer.logStringToSend.Write([]byte(StringBuilder(s.nozzleQueue.Pop(), s.verboseLogMessages, s.includeOnlyMatchingFilter, s.excludeAlwaysMatchingFilter)))
186232
buffer.logEventsInCurrentBuffer++
187233

188234
}
189-
func ParseCustomMetadata(customMetadata string) map[string]string {
190-
cMetadataArray := strings.Split(customMetadata, ",")
191-
customMetadataMap := make(map[string]string)
192-
for i := 0; i < len(cMetadataArray); i++ {
193-
customMetadataMap[strings.Split(cMetadataArray[i], ":")[0]] = strings.Split(cMetadataArray[i], ":")[1]
235+
func ParseCustomInput(customInput string) map[string]string {
236+
cInputArray := strings.Split(customInput, ",")
237+
customInputMap := make(map[string]string)
238+
for i := 0; i < len(cInputArray); i++ {
239+
customInputMap[strings.Split(cInputArray[i], ":")[0]] = strings.Split(cInputArray[i], ":")[1]
194240
}
195-
return customMetadataMap
241+
return customInputMap
196242
}
197243

198244
func (s *SumoLogicAppender) SendToSumo(logStringToSend string) {
@@ -219,7 +265,7 @@ func (s *SumoLogicAppender) SendToSumo(logStringToSend string) {
219265
}
220266

221267
if s.customMetadata != "" {
222-
customMetadataMap := ParseCustomMetadata(s.customMetadata)
268+
customMetadataMap := ParseCustomInput(s.customMetadata)
223269
for key, value := range customMetadataMap {
224270
request.Header.Add(key, value)
225271
}
@@ -253,6 +299,12 @@ func (s *SumoLogicAppender) SendToSumo(logStringToSend string) {
253299
if s.sumoCategory != "" {
254300
request.Header.Add("X-Sumo-Category", s.sumoCategory)
255301
}
302+
if s.customMetadata != "" {
303+
customMetadataMap := ParseCustomInput(s.customMetadata)
304+
for key, value := range customMetadataMap {
305+
request.Header.Add(key, value)
306+
}
307+
}
256308
//checking the timer before POST (retry intent)
257309
for time.Since(s.timerBetweenPost) < s.sumoPostMinimumDelay {
258310
logging.Trace.Println("Delaying Post because minimum post timer not expired")

0 commit comments

Comments
 (0)