added new logging setup

This commit is contained in:
Todd Gruben 2015-02-05 21:09:07 +00:00
parent d8d5a8dfcb
commit 377d850534
23 changed files with 169 additions and 153 deletions

View file

@ -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)
}

View file

@ -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"
}
}

View file

@ -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)
}
}

View file

@ -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 {

View file

@ -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}}}

View file

@ -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

View file

@ -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(`
<seelog>
<outputs>
<rollingfile type="size" filename="%s" maxsize="524288000" maxrolls="4" />
</outputs>
</seelog>
`, 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 {

View file

@ -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}
}

View file

@ -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}

View file

@ -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",

View file

@ -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)
}
}
}

View file

@ -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 {

View file

@ -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()

View file

@ -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)

View file

@ -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)

View file

@ -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()
}

View file

@ -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)

View file

@ -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())

View file

@ -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")
}
}

View file

@ -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

View file

@ -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) {

View file

@ -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
}

View file

@ -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
}