diff --git a/commands/pilosa-cruncher/cruncher.go b/commands/pilosa-cruncher/cruncher.go index 3158f31ce..deab217f3 100644 --- a/commands/pilosa-cruncher/cruncher.go +++ b/commands/pilosa-cruncher/cruncher.go @@ -2,8 +2,8 @@ package main import ( "flag" + log "github.com/cihub/seelog" "github.com/mitchellh/panicwrap" - "log" "os" "pilosa/core" "pilosa/cruncher" @@ -32,19 +32,20 @@ func main() { if *cpuprofile != "" { f, err := os.Create(*cpuprofile) if err != nil { - log.Fatal(err) + log.Warn(err) + os.Exit(1) } pprof.StartCPUProfile(f) defer pprof.StopCPUProfile() } cruncher := cruncher.NewCruncher() cruncher.Run() - log.Println("STOP") + log.Warn("STOP") } func panicHandler(output string) { // output contains the full output (including stack traces) of the // panic. Put it in a file or something. - log.Printf("The child panicked:\n\n%s\n", output) + log.Warn("The child panicked:\n\n", output) os.Exit(1) } diff --git a/config/config.go b/config/config.go index fec9eb0b1..d0ba69e0d 100644 --- a/config/config.go +++ b/config/config.go @@ -2,8 +2,8 @@ package config import ( "errors" + log "github.com/cihub/seelog" "io/ioutil" - "log" "os" "sync" @@ -77,7 +77,7 @@ func (self *Config) load() error { if config_file == "" { config_file = os.Getenv("PILOSA_CONFIG") if config_file == "" { - log.Println("PILOSA_CONFIG not set, defaulting to pilosa.yaml") + log.Warn("PILOSA_CONFIG not set, defaulting to pilosa.yaml") config_file = "pilosa.yaml" } } diff --git a/core/etcd.go b/core/etcd.go index a5ac92386..302510ac6 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -3,7 +3,7 @@ package core import ( "errors" "fmt" - "log" + log "github.com/cihub/seelog" "pilosa/config" "pilosa/db" "pilosa/util" @@ -23,7 +23,7 @@ type TopologyMapper struct { } func (self *TopologyMapper) Setup() { - log.Println(self.namespace + "/db") + log.Warn(self.namespace + "/db") db_path := self.namespace + "/db" resp, err := self.service.Etcd.Get(db_path, false, true) if err != nil { @@ -31,17 +31,17 @@ func (self *TopologyMapper) Setup() { if ok && ee.ErrorCode == 100 { // node does not exist resp, err = self.service.Etcd.CreateDir(db_path, 0) if err != nil { - log.Fatal(err) + log.Warn(err) } } else { - log.Fatal(err) + log.Warn(err) } } //need to lock the world for _, node := range flatten(resp.Node) { err := self.handlenode(node) if err != nil { - log.Println(err) + log.Warn(err) } } @@ -55,9 +55,9 @@ func (self *TopologyMapper) Run() { // TODO: use modindex to make sure watch catches everything for { ns := self.namespace + "/db" - log.Println(" ETCD watcher:", ns) + log.Warn(" ETCD watcher:", ns) resp, err := self.service.Etcd.Watch(ns, 0, true, receiver, stop) - log.Println("TopologyMapper ETCD watcher", resp, err) + log.Warn("TopologyMapper ETCD watcher", resp, err) } }() go func() { @@ -139,7 +139,7 @@ func (self *TopologyMapper) MakeFragments(db string, slice_int int) error { for _, frame := range frames_to_create { err := self.AllocateFragment(p.Key, db, frame, slice_int) if err != nil { - log.Println(err) + log.Warn(err) } } } @@ -156,13 +156,13 @@ func (self *TopologyMapper) AllocateFragment(process_guid, db, frame string, sli fuid := util.SUUID_to_Hex(util.Id()) fragment_key := fmt.Sprintf("%s/db/%s/frame/%s/slice/%d/fragment/%s/process", self.namespace, db, frame, slice_int, fuid) // need to check value to see how many we have left - log.Println("ALLOC:", process_guid, len(process_guid)) + log.Warn("ALLOC:", process_guid, len(process_guid)) if len(process_guid) > 1 { _, err := self.service.Etcd.Set(fragment_key, process_guid, 0) if err != nil { return err } - log.Printf("Fragment sent to etcd: %s(%s)", fragment_key, process_guid) + log.Warn("Fragment sent to etcd:", fragment_key, process_guid) } return nil @@ -189,7 +189,7 @@ func (self *TopologyMapper) handlenode(node *etcd.Node) error { process_uuid, err = util.ParseGUID(node.Value) if err != nil { - log.Println("Bad Process Guid", key) + log.Warn("Bad Process Guid", key) return errors.New("No Process Id") } } else { @@ -238,7 +238,7 @@ func (self *TopologyMapper) handlenode(node *etcd.Node) error { } if err != nil { - log.Println("Bad UUID:", process_uuid, key) + log.Warn("Bad UUID:", process_uuid, key) return err } process = db.NewProcess(&process_uuid) @@ -252,7 +252,7 @@ func (self *TopologyMapper) handlenode(node *etcd.Node) error { return err } func (self *TopologyMapper) remove_fragment(node *etcd.Node) error { - log.Println(" hot remove_fragment (Not Supported yet):", node) + log.Warn(" hot remove_fragment (Not Supported yet):", node) /* key := node.Key[len(self.namespace)+1:] bits := strings.Split(key, "/") @@ -395,7 +395,7 @@ func (self *ProcessMapper) getnode(u *util.GUID) *Node { func crash_on_error(err error) { if err != nil { - log.Fatal(err) + log.Warn(err) } } @@ -437,8 +437,8 @@ func (self *ProcessMapper) Run() { path := self.namespace + "/process" self_path := path + "/" + id_string - log.Println("Writing configuration to etcd...") - log.Println(self_path) + log.Warn("Writing configuration to etcd...") + log.Warn(self_path) var err error _, err = self.service.Etcd.Set(self_path+"/port_tcp", strconv.Itoa(config.GetInt("port_tcp")), 0) @@ -453,7 +453,7 @@ func (self *ProcessMapper) Run() { err := self.handlenode(node) if err != nil { out := spew.Sdump(node) - log.Println(err, out) + log.Warn(err, out) } } diff --git a/core/http.go b/core/http.go index cea31c3c5..396c760c8 100644 --- a/core/http.go +++ b/core/http.go @@ -8,7 +8,8 @@ import ( "encoding/json" "fmt" "io/ioutil" - "log" + //"log" + log "github.com/cihub/seelog" "net/http" "net/http/httputil" "pilosa/config" @@ -65,7 +66,7 @@ func (self *RequestLogger) Handle(w http.ResponseWriter, r *http.Request) { // Grab a dump of the incoming request dump, err := httputil.DumpRequest(r, true /*dump the body also*/) if err != nil { - log.Println("Dump Failure", err) + log.Warn("Dump Failure", err) } self.handler(w, r) @@ -100,8 +101,8 @@ func flush(requests []LogRecord, id string, records_to_dump int) { dest := genFileName(id) w, err := util.Create(dest) if err != nil { - log.Println("Error opening outfile ", dest) - log.Println(err) + log.Warn("Error opening outfile ", dest) + log.Warn(err) return } defer w.Close() @@ -134,7 +135,7 @@ func Logger(in chan []byte, end chan bool, id string, flusher chan bool) { } case <-end: flush(buffer, id, i) - log.Println("Shutdown Logger") + log.Info("Shutdown Logger") return } @@ -144,7 +145,7 @@ func Logger(in chan []byte, end chan bool, id string, flusher chan bool) { func (self *WebService) Run() { port_string := strconv.Itoa(config.GetInt("port_http")) - log.Printf("Serving HTTP on port %s...\n", port_string) + log.Info("Serving HTTP on port:", port_string) logger_chan := make(chan []byte, 1024) flusher := make(chan bool) mux := http.NewServeMux() @@ -259,8 +260,8 @@ func (self *WebService) HandleBatch(w http.ResponseWriter, r *http.Request) { encoder := json.NewEncoder(w) err = encoder.Encode(results) if err != nil { - log.Println("Error Batch results") - log.Println(spew.Sdump(r.Form)) + log.Warn("Error Batch results") + log.Warn(spew.Sdump(r.Form)) err = encoder.Encode("Bad Batch Request") } @@ -299,7 +300,7 @@ func (self *WebService) HandleQuery(w http.ResponseWriter, r *http.Request) { results, err := self.service.Executor.RunPQL(database_name, pql) if err != nil { - log.Println("PQL Exec Error:", err.Error(), database_name, pql) + log.Warn("PQL Exec Error:", err.Error(), database_name, pql) http.Error(w, "Error encoding: "+err.Error(), http.StatusInternalServerError) return } @@ -321,7 +322,7 @@ func (self *WebService) HandleQuery(w http.ResponseWriter, r *http.Request) { } if results == nil { - log.Println("Empty results:", database_name, pql) + log.Warn("Empty results:", database_name, pql) http.Error(w, "Error encoding: "+err.Error(), http.StatusInternalServerError) return } @@ -333,7 +334,7 @@ func (self *WebService) HandleQuery(w http.ResponseWriter, r *http.Request) { encoder := json.NewEncoder(w) err = encoder.Encode(results) if err != nil { - log.Println("Encode Error :", database_name, pql, err.Error()) + log.Warn("Encode Error :", database_name, pql, err.Error()) return } @@ -382,7 +383,7 @@ func bitmaps(frame string, obj JsonObject) chan uint64 { for i, id := range index.GetTimeIds(base_id, atime, quantum) { c <- id if i > 10 { - log.Println("TO MANY TIMEIDS", base_id, atime, quantum) + log.Warn("TO MANY TIMEIDS", base_id, atime, quantum) break } } @@ -433,7 +434,7 @@ func (self *WebService) HandleBit(w http.ResponseWriter, r *http.Request, ToSet //[{ "db": "3", "frame":"brand.","profile_id": 122,"filter":0, "bitmap_id":123}] var results []SBResult if len(args) > 4096 { - log.Println("Request too large:", len(args)) + log.Warn("Request too large:", len(args)) http.Error(w, "Request To large", http.StatusBadRequest) return } @@ -502,7 +503,7 @@ func (self *WebService) HandleBit(w http.ResponseWriter, r *http.Request, ToSet //pql := fmt.Sprintf("set(%d, %s, %d, %d)", bitmap_id, frame, filter, profile_id) if err != nil { - log.Println("Error running set_bit", dbs, frame, profile_id, ToSet) + log.Warn("Error running set_bit", dbs, frame, profile_id, ToSet) http.Error(w, err.Error(), http.StatusInternalServerError) return } @@ -514,8 +515,7 @@ func (self *WebService) HandleBit(w http.ResponseWriter, r *http.Request, ToSet encoder := json.NewEncoder(w) err = encoder.Encode(results) if err != nil { - log.Println("JSON SetBit ERROR:", err, ToSet) - //log.Println("Error encoding set_bit", spew.Sdump(results)) + log.Warn("JSON SetBit ERROR:", err, ToSet) //http.Error(w, "Error econding set_bit", http.StatusInternalServerError) return } @@ -554,7 +554,7 @@ func (self *WebService) HandleStats(w http.ResponseWriter, r *http.Request) { err := encoder.Encode(m) if err != nil { - log.Println("Error encoding stats") + log.Warn("Error encoding stats") http.Error(w, "Error econding stats", http.StatusMethodNotAllowed) } } @@ -694,7 +694,7 @@ func (self *WebService) streamer(writer func(map[string]interface{}) error) { "host": host, }) if err != nil { - log.Println("stopping") + log.Info("stopping") notify.Stop("inbox", inbox) notify.Stop("outbox", outbox) drain(inbox) @@ -708,14 +708,14 @@ func (self *WebService) HandleListenWS(w http.ResponseWriter, r *http.Request) { defer func() { err := recover() out := spew.Sdump(err) - log.Println(out) + log.Info(out) }() ws, err := websocket.Upgrade(w, r, nil, 1024, 1024) if _, ok := err.(websocket.HandshakeError); ok { http.Error(w, "Not a websocket handshake", 400) return } else if err != nil { - log.Println(err) + log.Warn(err) return } self.streamer(func(data map[string]interface{}) error { diff --git a/core/query.go b/core/query.go index 0af9ee360..4c0ab5039 100644 --- a/core/query.go +++ b/core/query.go @@ -1,7 +1,7 @@ package core import ( - "log" + log "github.com/cihub/seelog" "pilosa/db" "pilosa/index" "pilosa/query" @@ -164,7 +164,7 @@ func (self *Service) StashQueryStepHandler(msg *db.Message) { switch val := value.(type) { case index.BitmapHandle: - log.Println("STASH ADDING HANDLE", val) + log.Info("STASH ADDING HANDLE", val) //not sure what to do here.... //result.Handles = append(result.Handles, val) case []byte: @@ -174,7 +174,7 @@ func (self *Service) StashQueryStepHandler(msg *db.Message) { case query.Stash: result.Stash = append(result.Stash, val.Stash...) default: - log.Println("UNEXCPECTED MESSAG", value) + log.Warn("UNEXCPECTED MESSAG", value) } } result_message := db.Message{Data: query.StashQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} diff --git a/core/remotebits.go b/core/remotebits.go index d673f1351..17b763a1b 100644 --- a/core/remotebits.go +++ b/core/remotebits.go @@ -3,7 +3,7 @@ package core import ( // "github.com/davecgh/go-spew/spew" "encoding/gob" - "log" + log "github.com/cihub/seelog" "pilosa/db" "pilosa/util" ) @@ -54,7 +54,7 @@ func (self *RemoteSetBit) Request() { QueryId: random_id, DestProcessId: *process, } - wait := len(request) / 100 + wait := len(request) if wait < 10 { wait = 10 } @@ -75,7 +75,7 @@ func (self *RemoteSetBit) MergeResults(local_results []SBResult) []SBResult { go func(task remote_task) { value, err := self.service.Hold.Get(&task.id, task.wait_time) //eiher need to be the frame process or the handler process? if value == nil { - log.Println("Bad RemoteSetBit Result:", err) + log.Warn("Bad RemoteSetBit Result:", err) empty := make([]SBResult, 0, 0) answers <- empty diff --git a/core/service.go b/core/service.go index 3e79f7256..c36c64571 100644 --- a/core/service.go +++ b/core/service.go @@ -2,7 +2,8 @@ package core import ( "fmt" - "log" + //"log" + log "github.com/cihub/seelog" "os" "os/signal" "pilosa/config" @@ -61,12 +62,24 @@ func (self *Service) PrepareLogging() { if base_path == "" { base_path = "/tmp" } - f, err := os.OpenFile(fmt.Sprintf("%s/%s.%s", base_path, self.name, self.Id), os.O_RDWR|os.O_CREATE|os.O_APPEND, 0666) - if err != nil { - log.Println("error opening file: %v", err) - } - //defer f.Close() - log.SetOutput(f) + /* + f, err := os.OpenFile(fmt.Sprintf("%s/%s.%s", base_path, self.name, self.Id), os.O_RDWR|os.O_CREATE|os.O_APPEND, 0666) + if err != nil { + log.Warn("error opening file: %v", err) + } + //defer f.Close() + log.SetOutput(f) + */ + fname := fmt.Sprintf("%s/%s.%s", base_path, self.name, self.Id) + prod_config := fmt.Sprintf(` + + + + + + `, fname) + logger, _ := log.LoggerFromConfigAsBytes([]byte(prod_config)) + log.ReplaceLogger(logger) } func (service *Service) init_id() { @@ -74,15 +87,15 @@ func (service *Service) init_id() { var err error id_string := config.GetString("id") if id_string == "" { - log.Println("Service id not configured, generating...") + log.Info("Service id not configured, generating...") id = util.RandomUUID() if err != nil { - log.Fatal("problem generating uuid") + log.Critical("problem generating uuid") } } else { id, err = util.ParseGUID(id_string) if err != nil { - log.Fatalf("Service id '%s' not valid", id_string) + log.Critical("Service id not valid:", id_string) } } service.Id = &id @@ -101,9 +114,9 @@ func (service *Service) GetSignals() (chan os.Signal, chan os.Signal) { } func (service *Service) Run() { - log.Println("Setup service...", service.version) + log.Info("Setup service...", service.version) service.TopologyMapper.Setup() - log.Println("Running service...", service.version) + log.Info("Running service...", service.version) go service.TopologyMapper.Run() go service.ProcessMapper.Run() go service.WebService.Run() @@ -116,10 +129,10 @@ func (service *Service) Run() { for { select { case <-sighup: - log.Println("SIGHUP! Reloading configuration...") + log.Info("SIGHUP! Reloading configuration...") // TODO: reload configuration case <-sigterm: - log.Println("SIGTERM! Cleaning up...") + log.Info("SIGTERM! Cleaning up...") service.Index.Shutdown() service.WebService.Shutdown() util.ShutdownStats() @@ -127,7 +140,7 @@ func (service *Service) Run() { return } } - log.Println("Service stopping") + log.Info("Service stopping") } type Message interface { diff --git a/core/topn.go b/core/topn.go index c6ca3b109..a0aee9b5c 100644 --- a/core/topn.go +++ b/core/topn.go @@ -2,7 +2,7 @@ package core import ( "encoding/gob" - "log" + log "github.com/cihub/seelog" "pilosa/db" "pilosa/index" "pilosa/query" @@ -102,7 +102,7 @@ func (self *Service) TopFillHandler(msg *db.Message) { //in order for this to ge topfill := msg.Data.(TopFill) topn, err := self.Index.TopFillBatch(topfill.Args) if err != nil { - log.Println("TopFillHandler:", err) + log.Warn("TopFillHandler:", err) } result_message := db.Message{Data: query.FillResult{&query.BaseQueryResult{Id: &topfill.QueryId, Data: topn}}} @@ -133,7 +133,7 @@ func GatherResults(tasks map[util.GUID]*Task, service *Service) map[uint64]uint6 go func(id util.GUID) { value, err := service.Hold.Get(&id, 10) //eiher need to be the frame process or the handler process? if value == nil { - log.Println("Bad TopN Result:", err) + log.Warn("Bad TopN Result:", err) empty := make([]index.Pair, 0, 0) answers <- empty @@ -173,7 +173,7 @@ func (self *Service) TopNQueryStepHandler(msg *db.Message) { if qs.Input == nil { topn, err := self.Index.TopNAll(qs.Location.FragmentId, qs.N*2, qs.Filters) if err != nil { - log.Println(spew.Sdump(err)) + log.Warn(spew.Sdump(err)) } topnPackage = TopNPackage{*qs.Location.ProcessId, qs.Location.FragmentId, topn, bh} } else { @@ -189,7 +189,7 @@ func (self *Service) TopNQueryStepHandler(msg *db.Message) { topn, err := self.Index.TopN(qs.Location.FragmentId, bh, qs.N*2, qs.Filters) if err != nil { - log.Println(spew.Sdump(err)) + log.Warn(spew.Sdump(err)) } topnPackage = TopNPackage{*qs.Location.ProcessId, qs.Location.FragmentId, topn, bh} } diff --git a/db/topology.go b/db/topology.go index b706fc6e9..1d5d8bcb7 100644 --- a/db/topology.go +++ b/db/topology.go @@ -2,7 +2,7 @@ package db import ( "errors" - "log" + log "github.com/cihub/seelog" "pilosa/config" "pilosa/util" "sync" @@ -304,7 +304,7 @@ func (d *Database) GetFrameSliceIntersect(frame *Frame, slice *Slice) (*FrameSli return frameslice, nil } } - log.Println("Missing FrameSliceIntersect:", d.Name, frame, slice) + log.Warn("Missing FrameSliceIntersect:", d.Name, frame, slice) return nil, FrameSliceIntersectDoesNotExistError } @@ -356,15 +356,14 @@ func (d *Database) GetFragmentForBitmap(slice *Slice, bitmap *Bitmap) (*Fragment //defer d.mutex.Unlock() frame, err := d.getFrame(bitmap.FrameType) if err != nil { - log.Println("Missing FrameType", bitmap.FrameType, d.Name, slice) - log.Println(err) + log.Warn("Missing FrameType", bitmap.FrameType, d.Name, slice) + log.Warn(err) return nil, err } - //log.Println(frame, slice) fsi, err := d.GetFrameSliceIntersect(frame, slice) if err != nil { - log.Println("Missing frame,slice", frame, slice) - log.Println(err) + log.Warn("Missing frame,slice", frame, slice) + log.Warn(err) return nil, err } return fsi.GetFragment(), nil @@ -373,8 +372,8 @@ func (d *Database) GetFragmentForBitmap(slice *Slice, bitmap *Bitmap) (*Fragment func (d *Database) GetFragmentForFrameSlice(frame *Frame, slice *Slice) (*Fragment, error) { fsi, err := d.GetFrameSliceIntersect(frame, slice) if err != nil { - log.Println("Missing frame,slice", frame, slice) - log.Println(err) + log.Warn("Missing frame,slice", frame, slice) + log.Warn(err) return nil, err } return fsi.GetFragment(), nil @@ -400,7 +399,7 @@ func (d *Database) GetFragmentFromProfile(frame string, profile_id uint64) (*Fra func (d *Database) getFragment(frame *Frame, slice *Slice) (*Fragment, error) { fsi, err := d.GetFrameSliceIntersect(frame, slice) if err != nil { - log.Println(err) + log.Warn(err) return nil, err } return fsi.GetFragment(), nil @@ -409,7 +408,7 @@ func (d *Database) getFragment(frame *Frame, slice *Slice) (*Fragment, error) { func (d *Database) addFragment(frame *Frame, slice *Slice, fragment_id util.SUUID) *Fragment { fsi, err := d.GetFrameSliceIntersect(frame, slice) if err != nil { - log.Println("database.addFragment", err) + log.Warn("database.addFragment", err) return nil } fragment := Fragment{id: fragment_id} diff --git a/deps.json b/deps.json index a662fa67b..696fe067c 100644 --- a/deps.json +++ b/deps.json @@ -79,6 +79,11 @@ "version": "master", "type": "git" }, + "seelog": { + "repo": "github.com/cihub/seelog", + "version": "master", + "type": "git" + }, "spew": { "repo": "github.com/davecgh/go-spew/spew", "version": "e762b3d1320b76030bd7f6cc2bfc3d9acce874c0", diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index 45017e54e..67013b8a1 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -1,7 +1,7 @@ package dispatch import ( - "log" + log "github.com/cihub/seelog" "pilosa/core" "pilosa/db" "pilosa/query" @@ -14,18 +14,18 @@ type Dispatch struct { } func (self *Dispatch) Init() error { - log.Println("Starting Dispatcher") + log.Warn("Starting Dispatcher") return nil } func (self *Dispatch) Close() { - log.Println("Shutting down Dispatcher") + log.Warn("Shutting down Dispatcher") } // The Local Route func (self *Dispatch) Run() { - log.Println("Dispatch Run...") + log.Warn("Dispatch Run...") for { message := self.service.Transport.Receive() switch data := message.Data.(type) { @@ -61,9 +61,8 @@ func (self *Dispatch) Run() { case core.BitsResponse: self.service.Hold.Set(data.ResultId(), data.ResultData(), 30) default: - println("Dispatch Unhandled") spew.Dump(data) - log.Println("Unprocessed message", data) + log.Warn("Unprocessed message", data) } } } diff --git a/executor/executor.go b/executor/executor.go index a322f000d..9b73ef30b 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -1,7 +1,7 @@ package executor import ( - "log" + log "github.com/cihub/seelog" "pilosa/config" "pilosa/core" "pilosa/db" @@ -17,12 +17,12 @@ type Executor struct { } func (self *Executor) Init() error { - log.Println("Starting Executor") + log.Warn("Starting Executor") return nil } func (self *Executor) Close() { - log.Println("Shutting down Executor") + log.Warn("Shutting down Executor") } func (self *Executor) NewJob(job *db.Message) { @@ -52,8 +52,8 @@ func (self *Executor) NewJob(job *db.Message) { // case query.MaskQueryStep: // self.service.MaskQueryStepHandler(job) default: - log.Println("unknown") - log.Println(spew.Sdump(job.Data)) + log.Warn("unknown") + log.Warn(spew.Sdump(job.Data)) } } @@ -101,7 +101,7 @@ func (self *Executor) runQuery(database *db.Database, qry *query.Query) error { if loc != nil { self.service.Transport.Send(msg, loc.ProcessId) } else { - log.Println("Problem with querystep(nil location)", spew.Sdump(step)) + log.Warn("Problem with querystep(nil location)", spew.Sdump(step)) } } } @@ -172,7 +172,7 @@ func (self *Executor) RunPQL(database_name string, pql string) (interface{}, err }(query_list[i], result) if err != nil { out := spew.Sdump(err) - log.Println(out) + log.Warn(out) } } for z := 0; z < x; z++ { @@ -186,7 +186,7 @@ func (self *Executor) RunPQL(database_name string, pql string) (interface{}, err } func (self *Executor) Run() { - log.Println("Executor Run...") + log.Warn("Executor Run...") } func NewExecutor(service *core.Service) *Executor { diff --git a/index/bitmap.go b/index/bitmap.go index 5b690243a..2fd864946 100644 --- a/index/bitmap.go +++ b/index/bitmap.go @@ -7,8 +7,8 @@ import ( "compress/gzip" "encoding/base64" "encoding/gob" + log "github.com/cihub/seelog" "io/ioutil" - "log" "github.com/yasushi-saito/rbtree" ) @@ -300,7 +300,7 @@ func Union(a_bm IBitmap, b_bm IBitmap) IBitmap { a = a.Next() b = b.Next() } else { - log.Println("NEVER SHOULD BE HERE") + log.Warn("NEVER SHOULD BE HERE") break } } @@ -358,7 +358,7 @@ func Difference(a_bm IBitmap, b_bm IBitmap) IBitmap { a = a.Next() b = b.Next() } else { - log.Println("NEVER SHOULD BE HERE") + log.Warn("NEVER SHOULD BE HERE") break } } @@ -427,7 +427,7 @@ func CreateRBBitmap() IBitmap { func (self *Bitmap) FromCompressString(str string) { compressed_data, err := base64.StdEncoding.DecodeString(str) if err != nil { - log.Println(err) + log.Warn(err) return } reader, _ := gzip.NewReader(bytes.NewReader(compressed_data)) @@ -468,7 +468,7 @@ func (b *Bitmap) ToBytes() []byte { obj := i.Item().(*Chunk) err := enc.Encode(obj) if err != nil { - log.Println(err) + log.Warn(err) } } return buf.Bytes() diff --git a/index/brand.go b/index/brand.go index 895ffe887..ecd645783 100644 --- a/index/brand.go +++ b/index/brand.go @@ -11,7 +11,7 @@ import ( import ( "encoding/json" "fmt" - "log" + log "github.com/cihub/seelog" "pilosa/config" "pilosa/util" "sort" @@ -109,7 +109,7 @@ func (self *Brand) cache_it(bm IBitmap, bitmap_id uint64, category uint64) { if bm.Count() >= self.threshold_value { self.bitmap_cache[bitmap_id] = &Rank{&Pair{bitmap_id, bm.Count()}, bm, category} if len(self.bitmap_cache) > self.threshold_length { - log.Printf("RANK: %d %d %d", len(self.bitmap_cache), self.threshold_length, self.threshold_value) + log.Info("RANK: %d %d %d", len(self.bitmap_cache), self.threshold_length, self.threshold_value) self.Rank() self.trim() } @@ -121,7 +121,7 @@ func (self *Brand) trim() { delete(self.bitmap_cache, k) } } - log.Printf("TRIM: %d %d", len(self.bitmap_cache), self.threshold_length) + log.Info("TRIM:", len(self.bitmap_cache), self.threshold_length) } @@ -236,7 +236,7 @@ func (self *Brand) TopN(src_bitmap IBitmap, n int, categories []uint64) []Pair { } func dump(r RankList, n int) { for i, v := range r { - log.Println(i, v) + log.Info(i, v) if i > n { return } @@ -366,18 +366,18 @@ func (self *Brand) getFileName() string { } func (self *Brand) Persist() error { - log.Println("Brand Persist:", self.getFileName()) + log.Info("Brand Persist:", self.getFileName()) self.storage.FlushBatch() asize := len(self.bitmap_cache) if asize == 0 { - log.Println("Nothing to save ", self.getFileName()) + log.Warn("Nothing to save ", self.getFileName()) return nil } w, err := util.Create(self.getFileName()) if err != nil { - log.Println("Error opening outfile ", self.getFileName()) - log.Println(err) + log.Warn("Error opening outfile ", self.getFileName()) + log.Warn(err) return err } defer w.Close() @@ -402,11 +402,11 @@ func (self *Brand) Persist() error { } func (self *Brand) Load(requestChan chan Command, f *Fragment) { - log.Println("Brand Load") + log.Warn("Brand Load") time.Sleep(time.Duration(rand.Intn(32)) * time.Second) //trying to avoid mass cassandra hit r, err := util.Open(self.getFileName()) if err != nil { - log.Println("NO Brand Init File:", self.getFileName()) + log.Warn("NO Brand Init File:", self.getFileName()) return } dec := json.NewDecoder(r) diff --git a/index/commands.go b/index/commands.go index 24b34809d..09f4a7e72 100644 --- a/index/commands.go +++ b/index/commands.go @@ -3,8 +3,8 @@ package index import ( "bytes" "compress/gzip" + log "github.com/cihub/seelog" "io/ioutil" - "log" "time" ) @@ -141,7 +141,7 @@ func (self *CmdGetBytes) Execute(f *Fragment) Calculation { //*Compress it if !ok { bm = NewBitmap() - log.Println("cache miss") + log.Warn("cache miss") } var b bytes.Buffer w := gzip.NewWriter(&b) diff --git a/index/fragment_container.go b/index/fragment_container.go index 48b7c9741..692f8ca16 100644 --- a/index/fragment_container.go +++ b/index/fragment_container.go @@ -6,7 +6,7 @@ import ( "encoding/gob" "errors" "fmt" - "log" + log "github.com/cihub/seelog" "pilosa/config" "pilosa/util" "strings" @@ -25,7 +25,7 @@ func lookup(stmt *sql.Stmt, tile_id uint64) int { var category int err := stmt.QueryRow(tile_id).Scan(&category) // WHERE number = 13 if err != nil { - log.Println(err.Error()) + log.Warn(err.Error()) return 0 } return category @@ -53,7 +53,7 @@ func init() { } func (self *FragmentContainer) Shutdown() { - log.Println("Container Shutdown Started") + log.Warn("Container Shutdown Started") var wg sync.WaitGroup wg.Add(len(self.fragments)) @@ -61,7 +61,7 @@ func (self *FragmentContainer) Shutdown() { v.exit <- &wg } wg.Wait() - log.Println("Container Shutdown Complete") + log.Warn("Container Shutdown Complete") } func (self *FragmentContainer) LoadBitmap(frag_id util.SUUID, bitmap_id uint64, compressed_bitmap string, filter uint64) { @@ -294,7 +294,7 @@ func (self *FragmentContainer) Clear(frag_id util.SUUID) (bool, error) { func (self *FragmentContainer) AddFragment(db string, frame string, slice int, id util.SUUID) { _, ok := self.fragments[id] if !ok { - log.Println("ADD FRAGMENT", frame, db, slice, util.SUUID_to_Hex(id)) + log.Warn("ADD FRAGMENT", frame, db, slice, util.SUUID_to_Hex(id)) f := NewFragment(id, db, slice, frame) loader := make(chan Command) self.fragments[id] = f @@ -348,7 +348,7 @@ func getStorage(db string, slice int, frame string, fid util.SUUID) Storage { func NewFragment(frag_id util.SUUID, db string, slice int, frame string) *Fragment { var impl Pilosa - log.Println(fmt.Sprintf("XXXXXXXXXXXXXXXXXXXXXXXXXXX(%s)", frame)) + log.Warn(fmt.Sprintf("XXXXXXXXXXXXXXXXXXXXXXXXXXX(%s)", frame)) if strings.HasSuffix(frame, ".n") { impl = NewBrand(db, frame, slice, getStorage(db, slice, frame, frag_id), 50000, 45000, 100) } else { @@ -464,7 +464,7 @@ func (self *Fragment) difference(bitmaps []BitmapHandle) BitmapHandle { func (self *Fragment) Persist() { err := self.impl.Persist() if err != nil { - log.Println("Error saving:", err) + log.Warn("Error saving:", err) } } func (self *Fragment) Load(loadChan chan Command) { @@ -492,7 +492,7 @@ func (self *Fragment) ServeFragment(loadChan chan Command) { case req := <-loadChan: self.processCommand(req) case wg := <-self.exit: - log.Println("Fragment Shutdown") + log.Warn("Fragment Shutdown") self.Persist() wg.Done() } diff --git a/index/general.go b/index/general.go index 25041e4b5..0784733e0 100644 --- a/index/general.go +++ b/index/general.go @@ -3,7 +3,7 @@ package index import ( "encoding/json" "fmt" - "log" + log "github.com/cihub/seelog" "pilosa/config" "pilosa/util" @@ -98,10 +98,10 @@ func (self *General) getFileName() string { } func (self *General) Persist() error { - log.Println("General Persist") + log.Warn("General Persist") w, err := util.Create(self.getFileName()) if err != nil { - log.Println("Error saving:", err) + log.Warn("Error saving:", err) return err } defer w.Close() @@ -119,10 +119,10 @@ func (self *General) Persist() error { } func (self *General) Load(requestChan chan Command, f *Fragment) { - log.Println("General Load") + log.Warn("General Load") r, err := util.Open(self.getFileName()) if err != nil { - log.Println("NO General Init File:", self.getFileName()) + log.Warn("NO General Init File:", self.getFileName()) return } @@ -131,7 +131,6 @@ func (self *General) Load(requestChan chan Command, f *Fragment) { if err := dec.Decode(&keys); err != nil { return - //log.Println("Bad mojo") } for _, k := range keys { request := NewLoadRequest(k) diff --git a/index/storage_cass.go b/index/storage_cass.go index 3c8b839ed..6871898e8 100644 --- a/index/storage_cass.go +++ b/index/storage_cass.go @@ -3,7 +3,7 @@ package index // #cgo CFLAGS:-mpopcnt import ( - "log" + log "github.com/cihub/seelog" "pilosa/config" "pilosa/util" @@ -51,7 +51,7 @@ func NewCassStorage() Storage { // cluster.CQLVersion = "3.0.0" session, err := cluster.CreateSession() if err != nil { - log.Fatal(err) + log.Warn(err) } obj.db = session @@ -140,7 +140,7 @@ func (self *CassandraStorage) EndBatch() { } */ } else { - log.Println("NIL BATCH") + log.Warn("NIL BATCH") } delta := time.Since(start) util.SendTimer("cassandra_storage_EndBatch", delta.Nanoseconds()) diff --git a/query/parser.go b/query/parser.go index b6a3dac22..e73da71dc 100644 --- a/query/parser.go +++ b/query/parser.go @@ -3,7 +3,7 @@ package query import ( "errors" "fmt" - "log" + log "github.com/cihub/seelog" "pilosa/index" "pilosa/util" "strconv" @@ -303,7 +303,7 @@ ArgLoop: case TYPE_RB: // default: - log.Println(spew.Sdump("unexpected", token)) + log.Warn(spew.Sdump("unexpected", token)) return nil, errors.New("BAD TOKEN") } } diff --git a/query/planner.go b/query/planner.go index 29c495823..7e03547ff 100644 --- a/query/planner.go +++ b/query/planner.go @@ -4,7 +4,7 @@ import ( "encoding/gob" "errors" "fmt" - "log" + log "github.com/cihub/seelog" "math/rand" "pilosa/db" "pilosa/index" @@ -149,7 +149,7 @@ func (qt *TopNQueryTree) getLocation(d *db.Database) (*db.Location, error) { slice := d.GetOrCreateSlice(qt.Slice) fragment, err := d.GetFragmentForFrameSlice(frame, slice) if err != nil { - log.Println("GetFragmentForFrameSliceFailed TopNQueryTree", frame, slice) + log.Warn("GetFragmentForFrameSliceFailed TopNQueryTree", frame, slice) return nil, err } qt.location = fragment.GetLocation() @@ -332,7 +332,7 @@ func (qt *GetQueryTree) getLocation(d *db.Database) (*db.Location, error) { slice := d.GetOrCreateSlice(qt.slice) // TODO: this should probably be just GetSlice (no create) fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) if err != nil { - log.Println("GetFragmenForBitmapFailed GetQueryTree", slice) + log.Warn("GetFragmenForBitmapFailed GetQueryTree", slice) return nil, err } return fragment.GetLocation(), nil @@ -369,7 +369,7 @@ func (qt *SetQueryTree) getLocation(d *db.Database) (*db.Location, error) { } fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) if err != nil { - log.Println("NOT FOUND:", slice, qt.bitmap) + log.Warn("NOT FOUND:", slice, qt.bitmap) return nil, NewFragmentNotFound(d.Name, qt.bitmap.FrameType, db.GetSlice(qt.profile_id)) } return fragment.GetLocation(), nil @@ -467,7 +467,7 @@ func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) var p Appendable if query.Operation == "stash" { - // log.Println("STASH:", n) + // log.Warn("STASH:", n) tree = &StashQueryTree{N: n} } else { tree = &CatQueryTree{N: n} @@ -574,7 +574,7 @@ func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) tree = &StashQueryTree{subqueries: subqueries} } else { //TODO return error gracefully - log.Println(spew.Sdump(query)) + log.Warn(spew.Sdump(query)) return nil, errors.New("BuildTree Issues") } } @@ -781,7 +781,7 @@ func (qt *MaskQueryTree) getLocation(d *db.Database) (*db.Location, error) { slice := d.GetOrCreateSlice(qt.slice) // TODO: this should probably be just GetSlice (no create) fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) if err != nil { - log.Println("GetFragmenForBitmapFailed GetQueryTree", slice) + log.Warn("GetFragmenForBitmapFailed GetQueryTree", slice) return nil, err } return fragment.GetLocation(), nil @@ -812,7 +812,7 @@ func (qt *RangeQueryTree) getLocation(d *db.Database) (*db.Location, error) { slice := d.GetOrCreateSlice(qt.slice) // TODO: this should probably be just GetSlice (no create) fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) if err != nil { - log.Println("GetFragmenForBitmapFailed GetQueryTree", slice) + log.Warn("GetFragmenForBitmapFailed GetQueryTree", slice) return nil, err } return fragment.GetLocation(), nil diff --git a/transport/tcp.go b/transport/tcp.go index 5bf113242..f0c98d9fe 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -3,7 +3,7 @@ package transport import ( "encoding/gob" "fmt" - "log" + log "github.com/cihub/seelog" "net" "pilosa/config" "pilosa/core" @@ -38,14 +38,14 @@ BeginManageConnection: if self.conn == nil { process, err := self.transport.service.ProcessMap.GetProcess(self.process) if err != nil { - log.Println("transport/tcp: error getting process, retrying in 2 seconds... ", self.process, err) + log.Warn("transport/tcp: error getting process, retrying in 2 seconds... ", self.process, err) time.Sleep(2 * time.Second) continue } host_string := fmt.Sprintf("%s:%d", process.Host(), process.PortTcp()) conn, err := net.Dial("tcp", host_string) if err != nil { - log.Println("transport/tcp: error dialing: ", host_string, " Retrying in 2 seconds...") + log.Warn("transport/tcp: error dialing: ", host_string, " Retrying in 2 seconds...") time.Sleep(2 * time.Second) continue } @@ -62,7 +62,7 @@ BeginManageConnection: var mess *db.Message err := decoder.Decode(&mess) if err != nil { - log.Println("transport/tcp: error decoding message: ", err.Error()) + log.Warn("transport/tcp: error decoding message: ", err.Error()) exit <- 1 return } @@ -74,7 +74,7 @@ BeginManageConnection: case message := <-self.outbox: err := encoder.Encode(message) if err != nil { - log.Println(err.Error()) + log.Warn(err.Error()) return } case message := <-self.inbox: @@ -108,7 +108,7 @@ type TcpTransport struct { } func (self *TcpTransport) Run() { - log.Println("Initializing TCP transport") + log.Warn("Initializing TCP transport") go self.listen() for { select { @@ -130,12 +130,12 @@ func (self *TcpTransport) listen() { port_string := fmt.Sprintf(":%d", self.port) l, e := net.Listen("tcp", port_string) if e != nil { - log.Fatal("Cannot bind to port! ", self.port) + log.Warn("Cannot bind to port! ", self.port) } for { conn, err := l.Accept() if err != nil { - log.Println("Error accepting, trying again in 2 sec... ", err) + log.Warn("Error accepting, trying again in 2 sec... ", err) time.Sleep(2 * time.Second) continue } @@ -149,7 +149,7 @@ func (self *TcpTransport) manage(conn *net.Conn) { } func (self *TcpTransport) Close() { - log.Println("Shutting down TCP transport") + log.Warn("Shutting down TCP transport") } func (self *TcpTransport) Send(message *db.Message, host *GUID) { diff --git a/util/id.go b/util/id.go index 8a3bb41bf..4c46eaf81 100644 --- a/util/id.go +++ b/util/id.go @@ -5,7 +5,7 @@ import ( "encoding/binary" "encoding/hex" "fmt" - "log" + log "github.com/cihub/seelog" "math/rand" "os" "strings" @@ -23,7 +23,7 @@ func init() { rand.Seed(time.Now().UTC().UnixNano()) f, err := os.Open("/dev/urandom") if err != nil { - log.Fatal(err) + log.Warn(err) } Random = f } diff --git a/util/statd.go b/util/statd.go index 3455db813..571ea3b56 100644 --- a/util/statd.go +++ b/util/statd.go @@ -1,7 +1,7 @@ package util import ( - "log" + log "github.com/cihub/seelog" "pilosa/config" "github.com/cactus/go-statsd-client/statsd" @@ -24,7 +24,7 @@ func init() { count = make(chan string, 100) end = make(chan bool) stat_config := config.GetStringDefault("statsd_server", "127.0.0.1:8125") - log.Println("New Stats", stat_config) + log.Warn("New Stats", stat_config) stats, _ := statsd.New(stat_config, "") go func() { for { @@ -34,7 +34,7 @@ func init() { case stat := <-count: stats.Inc(stat, 1, 1.0) case <-end: - log.Println("DONE Stats") + log.Warn("DONE Stats") return } @@ -52,6 +52,6 @@ func SendInc(stat string) { } func ShutdownStats() { - log.Println("Shutdown Stats") + log.Warn("Shutdown Stats") end <- true }