diff --git a/ctl/backup.go b/ctl/backup.go index 3f42c0a61..ffcdd21a0 100644 --- a/ctl/backup.go +++ b/ctl/backup.go @@ -18,6 +18,7 @@ import ( "github.com/molecula/featurebase/v3/encoding/proto" "github.com/molecula/featurebase/v3/server" "github.com/pkg/errors" + "github.com/ricochet2200/go-disk-usage/du" "golang.org/x/sync/errgroup" ) @@ -145,6 +146,11 @@ func (cmd *BackupCommand) Run(ctx context.Context) (err error) { return err } + // Ensure there is enough free space + if err := cmd.checkFreeSpace(ctx); err != nil { + return fmt.Errorf("not enough disk space available: %w", err) + } + // Backup schema. if err := cmd.backupSchema(ctx, schema); err != nil { return fmt.Errorf("cannot back up schema: %w", err) @@ -280,6 +286,25 @@ func (cmd *BackupCommand) backupIndexData(ctx context.Context, ii *pilosa.IndexI return g.Wait() } +// checkFreeSpace checks if there is enough space in the output directory to backup data +func (cmd *BackupCommand) checkFreeSpace(ctx context.Context) (err error) { + freeSpace := du.NewDiskUsage(cmd.OutputDir).Free() + var usage pilosa.DiskUsage + if cmd.Index == "" { + usage, err = cmd.client.GetDiskUsage(ctx) + } else { + usage, err = cmd.client.GetIndexUsage(ctx, cmd.Index) + } + if err != nil { + return fmt.Errorf("getting size of data to be backed up: %s", err) + } + + if freeSpace < uint64(usage.Usage) { + return fmt.Errorf("not enough disk space available, free: %v, index usage: %v", freeSpace, usage.Usage) + } + return nil +} + // backupShard backs up a single shard from a single index. func (cmd *BackupCommand) backupShard(ctx context.Context, indexName string, shard uint64) (err error) { nodes, err := cmd.client.FragmentNodes(ctx, indexName, shard) @@ -309,6 +334,7 @@ func (cmd *BackupCommand) backupShardNode(ctx context.Context, indexName string, pilosa.WithClientRetryPeriod(cmd.RetryPeriod), pilosa.WithSerializer(proto.Serializer{})) rc, err := client.ShardReader(ctx, indexName, shard) + if err != nil { return fmt.Errorf("fetching shard reader: %w", err) } diff --git a/go.mod b/go.mod index b237c5b1f..a1fdf677e 100644 --- a/go.mod +++ b/go.mod @@ -40,6 +40,7 @@ require ( github.com/prometheus/prom2json v1.3.1 github.com/rakyll/statik v0.1.7 github.com/remyoudompheng/bigfft v0.0.0-20200410134404-eec4a21b6bb0 // indirect + github.com/ricochet2200/go-disk-usage/du v0.0.0-20210707232629-ac9918953285 github.com/satori/go.uuid v1.2.0 github.com/segmentio/kafka-go v0.4.29 github.com/shirou/gopsutil/v3 v3.22.5 diff --git a/go.sum b/go.sum index d42e374c7..4b4928098 100644 --- a/go.sum +++ b/go.sum @@ -923,6 +923,8 @@ github.com/rakyll/statik v0.1.7/go.mod h1:AlZONWzMtEnMs7W4e/1LURLiI49pIMmp6V9Ung github.com/rcrowley/go-metrics v0.0.0-20181016184325-3113b8401b8a/go.mod h1:bCqnVzQkZxMG4s8nGwiZ5l3QUCyqpo9Y+/ZMZ9VjZe4= github.com/remyoudompheng/bigfft v0.0.0-20200410134404-eec4a21b6bb0 h1:OdAsTTz6OkFY5QxjkYwrChwuRruF69c169dPK26NUlk= github.com/remyoudompheng/bigfft v0.0.0-20200410134404-eec4a21b6bb0/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo= +github.com/ricochet2200/go-disk-usage/du v0.0.0-20210707232629-ac9918953285 h1:d54EL9l+XteliUfUCGsEwwuk65dmmxX85VXF+9T6+50= +github.com/ricochet2200/go-disk-usage/du v0.0.0-20210707232629-ac9918953285/go.mod h1:fxIDly1xtudczrZeOOlfaUvd2OPb2qZAPuWdU2BsBTk= github.com/rogpeppe/fastuuid v0.0.0-20150106093220-6724a57986af/go.mod h1:XWv6SoW27p1b0cqNHllgS5HIMJraePCO15w5zCzIWYg= github.com/rogpeppe/fastuuid v1.2.0/go.mod h1:jVj6XXZzXRy/MSR5jhDC/2q6DgLz+nrA6LYCDYWNEvQ= github.com/rogpeppe/go-internal v1.1.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4= diff --git a/http_handler.go b/http_handler.go index a60ea5e25..dd6f90f80 100644 --- a/http_handler.go +++ b/http_handler.go @@ -538,6 +538,7 @@ func newRouter(handler *Handler) http.Handler { // other ones router.HandleFunc("/internal/mem-usage", handler.chkAuthZ(handler.handleGetMemUsage, authz.Read)).Methods("GET").Name("GetUsage") router.HandleFunc("/internal/disk-usage", handler.chkAuthZ(handler.handleGetDiskUsage, authz.Read)).Methods("GET").Name("GetUsage") + router.HandleFunc("/internal/disk-usage/{index}", handler.chkAuthZ(handler.handleGetDiskUsage, authz.Read)).Methods("GET").Name("GetUsage") router.HandleFunc("/internal/fragment/block/data", handler.chkAuthN(handler.handleGetFragmentBlockData)).Methods("GET").Name("GetFragmentBlockData") router.HandleFunc("/internal/fragment/blocks", handler.chkAuthN(handler.handleGetFragmentBlocks)).Methods("GET").Name("GetFragmentBlocks") router.HandleFunc("/internal/fragment/data", handler.chkAuthN(handler.handleGetFragmentData)).Methods("GET").Name("GetFragmentData") @@ -1163,7 +1164,13 @@ func (h *Handler) handleGetDiskUsage(w http.ResponseWriter, r *http.Request) { return } - use, err := GetDiskUsage(h.api.server.dataDir) + u := h.api.server.dataDir + indexName, ok := mux.Vars(r)["index"] + if ok { + u = fmt.Sprintf("%s/indexes/%s", u, indexName) + } + + use, err := GetDiskUsage(u) if err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) return diff --git a/internal_client.go b/internal_client.go index 0056e4f31..f834f2ed0 100644 --- a/internal_client.go +++ b/internal_client.go @@ -2477,3 +2477,59 @@ func (c *InternalClient) OAuthConfig() (rsp oauth2.Config, err error) { return rsp, nil } + +// GetDiskUsage gets the size of data directory across all nodes. +func (c *InternalClient) GetDiskUsage(ctx context.Context) (DiskUsage, error) { + span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.GetDiskUsage") + defer span.Finish() + return c.getDiskUsage(ctx, "") +} + +// GetIndexUsage gets the size of an index across all nodes. +func (c *InternalClient) GetIndexUsage(ctx context.Context, index string) (DiskUsage, error) { + span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.GetIndexUsage") + defer span.Finish() + return c.getDiskUsage(ctx, index) +} + +// getDiskUsage returns size of data directory if index is zero value. +func (c *InternalClient) getDiskUsage(ctx context.Context, index string) (DiskUsage, error) { + nodes, err := c.Nodes(ctx) + if err != nil { + return DiskUsage{}, fmt.Errorf("getting nodes: %s", err) + } + + var sum DiskUsage + for _, node := range nodes { + path := "/internal/disk-usage" + if index != "" { + path = path + "/" + index + } + u := uriPathToURL(&node.URI, path) + + req, err := http.NewRequest("GET", u.String(), nil) + if err != nil { + return DiskUsage{}, errors.Wrap(err, "creating request") + } + + req.Header.Set("User-Agent", "pilosa/"+Version) + req.Header.Set("Accept", "application/json") + AddAuthToken(ctx, &req.Header) + + // Execute request. + resp, err := c.executeRequest(req.WithContext(ctx)) + if err != nil { + return DiskUsage{}, err + } + defer resp.Body.Close() + + var rsp DiskUsage + if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil { + return DiskUsage{}, fmt.Errorf("json decode: %s", err) + } + + sum.Usage += rsp.Usage + } + + return sum, nil +} diff --git a/internal_client_test.go b/internal_client_test.go index d768b465c..ee3040d32 100644 --- a/internal_client_test.go +++ b/internal_client_test.go @@ -21,7 +21,9 @@ import ( "github.com/molecula/featurebase/v3/pql" "github.com/molecula/featurebase/v3/server" "github.com/molecula/featurebase/v3/test" + "github.com/molecula/featurebase/v3/vprint" "github.com/pkg/errors" + "github.com/ricochet2200/go-disk-usage/du" ) // Test distributed TopN Row count across 3 nodes. @@ -1029,6 +1031,16 @@ func TestClient_ImportValue(t *testing.T) { t.Fatalf("unexpected values: got max=%v, count=%v; expected max=40, cnt=1", vc.Val, vc.Count) } + // Calculate Data Usage before Import + preDUsage, err := c.GetDiskUsage(context.Background()) + if err != nil { + t.Fatal(err) + } + preIUsage, err := c.GetIndexUsage(context.Background(), "i") + if err != nil { + t.Fatal(err) + } + // Send import request. req = &pilosa.ImportValueRequest{ Index: "i", @@ -1040,6 +1052,28 @@ func TestClient_ImportValue(t *testing.T) { t.Fatal(err) } + // Check Equivalent growth in data directory and index + postDUsage, err := c.GetDiskUsage(context.Background()) + if err != nil { + t.Fatal(err) + } + postIUsage, err := c.GetIndexUsage(context.Background(), "i") + if err != nil { + t.Fatal(err) + } + + freeSpace := du.NewDiskUsage("/dev/disk1s5").Free() + vprint.VV("freespace: %+v", freeSpace) + + if (postDUsage.Usage - preDUsage.Usage) != (postIUsage.Usage - preIUsage.Usage) { + t.Errorf("expected size of data directory to grow the same amount as size of index: Before Import: disk usage: %v, index usage: %v, After Import: disk usage: %v, index usage: %v", preDUsage.Usage, preIUsage.Usage, postDUsage.Usage, preIUsage.Usage) + return + } + if postDUsage.Usage <= postIUsage.Usage { + t.Errorf("expected disk usage to be greater than index usage") + return + } + // Verify Sum. if resp, err := c.Query(context.Background(), "i", &pilosa.QueryRequest{Query: `Sum(field=f)`}); err != nil { t.Fatal(err)