Skip to content

Commit 0c0a182

Browse files
committed
Merge branch 'feature/SCFF-46' into develop
2 parents 5caac17 + 4262408 commit 0c0a182

11 files changed

Lines changed: 242 additions & 279 deletions

File tree

.gitignore

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
product/
22
release/
33
firehose-to-sumologic.zip
4-
my.db
4+
event.db
55

caching/caching_boltdb.go

Lines changed: 3 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@ func NewCachingBolt(gcfClientSet *cfClient.Client, boltDatabasePath string) Cach
3434
}
3535

3636
func (c *CachingBolt) CreateBucket() {
37+
//start to write inside the db
3738
c.Appdb.Update(func(tx *bolt.Tx) error {
3839
_, err := tx.CreateBucketIfNotExists([]byte("AppBucket"))
3940
if err != nil {
@@ -65,7 +66,6 @@ func (c *CachingBolt) fillDatabase(listApps []App) {
6566
if err != nil {
6667
return fmt.Errorf("create bucket: %s", err)
6768
}
68-
6969
serialize, err := json.Marshal(app)
7070

7171
if err != nil {
@@ -89,6 +89,7 @@ func (c *CachingBolt) GetAppByGuid(appGuid string) []App {
8989
if err != nil {
9090
return apps
9191
}
92+
9293
apps = append(apps, App{
9394
app.Name,
9495
app.Guid,
@@ -98,26 +99,24 @@ func (c *CachingBolt) GetAppByGuid(appGuid string) []App {
9899
app.SpaceData.Entity.OrgData.Entity.Guid,
99100
c.isOptOut(app.Environment),
100101
})
102+
101103
c.fillDatabase(apps)
102104
return apps
103105

104106
}
105107

106108
func (c *CachingBolt) GetAllApp() []App {
107-
108109
var apps []App
109110

110111
defer func() {
111112
if r := recover(); r != nil {
112113
//logging.LogError("Recovered in caching.GetAllApp()", r)
113114
}
114115
}()
115-
116116
cfApps, err := c.GcfClient.ListApps()
117117
if err != nil {
118118
return apps
119119
}
120-
121120
for _, app := range cfApps {
122121
//fmt.Printf("App [%s] Found... \n", app.Name)
123122
apps = append(apps, App{
@@ -130,15 +129,13 @@ func (c *CachingBolt) GetAllApp() []App {
130129
c.isOptOut(app.Environment),
131130
})
132131
}
133-
134132
c.fillDatabase(apps)
135133
//fmt.Printf("Found [%d] Apps!", len(apps))
136134

137135
return apps
138136
}
139137

140138
func (c *CachingBolt) GetAppInfo(appGuid string) App {
141-
142139
var d []byte
143140
var app App
144141
c.Appdb.View(func(tx *bolt.Tx) error {

eventRouting/eventRouting_suite_test.go

Lines changed: 0 additions & 13 deletions
This file was deleted.

eventRouting/eventrouting.go

Lines changed: 5 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -5,13 +5,10 @@ import (
55
"sort"
66
"strings"
77
"sync"
8-
"time"
98

109
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/caching"
1110
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/eventQueue"
1211
fevents "bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/events"
13-
"bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/sumoCFFirehose"
14-
"github.com/Sirupsen/logrus"
1512
"github.com/cloudfoundry/sonde-go/events"
1613
)
1714

@@ -20,16 +17,14 @@ type EventRouting struct {
2017
selectedEvents map[string]bool
2118
selectedEventsCount map[string]uint64
2219
mutex *sync.Mutex
23-
sLAppender sumoCFFirehose.SumoLogicAppender //**
2420
queue *eventQueue.Queue
2521
}
2622

27-
func NewEventRouting(caching caching.Caching, sLAppender sumoCFFirehose.SumoLogicAppender, queue *eventQueue.Queue) *EventRouting {
23+
func NewEventRouting(caching caching.Caching, queue *eventQueue.Queue) *EventRouting {
2824
return &EventRouting{
2925
CachingClient: caching,
3026
selectedEvents: make(map[string]bool),
3127
selectedEventsCount: make(map[string]uint64),
32-
sLAppender: sLAppender, //**
3328
queue: queue,
3429
mutex: &sync.Mutex{},
3530
}
@@ -118,7 +113,7 @@ func GetListAuthorizedEventEvents() (authorizedEvents string) {
118113
return strings.Join(arrEvents, ", ")
119114
}
120115

121-
func (e *EventRouting) GetTotalCountOfSelectedEvents() uint64 {
116+
/*func (e *EventRouting) GetTotalCountOfSelectedEvents() uint64 {
122117
var total = uint64(0)
123118
for _, count := range e.GetSelectedEventsCount() {
124119
total += count
@@ -148,9 +143,9 @@ func (e *EventRouting) LogEventTotals(logTotalsTime time.Duration) {
148143
e.queue.Push(event)
149144
}
150145
}()
151-
}
146+
}*/
152147

153-
func (e *EventRouting) getEventTotals(totalElapsedTime float64, elapsedTime float64, lastCount uint64) (*fevents.Event, uint64) {
148+
/*func (e *EventRouting) getEventTotals(totalElapsedTime float64, elapsedTime float64, lastCount uint64) (*fevents.Event, uint64) {
154149
e.mutex.Lock()
155150
defer e.mutex.Unlock()
156151
totalCount := e.GetTotalCountOfSelectedEvents()
@@ -172,3 +167,4 @@ func (e *EventRouting) getEventTotals(totalElapsedTime float64, elapsedTime floa
172167
event.AnnotateWithMetaData(map[string]string{})
173168
return event, totalCount
174169
}
170+
*/

events/events_suite_test.go

Lines changed: 0 additions & 41 deletions
This file was deleted.

events/events_test.go

Lines changed: 0 additions & 71 deletions
This file was deleted.

main.go

Lines changed: 42 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -19,25 +19,25 @@ import (
1919
)
2020

2121
var (
22-
apiEndpoint = kingpin.Flag("api-endpoint", "Api Endpoint").OverrideDefaultFromEnvar("API_ENDPOINT").String()
23-
sumoEndpoint = kingpin.Flag("sumo-endpoint", "Sumo Endpoint").OverrideDefaultFromEnvar("SUMO_ENDPOINT").String()
24-
dopplerEndpoint = kingpin.Flag("doppler-endpoint", "Overwrite default doppler endpoint return by /v2/info").OverrideDefaultFromEnvar("DOPPLER_ENDPOINT").String()
25-
subscriptionId = kingpin.Flag("subscription-id", "Id for the subscription.").Default("firehose").OverrideDefaultFromEnvar("FIREHOSE_SUBSCRIPTION_ID").String()
26-
user = kingpin.Flag("cloudfoundry-user", "Cloudfoundry User").OverrideDefaultFromEnvar("CLOUDFOUNDRY_USER").String() //user created in CF, authorized to connect the firehose
27-
password = kingpin.Flag("cloudfoundry-password", "Cloudfoundry Password").OverrideDefaultFromEnvar("CLOUDFOUNDRY_PASSWORD").String() // password created along with the firehose_user //kingpin.Flag("skip-ssl-validation", "Please don't").Default("false").OverrideDefaultFromEnvar("SKIP_SSL_VALIDATION").Bool()
28-
keepAlive, errK = time.ParseDuration("25s") //default
29-
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()
30-
boltDatabasePath = "my.db" //default
31-
tickerTime, errT = time.ParseDuration("60s") //Default
22+
apiEndpoint = kingpin.Flag("api-endpoint", "CF API Endpoint").OverrideDefaultFromEnvar("API_ENDPOINT").String()
23+
sumoEndpoint = kingpin.Flag("sumo-endpoint", "Sumo Logic Endpoint").OverrideDefaultFromEnvar("SUMO_ENDPOINT").String()
24+
//dopplerEndpoint = kingpin.Flag("doppler-endpoint", "Overwrite default doppler endpoint return by /v2/info").OverrideDefaultFromEnvar("DOPPLER_ENDPOINT").String()
25+
subscriptionId = kingpin.Flag("subscription-id", "Cloud Foundry ID for the subscription.").Default("firehose").OverrideDefaultFromEnvar("FIREHOSE_SUBSCRIPTION_ID").String()
26+
user = kingpin.Flag("cloudfoundry-user", "Cloud Foundry User").OverrideDefaultFromEnvar("CLOUDFOUNDRY_USER").String() //user created in CF, authorized to connect the firehose
27+
password = kingpin.Flag("cloudfoundry-password", "Cloud Foundry Password").OverrideDefaultFromEnvar("CLOUDFOUNDRY_PASSWORD").String() // password created along with the firehose_user //kingpin.Flag("skip-ssl-validation", "Please don't").Default("false").OverrideDefaultFromEnvar("SKIP_SSL_VALIDATION").Bool()
28+
keepAlive, errK = time.ParseDuration("25s") //default Error, ContainerMetric, HttpStart, HttpStop, HttpStartStop, LogMessage, ValueMetric, CounterEvent
29+
wantedEvents = kingpin.Flag("events", fmt.Sprintf("Comma separated list of events you would like. Valid options are %s", eventRouting.GetListAuthorizedEventEvents())).Default("Error, ContainerMetric, HttpStart, HttpStop, HttpStartStop, LogMessage, ValueMetric, CounterEvent").OverrideDefaultFromEnvar("EVENTS").String()
30+
boltDatabasePath = "event.db"
31+
tickerTime = kingpin.Flag("nozzle-polling-period", "Nozzle Polling Period").Default("15s").OverrideDefaultFromEnvar("NOZZLE_POLLING_PERIOD").Duration()
3232
eventsBatchSize = kingpin.Flag("log-events-batch-size", "Log Events Batch Size").OverrideDefaultFromEnvar("LOG_EVENTS_BATCH_SIZE").Int()
33-
sumoPostMinimumDelay = kingpin.Flag("sumo-post-minimum-delay", "Sumo Post Minimum Delay").OverrideDefaultFromEnvar("SUMO_POST_MINIMUM_DELAY").Duration()
34-
sumoCategory = kingpin.Flag("sumo-category", "Sumo Category").Default("").OverrideDefaultFromEnvar("SUMO_CATEGORY").String()
35-
sumoName = kingpin.Flag("sumo-name", "Sumo Name").Default("").OverrideDefaultFromEnvar("SUMO_NAME").String()
36-
sumoClient = kingpin.Flag("sumo-client", "Sumo Client").Default("").OverrideDefaultFromEnvar("SUMO_CLIENT").String()
33+
sumoPostMinimumDelay = kingpin.Flag("sumo-post-minimum-delay", "Sumo Logic HTTP Post Minimum Delay").OverrideDefaultFromEnvar("SUMO_POST_MINIMUM_DELAY").Duration()
34+
sumoCategory = kingpin.Flag("sumo-category", "Sumo Logic Category").Default("").OverrideDefaultFromEnvar("SUMO_CATEGORY").String()
35+
sumoName = kingpin.Flag("sumo-name", "Sumo Logic Name").Default("").OverrideDefaultFromEnvar("SUMO_NAME").String()
36+
sumoHost = kingpin.Flag("sumo-host", "Sumo Logic Host").Default("").OverrideDefaultFromEnvar("SUMO_HOST").String()
3737
)
3838

3939
var (
40-
version = "0.0.0"
40+
version = "0.1"
4141
)
4242

4343
func main() {
@@ -49,15 +49,26 @@ func main() {
4949

5050
runtime.GOMAXPROCS(1)
5151

52-
logging.Info.Println("Configurations set:")
53-
logging.Info.Println("Api Endpoint: " + *apiEndpoint)
54-
logging.Info.Println("Sumo Endpoint: " + *sumoEndpoint)
55-
logging.Info.Println("Cloudfoundry Doppler Endpoint: " + *dopplerEndpoint)
56-
logging.Info.Println("Cloudfoundry Nozzle Subscription ID: " + *subscriptionId)
57-
logging.Info.Println("Cloudfoundry User: " + *user)
58-
59-
logging.Info.Printf("Events Batch Size: [%d]\n", *eventsBatchSize)
60-
logging.Info.Println("Starting firehose-to-sumo " + version)
52+
logging.Info.Println("Set Configurations:")
53+
logging.Info.Println("CF API Endpoint: t" + *apiEndpoint)
54+
logging.Info.Println("Sumo Logic Endpoint: " + *sumoEndpoint)
55+
//logging.Info.Println("Cloud foundry Doppler Endpoint: " + *dopplerEndpoint) //TODO
56+
logging.Info.Println("Cloud Foundry Nozzle Subscription ID: " + *subscriptionId)
57+
logging.Info.Println("Cloud Foundry User: " + *user)
58+
logging.Info.Println("Events Selected: " + *wantedEvents)
59+
logging.Info.Printf("Nozzle Polling Period: %v\n", *tickerTime)
60+
logging.Info.Printf("Log Events Batch Size: [%d]\n", *eventsBatchSize)
61+
logging.Info.Printf("Sumo Logic HTTP Post Minimum Delay: %v\n", *sumoPostMinimumDelay)
62+
if *sumoName != "" {
63+
logging.Info.Println("Sumo Logic Name: " + *sumoName)
64+
}
65+
if *sumoHost != "" {
66+
logging.Info.Println("Sumo Logic Host: " + *sumoHost)
67+
}
68+
if *sumoCategory != "" {
69+
logging.Info.Println("Sumo Logic Category: " + *sumoCategory)
70+
}
71+
logging.Info.Println("Starting Sumo Logic Nozzle " + version)
6172

6273
c := cfclient.Config{
6374
ApiAddress: *apiEndpoint,
@@ -67,10 +78,6 @@ func main() {
6778
}
6879
cfClient, _ := cfclient.NewClient(&c)
6980

70-
if len(*dopplerEndpoint) > 0 {
71-
cfClient.Endpoint.DopplerEndpoint = *dopplerEndpoint
72-
}
73-
7481
//Creating Caching
7582
var cachingClient caching.Caching
7683
if caching.IsNeeded(*wantedEvents) {
@@ -81,11 +88,11 @@ func main() {
8188

8289
logging.Info.Println("Creating queue")
8390
queue := eventQueue.NewQueue(make([]*events.Event, 100))
84-
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 1000, &queue, *eventsBatchSize, *sumoPostMinimumDelay, *sumoCategory, *sumoName, *sumoClient)
91+
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 5000, &queue, *eventsBatchSize, *sumoPostMinimumDelay, *sumoCategory, *sumoName, *sumoHost)
8592
go loggingClientSumo.Start() //multi
8693

8794
logging.Info.Println("Creating Events")
88-
events := eventRouting.NewEventRouting(cachingClient, *loggingClientSumo, &queue)
95+
events := eventRouting.NewEventRouting(cachingClient, &queue)
8996
err := events.SetupEventRouting(*wantedEvents)
9097
if err != nil {
9198
logging.Error.Fatal("Error setting up event routing: ", err)
@@ -100,8 +107,12 @@ func main() {
100107
apps := cachingClient.GetAllApp()
101108
logging.Info.Printf("Done filling cache! Found [%d] Apps \n", len(apps))
102109

110+
logging.Info.Println("Apps found: ")
111+
for i := 0; i < len(apps); i++ {
112+
logging.Info.Printf("[%d] "+apps[i].Name+" GUID: "+apps[i].Guid, i+1)
113+
}
103114
//Let's start the goRoutine
104-
cachingClient.PerformPoollingCaching(tickerTime)
115+
cachingClient.PerformPoollingCaching(*tickerTime)
105116

106117
firehoseConfig := &firehoseclient.FirehoseConfig{
107118
TrafficControllerURL: cfClient.Endpoint.DopplerEndpoint,

manifest.yml

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,3 +4,5 @@ env:
44
FIREHOSE_SUBSCRIPTION_ID: firehose-to-sumologic
55
LOG_EVENTS_BATCHSIZE: 200
66
SUMO_POST_MINIMUM_DELAY: 200ms
7+
F2S_DISABLE_LOGGING: true
8+

0 commit comments

Comments
 (0)