diff --git a/core/http.go b/core/http.go index 396c760c8..8c8f90cf3 100644 --- a/core/http.go +++ b/core/http.go @@ -298,6 +298,7 @@ func (self *WebService) HandleQuery(w http.ResponseWriter, r *http.Request) { } _, bits := r.Form["bits"] + log.Info("PQL:", database_name, pql) results, err := self.service.Executor.RunPQL(database_name, pql) if err != nil { log.Warn("PQL Exec Error:", err.Error(), database_name, pql) diff --git a/core/remotebits.go b/core/remotebits.go index 17b763a1b..bd1b75707 100644 --- a/core/remotebits.go +++ b/core/remotebits.go @@ -54,7 +54,7 @@ func (self *RemoteSetBit) Request() { QueryId: random_id, DestProcessId: *process, } - wait := len(request) + wait := len(request) * 10 if wait < 10 { wait = 10 } diff --git a/core/service.go b/core/service.go index c36c64571..3c9e8b6f5 100644 --- a/core/service.go +++ b/core/service.go @@ -71,13 +71,15 @@ func (self *Service) PrepareLogging() { log.SetOutput(f) */ fname := fmt.Sprintf("%s/%s.%s", base_path, self.name, self.Id) + // + log_level := config.GetStringDefault("log_level", "info") prod_config := fmt.Sprintf(` - + - `, fname) + `, log_level, fname) logger, _ := log.LoggerFromConfigAsBytes([]byte(prod_config)) log.ReplaceLogger(logger) } diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index 67013b8a1..b3cdabd5b 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -30,10 +30,12 @@ func (self *Dispatch) Run() { message := self.service.Transport.Receive() switch data := message.Data.(type) { case core.BatchRequest: + log.Trace("Dispatch.Run BatchRequest") response := db.Message{Data: core.BatchResponse{Id: data.Id}} self.service.Index.LoadBitmap(data.Fragment_id, data.Bitmap_id, data.Compressed_bitmap, data.Filter) self.service.Transport.Send(&response, data.Source) case core.BitsRequest: + log.Trace("Dispatch.Run BitsRequest") var results []core.SBResult result := false @@ -50,13 +52,17 @@ func (self *Dispatch) Run() { response := db.Message{Data: core.BitsResponse{Id: &data.QueryId, Items: results}} self.service.Transport.Send(&response, &data.ReturnProcessId) case core.PingRequest: + log.Trace("Dispatch.Run Ping") pong := db.Message{Data: core.PongRequest{Id: data.Id}} self.service.Transport.Send(&pong, data.Source) case db.HoldResult: + log.Trace("Dispatch.Run HoldResult") self.service.Hold.Set(data.ResultId(), data.ResultData(), 30) case query.PortableQueryStep: + log.Trace("Dispatch.Run PortableQueryStep") go self.service.Executor.NewJob(message) case core.TopFill: + log.Trace("Dispatch.Run TopFill") go self.service.TopFillHandler(message) case core.BitsResponse: self.service.Hold.Set(data.ResultId(), data.ResultData(), 30) diff --git a/executor/executor.go b/executor/executor.go index 9b73ef30b..699415e4b 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -17,15 +17,16 @@ type Executor struct { } func (self *Executor) Init() error { - log.Warn("Starting Executor") + log.Trace("Executor.Init()") return nil } func (self *Executor) Close() { - log.Warn("Shutting down Executor") + log.Trace("Executor.Close()") } func (self *Executor) NewJob(job *db.Message) { + log.Trace("NewJob", job) switch job.Data.(type) { case query.CountQueryStep: self.service.CountQueryStepHandler(job) @@ -73,6 +74,7 @@ func (self *Executor) RunQueryTest(database_name string, pql string) string { } func (self *Executor) runQuery(database *db.Database, qry *query.Query) error { + log.Trace("Executor.runQuery", database, qry) process, err := self.service.GetProcess() if err != nil { return err @@ -109,6 +111,7 @@ func (self *Executor) runQuery(database *db.Database, qry *query.Query) error { } func (self *Executor) RunPQL(database_name string, pql string) (interface{}, error) { + log.Trace("Executor.RunPQL", database_name, pql) database := self.service.Cluster.GetOrCreateDatabase(database_name) // see if the outer query function is a custom query @@ -190,5 +193,6 @@ func (self *Executor) Run() { } func NewExecutor(service *core.Service) *Executor { + log.Trace("NewExector") return &Executor{service, make(chan *db.Message)} } diff --git a/hold/hold.go b/hold/hold.go index 660d44978..2dd94bebd 100644 --- a/hold/hold.go +++ b/hold/hold.go @@ -2,6 +2,7 @@ package hold import ( "errors" + log "github.com/cihub/seelog" . "pilosa/util" "time" ) @@ -24,11 +25,13 @@ type Holder struct { //var Hold Holder func (self *Holder) DelChan(id *GUID) { + log.Trace("Holder.DelChan", id) req := delhold{id} self.delchan <- req } func (self *Holder) GetChan(id *GUID) holdchan { + log.Trace("Holder.GetChan", id) reply := make(chan holdchan) req := gethold{id, reply} self.getchan <- req @@ -36,6 +39,7 @@ func (self *Holder) GetChan(id *GUID) holdchan { } func (self *Holder) Get(id *GUID, timeout int) (interface{}, error) { + log.Trace("Holder.Get", id, timeout) ch := self.GetChan(id) select { case val := <-ch: @@ -47,6 +51,7 @@ func (self *Holder) Get(id *GUID, timeout int) (interface{}, error) { } func (self *Holder) Set(id *GUID, value interface{}, timeout int) { + log.Trace("Holder.Set", id, value, timeout) ch := self.GetChan(id) go func() { select { diff --git a/index/storage_cass.go b/index/storage_cass.go index 6871898e8..fbb34a58b 100644 --- a/index/storage_cass.go +++ b/index/storage_cass.go @@ -104,6 +104,7 @@ func (c *CassandraStorage) Fetch(bitmap_id uint64, db string, frame string, slic } delta := time.Since(start) util.SendTimer("cassandra_storage_Fetch", delta.Nanoseconds()) + util.SendInc("cassandra_storage_Read") bitmap.SetCount(uint64(count)) return bitmap, uint64(filter) } @@ -144,6 +145,7 @@ func (self *CassandraStorage) EndBatch() { } delta := time.Since(start) util.SendTimer("cassandra_storage_EndBatch", delta.Nanoseconds()) + util.SendInc("cassandra_storage_Write") } diff --git a/query/lexer.go b/query/lexer.go index c84c90e7a..ffe95d177 100644 --- a/query/lexer.go +++ b/query/lexer.go @@ -3,6 +3,7 @@ package query import ( "errors" "fmt" + log "github.com/cihub/seelog" "strings" "unicode" "unicode/utf8" @@ -44,6 +45,7 @@ type Lexer struct { } func (lexer *Lexer) emit(typ int) { + log.Trace("Lexer.emit", typ) if lexer.start < lexer.pos { lexer.ch <- Token{lexer.text[lexer.start:lexer.pos], typ} } @@ -51,6 +53,7 @@ func (lexer *Lexer) emit(typ int) { } func (lexer *Lexer) acceptUntil(chars string, consume bool) (rune, error) { + log.Trace("Lexer.acceptUntil", chars, consume) start_pos := lexer.pos for { @@ -77,6 +80,7 @@ func (lexer *Lexer) acceptUntil(chars string, consume bool) (rune, error) { } func (lexer *Lexer) acceptRun(valid string) { + log.Trace("Lexer.acceptRun", valid) for strings.IndexRune(valid, lexer.next()) >= 0 { } lexer.backup() @@ -84,6 +88,7 @@ func (lexer *Lexer) acceptRun(valid string) { // next returns the next rune in the input. func (lexer *Lexer) next() (runey rune) { + log.Trace("Lexer.next", runey) if lexer.pos >= len(lexer.text) { lexer.width = 0 return 0 @@ -95,18 +100,21 @@ func (lexer *Lexer) next() (runey rune) { // ignore skips over the pending input before this point. func (lexer *Lexer) ignore() { + log.Trace("Lexer.ignore") lexer.start = lexer.pos } // backup steps back one rune. // Can be called only once per call of next. func (lexer *Lexer) backup() { + log.Trace("Lexer.backup") lexer.pos -= lexer.width } // peek returns but does not consume // the next rune in the input. func (lexer *Lexer) peek() rune { + log.Trace("Lexer.peek") for { next_rune := lexer.next() // ignore spaces @@ -119,6 +127,7 @@ func (lexer *Lexer) peek() rune { } func stateError(err error) func(lexer *Lexer) statefn { + log.Trace("stateError", err) return func(lexer *Lexer) statefn { lexer.ch <- Token{err.Error(), TYPE_ERROR} close(lexer.ch) @@ -127,6 +136,7 @@ func stateError(err error) func(lexer *Lexer) statefn { } func stateFunc(lexer *Lexer) statefn { + log.Trace("stateFunc", lexer) _, err := lexer.acceptUntil("(", true) if err != nil { return stateError(err) @@ -136,6 +146,7 @@ func stateFunc(lexer *Lexer) statefn { } func stateLP(lexer *Lexer) statefn { + log.Trace("stateLP", lexer) lexer.pos += 1 lexer.emit(TYPE_LP) // handle multiple LPs @@ -146,6 +157,7 @@ func stateLP(lexer *Lexer) statefn { } func stateLB(lexer *Lexer) statefn { + log.Trace("stateLB", lexer) lexer.acceptUntil("[", true) lexer.next() lexer.emit(TYPE_LB) @@ -177,6 +189,7 @@ func stateLB(lexer *Lexer) statefn { } func stateRB(lexer *Lexer) statefn { + log.Trace("stateRB", lexer) lexer.pos += 1 lexer.emit(TYPE_RB) @@ -191,6 +204,7 @@ func stateRB(lexer *Lexer) statefn { } func stateArgs(lexer *Lexer) statefn { + log.Trace("stateArgs", lexer) r, err := lexer.acceptUntil("(),=[]", false) if err != nil { return stateError(err) @@ -215,6 +229,7 @@ func stateArgs(lexer *Lexer) statefn { } func stateKeyword(lexer *Lexer) statefn { + log.Trace("stateKeyword", lexer) _, err := lexer.acceptUntil("=", true) if err != nil { return stateError(err) @@ -224,6 +239,7 @@ func stateKeyword(lexer *Lexer) statefn { } func stateEquals(lexer *Lexer) statefn { + log.Trace("stateEquals", lexer) e := lexer.next() if e != '=' { return stateError(errors.New("Expecting '='!")) @@ -233,6 +249,7 @@ func stateEquals(lexer *Lexer) statefn { } func stateValue(lexer *Lexer) statefn { + log.Trace("stateValue", lexer) r, err := lexer.acceptUntil("(),[", false) if err != nil { return stateError(err) @@ -259,6 +276,7 @@ func stateValue(lexer *Lexer) statefn { } func stateRP(lexer *Lexer) statefn { + log.Trace("stateRP", lexer) lexer.pos += 1 lexer.emit(TYPE_RP) @@ -275,23 +293,27 @@ func stateRP(lexer *Lexer) statefn { } func stateRPComma(lexer *Lexer) statefn { + log.Trace("stateRPComma", lexer) lexer.pos += 1 lexer.emit(TYPE_COMMA) return stateValue } func stateComma(lexer *Lexer) statefn { + log.Trace("stateComma", lexer) lexer.pos += 1 lexer.emit(TYPE_COMMA) return stateArgs } func stateEOF(lexer *Lexer) statefn { + log.Trace("stateEOF", lexer) close(lexer.ch) return nil } func (lexer *Lexer) Lex() (tokens []Token, err error) { + log.Trace("Lexer.Lex", tokens, err) defer func() { if r := recover(); r != nil { var ok bool @@ -321,6 +343,7 @@ func (lexer *Lexer) Lex() (tokens []Token, err error) { } func Lex(input string) ([]Token, error) { + log.Trace("Lex", input) lexer := Lexer{input, 0, 0, 0, TYPE_FUNC, make(chan Token)} return lexer.Lex() } diff --git a/query/parser.go b/query/parser.go index e73da71dc..16af7a979 100644 --- a/query/parser.go +++ b/query/parser.go @@ -40,6 +40,7 @@ bm := db.Bitmap{bitmap_id, frame_type, filter} return []QueryInput{&bm}, uint64(profile_id), 0 */ func (self *QueryParser) next() *Token { + log.Trace("QueryParser.next") self.pos += 1 if self.pos > len(self.tokens) { return nil @@ -48,16 +49,19 @@ func (self *QueryParser) next() *Token { } func (self *QueryParser) peek() *Token { + log.Trace("QueryParser.peek") token := self.next() self.backup() return token } func (self *QueryParser) backup() { + log.Trace("QueryParser.backup") self.pos -= 1 } func (self *QueryParser) Parse() (query *Query, err error) { + log.Trace("QueryParser.Parse", query, err) defer func() { if r := recover(); r != nil { var ok bool @@ -341,6 +345,7 @@ ArgLoop: } func Parse(tokens []Token) (*Query, error) { + log.Trace("Parse", tokens) parser := QueryParser{tokens, 0} return parser.Parse() } diff --git a/query/planner.go b/query/planner.go index 7e03547ff..c2abfcf34 100644 --- a/query/planner.go +++ b/query/planner.go @@ -66,17 +66,22 @@ type BaseQueryStep struct { } func (self *BaseQueryStep) GetId() *util.GUID { + log.Trace("BaseQueryStep", *self.Id) return self.Id } func (self *BaseQueryStep) GetLocation() *db.Location { + log.Trace("BaseQueryStep.GetLocation", *self.Location) return self.Location } func (self *BaseQueryStep) LocIsDest() bool { + log.Trace("BaseQueryStep.LocIsDest") if self.Location.ProcessId == self.Destination.ProcessId && self.Location.FragmentId == self.Destination.FragmentId { + log.Trace("BaseQueryStep.LocIsDest Return true") return true } + log.Trace("BaseQueryStep.LocIsDest Return false") return false } @@ -86,10 +91,13 @@ type BaseQueryResult struct { } func (self *BaseQueryResult) ResultId() *util.GUID { + log.Trace("BaseQueryStep.ResultId", self) + return self.Id } func (self *BaseQueryResult) ResultData() interface{} { + log.Trace("BaseQueryStep.ResultData", self) return self.Data } @@ -113,6 +121,7 @@ type CountQueryTree struct { // Uses consistent hashing function to select node containing data for GET operation func (qt *CountQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("CountQueryTree.getLocation", d, qt) return qt.subquery.getLocation(d) } @@ -143,6 +152,7 @@ type TopNQueryTree struct { // Uses consistent hashing function to select node containing data for GET operation func (qt *TopNQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("TopNQueryTree.getLocation", d, qt) var err error if qt.location == nil { frame := d.GetOrCreateFrame(qt.Frame) @@ -178,6 +188,7 @@ type UnionQueryTree struct { // Uses consistent hashing function to select node containing data for GET operation func (qt *UnionQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("UnionQueryTree.getLocation", d, qt) var err error if qt.location == nil { subqueryLength := len(qt.subqueries) @@ -210,6 +221,7 @@ type IntersectQueryTree struct { // Uses consistent hashing function to select node containing data for GET operation func (qt *IntersectQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("IntersectQueryTree.getLocation", d, qt) var err error if qt.location == nil { subqueryLength := len(qt.subqueries) @@ -242,6 +254,7 @@ type DifferenceQueryTree struct { // Uses consistent hashing function to select node containing data for GET operation func (qt *DifferenceQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("DifferenceQueryTree.getLocation", d, qt) var err error if qt.location == nil { subqueryLength := len(qt.subqueries) @@ -278,11 +291,13 @@ type CatQueryTree struct { } func (qt *CatQueryTree) Append(subtree QueryTree) { + log.Trace("CatQueryTree.Append", qt, subtree) qt.subqueries = append(qt.subqueries, subtree) } // Uses consistent hashing function to select node containing data for GET operation func (qt *CatQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("CatQueryTree.getLocation", d) var err error if qt.location == nil { subqueryLength := len(qt.subqueries) @@ -305,6 +320,7 @@ type AllQueryTree struct { // Uses consistent hashing function to select node containing data for GET operation func (qt *AllQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("AllQueryTree.getLocation", d) return nil, nil } @@ -329,6 +345,7 @@ type GetQueryTree struct { // Uses consistent hashing function to select node containing data for GET operation func (qt *GetQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("GetQueryTree.getLocation", d) slice := d.GetOrCreateSlice(qt.slice) // TODO: this should probably be just GetSlice (no create) fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) if err != nil { @@ -359,6 +376,7 @@ type SetQueryTree struct { // Uses consistent hashing function to select node containing data for GET operation func (qt *SetQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("SetQueryTree.getLocation", d) // check here for supported frames if !d.IsValidFrame(qt.bitmap.FrameType) { return nil, NewInvalidFrame(d.Name, qt.bitmap.FrameType) @@ -420,6 +438,7 @@ type QueryTree interface { } func validateRange(Args map[string]interface{}) error { + log.Trace("validateRange", Args) _, ok := Args["id"] if !ok { return errors.New("missing bitmap id") @@ -441,6 +460,7 @@ func validateRange(Args map[string]interface{}) error { // Builds QueryTree object from Query. Pass slice=-1 to perform operation on all slices func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { + log.Trace("QueryPlanner.buildTree", query, slice) var tree QueryTree // handle SET operation regardless of the slice @@ -583,6 +603,7 @@ func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) // Produces flattened QueryPlan from QueryTree input func (self *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Location) (*QueryPlan, error) { + log.Trace("QueryPlanner.flatten", self, qt, id, location) plan := QueryPlan{} if cat, ok := qt.(*CatQueryTree); ok { inputs := make([]*util.GUID, len(cat.subqueries)) @@ -750,6 +771,7 @@ func (self *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Loca // Transforms Query into QueryTree and flattens to QueryPlan object func (self *QueryPlanner) Plan(query *Query, id *util.GUID, destination *db.Location) (*QueryPlan, error) { + log.Trace("QueryPlanner.Plan", self, query, id, destination) queryTree, err := self.buildTree(query, -1) if err != nil { return nil, err @@ -809,6 +831,7 @@ type RangeQueryTree struct { } func (qt *RangeQueryTree) getLocation(d *db.Database) (*db.Location, error) { + log.Trace("RangeQueryTree", qt, d) slice := d.GetOrCreateSlice(qt.slice) // TODO: this should probably be just GetSlice (no create) fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) if err != nil { diff --git a/query/query.go b/query/query.go index 8d9971b65..aa7a7819a 100644 --- a/query/query.go +++ b/query/query.go @@ -1,6 +1,7 @@ package query import ( + log "github.com/cihub/seelog" "pilosa/db" "pilosa/util" "strings" @@ -28,6 +29,7 @@ type Query struct { } func QueryPlanForPQL(database *db.Database, pql string, destination *db.Location) (*QueryPlan, error) { + log.Trace("QueryPlanFOrPQL", database, pql, destination) tokens, err := Lex(pql) if err != nil { return nil, err @@ -36,6 +38,7 @@ func QueryPlanForPQL(database *db.Database, pql string, destination *db.Location } func QueryForPQL(pql string) (*Query, error) { + log.Trace("QueryForPQL", pql) tokens, err := Lex(pql) if err != nil { return nil, err @@ -44,6 +47,7 @@ func QueryForPQL(pql string) (*Query, error) { } func QueryForTokens(tokens []Token) (*Query, error) { + log.Trace("QueryForTokens", tokens) query, err := Parse(tokens) if err != nil { return nil, err @@ -52,6 +56,7 @@ func QueryForTokens(tokens []Token) (*Query, error) { } func QueryPlanForTokens(database *db.Database, tokens []Token, destination *db.Location) (*QueryPlan, error) { + log.Trace("QueryPlanForTokens", database, tokens, destination) query, err := QueryForTokens(tokens) if err != nil { return nil, err @@ -60,6 +65,7 @@ func QueryPlanForTokens(database *db.Database, tokens []Token, destination *db.L } func QueryPlanForQuery(database *db.Database, query *Query, destination *db.Location) (*QueryPlan, error) { + log.Trace("QueryPlanForQuery", database, query, destination) query_planner := QueryPlanner{Database: database, Query: query} id := util.RandomUUID() query_plan, err := query_planner.Plan(query, &id, destination) @@ -70,6 +76,7 @@ func QueryPlanForQuery(database *db.Database, query *Query, destination *db.Loca } func TokensToFilterStrings(tokens []Token) (string, []string) { + log.Trace("TokensToFilterStrings", tokens) var whole []string var filter string var filters []string diff --git a/transport/tcp.go b/transport/tcp.go index f0c98d9fe..1b2e98d93 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -153,12 +153,14 @@ func (self *TcpTransport) Close() { } func (self *TcpTransport) Send(message *db.Message, host *GUID) { + log.Trace("TcpTransport.Send", message, host) envelope := db.Envelope{message, host} notify.Post("outbox", &envelope) self.outbox <- envelope } func (self *TcpTransport) Receive() *db.Message { + log.Trace("TcpTransport.Receive") message := <-self.inbox notify.Post("inbox", message) return message