Merge pull request #907 from alanbernstein/cluster-data-endpoint

Add single-node 'bytesOnDisk' object to /status response
This commit is contained in:
alanbernstein 2020-10-01 10:42:27 -05:00 committed by GitHub
commit 11e4bbfc67
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
3 changed files with 116 additions and 0 deletions

66
api.go
View file

@ -25,6 +25,8 @@ import (
"io/ioutil"
"math"
"net/url"
"os"
"path"
"sort"
"strconv"
"strings"
@ -780,6 +782,70 @@ func (api *API) Node() *Node {
return &node
}
// Usage gets the disk usage per index
func (api *API) Usage() (map[string]int64, int64, error) {
indexSizes := make(map[string]int64)
var totalSize int64
dirName, err := expandDirName(api.server.dataDir)
if err != nil {
return indexSizes, totalSize, errors.Wrap(err, "expanding data directory")
}
dir, err := os.Open(dirName)
if err != nil {
return indexSizes, totalSize, errors.Wrap(err, "opening data directory")
}
defer dir.Close()
files, err := dir.Readdir(-1)
if err != nil {
return indexSizes, totalSize, errors.Wrap(err, "reading data directory")
}
for _, file := range files {
if !file.IsDir() {
continue
}
fullName := path.Join(dirName, file.Name())
indexSizes[file.Name()], err = diskUsage(fullName)
if err != nil {
break
}
totalSize += indexSizes[file.Name()]
}
return indexSizes, totalSize, nil
}
func diskUsage(fname string) (int64, error) {
var size int64
dir, err := os.Open(fname)
if err != nil {
return 0, errors.Wrap(err, "opening data subdirectory")
}
defer dir.Close()
files, err := dir.Readdir(-1)
if err != nil {
return 0, errors.Wrap(err, "reading data subdirectory")
}
for _, file := range files {
if file.IsDir() {
sz, err := diskUsage(path.Join(fname, file.Name()))
if err != nil {
return 0, err
}
size += sz
} else {
size += file.Size()
}
}
return size, nil
}
// RecalculateCaches forces all TopN caches to be updated.
// This is done internally within a TopN query, but a user may want to do it ahead of time?
func (api *API) RecalculateCaches(ctx context.Context) error {

View file

@ -392,6 +392,8 @@ func newRouter(handler *Handler) http.Handler {
router.HandleFunc("/queries", handler.handleGetActiveQueries).Methods("GET").Name("GetActiveQueries")
router.HandleFunc("/version", handler.handleGetVersion).Methods("GET").Name("GetVersion")
router.HandleFunc("/ui/usage", handler.handleGetUsage).Methods("GET").Name("GetUsage")
// /internal endpoints are for internal use only; they may change at any time.
// DO NOT rely on these for external applications!
router.HandleFunc("/internal/cluster/message", handler.handlePostClusterMessage).Methods("POST").Name("PostClusterMessage")
@ -657,6 +659,40 @@ func (h *Handler) handlePostSchema(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusNoContent)
}
// handleGetUsage handles GET /ui/usage requests.
func (h *Handler) handleGetUsage(w http.ResponseWriter, r *http.Request) {
if !validHeaderAcceptJSON(r.Header) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
usageIndexes, usageTotal, err := h.api.Usage()
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
}
disk := diskUsage{
Total: usageTotal,
Indexes: usageIndexes,
}
usage := getUsageResponse{
Disk: disk,
}
w.Header().Set("Content-Type", "application/json")
if err := json.NewEncoder(w).Encode(usage); err != nil {
h.logger.Printf("write status response error: %s", err)
}
}
type getUsageResponse struct {
Disk diskUsage `json:"bytesOnDisk"`
}
type diskUsage struct {
Total int64 `json:"total"`
Indexes map[string]int64 `json:"indexes"`
}
// handleGetStatus handles GET /status requests.
func (h *Handler) handleGetStatus(w http.ResponseWriter, r *http.Request) {
if !validHeaderAcceptJSON(r.Header) {

View file

@ -381,6 +381,20 @@ func TestHandler_Endpoints(t *testing.T) {
}
})
t.Run("UI/usage", func(t *testing.T) {
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/ui/usage", nil))
if w.Code != gohttp.StatusOK {
t.Fatalf("unexpected status code: %d", w.Code)
}
ret := mustJSONDecode(t, w.Body)
usage := ret["bytesOnDisk"].(map[string]interface{})
indexes := usage["indexes"].(map[string]interface{})
if len(indexes) != 2 {
t.Fatalf("wrong length index size list: %#v", indexes)
}
})
t.Run("Metrics", func(t *testing.T) {
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/metrics", nil))