Skip to content

Commit d9310ca

Browse files
committed
SCFF-30 unit test for StringBuilder method in Appender
1 parent ea72eb9 commit d9310ca

5 files changed

Lines changed: 102 additions & 38 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: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ import (
1010
func TestQueueFIFO(t *testing.T) {
1111

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(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]", "")
5151
}

main.go

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -28,7 +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()
31+
eventsBatch = kingpin.Flag("event-amount", "Events amount").Int()
3232
)
3333

3434
var (
@@ -42,10 +42,6 @@ func main() {
4242
fmt.Println("this is the sumo endpoint")
4343
fmt.Println(sumoEndpoint)
4444

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

5147
c := cfclient.Config{
@@ -68,6 +64,10 @@ func main() {
6864
cachingClient = caching.NewCachingEmpty()
6965
}
7066

67+
//Creating queue
68+
queue := eventQueue.NewQueue(make([]*eventQueue.Node, 100))
69+
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 1000, *queue, *eventsBatch)
70+
7171
//Creating Events
7272
events := eventRouting.NewEventRouting(cachingClient, *loggingClientSumo, *queue)
7373
err := events.SetupEventRouting(*wantedEvents)

sumoCFFirehose/sumoLogicAppender.go

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

21-
func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQueue eventQueue.Queue, eventsAmount int) *SumoLogicAppender {
21+
func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQueue eventQueue.Queue, eventsBatch int) *SumoLogicAppender {
2222
return &SumoLogicAppender{
2323
url: urlValue,
2424
connectionTimeout: connectionTimeoutValue,
2525
httpClient: http.Client{Timeout: time.Duration(connectionTimeoutValue * int(time.Millisecond))},
2626
nozzleQueue: nozzleQueue,
27-
eventsAmount: eventsAmount,
27+
eventsBatch: eventsBatch,
2828
}
2929
}
3030

@@ -45,22 +45,30 @@ func (s *SumoLogicAppender) Connect() bool {
4545
return success
4646
}
4747

48-
func (s *SumoLogicAppender) AppendLogs() {
49-
fmt.Println("i'm in appendLogs")
50-
// the appender calls for the next message in the queue and parse it to a string
48+
func StringBuilder(queue eventQueue.Queue) string {
5149
buf := new(bytes.Buffer)
52-
for i := 0; i <= s.eventsAmount; i++ { //Pop eventsAmount from queue
53-
event := s.nozzleQueue.Pop().GetNodeEvent()
50+
for queue.GetCount() > 0 { //Pop eventsBatch from queue
51+
event := queue.Pop().GetNodeEvent()
5452
if event.Fields["message_type"] == nil {
55-
return
53+
return ""
5654
}
5755
if event.Fields["message_type"] == "" {
58-
return
56+
return ""
5957
}
6058
message := time.Unix(0, event.Fields["timestamp"].(int64)*int64(time.Nanosecond)).String() + "\t" + event.Fields["message_type"].(string) + "\t" + event.Msg + "\n"
6159
buf.WriteString(message)
6260
}
63-
s.SendToSumo(buf.String())
61+
62+
return buf.String()
63+
}
64+
65+
func (s *SumoLogicAppender) AppendLogs(queue eventQueue.Queue) {
66+
// the appender calls for the next message in the queue and parse it to a string
67+
if queue.GetCount() == s.eventsBatch { //whent the batch limit is met, call stringBuilder
68+
logMessage := StringBuilder(queue)
69+
s.SendToSumo(logMessage)
70+
}
71+
6472
}
6573

6674
func (s *SumoLogicAppender) SendToSumo(log string) {
Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
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 := StringBuilder(queue)
52+
assert.Equal(t, finalString, "2016-12-12 16:02:41.828366387 -0300 CLST"+"\t"+"OUT"+"\t"+"index [01]"+"\n"+
53+
"2016-12-12 16:02:42.844737993 -0300 CLST"+"\t"+"OUT"+"\t"+"index [02]"+"\n"+
54+
"2016-12-12 16:02:43.862436654 -0300 CLST"+"\t"+"OUT"+"\t"+"index [03]"+"\n", "")
55+
56+
}

0 commit comments

Comments
 (0)