diff --git a/api.go b/api.go index 9140153a1..2bd8eb2e4 100644 --- a/api.go +++ b/api.go @@ -161,7 +161,7 @@ func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, er } if !req.Remote { - defer api.tracker.Finish(api.tracker.Start(req.Query, api.server.nodeID, req.Index, start)) + defer api.tracker.Finish(api.tracker.Start(req.Query, req.SQLQuery, api.server.nodeID, req.Index, start)) } // TODO can we get rid of exec options and pass the QueryRequest directly to executor? execOpts := &execOptions{ diff --git a/handler.go b/handler.go index d8e370d56..b50e6e946 100644 --- a/handler.go +++ b/handler.go @@ -29,6 +29,9 @@ type QueryRequest struct { // The query string to parse and execute. Query string + // The SQL source query, if applicable. + SQLQuery string + // The shards to include in the query execution. // If empty, all shards are included. Shards []uint64 diff --git a/http/handler.go b/http/handler.go index bc2fbb1ab..a6dfd5828 100644 --- a/http/handler.go +++ b/http/handler.go @@ -1182,7 +1182,7 @@ func (h *Handler) handleGetActiveQueries(w http.ResponseWriter, r *http.Request) } } for i, q := range queries { - _, err := fmt.Fprintf(w, "%*s%q\n", -(maxlen + 2), durations[i], q.Query) + _, err := fmt.Fprintf(w, "%*s%q\n", -(maxlen + 2), durations[i], q.PQL) if err != nil { h.logger.Printf("sending GetActiveQueries response: %s", err) return diff --git a/server/handler_test.go b/server/handler_test.go index 67b9af8e5..705cfdb5c 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -37,6 +37,7 @@ import ( "github.com/pilosa/pilosa/v2/encoding/proto" "github.com/pilosa/pilosa/v2/http" "github.com/pilosa/pilosa/v2/pql" + pb "github.com/pilosa/pilosa/v2/proto" "github.com/pilosa/pilosa/v2/server" "github.com/pilosa/pilosa/v2/test" ) @@ -1502,6 +1503,16 @@ func TestQueryHistory(t *testing.T) { test.Do(t, "POST", cmd.URL()+"/index/i0", "") test.Do(t, "POST", cmd.URL()+"/index/i0/field/f0", "") + + gh := server.NewGRPCHandler(cmd.API) + _, err = gh.QuerySQLUnary(context.Background(), &pb.QuerySQLRequest{ + Sql: `select * from i0`, + }) + + if err != nil { + t.Fatalf("QuerySQLUnary failed: %v", err) + } + test.Do(t, "POST", cmd.URL()+"/index/i0/query", "Set(0, f0=0)") test.Do(t, "POST", cmd.URL()+"/index/i0/query", "Set(3000000, f0=0)") test.Do(t, "POST", cmd.URL()+"/index/i0/query", "TopN(f0)") @@ -1511,7 +1522,7 @@ func TestQueryHistory(t *testing.T) { t.Fatalf("unexpected status code: %d", w.Code) } - ret := make([]pilosa.PastQueryStatus, 3) + ret := make([]pilosa.PastQueryStatus, 4) b, err := ioutil.ReadAll(w.Body) if err != nil { t.Fatalf("reading: %v", err) @@ -1522,10 +1533,10 @@ func TestQueryHistory(t *testing.T) { } // verify result length - if len(ret) != 3 { + if len(ret) != 4 { // each set query executes on both nodes once // topn query gets added to history on node0 once, node1 twice - t.Fatalf("expected list of length 3, got %d", len(ret)) + t.Fatalf("expected list of length 4, got %d\n%+v", len(ret), ret) } // verify sort order @@ -1543,8 +1554,14 @@ func TestQueryHistory(t *testing.T) { if ret[0].Node != cluster.GetNode(0).Server.NodeID() { t.Fatalf("response value for 'Node' was '%s', expected '%s'", ret[0].Node, cluster.GetNode(0).Server.NodeID()) } - if ret[0].Query != "TopN(f0)" { - t.Fatalf("response value for 'Query' was '%s', expected 'TopN(f0)'", ret[0].Query) + if ret[3].PQL != "Extract(All(),Rows(f0))" { + t.Fatalf("response value for 'PQL' was '%s', expected 'Extract(All(),Rows(f0))'", ret[0].PQL) + } + if ret[3].SQL != "select * from i0" { + t.Fatalf("response value for 'SQL' was '%s', expected 'select * from i0'", ret[0].SQL) + } + if ret[0].PQL != "TopN(f0)" { + t.Fatalf("response value for 'PQL' was '%s', expected 'TopN(f0)'", ret[0].PQL) } } diff --git a/sql/mapper.go b/sql/mapper.go index a43dba06d..eb6c23f0c 100644 --- a/sql/mapper.go +++ b/sql/mapper.go @@ -40,6 +40,7 @@ type MappedSQL struct { Statement sqlparser.Statement Mask QueryMask Tables []string + SQL string } // Mapper is responsible for mapping a SQL query to structure representation @@ -109,5 +110,6 @@ func (m *Mapper) MapSQL(sql string) (*MappedSQL, error) { Statement: stmt, Mask: qm, Tables: tableNames, + SQL: sql, }, nil } diff --git a/sql/select.go b/sql/select.go index b54eee0b4..b403f80c9 100644 --- a/sql/select.go +++ b/sql/select.go @@ -43,6 +43,7 @@ func NewSelectHandler(api *pilosa.API) *SelectHandler { // Handle executes mapped SQL func (s *SelectHandler) Handle(ctx context.Context, mapped *MappedSQL) (pproto.ToRowser, error) { stmt, ok := mapped.Statement.(*sqlparser.Select) + if !ok { return nil, fmt.Errorf("statement is not type select: %T", mapped.Statement) } @@ -50,7 +51,7 @@ func (s *SelectHandler) Handle(ctx context.Context, mapped *MappedSQL) (pproto.T if err != nil { return nil, errors.Wrap(err, "mapping select") } - return s.execMappingResult(ctx, mr) + return s.execMappingResult(ctx, mr, mapped.SQL) } func (s *SelectHandler) mapSelect(ctx context.Context, selectStmt *sqlparser.Select, qm QueryMask) (*MappingResult, error) { @@ -74,12 +75,12 @@ func (s *SelectHandler) mapSelect(ctx context.Context, selectStmt *sqlparser.Sel return mr, nil } -func (s *SelectHandler) execMappingResult(ctx context.Context, mr *MappingResult) (pproto.ToRowser, error) { +func (s *SelectHandler) execMappingResult(ctx context.Context, mr *MappingResult, sql string) (pproto.ToRowser, error) { if mr.Query == "" { return nil, errors.New("no pql query created") } - resp, err := s.api.Query(ctx, &pilosa.QueryRequest{Index: mr.IndexName, Query: mr.Query}) + resp, err := s.api.Query(ctx, &pilosa.QueryRequest{Index: mr.IndexName, Query: mr.Query, SQLQuery: sql}) if err != nil { return nil, errors.Wrap(err, "doing pql query") } diff --git a/tracker.go b/tracker.go index 359850005..bec0cf4db 100644 --- a/tracker.go +++ b/tracker.go @@ -21,14 +21,16 @@ import ( ) type ActiveQueryStatus struct { - Query string `json:"query"` + PQL string `json:"PQL"` + SQL string `json:"SQL,omitempty"` Node string `json:"node"` Index string `json:"index"` Age time.Duration `json:"age"` } type PastQueryStatus struct { - Query string `json:"query"` + PQL string `json:"PQL"` + SQL string `json:"SQL,omitempty"` Node string `json:"nodeID"` Index string `json:"index"` Start time.Time `json:"start"` @@ -36,14 +38,16 @@ type PastQueryStatus struct { } type activeQuery struct { - query string + PQL string + SQL string node string index string started time.Time } type pastQuery struct { - query string + PQL string + SQL string node string index string started time.Time @@ -123,7 +127,7 @@ func newQueryTracker(historyLength int) *queryTracker { select { case update := <-updates: if update.end { - pq := pastQuery{update.q.query, update.q.node, update.q.index, update.q.started, update.endTime.Sub(update.q.started)} + pq := pastQuery{update.q.PQL, update.q.SQL, update.q.node, update.q.index, update.q.started, update.endTime.Sub(update.q.started)} tracker.history.add(pq) delete(activeQueries, update.q) } else { @@ -146,8 +150,8 @@ func newQueryTracker(historyLength int) *queryTracker { return tracker } -func (t *queryTracker) Start(query, nodeID, index string, start time.Time) *activeQuery { - q := &activeQuery{query, nodeID, index, start} +func (t *queryTracker) Start(pql, sql, nodeID, index string, start time.Time) *activeQuery { + q := &activeQuery{pql, sql, nodeID, index, start} t.updates <- queryStatusUpdate{q, false, time.Time{}} return q } @@ -166,9 +170,9 @@ func (t *queryTracker) ActiveQueries() []ActiveQueryStatus { return true case queries[i].started.After(queries[j].started): return false - case queries[i].query < queries[j].query: + case queries[i].PQL < queries[j].PQL: return true - case queries[i].query > queries[j].query: + case queries[i].PQL > queries[j].PQL: return false default: return false @@ -177,7 +181,7 @@ func (t *queryTracker) ActiveQueries() []ActiveQueryStatus { now := time.Now() out := make([]ActiveQueryStatus, len(queries)) for i, v := range queries { - out[i] = ActiveQueryStatus{v.query, v.node, v.index, now.Sub(v.started)} + out[i] = ActiveQueryStatus{v.PQL, v.SQL, v.node, v.index, now.Sub(v.started)} } return out } @@ -186,7 +190,7 @@ func (t *queryTracker) PastQueries() []PastQueryStatus { queries := t.history.slice() out := make([]PastQueryStatus, len(queries)) for i, v := range queries { - out[i] = PastQueryStatus{v.query, v.node, v.index, v.started, v.runtime} + out[i] = PastQueryStatus{v.PQL, v.SQL, v.node, v.index, v.started, v.runtime} } return out diff --git a/tracker_test.go b/tracker_test.go index a5d8df418..569a7b609 100644 --- a/tracker_test.go +++ b/tracker_test.go @@ -47,11 +47,11 @@ func TestRingBuffer(t *testing.T) { } for n, q := range buffer.slice() { - if q.query != tests[k].queries[n] { - t.Fatalf("test[%d], buffer[%d] expected querystring '%s', found '%s'", k, n, tests[k].queries[n], q.query) + if q.PQL != tests[k].queries[n] { + t.Fatalf("test[%d], buffer[%d] expected querystring '%s', found '%s'", k, n, tests[k].queries[n], q.PQL) } } - buffer.add(pastQuery{query: fmt.Sprintf("%d", k)}) + buffer.add(pastQuery{PQL: fmt.Sprintf("%d", k)}) } } @@ -63,13 +63,13 @@ func TestQueryTracker(t *testing.T) { t.Fatalf("expected no active queries; found %v", queries) } - qs := tracker.Start("test query", "node0", "i", time.Now()) + qs := tracker.Start("test query", "test SQL", "node0", "i", time.Now()) var queries []ActiveQueryStatus for len(queries) < 1 { queries = tracker.ActiveQueries() } - if len(queries) > 1 || queries[0].Query != "test query" { + if len(queries) > 1 || queries[0].PQL != "test query" { t.Fatalf("unexpected queries: %v", queries) }