diff --git a/api.go b/api.go index dfbd33e01..9d6916558 100644 --- a/api.go +++ b/api.go @@ -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 } diff --git a/client.go b/client.go index 42ec51ba0..0cae10062 100644 --- a/client.go +++ b/client.go @@ -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 +} diff --git a/http/client.go b/http/client.go index 5f4e25b8a..3e2df94b4 100644 --- a/http/client.go +++ b/http/client.go @@ -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() diff --git a/http/handler.go b/http/handler.go index 98634111d..bb8cdf9f6 100644 --- a/http/handler.go +++ b/http/handler.go @@ -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) { diff --git a/server/handler_test.go b/server/handler_test.go index 47e4698b0..4096dc970 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -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) + } } })