Skip to content

Commit 7d7e07b

Browse files
committed
Merged in feature/SCFF-50 (pull request #19)
Feature/SCFF-50
2 parents 8ccb76b + 9c50213 commit 7d7e07b

2 files changed

Lines changed: 58 additions & 40 deletions

File tree

sumoCFFirehose/sumoLogicAppender.go

Lines changed: 15 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,6 @@ func newBuffer() SumoBuffer {
5858
}
5959

6060
func (s *SumoLogicAppender) Start() {
61-
//var wg sync.WaitGroup
6261
s.timerBetweenPost = time.Now()
6362
runtime.GOMAXPROCS(1)
6463
Buffer := newBuffer()
@@ -74,7 +73,7 @@ func (s *SumoLogicAppender) Start() {
7473
if time.Since(Buffer.timerIdlebuffer).Seconds() >= 10 && Buffer.logEventsInCurrentBuffer > 0 {
7574
logging.Info.Println("Sending current batch of logs after timer exceeded limit")
7675
//wg.Add(1)
77-
go s.SendToSumo(Buffer.logStringToSend.String() /*, &wg*/)
76+
go s.SendToSumo(Buffer.logStringToSend.String())
7877
Buffer = newBuffer()
7978
Buffer.timerIdlebuffer = time.Now()
8079
continue
@@ -90,8 +89,8 @@ func (s *SumoLogicAppender) Start() {
9089
s.AppendLogs(&Buffer)
9190
Buffer.timerIdlebuffer = time.Now()
9291
}
93-
//wg.Add(1)
94-
go s.SendToSumo(Buffer.logStringToSend.String() /*, &wg*/)
92+
93+
go s.SendToSumo(Buffer.logStringToSend.String())
9594
Buffer = newBuffer()
9695
} else {
9796
logging.Trace.Println("Pushing Logs to Buffer: ")
@@ -102,13 +101,12 @@ func (s *SumoLogicAppender) Start() {
102101
}
103102
}
104103
}
105-
//wg.Wait()
104+
106105
}
107106

108107
}
109108

110109
func StringBuilder(event *events.Event, verboseLogMessages bool) string {
111-
112110
eventType := event.Type
113111
var msg []byte
114112
switch eventType {
@@ -189,16 +187,12 @@ func (s *SumoLogicAppender) AppendLogs(buffer *SumoBuffer) {
189187

190188
}
191189

192-
func (s *SumoLogicAppender) SendToSumo(logStringToSend string /*, wg *sync.WaitGroup*/) {
190+
func (s *SumoLogicAppender) SendToSumo(logStringToSend string) {
193191
if logStringToSend != "" {
194192
var buf bytes.Buffer
195193
g := gzip.NewWriter(&buf)
196194
g.Write([]byte(logStringToSend))
197195
g.Close()
198-
for time.Since(s.timerBetweenPost) < s.sumoPostMinimumDelay {
199-
logging.Info. /*Trace*/ Println("Delaying post to honor minimum post delay")
200-
time.Sleep(100 * time.Millisecond)
201-
}
202196
request, err := http.NewRequest("POST", s.url, &buf)
203197
if err != nil {
204198
logging.Error.Printf("http.NewRequest() error: %v\n", err)
@@ -216,7 +210,11 @@ func (s *SumoLogicAppender) SendToSumo(logStringToSend string /*, wg *sync.WaitG
216210
if s.sumoCategory != "" {
217211
request.Header.Add("X-Sumo-Category", s.sumoCategory)
218212
}
219-
213+
//checking the timer before first POST intent
214+
for time.Since(s.timerBetweenPost) < s.sumoPostMinimumDelay {
215+
logging.Trace.Println("Delaying Post because minimum post timer not expired")
216+
time.Sleep(100 * time.Millisecond)
217+
}
220218
response, err := s.httpClient.Do(request)
221219

222220
if (err != nil) || (response.StatusCode != 200 && response.StatusCode != 302 && response.StatusCode < 500) {
@@ -226,7 +224,6 @@ func (s *SumoLogicAppender) SendToSumo(logStringToSend string /*, wg *sync.WaitG
226224
statusCode := 0
227225
err := Retry(func(attempt int) (bool, error) {
228226
var errRetry error
229-
//create again request
230227
request, err := http.NewRequest("POST", s.url, &buf)
231228
if err != nil {
232229
logging.Error.Printf("http.NewRequest() error: %v\n", err)
@@ -242,6 +239,11 @@ func (s *SumoLogicAppender) SendToSumo(logStringToSend string /*, wg *sync.WaitG
242239
if s.sumoCategory != "" {
243240
request.Header.Add("X-Sumo-Category", s.sumoCategory)
244241
}
242+
//checking the timer before POST (retry intent)
243+
for time.Since(s.timerBetweenPost) < s.sumoPostMinimumDelay {
244+
logging.Trace.Println("Delaying Post because minimum post timer not expired")
245+
time.Sleep(100 * time.Millisecond)
246+
}
245247
response, errRetry = s.httpClient.Do(request)
246248
if errRetry != nil {
247249
logging.Error.Printf("http.Do() error: %v\n", errRetry)
@@ -277,7 +279,6 @@ func (s *SumoLogicAppender) SendToSumo(logStringToSend string /*, wg *sync.WaitG
277279
if response != nil {
278280
defer response.Body.Close()
279281
}
280-
//wg.Done()
281282
}
282283

283284
}

sumoCFFirehose/sumoLogicAppender_test.go

Lines changed: 43 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -9,34 +9,51 @@ import (
99
. "bitbucket.org/mcplusa-ondemand/firehose-to-sumologic/events"
1010
)
1111

12-
func testAppenderStringBuilder(t *testing.T) {
12+
func TestAppenderStringBuilder(t *testing.T) {
1313
event1 := Event{
1414
Fields: map[string]interface{}{
15-
"timestamp": "1481569361828366387",
16-
"message_type": "OUT",
17-
"cf_app_id": "011",
15+
"deployment": "cf",
16+
"ip": "10.193.166.33",
17+
"job": "cloud_controller",
18+
"job_index": "c82feee9-2159-4b05-b669-a9929eb59017",
19+
"name": "requests.completed",
20+
"origin": "cc",
21+
"unit": "counter",
22+
"value": 558108,
1823
},
19-
Msg: "index [01]",
24+
Msg: "",
25+
Type: "ValueMetric",
2026
}
2127

2228
event2 := Event{
2329
Fields: map[string]interface{}{
24-
"timestamp": "1481569362844737993",
25-
"message_type": "OUT",
26-
"cf_app_id": "022",
30+
"delta": 9,
31+
"deployment": "cf-redis",
32+
"ip": "10.193.166.84",
33+
"job": "dedicated-node",
34+
"job_index": "8081eca4-9e27-49cb-83ce-948e703c0939",
35+
"name": "dropsondeMarshaller.sentEnvelopes",
36+
"origin": "MetronAgent",
37+
"total": 10249446,
2738
},
28-
Msg: "index [02]",
39+
Msg: "",
40+
Type: "CounterEvent",
2941
}
3042

3143
event3 := Event{
3244
Fields: map[string]interface{}{
33-
"timestamp": "1481569363862436654",
34-
"message_type": "OUT",
35-
"cf_app_id": "033",
45+
"delta": 582,
46+
"deployment": "cf-redis",
47+
"ip": "10.193.166.84",
48+
"job": "dedicated-node",
49+
"job_index": "23f9be01-bd83-4967-acba-69fc649f4ee6",
50+
"name": "dropsondeAgentListener.receivedByteCount",
51+
"origin": "MetronAgent",
52+
"total": 639557085,
3653
},
37-
Msg: "index [03]",
54+
Msg: "",
55+
Type: "CounterEvent",
3856
}
39-
4057
queue := Queue{
4158
Events: make([]*Event, 3),
4259
}
@@ -48,46 +65,46 @@ func testAppenderStringBuilder(t *testing.T) {
4865
for queue.GetCount() > 0 {
4966
finalString = finalString + StringBuilder(queue.Pop(), true)
5067
}
51-
assert.Equal(t, finalString, "2016-12-12 16:02:41.828366387 -0300 CLST"+"\t"+"OUT"+"\t"+"index [01]"+"\n"+
52-
"2016-12-12 16:02:42.844737993 -0300 CLST"+"\t"+"OUT"+"\t"+"index [02]"+"\n"+
53-
"2016-12-12 16:02:43.862436654 -0300 CLST"+"\t"+"OUT"+"\t"+"index [03]"+"\n", "")
68+
assert.Equal(t, finalString, "{\"Fields\":{\"deployment\":\"cf\",\"ip\":\"10.193.166.33\",\"job\":\"cloud_controller\",\"job_index\":\"c82feee9-2159-4b05-b669-a9929eb59017\",\"name\":\"requests.completed\",\"origin\":\"cc\",\"unit\":\"counter\",\"value\":558108},\"Msg\":\"\",\"Type\":\"ValueMetric\"}\n"+
69+
"{\"Fields\":{\"delta\":9,\"deployment\":\"cf-redis\",\"ip\":\"10.193.166.84\",\"job\":\"dedicated-node\",\"job_index\":\"8081eca4-9e27-49cb-83ce-948e703c0939\",\"name\":\"dropsondeMarshaller.sentEnvelopes\",\"origin\":\"MetronAgent\",\"total\":10249446},\"Msg\":\"\",\"Type\":\"CounterEvent\"}\n"+
70+
"{\"Fields\":{\"delta\":582,\"deployment\":\"cf-redis\",\"ip\":\"10.193.166.84\",\"job\":\"dedicated-node\",\"job_index\":\"23f9be01-bd83-4967-acba-69fc649f4ee6\",\"name\":\"dropsondeAgentListener.receivedByteCount\",\"origin\":\"MetronAgent\",\"total\":639557085},\"Msg\":\"\",\"Type\":\"CounterEvent\"}\n", "")
5471
}
5572

56-
func testStringBuilderVerboseLogsFalse(t *testing.T) {
73+
func TestStringBuilderVerboseLogsFalse(t *testing.T) {
5774
eventVerboseLogMessage := Event{
5875
Fields: map[string]interface{}{
5976
"message_type": "OUT",
60-
"source_instance": "0",
77+
"source_instance": 0,
6178
"deployment": "cf",
6279
"ip": "10.193.166.47",
6380
"job": "diego_cell",
6481
"job_index": "c62aebe5-16b8-43f5-a589-1267e09b9537",
6582
"cf_ignored_app": "false",
66-
"timestamp": "1483629662001580713",
83+
"timestamp": int64(1483629662001580713),
6784
"source_type": "APP",
6885
"origin": "rep",
6986
"cf_app_id": "7833dc75-4484-409c-9b74-90b6454906c6",
7087
},
7188
Msg: "Triggering 'app usage events fetcher'",
7289
Type: "LogMessage",
7390
}
74-
finalMessage := StringBuilder(&eventVerboseLogMessage, false)
7591

76-
assert.False(t, assert.Contains(t, finalMessage, "source_type", ""), "should be false")
92+
finalMessage := StringBuilder(&eventVerboseLogMessage, false)
93+
assert.NotContains(t, finalMessage, "source_type", "dsds")
7794

7895
}
7996

80-
func testStringBuilderVerboseLogsTrue(t *testing.T) {
97+
func TestStringBuilderVerboseLogsTrue(t *testing.T) {
8198
eventVerboseLogMessage := Event{
8299
Fields: map[string]interface{}{
83100
"message_type": "OUT",
84-
"source_instance": "0",
101+
"source_instance": 0,
85102
"deployment": "cf",
86103
"ip": "10.193.166.47",
87104
"job": "diego_cell",
88105
"job_index": "c62aebe5-16b8-43f5-a589-1267e09b9537",
89106
"cf_ignored_app": "false",
90-
"timestamp": "1483629662001580713",
107+
"timestamp": int64(1483629662001580713),
91108
"source_type": "APP",
92109
"origin": "rep",
93110
"cf_app_id": "7833dc75-4484-409c-9b74-90b6454906c6",
@@ -97,6 +114,6 @@ func testStringBuilderVerboseLogsTrue(t *testing.T) {
97114
}
98115
finalMessage := StringBuilder(&eventVerboseLogMessage, true)
99116

100-
assert.True(t, assert.Contains(t, finalMessage, "source_type", ""), "should be true")
117+
assert.Contains(t, finalMessage, "source_type", "dsds")
101118

102119
}

0 commit comments

Comments
 (0)