Skip to content

Commit 6628ec9

Browse files
committed
SCFF-30 remove the queue from Appender
1 parent 12ca162 commit 6628ec9

3 files changed

Lines changed: 22 additions & 20 deletions

File tree

eventQueue/eventQueue_test.go

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ import (
88
)
99

1010
func TestQueueFIFO(t *testing.T) {
11-
11+
assert := assert.New(t)
1212
node1 := Node{
1313
Event: Event{
1414
Fields: map[string]interface{}{
@@ -45,7 +45,7 @@ func TestQueueFIFO(t *testing.T) {
4545
queue.Push(&node2)
4646
queue.Push(&node3)
4747

48-
assert.Equal(t, queue.Pop().Event.Msg, "index [01]", "")
49-
assert.Equal(t, queue.Pop().Event.Msg, "index [02]", "")
50-
assert.Equal(t, queue.Pop().Event.Msg, "index [03]", "")
48+
assert.Equal(queue.Pop().Event.Msg, "index [01]", "")
49+
assert.Equal(queue.Pop().Event.Msg, "index [02]", "")
50+
assert.Equal(queue.Pop().Event.Msg, "index [03]", "")
5151
}

sumoCFFirehose/sumoLogicAppender.go

Lines changed: 14 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -45,29 +45,28 @@ func (s *SumoLogicAppender) Connect() bool {
4545
return success
4646
}
4747

48-
func StringBuilder(queue eventQueue.Queue) string {
48+
func StringBuilder(node *eventQueue.Node) string {
4949
buf := new(bytes.Buffer)
50-
for queue.GetCount() > 0 { //Pop eventsBatch from queue
51-
event := queue.Pop().GetNodeEvent()
52-
if event.Fields["message_type"] == nil {
53-
return ""
54-
}
55-
if event.Fields["message_type"] == "" {
56-
return ""
57-
}
58-
message := time.Unix(0, event.Fields["timestamp"].(int64)*int64(time.Nanosecond)).String() + "\t" + event.Fields["message_type"].(string) + "\t" + event.Msg + "\n"
59-
buf.WriteString(message)
50+
if node.Event.Fields["message_type"] == nil {
51+
return ""
6052
}
53+
if node.Event.Fields["message_type"] == "" {
54+
return ""
55+
}
56+
message := time.Unix(0, node.Event.Fields["timestamp"].(int64)*int64(time.Nanosecond)).String() + "\t" + node.Event.Fields["message_type"].(string) + "\t" + node.Event.Msg + "\n"
57+
buf.WriteString(message)
6158

6259
return buf.String()
6360
}
6461

65-
func (s *SumoLogicAppender) AppendLogs(queue eventQueue.Queue) {
62+
func (s *SumoLogicAppender) AppendLogs() {
6663
// the appender calls for the next message in the queue and parse it to a string
67-
if queue.GetCount() > s.eventsBatchSize { //whent the batch limit is met, call stringBuilder
68-
logMessage := StringBuilder(queue)
69-
s.SendToSumo(logMessage)
64+
logMessage := ""
65+
if s.nozzleQueue.GetCount() > s.eventsBatchSize { //when the batch limit is met, call stringBuilder
66+
logMessage = logMessage + StringBuilder(s.nozzleQueue.Pop())
67+
7068
}
69+
s.SendToSumo(logMessage)
7170

7271
}
7372

sumoCFFirehose/sumoLogicAppender_test.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,10 @@ func testAppenderStringBuilder(t *testing.T) {
4848
queue.Push(&node2)
4949
queue.Push(&node3)
5050

51-
finalString := StringBuilder(queue)
51+
finalString := ""
52+
for queue.GetCount() > 0 {
53+
finalString = finalString + StringBuilder(queue.Pop())
54+
}
5255
assert.Equal(t, finalString, "2016-12-12 16:02:41.828366387 -0300 CLST"+"\t"+"OUT"+"\t"+"index [01]"+"\n"+
5356
"2016-12-12 16:02:42.844737993 -0300 CLST"+"\t"+"OUT"+"\t"+"index [02]"+"\n"+
5457
"2016-12-12 16:02:43.862436654 -0300 CLST"+"\t"+"OUT"+"\t"+"index [03]"+"\n", "")

0 commit comments

Comments
 (0)