mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-09 04:17:51 +00:00
Add basic implementation of query history endpoint
This commit is contained in:
parent
55d0c12703
commit
c94d242097
3 changed files with 86 additions and 9 deletions
12
api.go
12
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: {},
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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"`
|
||||
}
|
||||
|
|
|
|||
68
tracker.go
68
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()
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue