diff --git a/commands/pilosa-nexter/nexter.go b/commands/pilosa-nexter/nexter.go index f2821c40e..25926f429 100644 --- a/commands/pilosa-nexter/nexter.go +++ b/commands/pilosa-nexter/nexter.go @@ -6,9 +6,11 @@ import ( "fmt" "log" "net/http" + "os" "strconv" "strings" + "github.com/cactus/go-statsd-client/statsd" "github.com/coreos/go-etcd/etcd" ) @@ -26,11 +28,15 @@ func (self *stringslice) Set(value string) error { var port uint var blocksize uint64 var etcd_nodes stringslice +var stats string +var logfile string func init() { flag.UintVar(&port, "port", 9000, "Port to run HTTP server on") flag.Uint64Var(&blocksize, "blocksize", 64, "Block size") flag.Var(&etcd_nodes, "etcd", "Etcd server") + flag.StringVar(&stats, "statsd", "127.0.0.1:8125", "Statsd server") + flag.StringVar(&logfile, "log", "/tmp/nexter.log", "Log file name") flag.Parse() } @@ -48,6 +54,7 @@ type DelReq struct { type Nexter struct { reqchan chan Req done chan bool + stats *statsd.Client } func (self *Nexter) countloop(ch chan uint64, id int, client *etcd.Client) { @@ -83,8 +90,11 @@ func (self *Nexter) countloop(ch chan uint64, id int, client *etcd.Client) { } _, err = client.CompareAndSwap(path, strconv.FormatUint(end, 10), 0, strconv.FormatUint(start, 10), 0) if err != nil { + log.Println("CAS failure", path, start, end) continue } + self.stats.Gauge("nexter."+strconv.Itoa(id), int64(end), 1.0) + log.Println("Allocated", id, start, end) for c := start; c < end; c += 1 { ch <- c } @@ -132,15 +142,27 @@ func (self *Nexter) Stop() { self.done <- true } -func NewNexter() *Nexter { - nexter := &Nexter{make(chan Req), make(chan bool)} +func NewNexter() (*Nexter, error) { + stats, err := statsd.New(stats, "") + if err != nil { + return nil, err + } + nexter := &Nexter{make(chan Req), make(chan bool), stats} go nexter.loop() - return nexter + return nexter, nil } func main() { + logf, err := os.OpenFile(logfile, os.O_RDWR|os.O_CREATE|os.O_APPEND, 0666) + if err != nil { + log.Println("Error opening file: %v", err) + } + log.SetOutput(logf) log.Println("Starting Nexter...") - nexter := NewNexter() + nexter, err := NewNexter() + if err != nil { + log.Fatal(err) + } http.HandleFunc("/nexter/", func(w http.ResponseWriter, r *http.Request) { url := r.URL.Path splits := strings.Split(url, "/") diff --git a/deps.json b/deps.json index db940d04e..feba75ef1 100644 --- a/deps.json +++ b/deps.json @@ -73,5 +73,11 @@ "repo": "github.com/gorilla/websocket", "version": "92334662baa9cbebc2e6e68b8d56bc1233f85a4c", "type": "git" + }, + "statsd": { + "repo": "github.com/cactus/go-statsd-client/statsd", + "version": "912f30c35e9cdf51f50bae24f071e227cca152fb", + "type": "git" } + }