Skip to content

Commit 305563b

Browse files
committed
SCFF-34 added logging and using Events instead Node in queue
1 parent e8aef1b commit 305563b

7 files changed

Lines changed: 152 additions & 123 deletions

File tree

eventQueue/eventQueue.go

Lines changed: 25 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -2,65 +2,55 @@ package eventQueue
22

33
import . "bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/events"
44

5-
//Node to put in queue
6-
type Node struct {
7-
Event Event
8-
}
9-
105
// Queue is a basic FIFO queue based on a circular list that resizes as needed.
116
type Queue struct {
12-
Nodes []*Node
13-
head int
14-
tail int
15-
count int
16-
}
17-
18-
func NewNode(event Event) *Node {
19-
return &Node{
20-
Event: event,
21-
}
7+
Events []*Event
8+
head int
9+
tail int
10+
count int
2211
}
2312

24-
func NewQueue(n []*Node) Queue {
13+
func NewQueue(n []*Event) Queue {
2514
return Queue{
26-
Nodes: n,
15+
Events: n,
2716
}
2817
}
2918

30-
func (q *Queue) GetNode() []*Node {
31-
return q.Nodes
19+
func (q *Queue) GetNode() []*Event {
20+
return q.Events
3221
}
3322

34-
func (n *Queue) GetCount() int {
35-
return n.count
23+
func (q *Queue) GetCount() int {
24+
return q.count
3625
}
3726

38-
func (n *Node) GetNodeEvent() Event {
39-
return n.Event
40-
}
27+
/*
28+
func (q *Queue) GetEvents() Event {
29+
return q.Events
30+
}*/
4131

4232
// Push adds a node to the queue.
43-
func (q *Queue) Push(n *Node) {
33+
func (q *Queue) Push(n *Event) {
4434
if q.head == q.tail && q.count > 0 {
45-
nodes := make([]*Node, len(q.Nodes)*2)
46-
copy(nodes, q.Nodes[q.head:])
47-
copy(nodes[len(q.Nodes)-q.head:], q.Nodes[:q.head])
35+
events := make([]*Event, len(q.Events)*2)
36+
copy(events, q.Events[q.head:])
37+
copy(events[len(q.Events)-q.head:], q.Events[:q.head])
4838
q.head = 0
49-
q.tail = len(q.Nodes)
50-
q.Nodes = nodes
39+
q.tail = len(q.Events)
40+
q.Events = events
5141
}
52-
q.Nodes[q.tail] = n
53-
q.tail = (q.tail + 1) % len(q.Nodes)
42+
q.Events[q.tail] = n
43+
q.tail = (q.tail + 1) % len(q.Events)
5444
q.count++
5545
}
5646

5747
// Pop removes and returns a node from the queue in first to last order.
58-
func (q *Queue) Pop() *Node {
48+
func (q *Queue) Pop() *Event {
5949
if q.count == 0 {
6050
return nil
6151
}
62-
node := q.Nodes[q.head]
63-
q.head = (q.head + 1) % len(q.Nodes)
52+
node := q.Events[q.head]
53+
q.head = (q.head + 1) % len(q.Events)
6454
q.count--
6555
return node
6656
}

eventQueue/eventQueue_test.go

Lines changed: 24 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -9,43 +9,39 @@ import (
99

1010
func TestQueueFIFO(t *testing.T) {
1111
assert := assert.New(t)
12-
node1 := Node{
13-
Event: Event{
14-
Fields: map[string]interface{}{
15-
"message_type": "OUT",
16-
"cf_app_id": "011",
17-
},
18-
Msg: "index [01]",
12+
event1 := Event{
13+
Fields: map[string]interface{}{
14+
"message_type": "OUT",
15+
"cf_app_id": "011",
1916
},
17+
Msg: "index [01]",
2018
}
21-
node2 := Node{
22-
Event: Event{
23-
Fields: map[string]interface{}{
24-
"message_type": "OUT",
25-
"cf_app_id": "022",
26-
},
27-
Msg: "index [02]",
19+
20+
event2 := Event{
21+
Fields: map[string]interface{}{
22+
"message_type": "OUT",
23+
"cf_app_id": "022",
2824
},
25+
Msg: "index [02]",
2926
}
30-
node3 := Node{
31-
Event: Event{
32-
Fields: map[string]interface{}{
33-
"message_type": "OUT",
34-
"cf_app_id": "033",
35-
},
36-
Msg: "index [03]",
27+
28+
event3 := Event{
29+
Fields: map[string]interface{}{
30+
"message_type": "OUT",
31+
"cf_app_id": "033",
3732
},
33+
Msg: "index [03]",
3834
}
3935

4036
queue := Queue{
41-
Nodes: make([]*Node, 3),
37+
Events: make([]*Event, 3),
4238
}
4339

44-
queue.Push(&node1)
45-
queue.Push(&node2)
46-
queue.Push(&node3)
40+
queue.Push(&event1)
41+
queue.Push(&event2)
42+
queue.Push(&event3)
4743

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]", "")
44+
assert.Equal(queue.Pop().Msg, "index [01]", "")
45+
assert.Equal(queue.Pop().Msg, "index [02]", "")
46+
assert.Equal(queue.Pop().Msg, "index [03]", "")
5147
}

eventRouting/eventrouting.go

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import (
1010
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/caching"
1111
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/eventQueue"
1212
fevents "bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/events"
13+
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/logging"
1314
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/sumoCFFirehose"
1415
"github.com/Sirupsen/logrus"
1516
"github.com/cloudfoundry/sonde-go/events"
@@ -76,8 +77,8 @@ func (e *EventRouting) RouteEvent(msg *events.Envelope) {
7677
e.selectedEventsCount["ignored_app_message"]++
7778
} else {
7879
//Push the event to the queue
79-
fmt.Println("pushing event to queue")
80-
e.queue.Push(eventQueue.NewNode(*event))
80+
logging.Info.Println("pushing event to queue")
81+
e.queue.Push(event)
8182
e.selectedEventsCount[eventType.String()]++
8283

8384
}
@@ -147,7 +148,7 @@ func (e *EventRouting) LogEventTotals(logTotalsTime time.Duration) {
147148
count = lastCount
148149

149150
//Push the event to the queue
150-
e.queue.Push(eventQueue.NewNode(*event))
151+
e.queue.Push(event)
151152
}
152153
}()
153154
}

logging/logging.go

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,45 @@
1+
package logging
2+
3+
import (
4+
"io"
5+
"log"
6+
)
7+
8+
var (
9+
Trace *log.Logger
10+
Info *log.Logger
11+
Warning *log.Logger
12+
Error *log.Logger
13+
)
14+
15+
func Init(
16+
traceHandle io.Writer,
17+
infoHandle io.Writer,
18+
warningHandle io.Writer,
19+
errorHandle io.Writer) {
20+
21+
Trace = log.New(traceHandle,
22+
"TRACE: ",
23+
log.Ldate|log.Ltime|log.Lshortfile)
24+
25+
Info = log.New(infoHandle,
26+
"INFO: ",
27+
log.Ldate|log.Ltime|log.Lshortfile)
28+
29+
Warning = log.New(warningHandle,
30+
"WARNING: ",
31+
log.Ldate|log.Ltime|log.Lshortfile)
32+
33+
Error = log.New(errorHandle,
34+
"ERROR: ",
35+
log.Ldate|log.Ltime|log.Lshortfile)
36+
}
37+
38+
/*func main() {
39+
Init(ioutil.Discard, os.Stdout, os.Stdout, os.Stderr)
40+
41+
Trace.Println("I have something standard to say")
42+
Info.Println("Special Information")
43+
Warning.Println("There is something you need to know about")
44+
Error.Println("Something has failed")
45+
}*/

main.go

Lines changed: 17 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -2,16 +2,17 @@ package main
22

33
import (
44
"fmt"
5-
"log"
5+
"io/ioutil"
66
"os"
7-
"time"
8-
97
"runtime"
8+
"time"
109

1110
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/caching"
1211
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/eventQueue"
1312
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/eventRouting"
13+
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/events"
1414
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/firehoseclient"
15+
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/logging"
1516
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/sumoCFFirehose"
1617
"github.com/cloudfoundry-community/go-cfclient"
1718
"gopkg.in/alecthomas/kingpin.v2"
@@ -38,11 +39,15 @@ var (
3839
)
3940

4041
func main() {
42+
//logging init
43+
logging.Init(ioutil.Discard, os.Stdout, os.Stdout, os.Stderr)
44+
4145
kingpin.Version(version)
4246
kingpin.Parse()
47+
4348
runtime.GOMAXPROCS(1)
4449

45-
fmt.Printf("Starting firehose-to-sumo %s \n", version)
50+
logging.Info.Println("Starting firehose-to-sumo " + version)
4651

4752
c := cfclient.Config{
4853
ApiAddress: apiEndpoint,
@@ -65,25 +70,25 @@ func main() {
6570
}
6671

6772
//Creating queue
68-
queue := eventQueue.NewQueue(make([]*eventQueue.Node, 100))
73+
queue := eventQueue.NewQueue(make([]*events.Event, 100))
6974
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 1000, &queue, *eventsBatchSize)
7075
go loggingClientSumo.Start() //multi
7176

7277
//Creating Events
7378
events := eventRouting.NewEventRouting(cachingClient, *loggingClientSumo, &queue)
7479
err := events.SetupEventRouting(*wantedEvents)
7580
if err != nil {
76-
log.Fatal("Error setting up event routing: ", err)
81+
logging.Error.Fatal("Error setting up event routing: ", err)
7782
os.Exit(1)
7883

7984
}
8085

8186
// Parse extra fields from cmd call
8287
cachingClient.CreateBucket()
8388
//Let's Update the database the first time
84-
fmt.Printf("Start filling app/space/org cache.\n")
89+
logging.Info.Printf("Start filling app/space/org cache.\n")
8590
apps := cachingClient.GetAllApp()
86-
fmt.Printf("Done filling cache! Found [%d] Apps \n", len(apps))
91+
logging.Info.Printf("Done filling cache! Found [%d] Apps \n", len(apps))
8792

8893
//Let's start the goRoutine
8994
cachingClient.PerformPoollingCaching(tickerTime)
@@ -95,15 +100,12 @@ func main() {
95100
FirehoseSubscriptionID: *subscriptionId,
96101
}
97102

98-
if /*loggingClientSumo.Connect() ||*/ debug {
99-
100-
fmt.Printf("Connected to Server! Connecting to Firehose... \n")
103+
logging.Info.Printf("Connecting to Firehose... \n")
101104

102-
firehoseClient := firehoseclient.NewFirehoseNozzle(cfClient, events, firehoseConfig)
103-
go firehoseClient.Start()
105+
firehoseClient := firehoseclient.NewFirehoseNozzle(cfClient, events, firehoseConfig)
106+
go firehoseClient.Start()
104107

105-
defer firehoseClient.Start()
106-
}
108+
defer firehoseClient.Start()
107109

108110
defer cachingClient.Close()
109111

0 commit comments

Comments
 (0)