-
Notifications
You must be signed in to change notification settings - Fork 71
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
1 parent
2c141b9
commit b330039
Showing
5 changed files
with
192 additions
and
12 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,148 @@ | ||
package librato | ||
|
||
import ( | ||
"bytes" | ||
"encoding/json" | ||
"net/http" | ||
"strings" | ||
"sync" | ||
"time" | ||
|
||
"github.com/nyaruka/courier/utils" | ||
"github.com/sirupsen/logrus" | ||
) | ||
|
||
// Default is our default librato collector | ||
var Default *Sender | ||
|
||
// NewSender creates a new librato Sender with the passed in parameters | ||
func NewSender(waitGroup *sync.WaitGroup, username string, token string, source string, timeout time.Duration) *Sender { | ||
return &Sender{ | ||
waitGroup: waitGroup, | ||
stop: make(chan bool), | ||
|
||
buffer: make(chan gauge, 1000), | ||
username: username, | ||
token: token, | ||
source: source, | ||
timeout: timeout, | ||
} | ||
} | ||
|
||
// AddGauge can be used to add a new gauge to be sent to librato | ||
func (c *Sender) AddGauge(name string, value float64) { | ||
// if no librato configured, return | ||
if c == nil { | ||
return | ||
} | ||
|
||
// our buffer is full, log an error but continue | ||
if len(c.buffer) >= cap(c.buffer) { | ||
logrus.Error("unable to add new gauges, buffer full, you may want to increase your buffer size or decrease your timeout") | ||
return | ||
} | ||
|
||
c.buffer <- gauge{Name: strings.ToLower(name), Value: value, MeasureTime: time.Now().Unix()} | ||
} | ||
|
||
// Start starts our librato sender, callers can use Stop to stop it | ||
func (c *Sender) Start() { | ||
if c == nil { | ||
return | ||
} | ||
|
||
go func() { | ||
c.waitGroup.Add(1) | ||
defer c.waitGroup.Done() | ||
for { | ||
select { | ||
case <-c.stop: | ||
c.flush(1000) | ||
logrus.WithField("comp", "librato").Info("stopped") | ||
return | ||
|
||
case <-time.After(c.timeout * time.Second): | ||
c.flush(300) | ||
} | ||
} | ||
}() | ||
} | ||
|
||
func (c *Sender) flush(count int) { | ||
if len(c.buffer) <= 0 { | ||
return | ||
} | ||
|
||
// build our payload | ||
reqPayload := &payload{ | ||
MeasureTime: time.Now().Unix(), | ||
Source: c.source, | ||
Gauges: make([]gauge, 0, len(c.buffer)), | ||
} | ||
|
||
// read up to our count of gauges | ||
for i := 0; i < count; i++ { | ||
select { | ||
case g := <-c.buffer: | ||
reqPayload.Gauges = append(reqPayload.Gauges, g) | ||
default: | ||
break | ||
} | ||
} | ||
|
||
// send it off | ||
encoded, err := json.Marshal(reqPayload) | ||
if err != nil { | ||
logrus.WithField("comp", "librato").WithError(err).Error("error encoding librato metrics") | ||
return | ||
} | ||
|
||
req, err := http.NewRequest("POST", "https://metrics-api.librato.com/v1/metrics", bytes.NewReader(encoded)) | ||
if err != nil { | ||
logrus.WithField("comp", "librato").WithError(err).Error("error sending librato metrics") | ||
return | ||
} | ||
req.SetBasicAuth(c.username, c.token) | ||
req.Header.Set("Content-Type", "application/json") | ||
_, err = utils.MakeHTTPRequest(req) | ||
|
||
if err != nil { | ||
logrus.WithField("comp", "librato").WithError(err).Error("error sending librato metrics") | ||
return | ||
} | ||
|
||
logrus.WithField("comp", "librato").WithField("body", string(encoded)).WithField("count", len(reqPayload.Gauges)).Debug("flushed to librato") | ||
} | ||
|
||
// Stop stops our sender, callers can use the WaitGroup used during initialization to block for stop | ||
func (c *Sender) Stop() { | ||
if c == nil { | ||
return | ||
} | ||
close(c.stop) | ||
} | ||
|
||
type gauge struct { | ||
Name string `json:"name"` | ||
Value float64 `json:"value"` | ||
MeasureTime int64 `json:"measure_time"` | ||
} | ||
|
||
type payload struct { | ||
MeasureTime int64 `json:"measure_time"` | ||
Source string `json:"source"` | ||
Gauges []gauge `json:"gauges"` | ||
} | ||
|
||
// Sender is responsible for collecting gauges and sending them in batches to our librato server | ||
type Sender struct { | ||
waitGroup *sync.WaitGroup | ||
stop chan bool | ||
|
||
buffer chan gauge | ||
|
||
username string | ||
token string | ||
source string | ||
timeout time.Duration | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters