diff --git a/api.go b/api.go index 4fc631278..2038679c1 100644 --- a/api.go +++ b/api.go @@ -158,7 +158,7 @@ func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, er if err != nil { return QueryResponse{}, errors.Wrap(err, "parsing") } - defer api.tracker.Finish(api.tracker.Start(req.Query)) + defer api.tracker.Finish(api.tracker.Start(req.Query, api.server.nodeID), time.Now()) execOpts := &execOptions{ Remote: req.Remote, Profile: req.Profile, @@ -2050,6 +2050,14 @@ func (api *API) ActiveQueries(ctx context.Context) ([]ActiveQueryStatus, error) return api.tracker.ActiveQueries(), nil } +func (api *API) PastQueries(ctx context.Context) ([]PastQueryStatus, error) { + if err := api.validate(apiPastQueries); err != nil { + return nil, errors.Wrap(err, "validating api method") + } + x := api.tracker.PastQueries() + return x, nil +} + // TranslateIndexDB is an internal function to load the index keys database // rd is a boltdb file. func (api *API) TranslateIndexDB(ctx context.Context, indexName string, partitionID int, rd io.Reader) error { @@ -2124,6 +2132,7 @@ const ( apiTransactions apiGetTransaction apiActiveQueries + apiPastQueries ) var methodsCommon = map[apiMethod]struct{}{ @@ -2164,4 +2173,5 @@ var methodsNormal = map[apiMethod]struct{}{ apiTransactions: {}, apiGetTransaction: {}, apiActiveQueries: {}, + apiPastQueries: {}, } diff --git a/http/handler.go b/http/handler.go index f9339a751..69af4d81a 100644 --- a/http/handler.go +++ b/http/handler.go @@ -397,8 +397,10 @@ func newRouter(handler *Handler) http.Handler { router.HandleFunc("/transaction/{id}/finish", handler.handlePostFinishTransaction).Methods("POST").Name("PostFinishTransaction") router.HandleFunc("/transactions", handler.handleGetTransactions).Methods("GET").Name("GetTransactions") router.HandleFunc("/queries", handler.handleGetActiveQueries).Methods("GET").Name("GetActiveQueries") + router.HandleFunc("/query-history", handler.handleGetPastQueries).Methods("GET").Name("GetPastQueries") router.HandleFunc("/version", handler.handleGetVersion).Methods("GET").Name("GetVersion") + // /ui endpoints are for UI use; they may change at any time. router.HandleFunc("/ui/usage", handler.handleGetUsage).Methods("GET").Name("GetUsage") router.HandleFunc("/ui/transaction", handler.handleGetTransactionList).Methods("GET").Name("GetTransactionList") router.HandleFunc("/ui/transaction/", handler.handleGetTransactionList).Methods("GET").Name("GetTransactionList") @@ -1184,6 +1186,19 @@ func (h *Handler) handleGetActiveQueries(w http.ResponseWriter, r *http.Request) } } +func (h *Handler) handleGetPastQueries(w http.ResponseWriter, r *http.Request) { + queries, err := h.api.PastQueries(r.Context()) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "application/json") + if err := json.NewEncoder(w).Encode(queries); err != nil { + h.logger.Printf("encoding GetActiveQueries response: %s", err) + } + +} + type postIndexAttrDiffRequest struct { Blocks []pilosa.AttrBlock `json:"blocks"` } diff --git a/tracker.go b/tracker.go index af5788137..66bf0e924 100644 --- a/tracker.go +++ b/tracker.go @@ -22,22 +22,40 @@ import ( type ActiveQueryStatus struct { Query string `json:"query"` + Node string `json:"node"` Age time.Duration `json:"age"` } +type PastQueryStatus struct { + Query string `json:"query"` + Node string `json:"node"` + Age time.Duration `json:"age"` + Runtime time.Duration `json:"runtime"` +} + type activeQuery struct { query string + node string started time.Time } +type pastQuery struct { + query string + node string + started time.Time + runtime time.Duration +} + type queryStatusUpdate struct { - q *activeQuery - end bool + q *activeQuery + end bool + endTime time.Time } type queryTracker struct { updates chan<- queryStatusUpdate checks chan<- chan<- []*activeQuery + history map[pastQuery]struct{} // TODO not in memory wg sync.WaitGroup stop chan struct{} } @@ -46,9 +64,11 @@ func newQueryTracker() *queryTracker { done := make(chan struct{}) updates := make(chan queryStatusUpdate, 128) checks := make(chan chan<- []*activeQuery) + history := make(map[pastQuery]struct{}) tracker := &queryTracker{ updates: updates, checks: checks, + history: history, stop: done, } tracker.wg.Add(1) @@ -64,6 +84,8 @@ func newQueryTracker() *queryTracker { delete(activeQueries, update.q) } else { activeQueries[update.q] = struct{}{} + pq := pastQuery{update.q.query, update.q.node, update.q.started, update.endTime.Sub(update.q.started)} // TODO move end time to api.go + tracker.history[pq] = struct{}{} } case check := <-checks: out := make([]*activeQuery, len(activeQueries)) @@ -82,15 +104,15 @@ func newQueryTracker() *queryTracker { return tracker } -func (t *queryTracker) Start(query string) *activeQuery { +func (t *queryTracker) Start(query, nodeID string) *activeQuery { now := time.Now() - q := &activeQuery{query, now} - t.updates <- queryStatusUpdate{q, false} + q := &activeQuery{query, nodeID, now} + t.updates <- queryStatusUpdate{q, false, time.Time{}} return q } -func (t *queryTracker) Finish(q *activeQuery) { - t.updates <- queryStatusUpdate{q, true} +func (t *queryTracker) Finish(q *activeQuery, endTime time.Time) { + t.updates <- queryStatusUpdate{q, true, endTime} } func (t *queryTracker) ActiveQueries() []ActiveQueryStatus { @@ -114,11 +136,41 @@ func (t *queryTracker) ActiveQueries() []ActiveQueryStatus { now := time.Now() out := make([]ActiveQueryStatus, len(queries)) for i, v := range queries { - out[i] = ActiveQueryStatus{v.query, now.Sub(v.started)} + out[i] = ActiveQueryStatus{v.query, v.node, now.Sub(v.started)} } return out } +func (t *queryTracker) PastQueries() []PastQueryStatus { + queries := make([]pastQuery, 0, len(t.history)) + for pq, _ := range t.history { + queries = append(queries, pq) + } + // TODO use sort.Sort + sort.Slice(queries, func(i, j int) bool { + switch { + case queries[i].started.Before(queries[j].started): + return true + case queries[i].started.After(queries[j].started): + return false + case queries[i].query < queries[j].query: + return true + case queries[i].query > queries[j].query: + return false + default: + return false + } + }) + + now := time.Now() + out := make([]PastQueryStatus, len(queries)) + for i, v := range queries { + out[i] = PastQueryStatus{v.query, v.node, now.Sub(v.started), v.runtime} + } + return out + +} + func (t *queryTracker) Stop() { close(t.stop) t.wg.Wait()