Skip to content

Commit ea72eb9

Browse files
committed
SCFF-30 created the Stringbuilder for the log entries using eventsAmount (kingpin)
1 parent 636fe17 commit ea72eb9

2 files changed

Lines changed: 17 additions & 19 deletions

File tree

main.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ var (
2828
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()
2929
boltDatabasePath = "my.db" //default
3030
tickerTime, errT = time.ParseDuration("60s") //Default
31+
eventsAmount = kingpin.Flag("event-amount", "Events amount").Int()
3132
)
3233

3334
var (
@@ -43,7 +44,7 @@ func main() {
4344

4445
//Creating queue
4546
queue := eventQueue.NewQueue(make([]*eventQueue.Node, 100))
46-
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 1000, *queue)
47+
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 1000, *queue, *eventsAmount)
4748

4849
fmt.Printf("Starting firehose-to-sumo %s \n", version)
4950

@@ -99,7 +100,6 @@ func main() {
99100

100101
firehoseClient := firehoseclient.NewFirehoseNozzle(cfClient, events, firehoseConfig)
101102
err = firehoseClient.Start()
102-
fmt.Printf("I created the Firehose... \n")
103103
if err != nil {
104104
fmt.Printf("Failed connecting to Firehose...Please check settings and try again! \n") //Log error
105105

sumoCFFirehose/sumoLogicAppender.go

Lines changed: 15 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -15,14 +15,16 @@ type SumoLogicAppender struct {
1515
connectionTimeout int //10000
1616
httpClient http.Client
1717
nozzleQueue eventQueue.Queue
18+
eventsAmount int
1819
}
1920

20-
func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQueue eventQueue.Queue) *SumoLogicAppender {
21+
func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQueue eventQueue.Queue, eventsAmount int) *SumoLogicAppender {
2122
return &SumoLogicAppender{
2223
url: urlValue,
2324
connectionTimeout: connectionTimeoutValue,
2425
httpClient: http.Client{Timeout: time.Duration(connectionTimeoutValue * int(time.Millisecond))},
2526
nozzleQueue: nozzleQueue,
27+
eventsAmount: eventsAmount,
2628
}
2729
}
2830

@@ -44,25 +46,21 @@ func (s *SumoLogicAppender) Connect() bool {
4446
}
4547

4648
func (s *SumoLogicAppender) AppendLogs() {
47-
// the appender calls for the next message in the queue and parse it to a string
4849
fmt.Println("i'm in appendLogs")
49-
event := s.nozzleQueue.Pop().GetNodeEvent()
50-
/*
51-
if event == nil {
50+
// the appender calls for the next message in the queue and parse it to a string
51+
buf := new(bytes.Buffer)
52+
for i := 0; i <= s.eventsAmount; i++ { //Pop eventsAmount from queue
53+
event := s.nozzleQueue.Pop().GetNodeEvent()
54+
if event.Fields["message_type"] == nil {
5255
return
53-
}*/
54-
55-
if event.Fields["message_type"] == nil {
56-
return
57-
}
58-
59-
if event.Fields["message_type"] == "" {
60-
return
56+
}
57+
if event.Fields["message_type"] == "" {
58+
return
59+
}
60+
message := time.Unix(0, event.Fields["timestamp"].(int64)*int64(time.Nanosecond)).String() + "\t" + event.Fields["message_type"].(string) + "\t" + event.Msg + "\n"
61+
buf.WriteString(message)
6162
}
62-
63-
Message := time.Unix(0, event.Fields["timestamp"].(int64)*int64(time.Nanosecond)).String() + "\t" + event.Fields["message_type"].(string) + "\t" + event.Msg + "\n"
64-
fmt.Println(Message)
65-
s.SendToSumo(Message)
63+
s.SendToSumo(buf.String())
6664
}
6765

6866
func (s *SumoLogicAppender) SendToSumo(log string) {

0 commit comments

Comments
 (0)