Skip to content

Commit 3f77cc0

Browse files
committed
Merged in feature/SCFF-34 (pull request #10)
SCFF-34 added logging and using Events instead Node in queue
2 parents e8aef1b + d867bc8 commit 3f77cc0

8 files changed

Lines changed: 167 additions & 133 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: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -76,10 +76,8 @@ func (e *EventRouting) RouteEvent(msg *events.Envelope) {
7676
e.selectedEventsCount["ignored_app_message"]++
7777
} else {
7878
//Push the event to the queue
79-
fmt.Println("pushing event to queue")
80-
e.queue.Push(eventQueue.NewNode(*event))
79+
e.queue.Push(event)
8180
e.selectedEventsCount[eventType.String()]++
82-
8381
}
8482
e.mutex.Unlock()
8583
}
@@ -147,7 +145,7 @@ func (e *EventRouting) LogEventTotals(logTotalsTime time.Duration) {
147145
count = lastCount
148146

149147
//Push the event to the queue
150-
e.queue.Push(eventQueue.NewNode(*event))
148+
e.queue.Push(event)
151149
}
152150
}()
153151
}

firehoseclient/firehoseclient.go

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2,10 +2,10 @@ package firehoseclient
22

33
import (
44
"crypto/tls"
5-
"fmt"
65
"time"
76

87
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/eventRouting"
8+
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/logging"
99
"github.com/cloudfoundry-community/go-cfclient"
1010
"github.com/cloudfoundry/noaa/consumer"
1111
"github.com/cloudfoundry/sonde-go/events"
@@ -43,11 +43,10 @@ func NewFirehoseNozzle(cfClient *cfclient.Client, eventRouting *eventRouting.Eve
4343
}
4444

4545
func (f *FirehoseNozzle) Start() error {
46-
fmt.Printf("Started the Nozzle... \n")
46+
logging.Info.Printf("Started the Nozzle... \n")
4747
f.consumeFirehose()
48-
fmt.Printf("consume the firehose... \n")
48+
logging.Info.Printf("consume the firehose... \n")
4949
err := f.routeEvent()
50-
fmt.Printf("route event... \n")
5150
return err
5251
}
5352

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: 31 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -2,24 +2,25 @@ 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"
1819
)
1920

2021
var (
21-
debug = true //debug", "Enable debug mode, print in console
22-
apiEndpoint = "https://api.bosh-lite.com"
22+
debug = true //debug", "Enable debug mode, print in console
23+
apiEndpoint = kingpin.Flag("api-endpoint", "Sumo Endpoint").String() //"https://api.bosh-lite.com"
2324
sumoEndpoint = kingpin.Flag("sumo-endpoint", "Sumo Endpoint").String()
2425
dopplerEndpoint = kingpin.Flag("doppler-endpoint", "Overwrite default doppler endpoint return by /v2/info").OverrideDefaultFromEnvar("DOPPLER_ENDPOINT").String()
2526
subscriptionId = kingpin.Flag("subscription-id", "Id for the subscription.").Default("firehose").OverrideDefaultFromEnvar("FIREHOSE_SUBSCRIPTION_ID").String()
@@ -38,14 +39,27 @@ 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("Configurations set:")
51+
logging.Info.Println("Api Endpoint: " + *apiEndpoint)
52+
logging.Info.Println("Sumo Endpoint: " + *sumoEndpoint)
53+
logging.Info.Println("Cloudfoundry Doppler Endpoint: " + *dopplerEndpoint)
54+
logging.Info.Println("Cloudfoundry Nozzle Subscription ID: " + *subscriptionId)
55+
logging.Info.Println("Cloudfoundry User: " + user)
56+
57+
logging.Info.Printf("Events Batch Size: [%d]\n", eventsBatchSize)
58+
59+
logging.Info.Println("Starting firehose-to-sumo " + version)
4660

4761
c := cfclient.Config{
48-
ApiAddress: apiEndpoint,
62+
ApiAddress: *apiEndpoint,
4963
Username: user,
5064
Password: password,
5165
SkipSslValidation: *skipSSLValidation,
@@ -64,26 +78,26 @@ func main() {
6478
cachingClient = caching.NewCachingEmpty()
6579
}
6680

67-
//Creating queue
68-
queue := eventQueue.NewQueue(make([]*eventQueue.Node, 100))
81+
logging.Info.Println("Creating queue")
82+
queue := eventQueue.NewQueue(make([]*events.Event, 100))
6983
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 1000, &queue, *eventsBatchSize)
7084
go loggingClientSumo.Start() //multi
7185

72-
//Creating Events
86+
logging.Info.Println("Creating Events")
7387
events := eventRouting.NewEventRouting(cachingClient, *loggingClientSumo, &queue)
7488
err := events.SetupEventRouting(*wantedEvents)
7589
if err != nil {
76-
log.Fatal("Error setting up event routing: ", err)
90+
logging.Error.Fatal("Error setting up event routing: ", err)
7791
os.Exit(1)
7892

7993
}
8094

8195
// Parse extra fields from cmd call
8296
cachingClient.CreateBucket()
8397
//Let's Update the database the first time
84-
fmt.Printf("Start filling app/space/org cache.\n")
98+
logging.Info.Printf("Start filling app/space/org cache.\n")
8599
apps := cachingClient.GetAllApp()
86-
fmt.Printf("Done filling cache! Found [%d] Apps \n", len(apps))
100+
logging.Info.Printf("Done filling cache! Found [%d] Apps \n", len(apps))
87101

88102
//Let's start the goRoutine
89103
cachingClient.PerformPoollingCaching(tickerTime)
@@ -95,15 +109,12 @@ func main() {
95109
FirehoseSubscriptionID: *subscriptionId,
96110
}
97111

98-
if /*loggingClientSumo.Connect() ||*/ debug {
112+
logging.Info.Printf("Connecting to Firehose... \n")
99113

100-
fmt.Printf("Connected to Server! Connecting to Firehose... \n")
114+
firehoseClient := firehoseclient.NewFirehoseNozzle(cfClient, events, firehoseConfig)
115+
go firehoseClient.Start()
101116

102-
firehoseClient := firehoseclient.NewFirehoseNozzle(cfClient, events, firehoseConfig)
103-
go firehoseClient.Start()
104-
105-
defer firehoseClient.Start()
106-
}
117+
defer firehoseClient.Start()
107118

108119
defer cachingClient.Close()
109120

0 commit comments

Comments
 (0)