Skip to content

Commit 2c426fc

Browse files
committed
SCFF-35 Added timer for sumo post minimum and Default values
1 parent 87bc0c0 commit 2c426fc

2 files changed

Lines changed: 42 additions & 33 deletions

File tree

main.go

Lines changed: 15 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -19,19 +19,20 @@ import (
1919
)
2020

2121
var (
22-
debug = true //debug", "Enable debug mode, print in console
23-
apiEndpoint = kingpin.Flag("api-endpoint", "Sumo Endpoint").String() //"https://api.bosh-lite.com"
24-
sumoEndpoint = kingpin.Flag("sumo-endpoint", "Sumo Endpoint").String()
25-
dopplerEndpoint = kingpin.Flag("doppler-endpoint", "Overwrite default doppler endpoint return by /v2/info").OverrideDefaultFromEnvar("DOPPLER_ENDPOINT").String()
26-
subscriptionId = kingpin.Flag("subscription-id", "Id for the subscription.").Default("firehose").OverrideDefaultFromEnvar("FIREHOSE_SUBSCRIPTION_ID").String()
27-
user = "firehose_user" //user created in CF, authorized to connect the firehose
28-
password = "firehose_password" // password created along with the firehose_user
29-
skipSSLValidation = kingpin.Flag("skip-ssl-validation", "Please don't").Default("false").OverrideDefaultFromEnvar("SKIP_SSL_VALIDATION").Bool()
30-
keepAlive, errK = time.ParseDuration("25s") //default
31-
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()
32-
boltDatabasePath = "my.db" //default
33-
tickerTime, errT = time.ParseDuration("60s") //Default
34-
eventsBatchSize = kingpin.Flag("log-events-batch-size", "Log Events Batch Size").Int()
22+
debug = true //debug", "Enable debug mode, print in console
23+
apiEndpoint = kingpin.Flag("api-endpoint", "Sumo Endpoint").String() //"https://api.bosh-lite.com"
24+
sumoEndpoint = kingpin.Flag("sumo-endpoint", "Sumo Endpoint").String()
25+
dopplerEndpoint = kingpin.Flag("doppler-endpoint", "Overwrite default doppler endpoint return by /v2/info").OverrideDefaultFromEnvar("DOPPLER_ENDPOINT").String()
26+
subscriptionId = kingpin.Flag("subscription-id", "Id for the subscription.").Default("firehose").OverrideDefaultFromEnvar("FIREHOSE_SUBSCRIPTION_ID").String()
27+
user = "firehose_user" //user created in CF, authorized to connect the firehose
28+
password = "firehose_password" // password created along with the firehose_user
29+
skipSSLValidation = kingpin.Flag("skip-ssl-validation", "Please don't").Default("false").OverrideDefaultFromEnvar("SKIP_SSL_VALIDATION").Bool()
30+
keepAlive, errK = time.ParseDuration("25s") //default
31+
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()
32+
boltDatabasePath = "my.db" //default
33+
tickerTime, errT = time.ParseDuration("60s") //Default
34+
eventsBatchSize = kingpin.Flag("log-events-batch-size", "Log Events Batch Size").Default("10").Int()
35+
sumoPostMinimumDelay = kingpin.Flag("sumo-Post-Minimum-Delay", "Sumo Post Minimum Delay").Default("10s").Duration()
3536
)
3637

