Collect size-on-disk usage data from all nodes

This commit is contained in:
Alan Bernstein 2020-10-06 02:45:30 -05:00 • committed by Jason E. Aten
parent e1e0b298f5
commit 2a44bdd25c
5 changed files with 108 additions and 35 deletions

62
api.go
View file

@ -793,26 +793,43 @@ 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)
// NodeUsage represents all usage measurements for one node.
type NodeUsage struct {
Disk DiskUsage `json:"bytesOnDisk"`
}
// DiskUsage represents the storage space used on disk by one node.
type DiskUsage struct {
Total int64 `json:"total"`
Indexes map[string]int64 `json:"indexes"`
}
// Usage gets the disk usage per index, in a map[nodeID]NodeUsage
func (api *API) Usage(ctx context.Context, remote bool) (map[string]NodeUsage, error) {
span, _ := tracing.StartSpanFromContext(ctx, "API.Usage")
defer span.Finish()
nodeUsages := make(map[string]NodeUsage)
var totalSize int64
// Open storage directory.
dirName, err := expandDirName(api.server.dataDir)
if err != nil {
return indexSizes, totalSize, errors.Wrap(err, "expanding data directory")
return nodeUsages, errors.Wrap(err, "expanding data directory")
}
dir, err := os.Open(dirName)
if err != nil {
return indexSizes, totalSize, errors.Wrap(err, "opening data directory")
return nodeUsages, 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")
return nodeUsages, errors.Wrap(err, "reading data directory")
}
// Read size on disk for each index directory.
indexSizes := make(map[string]int64)
for _, file := range files {
if !file.IsDir() {
continue
@ -821,17 +838,40 @@ func (api *API) Usage() (map[string]int64, int64, error) {
continue
}
fullName := path.Join(dirName, file.Name())
indexSizes[file.Name()], err = diskUsage(fullName)
indexSizes[file.Name()], err = directoryUsage(fullName)
if err != nil {
break
return nodeUsages, errors.Wrap(err, "getting disk usage")
}
totalSize += indexSizes[file.Name()]
}
return indexSizes, totalSize, nil
// Insert into result.
nodeUsage := NodeUsage{
Disk: DiskUsage{
Total: totalSize,
Indexes: indexSizes,
},
}
nodeUsages[api.server.nodeID] = nodeUsage
// Collect size on disk from remote nodes
if !remote {
nodes := api.cluster.Nodes()
for _, node := range nodes {
if node.ID == api.server.nodeID {
continue
}
nodeUsage, err := api.server.defaultClient.GetNodeUsage(ctx, &node.URI)
if err != nil {
return nil, errors.Wrapf(err, "collecting disk usage from %s", node.URI)
}
nodeUsages[node.ID] = nodeUsage[node.ID]
}
}
return nodeUsages, nil
}
func diskUsage(fname string) (int64, error) {
func directoryUsage(fname string) (int64, error) {
var size int64
dir, err := os.Open(fname)
@ -847,7 +887,7 @@ func diskUsage(fname string) (int64, error) {
for _, file := range files {
if file.IsDir() {
sz, err := diskUsage(path.Join(fname, file.Name()))
sz, err := directoryUsage(path.Join(fname, file.Name()))
if err != nil {
return 0, err
}

View file

@ -81,6 +81,8 @@ type InternalClient interface {
FinishTransaction(ctx context.Context, id string) (*Transaction, error)
Transactions(ctx context.Context) (map[string]*Transaction, error)
GetTransaction(ctx context.Context, id string) (*Transaction, error)
GetNodeUsage(ctx context.Context, uri *URI) (map[string]NodeUsage, error)
}
//===============
@ -227,3 +229,7 @@ func (n nopInternalClient) Transactions(ctx context.Context) (map[string]*Transa
func (n nopInternalClient) GetTransaction(ctx context.Context, id string) (*Transaction, error) {
return nil, nil
}
func (n nopInternalClient) GetNodeUsage(ctx context.Context, uri *URI) (map[string]NodeUsage, error) {
return nil, nil
}

View file

@ -1247,6 +1247,37 @@ func (c *InternalClient) TranslateIDsNode(ctx context.Context, uri *pilosa.URI,
return tkresp.Keys, nil
}
// GetNodeUsage retrieves the size-on-disk information for the specified node.
func (c *InternalClient) GetNodeUsage(ctx context.Context, uri *pilosa.URI) (map[string]pilosa.NodeUsage, error) {
u := uri.Path("/ui/usage?remote=true")
req, err := http.NewRequest("GET", u, nil)
if err != nil {
return nil, errors.Wrap(err, "creating request")
}
req.Header.Set("Accept", "application/json")
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
// Execute request against the host.
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, err
}
defer resp.Body.Close()
// Read body and unmarshal response.
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return nil, errors.Wrap(err, "reading")
}
nodeUsages := make(map[string]pilosa.NodeUsage) // map of size 1
if err := json.Unmarshal(body, &nodeUsages); err != nil {
return nil, fmt.Errorf("unmarshal response: %s", err)
}
return nodeUsages, nil
}
func (c *InternalClient) Transactions(ctx context.Context) (map[string]*pilosa.Transaction, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.Transactions")
defer span.Finish()

View file

@ -665,34 +665,25 @@ func (h *Handler) handleGetUsage(w http.ResponseWriter, r *http.Request) {
http.Error(w, "JSON only acceptable response", http.StatusNotAcceptable)
return
}
usageIndexes, usageTotal, err := h.api.Usage()
q := r.URL.Query()
remoteStr := q.Get("remote")
var remote bool
if remoteStr == "true" {
remote = true
}
nodeUsages, err := h.api.Usage(r.Context(), remote)
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 {
if err := json.NewEncoder(w).Encode(nodeUsages); 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

@ -385,13 +385,18 @@ func TestHandler_Endpoints(t *testing.T) {
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/ui/usage", nil))
if w.Code != gohttp.StatusOK {
fmt.Printf("%+v\n", w.Body)
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)
nodeUsages := make(map[string]pilosa.NodeUsage)
if err := json.Unmarshal(w.Body.Bytes(), &nodeUsages); err != nil {
t.Fatalf("unmarshal")
}
for _, nodeUsage := range nodeUsages {
if len(nodeUsage.Disk.Indexes) != 2 {
t.Fatalf("wrong length index size list: %#v", nodeUsage.Disk.Indexes)
}
}
})