Skip to content

Commit 314ee71

Browse files
committed
SCFF-46 logging app name list, fixed error in HTTP post timeout, check not send empty post
1 parent 5caac17 commit 314ee71

2 files changed

Lines changed: 111 additions & 106 deletions

File tree

main.go

Lines changed: 32 additions & 28 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
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 = "my.db" //TODO remove once database code is removed
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+
sumoClient = kingpin.Flag("sumo-client", "Sumo Logic Client").Default("").OverrideDefaultFromEnvar("SUMO_CLIENT").String()
3737
)
3838

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

4343
func main() {
@@ -49,15 +49,15 @@ 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)
52+
logging.Info.Println("Set Configurations:")
53+
logging.Info.Println("CF API Endpoint: " + *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)
5858

59-
logging.Info.Printf("Events Batch Size: [%d]\n", *eventsBatchSize)
60-
logging.Info.Println("Starting firehose-to-sumo " + version)
59+
logging.Info.Printf("Log Events Batch Size: [%d]\n", *eventsBatchSize)
60+
logging.Info.Println("Starting Sumo Logic Nozzle " + version)
6161

6262
c := cfclient.Config{
6363
ApiAddress: *apiEndpoint,
@@ -67,9 +67,9 @@ func main() {
6767
}
6868
cfClient, _ := cfclient.NewClient(&c)
6969

70-
if len(*dopplerEndpoint) > 0 {
70+
/*if len(*dopplerEndpoint) > 0 {
7171
cfClient.Endpoint.DopplerEndpoint = *dopplerEndpoint
72-
}
72+
}*/ //TODO
7373

7474
//Creating Caching
7575
var cachingClient caching.Caching
@@ -81,7 +81,7 @@ func main() {
8181

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

8787
logging.Info.Println("Creating Events")
@@ -99,9 +99,13 @@ func main() {
9999
logging.Info.Printf("Start filling app/space/org cache.\n")
100100
apps := cachingClient.GetAllApp()
101101
logging.Info.Printf("Done filling cache! Found [%d] Apps \n", len(apps))
102-
102+
//TODO log out each apps name
103+
logging.Info.Println("Apps founded: ")
104+
for i := 0; i < len(apps); i++ {
105+
logging.Info.Printf("[%d] "+apps[i].Name, i+1)
106+
}
103107
//Let's start the goRoutine
104-
cachingClient.PerformPoollingCaching(tickerTime)
108+
cachingClient.PerformPoollingCaching(*tickerTime)
105109

106110
firehoseConfig := &firehoseclient.FirehoseConfig{
107111
TrafficControllerURL: cfClient.Endpoint.DopplerEndpoint,

sumoCFFirehose/sumoLogicAppender.go

Lines changed: 79 additions & 78 deletions
Original file line numberDiff line numberDiff line change
@@ -120,92 +120,93 @@ func (s *SumoLogicAppender) AppendLogs(buffer *SumoBuffer) {
120120
}
121121

122122
func (s *SumoLogicAppender) SendToSumo(logStringToSend string) {
123-
var buf bytes.Buffer
124-
g := gzip.NewWriter(&buf)
125-
g.Write([]byte(logStringToSend))
126-
g.Close()
127-
128-
for time.Since(s.timerBetweenPost) < s.sumoPostMinimumDelay {
129-
logging.Trace.Println("Delaying post to honor minimum post delay")
130-
time.Sleep(100 * time.Millisecond)
131-
}
123+
if logStringToSend != "" {
124+
var buf bytes.Buffer
125+
g := gzip.NewWriter(&buf)
126+
g.Write([]byte(logStringToSend))
127+
g.Close()
128+
for time.Since(s.timerBetweenPost) < s.sumoPostMinimumDelay {
129+
logging.Trace.Println("Delaying post to honor minimum post delay")
130+
time.Sleep(100 * time.Millisecond)
131+
}
132132

133-
request, err := http.NewRequest("POST", s.url, &buf)
134-
if err != nil {
135-
logging.Error.Printf("http.NewRequest() error: %v\n", err)
136-
return
137-
}
133+
request, err := http.NewRequest("POST", s.url, &buf)
134+
if err != nil {
135+
logging.Error.Printf("http.NewRequest() error: %v\n", err)
136+
return
137+
}
138138

139-
request.Header.Add("Content-Encoding", "gzip")
139+
request.Header.Add("Content-Encoding", "gzip")
140140

141-
if s.sumoName != "" {
142-
request.Header.Add("X-Sumo-Name", s.sumoName)
143-
}
144-
if s.sumoClient != "" {
145-
request.Header.Add("X-Sumo-Client", s.sumoClient)
146-
}
147-
if s.sumoCategory != "" {
148-
request.Header.Add("X-Sumo-Category", s.sumoCategory)
149-
}
141+
if s.sumoName != "" {
142+
request.Header.Add("X-Sumo-Name", s.sumoName)
143+
}
144+
if s.sumoClient != "" {
145+
request.Header.Add("X-Sumo-Client", s.sumoClient)
146+
}
147+
if s.sumoCategory != "" {
148+
request.Header.Add("X-Sumo-Category", s.sumoCategory)
149+
}
150150

151-
response, err := s.httpClient.Do(request)
152-
153-
if (err != nil) || (response.StatusCode != 200 && response.StatusCode != 302 && response.StatusCode < 500) {
154-
logging.Info.Println("Endpoint dropped the post send")
155-
logging.Info.Println("Waiting for 300 ms to retry")
156-
time.Sleep(300 * time.Millisecond)
157-
statusCode := 0
158-
err := Retry(func(attempt int) (bool, error) {
159-
var errRetry error
160-
//create again request
161-
request, err := http.NewRequest("POST", s.url, &buf)
162-
if err != nil {
163-
logging.Error.Printf("http.NewRequest() error: %v\n", err)
164-
}
165-
request.Header.Add("Content-Encoding", "gzip")
151+
response, err := s.httpClient.Do(request)
166152

167-
if s.sumoName != "" {
168-
request.Header.Add("X-Sumo-Name", s.sumoName)
169-
}
170-
if s.sumoClient != "" {
171-
request.Header.Add("X-Sumo-Client", s.sumoClient)
172-
}
173-
if s.sumoCategory != "" {
174-
request.Header.Add("X-Sumo-Category", s.sumoCategory)
175-
}
176-
response, errRetry = s.httpClient.Do(request)
177-
if errRetry != nil {
178-
logging.Error.Printf("http.Do() error: %v\n", errRetry)
179-
logging.Info.Println("Waiting for 300 ms to retry after error")
180-
time.Sleep(300 * time.Millisecond)
181-
return attempt < 5, errRetry
182-
} else if response.StatusCode != 200 && response.StatusCode != 302 && response.StatusCode < 500 {
183-
logging.Info.Println("Endpoint dropped the post send again")
184-
logging.Info.Println("Waiting for 300 ms to retry after a retry ...")
185-
statusCode = response.StatusCode
186-
time.Sleep(300 * time.Millisecond)
153+
if (err != nil) || (response.StatusCode != 200 && response.StatusCode != 302 && response.StatusCode < 500) {
154+
logging.Info.Println("Endpoint dropped the post send")
155+
logging.Info.Println("Waiting for 300 ms to retry")
156+
time.Sleep(300 * time.Millisecond)
157+
statusCode := 0
158+
err := Retry(func(attempt int) (bool, error) {
159+
var errRetry error
160+
//create again request
161+
request, err := http.NewRequest("POST", s.url, &buf)
162+
if err != nil {
163+
logging.Error.Printf("http.NewRequest() error: %v\n", err)
164+
}
165+
request.Header.Add("Content-Encoding", "gzip")
166+
167+
if s.sumoName != "" {
168+
request.Header.Add("X-Sumo-Name", s.sumoName)
169+
}
170+
if s.sumoClient != "" {
171+
request.Header.Add("X-Sumo-Client", s.sumoClient)
172+
}
173+
if s.sumoCategory != "" {
174+
request.Header.Add("X-Sumo-Category", s.sumoCategory)
175+
}
176+
response, errRetry = s.httpClient.Do(request)
177+
if errRetry != nil {
178+
logging.Error.Printf("http.Do() error: %v\n", errRetry)
179+
logging.Info.Println("Waiting for 300 ms to retry after error")
180+
time.Sleep(300 * time.Millisecond)
181+
return attempt < 5, errRetry
182+
} else if response.StatusCode != 200 && response.StatusCode != 302 && response.StatusCode < 500 {
183+
logging.Info.Println("Endpoint dropped the post send again")
184+
logging.Info.Println("Waiting for 300 ms to retry after a retry ...")
185+
statusCode = response.StatusCode
186+
time.Sleep(300 * time.Millisecond)
187+
return attempt < 5, errRetry
188+
} else if response.StatusCode == 200 {
189+
logging.Info. /*Trace*/ Println("Post of logs successful after retry...")
190+
s.timerBetweenPost = time.Now()
191+
statusCode = response.StatusCode
192+
return true, err
193+
}
187194
return attempt < 5, errRetry
188-
} else if response.StatusCode == 200 {
189-
logging.Info. /*Trace*/ Println("Post of logs successful after retry...")
190-
s.timerBetweenPost = time.Now()
191-
statusCode = response.StatusCode
192-
return true, err
195+
})
196+
if err != nil {
197+
logging.Error.Println("Error, Not able to post after retry")
198+
logging.Error.Printf("http.Do() error: %v\n", err)
199+
return
200+
} else if statusCode != 200 {
201+
logging.Error.Printf("Not able to post after retry, with status code: %d", statusCode)
193202
}
194-
return attempt < 5, errRetry
195-
})
196-
if err != nil {
197-
logging.Error.Println("Error, Not able to post after retry")
198-
logging.Error.Printf("http.Do() error: %v\n", err)
199-
return
200-
} else if statusCode != 200 {
201-
logging.Error.Printf("Not able to post after retry, with status code: %d", statusCode)
203+
} else if response.StatusCode == 200 {
204+
logging.Info. /*Trace*/ Println("Post of logs successful")
205+
s.timerBetweenPost = time.Now()
206+
}
207+
if response != nil {
208+
defer response.Body.Close()
202209
}
203-
} else if response.StatusCode == 200 {
204-
logging.Info. /*Trace*/ Println("Post of logs successful")
205-
s.timerBetweenPost = time.Now()
206-
}
207-
if response != nil {
208-
defer response.Body.Close()
209210
}
210211

211212
}

0 commit comments

Comments
 (0)