Skip to content

Commit ec11ea1

Browse files
committed
SCFF-35 created the two timers regarding the buffer content and the post to sumo
1 parent 9ea8d7f commit ec11ea1

1 file changed

Lines changed: 17 additions & 17 deletions

File tree

sumoCFFirehose/sumoLogicAppender.go

Lines changed: 17 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -19,12 +19,13 @@ type SumoLogicAppender struct {
1919
nozzleQueue *eventQueue.Queue
2020
eventsBatchSize int
2121
sumoPostMinimumDelay time.Duration
22+
timerBetweenPost time.Time
2223
}
2324

2425
type SumoBuffer struct {
2526
logStringToSend *bytes.Buffer
2627
logEventsInCurrentBuffer int
27-
timerPostMinimum time.Time
28+
timerIdlebuffer time.Time
2829
channelMessage chan string
2930
}
3031

@@ -43,21 +44,21 @@ func newBuffer() SumoBuffer {
4344
return SumoBuffer{
4445
logStringToSend: bytes.NewBufferString(""),
4546
logEventsInCurrentBuffer: 0,
46-
timerPostMinimum: time.Now(),
4747
channelMessage: make(chan string),
4848
}
4949
}
5050

5151
func (s *SumoLogicAppender) Start() {
52+
s.timerBetweenPost = time.Now()
5253
runtime.GOMAXPROCS(1)
53-
timer := time.Now()
54-
Buffer := newBuffer() //creating Buffer
54+
Buffer := newBuffer() //creating Buffer
55+
Buffer.timerIdlebuffer = time.Now() //starting timerPostMinimum
5556
msgFromChannel := ""
5657
logging.Info.Println("Starting Appender Worker")
5758
for {
5859
fmt.Println("first for")
59-
time.Sleep(300 * time.Millisecond) //delay
60-
if Buffer.logEventsInCurrentBuffer >= s.eventsBatchSize || msgFromChannel == "Buffer in use" { //if buffer is full, create a new one
60+
time.Sleep(300 * time.Millisecond) //delay
61+
if Buffer.logEventsInCurrentBuffer >= s.eventsBatchSize || msgFromChannel == "Buffer being sent" { //if buffer is full, create a new one
6162
fmt.Println("Creating new Buffer")
6263
fmt.Println(Buffer.logEventsInCurrentBuffer)
6364
fmt.Println(msgFromChannel)
@@ -68,20 +69,23 @@ func (s *SumoLogicAppender) Start() {
6869
for s.nozzleQueue.GetCount() != 0 && Buffer.logEventsInCurrentBuffer < s.eventsBatchSize {
6970
fmt.Println("second for")
7071
s.AppendLogs(&Buffer) //this method POP an event from queue to Buffer
71-
timer = time.Now() //reset timer
72+
Buffer.timerIdlebuffer = time.Now() //reset buffer timer, buffer has new content
7273
if Buffer.logEventsInCurrentBuffer == s.eventsBatchSize { //if buffer is full, send logs to sumo
7374
logging.Info.Println("Batch Size complete")
7475
break
75-
} else if time.Since(timer).Seconds() >= 10 { // else if timer is up, send existing logs to sumo
76+
} else if time.Since(Buffer.timerIdlebuffer).Seconds() >= 10 { // else if buffer timer is up, send existing logs to sumo
7677
logging.Info.Println("Sending current batch of logs after timer exceeded limit")
7778
break
7879
}
7980
}
80-
//if batch size is met, send to sumo
81+
//wait period between posts to Sumo
82+
for time.Since(s.timerBetweenPost) < s.sumoPostMinimumDelay {
83+
time.Sleep(30 * time.Millisecond) // wait to retry
84+
}
8185
go s.SendToSumo(&Buffer)
8286
msgFromChannel = <-Buffer.channelMessage
8387

84-
timer = time.Now() //reset timer
88+
Buffer.timerIdlebuffer = time.Now() //reset Buffer timer
8589
}
8690

8791
}
@@ -109,13 +113,10 @@ func (s *SumoLogicAppender) AppendLogs(buffer *SumoBuffer) {
109113
}
110114

111115
func (s *SumoLogicAppender) SendToSumo(buffer *SumoBuffer) {
112-
buffer.channelMessage <- "Buffer in use"
113-
//wait period between posts to Sumo
116+
buffer.channelMessage <- "Buffer being sent"
114117
fmt.Println(buffer.logStringToSend.String())
115118
fmt.Println("......................")
116-
for time.Since(buffer.timerPostMinimum) < s.sumoPostMinimumDelay {
117-
time.Sleep(30 * time.Millisecond) // wait to retry
118-
}
119+
119120
logging.Trace.Println("Sending logs to Sumologic...")
120121
request, err := http.NewRequest("POST", s.url, buffer.logStringToSend)
121122
if err != nil {
@@ -130,9 +131,8 @@ func (s *SumoLogicAppender) SendToSumo(buffer *SumoBuffer) {
130131
return
131132
} else {
132133
logging.Trace.Println("Do(Request) successful")
133-
buffer.timerPostMinimum = time.Now() //reset timer post minimum
134+
s.timerBetweenPost = time.Now() //reset timer post minimum
134135
}
135136

136137
defer response.Body.Close()
137-
buffer.channelMessage <- "Buffer finished"
138138
}

0 commit comments

Comments
 (0)