Skip to content

Commit b1962d9

Browse files
committed
Merged in feature/SCFF-30 (pull request #8)
Feature/SCFF-30
2 parents 636fe17 + 0a1a96e commit b1962d9

5 files changed

Lines changed: 108 additions & 44 deletions

File tree

eventQueue/eventQueue.go

Lines changed: 15 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -6,53 +6,53 @@ import (
66

77
//Node to put in queue
88
type Node struct {
9-
event Event
9+
Event Event
1010
}
1111

1212
// Queue is a basic FIFO queue based on a circular list that resizes as needed.
1313
type Queue struct {
14-
nodes []*Node
14+
Nodes []*Node
1515
head int
1616
tail int
1717
count int
1818
}
1919

2020
func NewNode(event Event) *Node {
2121
return &Node{
22-
event: event,
22+
Event: event,
2323
}
2424
}
2525

2626
func NewQueue(n []*Node) *Queue {
2727
return &Queue{
28-
nodes: n,
28+
Nodes: n,
2929
}
3030
}
3131

3232
func (q *Queue) GetNode() []*Node {
33-
return q.nodes
33+
return q.Nodes
3434
}
3535

3636
func (n *Queue) GetCount() int {
3737
return n.count
3838
}
3939

4040
func (n *Node) GetNodeEvent() Event {
41-
return n.event
41+
return n.Event
4242
}
4343

4444
// Push adds a node to the queue.
4545
func (q *Queue) Push(n *Node) {
4646
if q.head == q.tail && q.count > 0 {
47-
nodes := make([]*Node, len(q.nodes)*2)
48-
copy(nodes, q.nodes[q.head:])
49-
copy(nodes[len(q.nodes)-q.head:], q.nodes[:q.head])
47+
nodes := make([]*Node, len(q.Nodes)*2)
48+
copy(nodes, q.Nodes[q.head:])
49+
copy(nodes[len(q.Nodes)-q.head:], q.Nodes[:q.head])
5050
q.head = 0
51-
q.tail = len(q.nodes)
52-
q.nodes = nodes
51+
q.tail = len(q.Nodes)
52+
q.Nodes = nodes
5353
}
54-
q.nodes[q.tail] = n
55-
q.tail = (q.tail + 1) % len(q.nodes)
54+
q.Nodes[q.tail] = n
55+
q.tail = (q.tail + 1) % len(q.Nodes)
5656
q.count++
5757
}
5858

@@ -61,8 +61,8 @@ func (q *Queue) Pop() *Node {
6161
if q.count == 0 {
6262
return nil
6363
}
64-
node := q.nodes[q.head]
65-
q.head = (q.head + 1) % len(q.nodes)
64+
node := q.Nodes[q.head]
65+
q.head = (q.head + 1) % len(q.Nodes)
6666
q.count--
6767
return node
6868
}

eventQueue/eventQueue_test.go

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

1010
func TestQueueFIFO(t *testing.T) {
11-
11+
assert := assert.New(t)
1212
node1 := Node{
13-
event: Event{
13+
Event: Event{
1414
Fields: map[string]interface{}{
1515
"message_type": "OUT",
1616
"cf_app_id": "011",
@@ -19,7 +19,7 @@ func TestQueueFIFO(t *testing.T) {
1919
},
2020
}
2121
node2 := Node{
22-
event: Event{
22+
Event: Event{
2323
Fields: map[string]interface{}{
2424
"message_type": "OUT",
2525
"cf_app_id": "022",
@@ -28,7 +28,7 @@ func TestQueueFIFO(t *testing.T) {
2828
},
2929
}
3030
node3 := Node{
31-
event: Event{
31+
Event: Event{
3232
Fields: map[string]interface{}{
3333
"message_type": "OUT",
3434
"cf_app_id": "033",
@@ -38,14 +38,14 @@ func TestQueueFIFO(t *testing.T) {
3838
}
3939

4040
queue := Queue{
41-
nodes: make([]*Node, 3),
41+
Nodes: make([]*Node, 3),
4242
}
4343

4444
queue.Push(&node1)
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
}

main.go

Lines changed: 5 additions & 5 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+
eventsBatchSize = kingpin.Flag("log-events-batch-size", "Log Events Batch Size").Int()
3132
)
3233

3334
var (
@@ -41,10 +42,6 @@ func main() {
4142
fmt.Println("this is the sumo endpoint")
4243
fmt.Println(sumoEndpoint)
4344

44-
//Creating queue
45-
queue := eventQueue.NewQueue(make([]*eventQueue.Node, 100))
46-
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 1000, *queue)
47-
4845
fmt.Printf("Starting firehose-to-sumo %s \n", version)
4946

5047
c := cfclient.Config{
@@ -67,6 +64,10 @@ func main() {
6764
cachingClient = caching.NewCachingEmpty()
6865
}
6966

67+
//Creating queue
68+
queue := eventQueue.NewQueue(make([]*eventQueue.Node, 100))
69+
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 1000, *queue, *eventsBatchSize)
70+
7071
//Creating Events
7172
events := eventRouting.NewEventRouting(cachingClient, *loggingClientSumo, *queue)
7273
err := events.SetupEventRouting(*wantedEvents)
@@ -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: 21 additions & 16 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+
eventsBatchSize int
1819
}
1920

20-
func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQueue eventQueue.Queue) *SumoLogicAppender {
21+
func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQueue eventQueue.Queue, eventsBatchSize 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+
eventsBatchSize: eventsBatchSize,
2628
}
2729
}
2830

@@ -43,26 +45,29 @@ func (s *SumoLogicAppender) Connect() bool {
4345
return success
4446
}
4547

48+
func StringBuilder(node *eventQueue.Node) string {
49+
buf := new(bytes.Buffer)
50+
if node.Event.Fields["message_type"] == nil {
51+
return ""
52+
}
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)
58+
59+
return buf.String()
60+
}
61+
4662
func (s *SumoLogicAppender) AppendLogs() {
4763
// the appender calls for the next message in the queue and parse it to a string
48-
fmt.Println("i'm in appendLogs")
49-
event := s.nozzleQueue.Pop().GetNodeEvent()
50-
/*
51-
if event == nil {
52-
return
53-
}*/
64+
logMessage := ""
65+
if s.nozzleQueue.GetCount() > s.eventsBatchSize { //when the batch limit is met, call stringBuilder
66+
logMessage = logMessage + StringBuilder(s.nozzleQueue.Pop())
5467

55-
if event.Fields["message_type"] == nil {
56-
return
57-
}
58-
59-
if event.Fields["message_type"] == "" {
60-
return
6168
}
69+
s.SendToSumo(logMessage)
6270

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)
6671
}
6772

6873
func (s *SumoLogicAppender) SendToSumo(log string) {
Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,59 @@
1+
package sumoCFFirehose
2+
3+
import (
4+
"testing"
5+
6+
. "bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/eventQueue"
7+
. "bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/events"
8+
"github.com/stretchr/testify/assert"
9+
)
10+
11+
func testAppenderStringBuilder(t *testing.T) {
12+
13+
node1 := Node{
14+
Event: Event{
15+
Fields: map[string]interface{}{
16+
"timestamp": "1481569361828366387",
17+
"message_type": "OUT",
18+
"cf_app_id": "011",
19+
},
20+
Msg: "index [01]",
21+
},
22+
}
23+
node2 := Node{
24+
Event: Event{
25+
Fields: map[string]interface{}{
26+
"timestamp": "1481569362844737993",
27+
"message_type": "OUT",
28+
"cf_app_id": "022",
29+
},
30+
Msg: "index [02]",
31+
},
32+
}
33+
node3 := Node{
34+
Event: Event{
35+
Fields: map[string]interface{}{
36+
"timestamp": "1481569363862436654",
37+
"message_type": "OUT",
38+
"cf_app_id": "033",
39+
},
40+
Msg: "index [03]",
41+
},
42+
}
43+
44+
queue := Queue{
45+
Nodes: make([]*Node, 3),
46+
}
47+
queue.Push(&node1)
48+
queue.Push(&node2)
49+
queue.Push(&node3)
50+
51+
finalString := ""
52+
for queue.GetCount() > 0 {
53+
finalString = finalString + StringBuilder(queue.Pop())
54+
}
55+
assert.Equal(t, finalString, "2016-12-12 16:02:41.828366387 -0300 CLST"+"\t"+"OUT"+"\t"+"index [01]"+"\n"+
56+
"2016-12-12 16:02:42.844737993 -0300 CLST"+"\t"+"OUT"+"\t"+"index [02]"+"\n"+
57+
"2016-12-12 16:02:43.862436654 -0300 CLST"+"\t"+"OUT"+"\t"+"index [03]"+"\n", "")
58+
59+
}

0 commit comments

Comments
 (0)