3738
var (
@@ -80,7 +81,7 @@ func main() {
8081

8182
logging.Info.Println("Creating queue")
8283
queue := eventQueue.NewQueue(make([]*events.Event, 100))
83-
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 1000, &queue, *eventsBatchSize)
84+
loggingClientSumo := sumoCFFirehose.NewSumoLogicAppender(*sumoEndpoint, 1000, &queue, *eventsBatchSize, *sumoPostMinimumDelay)
8485
go loggingClientSumo.Start() //multi
8586

8687
logging.Info.Println("Creating Events")

sumoCFFirehose/sumoLogicAppender.go

Lines changed: 27 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -18,9 +18,11 @@ type SumoLogicAppender struct {
1818
eventsBatchSize int
1919
logEventsInCurrentBuffer int
2020
logStringToSend *bytes.Buffer
21+
sumoPostMinimumDelay time.Duration
22+
timerPostMinimum time.Time
2123
}
2224

23-
func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQueue *eventQueue.Queue, eventsBatchSize int) *SumoLogicAppender {
25+
func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQueue *eventQueue.Queue, eventsBatchSize int, sumoPostMinimumDelay time.Duration) *SumoLogicAppender {
2426
return &SumoLogicAppender{
2527
url: urlValue,
2628
connectionTimeout: connectionTimeoutValue,
@@ -29,14 +31,16 @@ func NewSumoLogicAppender(urlValue string, connectionTimeoutValue int, nozzleQue
2931
eventsBatchSize: eventsBatchSize,
3032
logEventsInCurrentBuffer: 0,
3133
logStringToSend: bytes.NewBufferString(""),
34+
sumoPostMinimumDelay: sumoPostMinimumDelay,
3235
}
3336
}
3437

3538
func (s *SumoLogicAppender) Start() {
3639
timer := time.Now()
40+
s.timerPostMinimum = time.Now() //starting timer for sumo post minimum
3741
logging.Info.Println("Starting Appender Worker")
3842
for {
39-
time.Sleep(300 * time.Millisecond)
43+
time.Sleep(300 * time.Millisecond) //delay
4044
// while queue is not empty && s.eventsBatchSize not completed, queue.POP (appendLogs)
4145
for s.nozzleQueue.GetCount() != 0 && s.logEventsInCurrentBuffer <= s.eventsBatchSize {
4246
s.AppendLogs() //this method POP an event from queue
@@ -75,24 +79,28 @@ func (s *SumoLogicAppender) AppendLogs() {
7579
}
7680

7781
func (s *SumoLogicAppender) SendToSumo(log *bytes.Buffer) {
78-
logging.Trace.Println("Sending logs to Sumologic...")
79-
request, err := http.NewRequest("POST", s.url, log)
80-
if err != nil {
81-
logging.Error.Printf("http.NewRequest() error: %v\n", err)
82-
return
83-
}
84-
//request.Header.Add("content-type", "application/json")
85-
//request.SetBasicAuth("admin", "admin")
86-
response, err := s.httpClient.Do(request)
82+
//wait period between posts to Sumo
83+
if time.Since(s.timerPostMinimum) >= s.sumoPostMinimumDelay {
84+
logging.Trace.Println("Sending logs to Sumologic...")
85+
request, err := http.NewRequest("POST", s.url, log)
86+
if err != nil {
87+
logging.Error.Printf("http.NewRequest() error: %v\n", err)
88+
return
89+
}
90+
//request.Header.Add("content-type", "application/json")
91+
//request.SetBasicAuth("admin", "admin")
92+
response, err := s.httpClient.Do(request)
8793

88-
if err != nil {
89-
logging.Error.Printf("http.Do() error: %v\n", err)
90-
return
91-
} else {
92-
logging.Trace.Println("Do(Request) successful")
94+
if err != nil {
95+
logging.Error.Printf("http.Do() error: %v\n", err)
96+
return
97+
} else {
98+
logging.Trace.Println("Do(Request) successful")
99+
}
100+
s.logEventsInCurrentBuffer = 0 // reset counter
101+
s.logStringToSend = bytes.NewBufferString("") //reset String
102+
defer response.Body.Close()
103+
s.timerPostMinimum = time.Now() //reset timer post minimum
93104
}
94-
s.logEventsInCurrentBuffer = 0 // reset counter
95-
s.logStringToSend = bytes.NewBufferString("") //reset String
96-
defer response.Body.Close()
97105

98106
}

0 commit comments

Comments
 (0)