From ac09c11badb630a375d65b930e606f17e3e14662 Mon Sep 17 00:00:00 2001 From: Nia Weiss Date: Tue, 26 Jan 2021 12:30:13 -0500 Subject: [PATCH 01/25] Invoke precalls directly in count operations This changes Count(Precall()) operations to execute the precall directly inside of the count operation, bypassing the transformation to a Precomputed() call. Eliminating the Precomputed() step causes Count(Distinct()) to work properly on negative integers. --- executor.go | 52 ++++++++++++++++++++++++++---------------------- executor_test.go | 24 ++++++++++++++++++++++ pql/ast.go | 9 ++++++--- 3 files changed, 58 insertions(+), 27 deletions(-) diff --git a/executor.go b/executor.go index b1a0ae1eb..7e046c11c 100644 --- a/executor.go +++ b/executor.go @@ -556,11 +556,22 @@ func (e *executor) execute(ctx context.Context, qcx *Qcx, index string, q *pql.Q // about the positive values, because only positive values // are valid column IDs. So we don't actually eat top-level // pre calls. - err := e.handlePreCallChildren(ctx, qcx, index, call, shards, opt) - if err != nil { - return nil, err + if call.Name == "Count" { + // Handle count specially, skipping the level directly underneath it. + for _, child := range call.Children { + err := e.handlePreCallChildren(ctx, qcx, index, child, shards, opt) + if err != nil { + return nil, err + } + } + } else { + err := e.handlePreCallChildren(ctx, qcx, index, call, shards, opt) + if err != nil { + return nil, err + } } var v interface{} + var err error // Top-level calls don't need to precompute cross-index things, // because we can just pick whatever index we want, but we // still need to handle them. Since everything else was @@ -4618,28 +4629,21 @@ func (e *executor) executeCount(ctx context.Context, qcx *Qcx, index string, c * child := c.Children[0] - // If the child is precomputed, we'll bypass mapreduce, ignore - // shards, and just count the number of bits. - if child.Name == "Precomputed" { - count := uint64(0) - for _, irow := range child.Precomputed { - switch row := irow.(type) { - case *Row: - for _, seg := range row.segments { - count += seg.n - } - case SignedRow: - for _, seg := range row.Pos.segments { - count += seg.n - } - for _, seg := range row.Neg.segments { - count += seg.n - } - default: - return 0, errors.Errorf("unexpected precomputed value type inside count: %+v", row) - } + // If the child is distinct/similar, execute it directly here and count the result. + if child.Type == pql.PrecallGlobal { + result, err := e.executeCall(ctx, qcx, index, child, shards, opt) + if err != nil { + return 0, err + } + + switch row := result.(type) { + case *Row: + return row.Count(), nil + case SignedRow: + return row.Pos.Count() + row.Neg.Count(), nil + default: + return 0, errors.Errorf("cannot count result of type %T from call %q", row, child.String()) } - return count, nil } // Execute calls in bulk on each remote node and merge. diff --git a/executor_test.go b/executor_test.go index 6a202ca3a..ae86ba7d6 100644 --- a/executor_test.go +++ b/executor_test.go @@ -6569,6 +6569,30 @@ func TestExecutor_Execute_CountDistinct(t *testing.T) { }) } +// Ensure that Count(Distinct()) works with negative numbers. +func TestExecutor_Execute_CountDistinctSigned(t *testing.T) { + c := test.MustRunCluster(t, 2) + defer c.Close() + + c.CreateField(t, "i", pilosa.IndexOptions{}, "ints", + pilosa.OptFieldTypeInt(-100, 100), + ) + + c.Query(t, "i", fmt.Sprintf(` + Set(0, ints=1) + Set(%d, ints=-1) + `, ShardWidth)) + + resp := c.Query(t, "i", "Count(Distinct(field=ints))") + cnt, ok := resp.Results[0].(uint64) + if !ok { + t.Fatalf("invalid response type, expected: uint64, got: %T", resp.Results[0]) + } + if cnt != 2 { + t.Fatalf("invalid result, expected: 2, got: %v", cnt) + } +} + // Ensure that a top-level, bare distinct on multiple nodes // is handled correctly. func TestExecutor_BareDistinct(t *testing.T) { diff --git a/pql/ast.go b/pql/ast.go index 54ae2dee7..a4b80f4ae 100644 --- a/pql/ast.go +++ b/pql/ast.go @@ -274,14 +274,17 @@ type callStackElem struct { type CallType byte const ( - // Normal calls can be executed per shard. + // PrecallNone calls can be executed per shard. PrecallNone = CallType(iota) - // PreCallGlobal indicates a call which must be run globally *before* + + // PrecallGlobal indicates a call which must be run globally *before* // distributing the call to other shards. Example: A Distinct query, // where every shard could potentially produce results for any shard, // so you have to produce the results up front. + // These are processed directly when inside of a count operation. PrecallGlobal - // PreCallPerNode indicates a call which needs to be run per-shard + + // PrecallPerNode indicates a call which needs to be run per-shard // in a way that lets it be done on each shard, but where it should // be done prior to spawning per-shard goroutines. Example: // A cross-index query, where each local shard may or may not need From 48ac989e6a5084e262f836649c7fd2e1de82091f Mon Sep 17 00:00:00 2001 From: Travis Date: Tue, 26 Jan 2021 15:41:47 -0600 Subject: [PATCH 02/25] Return zero-bit row (with Index/Field) instead of nil in executeDistinctShardSet --- executor.go | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/executor.go b/executor.go index b1a0ae1eb..9fa829aac 100644 --- a/executor.go +++ b/executor.go @@ -143,7 +143,6 @@ func (e *executor) Close() error { // Execute executes a PQL query. func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shards []uint64, opt *execOptions) (QueryResponse, error) { - span, ctx := tracing.StartSpanFromContext(ctx, "Executor.Execute") span.LogKV("pql", q.String()) defer span.Finish() @@ -1505,7 +1504,18 @@ func executeDistinctShardSet(ctx context.Context, qcx *Qcx, idx *Index, fieldNam fragData, _, err := tx.ContainerIterator(index, fieldName, "standard", shard, 0) switch errors.Cause(err) { case ViewNotFound, FragmentNotFound: - return nil, nil + // It may seem reasonable to return `nil` here in the case where the + // fragment for this shard does not exist. The problem with doing that + // is that if this operation is being performed on a remote node, then + // this result is going to get serialized as a QueryResponse and sent + // back to the original, non-remote node. When this happens, the + // encodeRow/decodeRow logic replaces `nil` with an empty Row. An empty + // Row will cause problems during the union step of the reduce phase if + // it is the "left" side of the union, because then the resulting Row + // after the union will have blank Index and Field values. Here, we + // ensure that we send a non-nil Row with valid Index and Field values + // so that the union step doesn't cause problems. + return &Row{Index: index, Field: fieldName}, nil case nil: default: return nil, errors.Wrap(err, "getting fragment data") From 0d97ec559b5f2b62188dfaee955e49d9400cf2dd Mon Sep 17 00:00:00 2001 From: Nia Weiss Date: Wed, 27 Jan 2021 10:29:40 -0500 Subject: [PATCH 03/25] switch signed count distinct test to use TestVariousQueries --- executor_test.go | 38 ++++++++++++-------------------------- 1 file changed, 12 insertions(+), 26 deletions(-) diff --git a/executor_test.go b/executor_test.go index ae86ba7d6..47e1e6d57 100644 --- a/executor_test.go +++ b/executor_test.go @@ -6569,30 +6569,6 @@ func TestExecutor_Execute_CountDistinct(t *testing.T) { }) } -// Ensure that Count(Distinct()) works with negative numbers. -func TestExecutor_Execute_CountDistinctSigned(t *testing.T) { - c := test.MustRunCluster(t, 2) - defer c.Close() - - c.CreateField(t, "i", pilosa.IndexOptions{}, "ints", - pilosa.OptFieldTypeInt(-100, 100), - ) - - c.Query(t, "i", fmt.Sprintf(` - Set(0, ints=1) - Set(%d, ints=-1) - `, ShardWidth)) - - resp := c.Query(t, "i", "Count(Distinct(field=ints))") - cnt, ok := resp.Results[0].(uint64) - if !ok { - t.Fatalf("invalid response type, expected: uint64, got: %T", resp.Results[0]) - } - if cnt != 2 { - t.Fatalf("invalid result, expected: 2, got: %v", cnt) - } -} - // Ensure that a top-level, bare distinct on multiple nodes // is handled correctly. func TestExecutor_BareDistinct(t *testing.T) { @@ -6842,6 +6818,8 @@ func TestMissingKeyRegression(t *testing.T) { func TestVariousQueries(t *testing.T) { for _, clusterSize := range []int{1, 3, 4, 7} { t.Run(fmt.Sprintf("%d-node", clusterSize), func(t *testing.T) { + t.Parallel() + variousQueries(t, clusterSize) }) } @@ -6921,8 +6899,7 @@ func variousQueries(t *testing.T, clusterSize int) { {Val: 0, Key: "userE"}, }) - // Create and populate "affinity" int field with negative, positive, zero and null values. - + // Create and populate "net_worth" int field with positive values. c.CreateField(t, "users", pilosa.IndexOptions{Keys: true, TrackExistence: true}, "net_worth", pilosa.OptFieldTypeInt(-100000000, 100000000)) c.ImportIntKey(t, "users", "net_worth", []test.IntKey{ {Val: 1, Key: "userA"}, @@ -7056,6 +7033,15 @@ toronto,2,11 }, csvVerifier: "-10\n-5\n0\n5\n10\n", }, + { + query: "Count(Distinct(field=affinity))", + qrVerifier: func(t *testing.T, resp pilosa.QueryResponse) { + if resp.Results[0].(uint64) != 5 { + t.Errorf("wrong number of values: %+v", resp.Results[0]) + } + }, + csvVerifier: "5\n", + }, { query: "Distinct(Row(affinity>=0),field=affinity)", qrVerifier: func(t *testing.T, resp pilosa.QueryResponse) { From 60b6eecc7a81ad7438c3de013ee7a117c1eb7dce Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Thu, 28 Jan 2021 16:36:45 -0700 Subject: [PATCH 04/25] Update keyed/unkeyed benchmarks --- scripts/bench_read.sh | 2 +- scripts/bench_write.sh | 27 +++++++++++++----------- scripts/etc/gloat/gh.issues.keyed.yml | 9 ++++++++ scripts/etc/gloat/gh.issues.unkeyed.yml | 9 ++++++++ scripts/etc/gloat/query.count.keyed.yml | 11 ++++++++++ scripts/populate_query_db.keyed.sh | 28 +++++++++++++++++++++++++ 6 files changed, 73 insertions(+), 13 deletions(-) create mode 100644 scripts/etc/gloat/gh.issues.keyed.yml create mode 100644 scripts/etc/gloat/gh.issues.unkeyed.yml create mode 100644 scripts/etc/gloat/query.count.keyed.yml create mode 100755 scripts/populate_query_db.keyed.sh diff --git a/scripts/bench_read.sh b/scripts/bench_read.sh index a070bfb94..f2ad4d96f 100755 --- a/scripts/bench_read.sh +++ b/scripts/bench_read.sh @@ -21,7 +21,7 @@ SHA=$(git -C $PILOSA_SRC rev-parse HEAD) # Format current date. DATE=$(date '+%Y%m%d') -for TYPE in row row-bsi row-range count intersect union difference xor groupby topk +for TYPE in row row-bsi row-range count count-keyed intersect union difference xor groupby topk do WORKFLOW_PATH="${BASH_SOURCE%/*}/etc/gloat/query.${TYPE}.yml" WORKFLOW_NAME="$(gloat workflow name $WORKFLOW_PATH)" diff --git a/scripts/bench_write.sh b/scripts/bench_write.sh index 02e8067ca..a5306debc 100755 --- a/scripts/bench_write.sh +++ b/scripts/bench_write.sh @@ -21,19 +21,22 @@ SHA=$(git -C $PILOSA_SRC rev-parse HEAD) # Format current date. DATE=$(date '+%Y%m%d') -WORKFLOW_PATH="${BASH_SOURCE%/*}/etc/gloat/gh.1m.yml" -WORKFLOW_NAME="$(gloat workflow name $WORKFLOW_PATH)" -TITLE="RBF vs Roaring, $WORKFLOW_NAME, $DATE ($SHA)" +for FILENAME in gh.1m.yml gh.issues.keyed.yml gh.issues.unkeyed.yml +do + WORKFLOW_PATH="${BASH_SOURCE%/*}/etc/gloat/${FILENAME}" + WORKFLOW_NAME="$(gloat workflow name $WORKFLOW_PATH)" + TITLE="RBF vs Roaring, $WORKFLOW_NAME, $DATE ($SHA)" -# Execute RBF/Roaring benchmark. -RBF_PATH=gloat/data/1m/rbf/${DATE}.tar.gz -TXSRC=rbf gloat run -v -o $RBF_PATH $WORKFLOW_PATH + # Execute RBF/Roaring benchmark. + RBF_PATH=gloat/data/1m/rbf/${DATE}.tar.gz + TXSRC=rbf gloat run -v -o $RBF_PATH $WORKFLOW_PATH -ROARING_PATH=gloat/data/1m/roaring/${DATE}.tar.gz -TXSRC=roaring gloat run -v -o $ROARING_PATH $WORKFLOW_PATH + ROARING_PATH=gloat/data/1m/roaring/${DATE}.tar.gz + TXSRC=roaring gloat run -v -o $ROARING_PATH $WORKFLOW_PATH -# Generate graph from results. -gloat graph -layout 2,5 -size 5120,820 -title "$TITLE" -name utime,stime,heap_alloc,heap_inuse,heap_objects,num_gc,rchar,wchar,syscr,syscw -series rbf,roaring -o /tmp/output.png $RBF_PATH $ROARING_PATH + # Generate graph from results. + gloat graph -layout 2,5 -size 5120,820 -title "$TITLE" -name utime,stime,heap_alloc,heap_inuse,heap_objects,num_gc,rchar,wchar,syscr,syscw -series rbf,roaring -o /tmp/output.png $RBF_PATH $ROARING_PATH -# Post graph to Slack with SHA. -curl -F file=@/tmp/output.png -F channels=C01HBFKRLGH -F "initial_comment=$TITLE" -H "Authorization: Bearer $SLACK_OAUTH_TOKEN" https://slack.com/api/files.upload + # Post graph to Slack with SHA. + curl -F file=@/tmp/output.png -F channels=C01HBFKRLGH -F "initial_comment=$TITLE" -H "Authorization: Bearer $SLACK_OAUTH_TOKEN" https://slack.com/api/files.upload +done \ No newline at end of file diff --git a/scripts/etc/gloat/gh.issues.keyed.yml b/scripts/etc/gloat/gh.issues.keyed.yml new file mode 100644 index 000000000..215f7f0d3 --- /dev/null +++ b/scripts/etc/gloat/gh.issues.keyed.yml @@ -0,0 +1,9 @@ +name: "GitHub Issues Import Load Testing (1 month, keyed)" + +main: "pilosa server --data-dir ${TMPDIR} --txsrc ${TXSRC}" +load: "molecula-consumer-github -i issues -r url --record-type issue --batch-size=100000 --start-time 2020-01-01T00:00:00Z --end-time 2020-01-13T23:00:00Z --cache-dir ~/.githubarchive" + +health_url: "http://localhost:10101/status" +vars_urls: + - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/gh.issues.unkeyed.yml b/scripts/etc/gloat/gh.issues.unkeyed.yml new file mode 100644 index 000000000..45bbb7650 --- /dev/null +++ b/scripts/etc/gloat/gh.issues.unkeyed.yml @@ -0,0 +1,9 @@ +name: "GitHub Issues Import Load Testing (1 month, unkeyed)" + +main: "pilosa server --data-dir ${TMPDIR} --txsrc ${TXSRC}" +load: "molecula-consumer-github -i issues -d id --record-type issue --batch-size=100000 --start-time 2020-01-01T00:00:00Z --end-time 2020-01-13T23:00:00Z --cache-dir ~/.githubarchive" + +health_url: "http://localhost:10101/status" +vars_urls: + - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/query.count.keyed.yml b/scripts/etc/gloat/query.count.keyed.yml new file mode 100644 index 000000000..3d0eb49b3 --- /dev/null +++ b/scripts/etc/gloat/query.count.keyed.yml @@ -0,0 +1,11 @@ +name: "Count() Load Testing w/ Keys" + +main: "pilosa server --data-dir ~/pilosa.query.keyed.${TXSRC} --txsrc ${TXSRC}" +load: "pilosa-bench -type count -rate 100 -n 3000" + +health_url: "http://localhost:10101/status" +health_regexp: "NORMAL" + +vars_urls: + - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/populate_query_db.keyed.sh b/scripts/populate_query_db.keyed.sh new file mode 100755 index 000000000..5476d54b8 --- /dev/null +++ b/scripts/populate_query_db.keyed.sh @@ -0,0 +1,28 @@ +#!/bin/bash +set -e + +# This script generates data query load testing to be run against. +# +# Environment variables: +# - TXSRC: Transaction store type ("roaring", "rbf") +# - CACHEDIR: Path to local GitHub Archive data, if available. + +# Require environment variables. +: "${TXSRC:?Must set TXSRC environment variable}" +: "${GHCACHEDIR:''}" + +echo "Starting pilosa" +pilosa server --data-dir ~/pilosa.query.keyed.${TXSRC} --txsrc ${TXSRC} & pid_pilosa=$! +sleep 5 + +echo "" +echo "Importing GitHub Archive" +molecula-consumer-github -i issues -r url --record-type issue --batch-size=100000 \ + --start-time 2020-01-01T00:00:00Z --end-time 2020-01-31T23:00:00Z \ + --cache-dir "$GHCACHEDIR" + +echo "" +echo "Import complete, shutting down pilosa" + +sleep 5 +kill $pid_pilosa From 4fba6bea82f00ac811093681130ccf83a79dee63 Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Fri, 29 Jan 2021 09:39:01 -0600 Subject: [PATCH 05/25] Prevent nil pointer exception during Distinct key translation --- executor.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/executor.go b/executor.go index 9fa829aac..9b4ad515e 100644 --- a/executor.go +++ b/executor.go @@ -6572,6 +6572,9 @@ func (e *executor) translateResult(ctx context.Context, index string, idx *Index if field.Keys() { rslt := result.Pos + if rslt == nil { + return &SignedRow{Pos: &Row{}}, nil + } other := &Row{Attrs: rslt.Attrs} for _, segment := range rslt.Segments() { keys, err := e.Cluster.translateIndexIDs(context.Background(), field.ForeignIndex(), segment.Columns()) From c4455acbd8df8bcdefc8c103f84673697f340fb3 Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Fri, 29 Jan 2021 09:39:31 -0600 Subject: [PATCH 06/25] Use inconsistent JSON schema to reach Distinct translation error condition --- executor_test.go | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/executor_test.go b/executor_test.go index 6a202ca3a..17ebbfc95 100644 --- a/executor_test.go +++ b/executor_test.go @@ -5397,6 +5397,19 @@ func TestExecutor_ForeignIndex(t *testing.T) { pilosa.OptFieldKeys(), ) + // stepchild/other field needs to have usesKeys=true + crashSchemaJson := `{"indexes": [{"name": "stepparent","createdAt": 1611247966371721700,"options": {"keys": true,"trackExistence": true},"shardWidth": 1048576},{"name": "stepchild","createdAt": 1611247953796662800,"options": {"keys": true,"trackExistence": true},"shardWidth": 1048576,"fields": [{"name": "parent_id","createdAt": 1611247953797265700,"options": {"type": "int","base": 0,"bitDepth": 28,"min": -9223372036854776000,"max": 9223372036854776000,"keys": false,"foreignIndex": "stepparent"}},{"name": "other","createdAt": 1611247953796814000,"options": {"type": "int","base": 0,"bitDepth": 17,"min": -9223372036854776000,"max": 9223372036854776000,"keys": true,"foreignIndex": ""}}]}]}` + + crashSchema := &pilosa.Schema{} + err := json.Unmarshal([]byte(crashSchemaJson), &crashSchema) + if err != nil { + t.Fatalf("json unmarshall: %v", err) + } + err = c.GetNode(0).API.ApplySchema(context.Background(), crashSchema, false) + if err != nil { + t.Fatalf("applying JSON schema: %v", err) + } + // Populate parent data. c.Query(t, "parent", fmt.Sprintf(` Set("one", general=1) @@ -5442,6 +5455,12 @@ func TestExecutor_ForeignIndex(t *testing.T) { t.Fatalf("unexpected keys: %v", row.Keys) } + crash := c.Query(t, "stepchild", `Distinct(Row(parent_id=3), field=other)`).Results[0].(pilosa.SignedRow) + if !sameStringSlice(crash.Pos.Keys, []string{}) { + // empty result; error condition does not require data + t.Fatalf("unexpected columns: %v", crash.Pos.Keys) + } + eq := c.Query(t, "child", `Row(parent_id=="one")`).Results[0].(*pilosa.Row) if !reflect.DeepEqual(eq.Columns(), []uint64{1, ShardWidth}) { t.Fatalf("unexpected columns: %v", eq.Columns()) From e397d35ed5c723d6a6539ccff582234b292e7f69 Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Wed, 13 Jan 2021 17:53:05 -0600 Subject: [PATCH 07/25] Include roaring field and key details in usage endpoint --- api.go | 38 ++++++++++++----- bluegreentx.go | 4 ++ bolt.go | 4 ++ catcher.go | 4 ++ rbf.go | 4 ++ rbf/tx.go | 22 +++++++++- rrtx.go | 4 ++ stattx.go | 4 ++ tx.go | 2 + txfactory.go | 113 ++++++++++++++++++++++++++++++++++++++++--------- 10 files changed, 168 insertions(+), 31 deletions(-) diff --git a/api.go b/api.go index 2bd8eb2e4..62b445484 100644 --- a/api.go +++ b/api.go @@ -821,25 +821,41 @@ type NodeUsage struct { // DiskUsage represents the storage space used on disk by one node. type DiskUsage struct { - Capacity uint64 `json:"capacity,omitempty"` - TotalUse int64 `json:"totalInUse"` - Indexes map[string]int64 `json:"indexes"` + Capacity uint64 `json:"capacity,omitempty"` + TotalUse uint64 `json:"totalInUse"` + IndexUsage map[string]IndexUsage `json:"indexes"` } -// Usage gets the disk usage per index, in a map[nodeID]NodeUsage +// IndexUsage represents the storage space used on disk by one index, on one node. +type IndexUsage struct { + Total uint64 `json:"total"` + IndexKeys uint64 `json:"indexKeys"` + FieldKeysTotal uint64 `json:"fieldKeysTotal"` + Fragments uint64 `json:"fragments"` + Fields map[string]FieldUsage `json:"fields"` +} + +// FieldUsage represents the storage space used on disk by one field, on one node +type FieldUsage struct { + Total uint64 `json:"total"` + Fragments uint64 `json:"fragments"` + Keys uint64 `json:"keys"` +} + +// Usage gets the disk usage, 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) - indexSizes, err := api.holder.Txf().IndexSizes() + indexDetails, err := api.holder.Txf().IndexUsageDetails() if err != nil { return nil, errors.Wrap(err, "getting index usage") } - var totalSize int64 - for _, s := range indexSizes { - totalSize += s + var totalSize uint64 + for _, s := range indexDetails { + totalSize += s.Total } capacity, err := api.server.systemInfo.DiskCapacity(api.holder.path) @@ -850,9 +866,9 @@ func (api *API) Usage(ctx context.Context, remote bool) (map[string]NodeUsage, e // Insert into result. nodeUsage := NodeUsage{ Disk: DiskUsage{ - Capacity: capacity, - TotalUse: totalSize, - Indexes: indexSizes, + Capacity: capacity, + TotalUse: totalSize, + IndexUsage: indexDetails, }, } nodeUsages[api.server.nodeID] = nodeUsage diff --git a/bluegreentx.go b/bluegreentx.go index 62f1c0bd8..e0d9ddee9 100644 --- a/bluegreentx.go +++ b/bluegreentx.go @@ -663,6 +663,10 @@ func (c *blueGreenTx) ContainerIterator(index, field, view string, shard uint64, return bgi, bfound, errB } +func (tx *blueGreenTx) GetFieldSizeBytes(index, field string) (uint64, error) { + return 0, nil +} + func NewBlueGreenIterator(tx *blueGreenTx, ait, bit roaring.ContainerIterator) *blueGreenIterator { return &blueGreenIterator{ tx: tx, diff --git a/bolt.go b/bolt.go index b8a3f25be..3b37a2bec 100644 --- a/bolt.go +++ b/bolt.go @@ -759,6 +759,10 @@ func (tx *BoltTx) ContainerIterator(index, field, view string, shard uint64, fir return bi, bytes.Equal(bi.lastKey, needle), nil } +func (tx *BoltTx) GetFieldSizeBytes(index, field string) (uint64, error) { + return 0, nil +} + // BoltIterator is the iterator returned from a BoltTx.ContainerIterator() call. // It implements the roaring.ContainerIterator interface. type BoltIterator struct { diff --git a/catcher.go b/catcher.go index 671ee7662..09f24ce93 100644 --- a/catcher.go +++ b/catcher.go @@ -317,3 +317,7 @@ func (c *catcherTx) ApplyFilter(index, field, view string, shard uint64, ckey ui func (c *catcherTx) GetSortedFieldViewList(idx *Index, shard uint64) (fvs []txkey.FieldView, err error) { return c.b.GetSortedFieldViewList(idx, shard) } + +func (tx *catcherTx) GetFieldSizeBytes(index, field string) (uint64, error) { + return 0, nil +} diff --git a/rbf.go b/rbf.go index dff5a1123..21c1a321a 100644 --- a/rbf.go +++ b/rbf.go @@ -435,6 +435,10 @@ func (tx *RBFTx) GetSortedFieldViewList(idx *Index, shard uint64) (fvs []txkey.F return tx.tx.GetSortedFieldViewList() } +func (tx *RBFTx) GetFieldSizeBytes(index, field string) (uint64, error) { + return 0, nil +} + // rbfName returns a NULL-separated key used for identifying bitmap maps in RBF. func rbfName(index, field, view string, shard uint64) string { return string(txkey.Prefix(index, field, view, shard)) diff --git a/rbf/tx.go b/rbf/tx.go index a952104ec..f89b38048 100644 --- a/rbf/tx.go +++ b/rbf/tx.go @@ -839,6 +839,25 @@ func (tx *Tx) inusePageSet() (map[uint32]struct{}, error) { return m, nil } +func (tx *Tx) GetFieldSizeBytes(index, field string) (uint64, error) { + + fmt.Printf("getting RBF field size %s/%s\n", index, field) + + var pgno uint32 + var parent uint32 + + var pageCount uint64 + + if err := tx.walkTree(pgno, parent, func(pgno, parent, typ uint32) error { + pageCount++ + return nil + }); err != nil { + return 0, err + } + + return uint64(pageCount * PageSize), nil +} + // walkTree recursively iterates over a page and all its children. func (tx *Tx) walkTree(pgno, parent uint32, fn func(pgno, parent, typ uint32) error) error { // Read page and iterate over children. @@ -964,6 +983,7 @@ func (tx *Tx) deallocateTree(pgno uint32) error { func (tx *Tx) readPage(pgno uint32) (_ []byte, isHeap bool, err error) { // Meta page is always cached on the transaction. + //fmt.Printf("readPage %d\n", pgno) if pgno == 0 { return tx.meta[:], false, nil } @@ -971,7 +991,7 @@ func (tx *Tx) readPage(pgno uint32) (_ []byte, isHeap bool, err error) { // Verify page number requested is within current size of database. pageN := readMetaPageN(tx.meta[:]) if pgno > pageN { - return nil, false, fmt.Errorf("rbf: page read out of bounds: pgno=%d max=%d", pgno, pageN) + return nil, false, fmt.Errorf("rbf: page read out of bounds: pgno=%d max=%d", pgno, pageN-1) } // Check if page has been updated in this tx. diff --git a/rrtx.go b/rrtx.go index 8c75a2924..ca89b6d3d 100644 --- a/rrtx.go +++ b/rrtx.go @@ -589,6 +589,10 @@ func (tx *RoaringTx) GetSortedFieldViewList(idx *Index, shard uint64) (fvs []txk return } +func (tx *RoaringTx) GetFieldSizeBytes(index, field string) (uint64, error) { + return 0, nil +} + //////// registrar and wrapper machinery // roaringRegistrar mirrors the machinery expected diff --git a/stattx.go b/stattx.go index 02b60d6e1..7009de9ac 100644 --- a/stattx.go +++ b/stattx.go @@ -681,3 +681,7 @@ func (c *statTx) Sn() int64 { func (c *statTx) GetSortedFieldViewList(idx *Index, shard uint64) (fvs []txkey.FieldView, err error) { return c.b.GetSortedFieldViewList(idx, shard) } + +func (tx *statTx) GetFieldSizeBytes(index, field string) (uint64, error) { + return 0, nil +} diff --git a/tx.go b/tx.go index 771e94ff8..f33648cb0 100644 --- a/tx.go +++ b/tx.go @@ -219,6 +219,8 @@ type Tx interface { // GetSortedFieldViewList gets the set of FieldView(s) GetSortedFieldViewList(idx *Index, shard uint64) (fvs []txkey.FieldView, err error) + + GetFieldSizeBytes(index, field string) (uint64, error) } // Closer is used by Finders diff --git a/txfactory.go b/txfactory.go index d8aaf3aee..5e2ad647b 100644 --- a/txfactory.go +++ b/txfactory.go @@ -593,40 +593,115 @@ func (f *TxFactory) DumpAll() { f.dbPerShard.DumpAll() } -func (f *TxFactory) IndexSizes() (index2bytes map[string]int64, err error) { - // Open storage directory. - index2bytes = make(map[string]int64) +func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, error) { + indexUsage := make(map[string]IndexUsage) dirName, err := expandDirName(f.holder.path) if err != nil { - return index2bytes, errors.Wrap(err, "expanding data directory") + return indexUsage, errors.Wrap(err, "expanding data directory") } idxs := f.holder.Indexes() + /* + qcx := f.NewQcx() + tx, finisher, err := qcx.GetTx(Txo{Write: !writable}) + if err != nil { + return indexUsage, errors.Wrap(err, "qcx.GetTx") + } + defer finisher(nil) + */ + for _, idx := range idxs { index := idx.name - fullName := path.Join(dirName, index) - roaringAndMeta, err := directoryUsage(fullName) - if err != nil { - return index2bytes, errors.Wrap(err, "getting disk usage for roaring and meta") + println(" i:" + index) + indexPath := path.Join(dirName, index) + + // field usage + fieldUsages := make(map[string]FieldUsage) + fragmentsTotal := uint64(0) + fieldKeysTotal := uint64(0) + flds := idx.Fields() + for _, fld := range flds { + field := fld.Name() + fUsage, err := f.FieldUsage(indexPath, fld) + if err != nil { + return indexUsage, errors.Wrapf(err, "getting disk usage for index (%s)", index) + } + fieldUsages[field] = fUsage + keysBytes := fieldUsages[field].Keys + fieldKeysTotal += keysBytes + fragmentsTotal += fieldUsages[field].Fragments + + // non-roaring field usage + /* + fieldBytes, err := tx.GetFieldSizeBytes(index, field) + if err != nil { + return indexUsage, errors.Wrapf(err, "getting disk usage for non-roaring fragments (%s)", field) + } + fieldUsages[field] = FieldUsage{ + Total: fieldBytes, + Fragments: fieldBytes - keysBytes, + Keys: keysBytes, + } + */ } - fullName += ".index.txstores@@@" - rbfOrLmdb, err := directoryUsage(fullName) - if err != nil { - return index2bytes, errors.Wrap(err, "getting disk usage for backend") + + // index keys usage + keysBytes := uint64(0) + if idx.keys { + keysPath := path.Join(indexPath, translateStoreDir) + keysBytes, err = directoryUsage(keysPath) + if err != nil { + return indexUsage, errors.Wrapf(err, "getting disk usage for index keys (%s)", index) + } + } + + indexUsage[index] = IndexUsage{ + Total: keysBytes + fieldKeysTotal + fragmentsTotal, + IndexKeys: keysBytes, + FieldKeysTotal: fieldKeysTotal, + Fragments: fragmentsTotal, + Fields: fieldUsages, } - index2bytes[index] = roaringAndMeta + rbfOrLmdb } - return index2bytes, nil + return indexUsage, nil } -func directoryUsage(fname string) (int64, error) { - if !DirExists(fname) { - return 0, nil +func (f *TxFactory) FieldUsage(indexPath string, fld *Field) (FieldUsage, error) { + fieldUsage := FieldUsage{} + + field := fld.name + println(" f:" + field) + + // roaring field usage + fieldPath := path.Join(indexPath, field) + fieldBytes, err := directoryUsage(fieldPath) + if err != nil { + return fieldUsage, errors.Wrapf(err, "getting disk usage for field (%s)", field) + } + keysBytes := int64(0) + if fld.usesKeys { + keysBytes, err = fileSize(fld.TranslateStorePath()) + if err != nil { + return fieldUsage, errors.Wrapf(err, "getting disk usage for field keys (%s)", field) + } + } + fieldUsage = FieldUsage{ + Total: fieldBytes, + Fragments: fieldBytes - uint64(keysBytes), + Keys: uint64(keysBytes), } - var size int64 + return fieldUsage, nil +} + +func directoryUsage(fname string) (uint64, error) { + if !DirExists(fname) { + return 0, errors.Errorf("directory does not exist (%s)", fname) + } + + var size uint64 dir, err := os.Open(fname) if err != nil { @@ -647,7 +722,7 @@ func directoryUsage(fname string) (int64, error) { } size += sz } else { - size += file.Size() + size += uint64(file.Size()) // NOTE this cast is safe for regular files, not necessarily others } } From 32a35805a409df8e50d510dc2c39e71898397757 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Thu, 14 Jan 2021 13:09:57 -0700 Subject: [PATCH 08/25] Add RBF index/field usage stats --- rbf.go | 2 +- rbf/tx.go | 36 ++++++++++++++++++++++-------------- server/server.go | 2 +- txfactory.go | 44 ++++++++++++++++++++++++-------------------- 4 files changed, 48 insertions(+), 36 deletions(-) diff --git a/rbf.go b/rbf.go index 21c1a321a..4fa2f8900 100644 --- a/rbf.go +++ b/rbf.go @@ -436,7 +436,7 @@ func (tx *RBFTx) GetSortedFieldViewList(idx *Index, shard uint64) (fvs []txkey.F } func (tx *RBFTx) GetFieldSizeBytes(index, field string) (uint64, error) { - return 0, nil + return tx.tx.GetSizeBytesWithPrefix(string(txkey.FieldPrefix(index, field))) } // rbfName returns a NULL-separated key used for identifying bitmap maps in RBF. diff --git a/rbf/tx.go b/rbf/tx.go index f89b38048..5e7c35001 100644 --- a/rbf/tx.go +++ b/rbf/tx.go @@ -839,23 +839,31 @@ func (tx *Tx) inusePageSet() (map[uint32]struct{}, error) { return m, nil } -func (tx *Tx) GetFieldSizeBytes(index, field string) (uint64, error) { - - fmt.Printf("getting RBF field size %s/%s\n", index, field) - - var pgno uint32 - var parent uint32 - - var pageCount uint64 - - if err := tx.walkTree(pgno, parent, func(pgno, parent, typ uint32) error { - pageCount++ - return nil - }); err != nil { +// GetSizeBytesWithPrefix returns the size of bitmaps with a given key prefix. +func (tx *Tx) GetSizeBytesWithPrefix(prefix string) (n uint64, err error) { + records, err := tx.RootRecords() + if err != nil { return 0, err } - return uint64(pageCount * PageSize), nil + // Loop over each bitmap in the database. + for itr := records.Iterator(); !itr.Done(); { + name, pgno := itr.Next() + + // Skip over any bitmaps that don't have a matching prefix. + if !strings.HasPrefix(name.(string), prefix) { + continue + } + + // Traverse the bitmap's b-tree and count the bytes for each page. + if err := tx.walkTree(pgno.(uint32), 0, func(pgno, parent, typ uint32) error { + n += PageSize + return nil + }); err != nil { + return 0, err + } + } + return n, nil } // walkTree recursively iterates over a page and all its children. diff --git a/server/server.go b/server/server.go index 530e4102a..51729a8b1 100644 --- a/server/server.go +++ b/server/server.go @@ -185,7 +185,7 @@ func (m *Command) Start() (err error) { return errors.Wrap(err, "opening server") } - m.logger.Printf("listening as %s\n", m.listenURI) + m.logger.Printf("listening! as %s\n", m.listenURI) go func() { if err := m.grpcServer.Serve(); err != nil { m.logger.Printf("grpc server error: %v", err) diff --git a/txfactory.go b/txfactory.go index 5e2ad647b..a7984e518 100644 --- a/txfactory.go +++ b/txfactory.go @@ -602,15 +602,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, error) { idxs := f.holder.Indexes() - /* - qcx := f.NewQcx() - tx, finisher, err := qcx.GetTx(Txo{Write: !writable}) - if err != nil { - return indexUsage, errors.Wrap(err, "qcx.GetTx") - } - defer finisher(nil) - */ - + qcx := f.NewQcx() for _, idx := range idxs { index := idx.name println(" i:" + index) @@ -632,18 +624,30 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, error) { fieldKeysTotal += keysBytes fragmentsTotal += fieldUsages[field].Fragments - // non-roaring field usage - /* - fieldBytes, err := tx.GetFieldSizeBytes(index, field) - if err != nil { - return indexUsage, errors.Wrapf(err, "getting disk usage for non-roaring fragments (%s)", field) + fieldUsage := FieldUsage{Keys: keysBytes} + + for _, shard := range fld.AvailableShards(true).Slice() { + if err := func() error { + tx, finisher, err := qcx.GetTx(Txo{Write: !writable, Index: idx, Shard: shard}) + if err != nil { + return errors.Wrap(err, "qcx.GetTx") + } + defer finisher(nil) + + // non-roaring field usage + fieldBytes, err := tx.GetFieldSizeBytes(index, field) + if err != nil { + return errors.Wrapf(err, "getting disk usage for non-roaring fragments (%s)", field) + } + fieldUsage.Total += fieldBytes + return nil + }(); err != nil { + return indexUsage, err } - fieldUsages[field] = FieldUsage{ - Total: fieldBytes, - Fragments: fieldBytes - keysBytes, - Keys: keysBytes, - } - */ + } + + fieldUsage.Fragments = fieldUsage.Total - keysBytes + fieldUsages[field] = fieldUsage } // index keys usage From 85fad859e2245dfdb770241ac97859e2bfd7f6e6 Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Fri, 15 Jan 2021 03:45:10 -0600 Subject: [PATCH 09/25] Correct some disk usage computations --- api.go | 10 ++-- server/handler_test.go | 10 +++- server/server.go | 2 +- txfactory.go | 110 +++++++++++++++++++++++++++-------------- 4 files changed, 86 insertions(+), 46 deletions(-) diff --git a/api.go b/api.go index 62b445484..913433807 100644 --- a/api.go +++ b/api.go @@ -816,7 +816,7 @@ func (api *API) Node() *Node { // NodeUsage represents all usage measurements for one node. type NodeUsage struct { - Disk DiskUsage `json:"bytesOnDisk"` + Disk DiskUsage `json:"diskUsage"` } // DiskUsage represents the storage space used on disk by one node. @@ -849,11 +849,11 @@ func (api *API) Usage(ctx context.Context, remote bool) (map[string]NodeUsage, e nodeUsages := make(map[string]NodeUsage) - indexDetails, err := api.holder.Txf().IndexUsageDetails() + indexDetails, nodeMetadataBytes, err := api.holder.Txf().IndexUsageDetails() if err != nil { - return nil, errors.Wrap(err, "getting index usage") + return nil, errors.Wrap(err, "getting node usage") } - var totalSize uint64 + totalSize := nodeMetadataBytes for _, s := range indexDetails { totalSize += s.Total } @@ -873,7 +873,7 @@ func (api *API) Usage(ctx context.Context, remote bool) (map[string]NodeUsage, e } nodeUsages[api.server.nodeID] = nodeUsage - // Collect size on disk from remote nodes + // Collect diskUsage from remote nodes if !remote { nodes := api.cluster.Nodes() for _, node := range nodes { diff --git a/server/handler_test.go b/server/handler_test.go index 705cfdb5c..5154349e7 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -396,10 +396,16 @@ func TestHandler_Endpoints(t *testing.T) { } for _, nodeUsage := range nodeUsages { - if len(nodeUsage.Disk.Indexes) != 2 { - t.Fatalf("wrong length index size list: %#v", nodeUsage.Disk.Indexes) + numIndexes := len(nodeUsage.Disk.IndexUsage) + if numIndexes != 2 { + t.Fatalf("wrong length index usage list: expected %d, got %d", 2, numIndexes) + } + numFields := len(nodeUsage.Disk.IndexUsage["i1"].Fields) + if numFields != len(i1.Fields()) { + t.Fatalf("wrong length field usage list: expected %d, got %d", len(i1.Fields()), numFields) } } + }) t.Run("UI/shard-distribution", func(t *testing.T) { diff --git a/server/server.go b/server/server.go index 51729a8b1..530e4102a 100644 --- a/server/server.go +++ b/server/server.go @@ -185,7 +185,7 @@ func (m *Command) Start() (err error) { return errors.Wrap(err, "opening server") } - m.logger.Printf("listening! as %s\n", m.listenURI) + m.logger.Printf("listening as %s\n", m.listenURI) go func() { if err := m.grpcServer.Serve(); err != nil { m.logger.Printf("grpc server error: %v", err) diff --git a/txfactory.go b/txfactory.go index a7984e518..a93fa9c3d 100644 --- a/txfactory.go +++ b/txfactory.go @@ -593,11 +593,13 @@ func (f *TxFactory) DumpAll() { f.dbPerShard.DumpAll() } -func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, error) { +// IndexUsageDetails computes the sum of filesizes used by the node, broken down +// by index, field, fragments and keys. +func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { indexUsage := make(map[string]IndexUsage) - dirName, err := expandDirName(f.holder.path) + holderPath, err := expandDirName(f.holder.path) if err != nil { - return indexUsage, errors.Wrap(err, "expanding data directory") + return indexUsage, 0, errors.Wrap(err, "expanding data directory") } idxs := f.holder.Indexes() @@ -605,8 +607,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, error) { qcx := f.NewQcx() for _, idx := range idxs { index := idx.name - println(" i:" + index) - indexPath := path.Join(dirName, index) + indexPath := path.Join(holderPath, index) // field usage fieldUsages := make(map[string]FieldUsage) @@ -615,16 +616,16 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, error) { flds := idx.Fields() for _, fld := range flds { field := fld.Name() - fUsage, err := f.FieldUsage(indexPath, fld) - if err != nil { - return indexUsage, errors.Wrapf(err, "getting disk usage for index (%s)", index) + if field == "_keys" { + continue + } + fUsage, err := f.fieldUsage(indexPath, fld) + if err != nil { + return indexUsage, 0, errors.Wrapf(err, "getting disk usage for index (%s)", index) } - fieldUsages[field] = fUsage - keysBytes := fieldUsages[field].Keys - fieldKeysTotal += keysBytes - fragmentsTotal += fieldUsages[field].Fragments - fieldUsage := FieldUsage{Keys: keysBytes} + // non-roaring field usage + fragmentUsage := uint64(0) for _, shard := range fld.AvailableShards(true).Slice() { if err := func() error { @@ -634,74 +635,107 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, error) { } defer finisher(nil) - // non-roaring field usage fieldBytes, err := tx.GetFieldSizeBytes(index, field) if err != nil { return errors.Wrapf(err, "getting disk usage for non-roaring fragments (%s)", field) } - fieldUsage.Total += fieldBytes + fragmentUsage += fieldBytes return nil }(); err != nil { - return indexUsage, err + return indexUsage, 0, err } } - fieldUsage.Fragments = fieldUsage.Total - keysBytes - fieldUsages[field] = fieldUsage + // add non-roaring to roaring + fUsage.Fragments += fragmentUsage + fUsage.Total += fragmentUsage + + // add to running total + fieldKeysTotal += fUsage.Keys + fragmentsTotal += fUsage.Fragments + + fieldUsages[field] = fUsage + } + + // index metadata, e.g. columnAttrs + indexMetaBytes, err := directoryUsage(indexPath, false) + if err != nil { + return indexUsage, 0, errors.Wrapf(err, "getting disk usage for index metadata (%s)", index) } // index keys usage - keysBytes := uint64(0) + indexKeysBytes := uint64(0) if idx.keys { keysPath := path.Join(indexPath, translateStoreDir) - keysBytes, err = directoryUsage(keysPath) + indexKeysBytes, err = directoryUsage(keysPath, true) if err != nil { - return indexUsage, errors.Wrapf(err, "getting disk usage for index keys (%s)", index) + return indexUsage, 0, errors.Wrapf(err, "getting disk usage for index keys (%s)", index) } } indexUsage[index] = IndexUsage{ - Total: keysBytes + fieldKeysTotal + fragmentsTotal, - IndexKeys: keysBytes, + Total: indexMetaBytes + indexKeysBytes + fieldKeysTotal + fragmentsTotal, + IndexKeys: indexKeysBytes, FieldKeysTotal: fieldKeysTotal, Fragments: fragmentsTotal, Fields: fieldUsages, } } - return indexUsage, nil + // node metadata, e.g. id allocator + nodeMetaBytes, err := directoryUsage(holderPath, false) + if err != nil { + return indexUsage, 0, errors.Wrapf(err, "getting disk usage for node metadata") + } + + return indexUsage, nodeMetaBytes, nil } -func (f *TxFactory) FieldUsage(indexPath string, fld *Field) (FieldUsage, error) { +// fieldUsage computes the sum of filesizes used by a field in +// the filesystem tree (roaring storage), broken down by keys and fragments. +func (f *TxFactory) fieldUsage(indexPath string, fld *Field) (FieldUsage, error) { fieldUsage := FieldUsage{} field := fld.name - println(" f:" + field) - // roaring field usage - fieldPath := path.Join(indexPath, field) - fieldBytes, err := directoryUsage(fieldPath) - if err != nil { - return fieldUsage, errors.Wrapf(err, "getting disk usage for field (%s)", field) - } + // row keys keysBytes := int64(0) + var err error if fld.usesKeys { keysBytes, err = fileSize(fld.TranslateStorePath()) if err != nil { return fieldUsage, errors.Wrapf(err, "getting disk usage for field keys (%s)", field) } } + + // field metadata, e.g. rowAttrs + fieldPath := path.Join(indexPath, field) + metaBytes, err := directoryUsage(fieldPath, false) // this includes keys + if err != nil { + return fieldUsage, errors.Wrapf(err, "getting disk usage for field meta (%s)", field) + } + + // fragment data + viewsPath := path.Join(fieldPath, "views") + fragmentBytes := uint64(0) + if dirExists(viewsPath) { + fragmentBytes, err = directoryUsage(viewsPath, true) + if err != nil { + return fieldUsage, errors.Wrapf(err, "getting disk usage for field fragments (%s)", field) + } + } + fieldUsage = FieldUsage{ - Total: fieldBytes, - Fragments: fieldBytes - uint64(keysBytes), + Total: metaBytes + fragmentBytes, + Fragments: fragmentBytes, Keys: uint64(keysBytes), } return fieldUsage, nil } -func directoryUsage(fname string) (uint64, error) { - if !DirExists(fname) { +func directoryUsage(fname string, recursive bool) (uint64, error) { + if !dirExists(fname) { return 0, errors.Errorf("directory does not exist (%s)", fname) } @@ -719,8 +753,8 @@ func directoryUsage(fname string) (uint64, error) { } for _, file := range files { - if file.IsDir() { - sz, err := directoryUsage(path.Join(fname, file.Name())) + if recursive && file.IsDir() { + sz, err := directoryUsage(path.Join(fname, file.Name()), true) if err != nil { return 0, err } From 703bd14048be033661b9f83ab18e1045f72c066c Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Mon, 18 Jan 2021 17:01:35 -0600 Subject: [PATCH 10/25] Fix errors in usage check --- txfactory.go | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/txfactory.go b/txfactory.go index a93fa9c3d..ba964c7cc 100644 --- a/txfactory.go +++ b/txfactory.go @@ -605,6 +605,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { idxs := f.holder.Indexes() qcx := f.NewQcx() + defer qcx.Abort() for _, idx := range idxs { index := idx.name indexPath := path.Join(holderPath, index) @@ -667,10 +668,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { indexKeysBytes := uint64(0) if idx.keys { keysPath := path.Join(indexPath, translateStoreDir) - indexKeysBytes, err = directoryUsage(keysPath, true) - if err != nil { - return indexUsage, 0, errors.Wrapf(err, "getting disk usage for index keys (%s)", index) - } + indexKeysBytes, _ = directoryUsage(keysPath, true) // if directory doesn't exist, size = 0 } indexUsage[index] = IndexUsage{ @@ -704,7 +702,8 @@ func (f *TxFactory) fieldUsage(indexPath string, fld *Field) (FieldUsage, error) if fld.usesKeys { keysBytes, err = fileSize(fld.TranslateStorePath()) if err != nil { - return fieldUsage, errors.Wrapf(err, "getting disk usage for field keys (%s)", field) + // if file doesn't exist, size = 0 + keysBytes = 0 } } From 5dc00883bde0df2400932dbcf9a026a6f4c4cfac Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Mon, 18 Jan 2021 17:12:08 -0600 Subject: [PATCH 11/25] Add more involved diskUsage test --- server/handler_test.go | 89 ++++++++++++++++++++++++++++++++++++++++++ test/pilosa.go | 7 ++++ 2 files changed, 96 insertions(+) diff --git a/server/handler_test.go b/server/handler_test.go index 5154349e7..593f801cc 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -35,6 +35,7 @@ import ( "github.com/pilosa/pilosa/v2" "github.com/pilosa/pilosa/v2/boltdb" "github.com/pilosa/pilosa/v2/encoding/proto" + "github.com/pilosa/pilosa/v2/gopsutil" "github.com/pilosa/pilosa/v2/http" "github.com/pilosa/pilosa/v2/pql" pb "github.com/pilosa/pilosa/v2/proto" @@ -1571,6 +1572,94 @@ func TestQueryHistory(t *testing.T) { } } +func Test_DiskUsage_Roaring(t *testing.T) { + + // roaring-usage-1: keys, existence + // roaring-usage-1/f1: set, no keys + // roaring-usage-1/f2: set, keys + // roaring-usage-2: no keys, no existence + // roaring-usage-2/g1: no keys, no existence + // roaring-usage-3: keys, no existence, no fields + schemaString := `{"indexes": [{"fields": [{"options": {"keys": false,"cacheSize": 50000,"cacheType": "ranked","type": "set"},"name": "f1"},{"options": {"keys": true,"cacheSize": 50000,"cacheType": "ranked","type": "set"},"name": "f2"}],"options": {"trackExistence": true,"keys": true},"name": "roaring-usage-1"},{"fields": [{"options": {"keys": true,"cacheSize": 50000,"cacheType": "ranked","type": "set"},"name": "g1"}],"options": {"trackExistence": false,"keys": false},"name": "roaring-usage-2"},{"fields": [],"options": {"trackExistence": false,"keys": true},"name": "roaring-usage-3"}]} +` + + txsrc := []string{"roaring", "rbf"} + + exp0f := []string{`{"Test_DiskUsage_Roaring__0":{"diskUsage":{"capacity":%d,"totalInUse":16562,"indexes":{}}},"Test_DiskUsage_Roaring__1":{"diskUsage":{"capacity":%[1]d,"totalInUse":16562,"indexes":{}}},"Test_DiskUsage_Roaring__2":{"diskUsage":{"capacity":%[1]d,"totalInUse":16562,"indexes":{}}}} +`, `{"Test_DiskUsage_Roaring__0":{"diskUsage":{"capacity":%d,"totalInUse":16562,"indexes":{}}},"Test_DiskUsage_Roaring__1":{"diskUsage":{"capacity":%[1]d,"totalInUse":16562,"indexes":{}}},"Test_DiskUsage_Roaring__2":{"diskUsage":{"capacity":%[1]d,"totalInUse":16562,"indexes":{}}}} +`} + + exp1f := []string{`{"Test_DiskUsage_Roaring__0":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}},"Test_DiskUsage_Roaring__1":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}},"Test_DiskUsage_Roaring__2":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}}} +`, `{"Test_DiskUsage_Roaring__0":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}},"Test_DiskUsage_Roaring__1":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}},"Test_DiskUsage_Roaring__2":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}}} +`} + + for n := 0; n < 2; n++ { + cluster := test.MustRunCluster(t, 3, []server.CommandOption{test.OptTxSrc(txsrc[n])}) + defer cluster.Close() + cmd := cluster.GetNode(0) + h := cmd.Handler.(*http.Handler).Handler + holder := cmd.Server.Holder() + + sysInfo := gopsutil.NewSystemInfo() + capacity, err := sysInfo.DiskCapacity(holder.Path()) + if err != nil { + t.Fatalf("unable to check disk capacity: %s", err) + } + + // check usage for empty cluster + exp0 := fmt.Sprintf(exp0f[n], capacity) + + 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) + } + body := w.Body.String() + if body != exp0 { + t.Fatalf("unexpected body:\n %s\nexpected:\n %s", body, exp0) + } + + // create schema + w = httptest.NewRecorder() + h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/schema", strings.NewReader(schemaString))) + + if w.Code != gohttp.StatusNoContent { + bod, err := ioutil.ReadAll(w.Result().Body) + if err != nil { + t.Errorf("reading body: %v", err) + } + t.Fatalf("unexpected code: %v, bod: %s", w.Code, bod) + } + + idx, err := cmd.API.Index(context.Background(), "roaring-usage-1") + if err != nil { + t.Fatalf("getting index: %v", err) + } + if idx.Name() != "roaring-usage-1" { + t.Fatalf("index did not get set, got %v", idx.Name()) + } + + // set some bits + test.MustNewHTTPRequest("POST", "/index/roaring-usage-1/query", strings.NewReader(`Set("col1", f2=10)`)) + test.MustNewHTTPRequest("POST", "/index/roaring-usage-1/query", strings.NewReader(`Set("col1", f2="row1")`)) + test.MustNewHTTPRequest("POST", "/index/roaring-usage-2/query", strings.NewReader(`Set(21, g1=42)`)) + // check usage with data + exp1 := fmt.Sprintf(exp1f[n], capacity) + + 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) + } + body = w.Body.String() + if body != exp1 { + t.Fatalf("unexpected body:\n %s\nexpected:\n %s", body, exp1) + } + } +} + func mustJSONDecode(t *testing.T, r io.Reader) (ret map[string]interface{}) { dec := json.NewDecoder(r) err := dec.Decode(&ret) diff --git a/test/pilosa.go b/test/pilosa.go index 264e3f4c9..7cdac8a10 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -41,6 +41,13 @@ type Command struct { commandOptions []server.CommandOption } +func OptTxSrc(src string) server.CommandOption { + return func(m *server.Command) error { + m.Config.Txsrc = src + return nil + } +} + func OptAllowedOrigins(origins []string) server.CommandOption { return func(m *server.Command) error { m.Config.Handler.AllowedOrigins = origins From 15443b371c9692f140d8a0f2b029debdc7718d90 Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Mon, 18 Jan 2021 22:40:00 -0600 Subject: [PATCH 12/25] Force consistent timestamp width in startup log --- holder.go | 6 ++---- server/handler_test.go | 47 ++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 49 insertions(+), 4 deletions(-) diff --git a/holder.go b/holder.go index 256871058..3e40ea14b 100644 --- a/holder.go +++ b/holder.go @@ -1277,10 +1277,8 @@ func (h *Holder) LoadNodeID() (string, error) { // Log startup time and version to $DATA_DIR/.startup.log func (h *Holder) logStartup() error { - time, err := time.Now().MarshalText() - if err != nil { - return errors.Wrap(err, "creating timestamp") - } + RFC3339NanoFixedWidth := "2006-01-02T15:04:05.000000 07:00" + time := time.Now().Format(RFC3339NanoFixedWidth) logLine := fmt.Sprintf("%s\t%s\n", time, Version) f, err := os.OpenFile(h.path+"/.startup.log", os.O_APPEND|os.O_WRONLY|os.O_CREATE, 0600) diff --git a/server/handler_test.go b/server/handler_test.go index 593f801cc..c17a61fd7 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -25,6 +25,8 @@ import ( "math" gohttp "net/http" "net/http/httptest" + "os" + "path/filepath" "reflect" "sort" "strings" @@ -1607,6 +1609,11 @@ func Test_DiskUsage_Roaring(t *testing.T) { } // check usage for empty cluster + dumpFile(holder.Path() + "/.id") + dumpFile(holder.Path() + "/.startup.log") + dumpFile(holder.Path() + "/.topology") + dumpFile(holder.Path() + "/idalloc.db") + dumpDir(holder.Path()) exp0 := fmt.Sprintf(exp0f[n], capacity) w := httptest.NewRecorder() @@ -1647,6 +1654,8 @@ func Test_DiskUsage_Roaring(t *testing.T) { // check usage with data exp1 := fmt.Sprintf(exp1f[n], capacity) + dumpDir(holder.Path()) + w = httptest.NewRecorder() h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/ui/usage", nil)) if w.Code != gohttp.StatusOK { @@ -1660,6 +1669,44 @@ func Test_DiskUsage_Roaring(t *testing.T) { } } +func dumpFile(pth string) { + file, err := os.Open(pth) + defer file.Close() + msg := pth + if err != nil { + fmt.Printf("\n", err) + return + } + b, err := ioutil.ReadAll(file) + if err != nil { + fmt.Printf("\n", err) + } + msg += fmt.Sprintf(" (%d bytes):", len(b)) + if len(b) <= 1000 { + msg += fmt.Sprintf(string(b)) + } + fmt.Printf(msg) + fmt.Printf("\n") +} + +func dumpDir(pth string) { + var files []string + + err := filepath.Walk(pth, func(path string, info os.FileInfo, err error) error { + files = append(files, path) + return nil + }) + if err != nil { + panic(err) + } + for _, pth2 := range files { + file, _ := os.Open(pth2) + b, _ := ioutil.ReadAll(file) + fmt.Printf("%10d %s\n", len(b), pth2) + file.Close() + } +} + func mustJSONDecode(t *testing.T, r io.Reader) (ret map[string]interface{}) { dec := json.NewDecoder(r) err := dec.Decode(&ret) From c81ea88847bf3fdb62cc0ef0070415f2cc507bf3 Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Tue, 19 Jan 2021 19:05:49 -0600 Subject: [PATCH 13/25] Fix total summation --- txfactory.go | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/txfactory.go b/txfactory.go index ba964c7cc..c1801fb5e 100644 --- a/txfactory.go +++ b/txfactory.go @@ -614,6 +614,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { fieldUsages := make(map[string]FieldUsage) fragmentsTotal := uint64(0) fieldKeysTotal := uint64(0) + fieldsTotal := uint64(0) flds := idx.Fields() for _, fld := range flds { field := fld.Name() @@ -654,6 +655,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { // add to running total fieldKeysTotal += fUsage.Keys fragmentsTotal += fUsage.Fragments + fieldsTotal += fUsage.Total fieldUsages[field] = fUsage } @@ -672,7 +674,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { } indexUsage[index] = IndexUsage{ - Total: indexMetaBytes + indexKeysBytes + fieldKeysTotal + fragmentsTotal, + Total: indexMetaBytes + indexKeysBytes + fieldsTotal, IndexKeys: indexKeysBytes, FieldKeysTotal: fieldKeysTotal, Fragments: fragmentsTotal, @@ -725,7 +727,7 @@ func (f *TxFactory) fieldUsage(indexPath string, fld *Field) (FieldUsage, error) } fieldUsage = FieldUsage{ - Total: metaBytes + fragmentBytes, + Total: metaBytes + fragmentBytes, // metaBytes includes keys Fragments: fragmentBytes, Keys: uint64(keysBytes), } From 5beb2664b0212853be75192ed575023d327a24ea Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Tue, 19 Jan 2021 19:31:12 -0600 Subject: [PATCH 14/25] Use simpler test --- server/handler_test.go | 142 ++--------------------------------------- 1 file changed, 6 insertions(+), 136 deletions(-) diff --git a/server/handler_test.go b/server/handler_test.go index c17a61fd7..0ffaf3d0e 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -25,8 +25,6 @@ import ( "math" gohttp "net/http" "net/http/httptest" - "os" - "path/filepath" "reflect" "sort" "strings" @@ -37,7 +35,6 @@ import ( "github.com/pilosa/pilosa/v2" "github.com/pilosa/pilosa/v2/boltdb" "github.com/pilosa/pilosa/v2/encoding/proto" - "github.com/pilosa/pilosa/v2/gopsutil" "github.com/pilosa/pilosa/v2/http" "github.com/pilosa/pilosa/v2/pql" pb "github.com/pilosa/pilosa/v2/proto" @@ -400,6 +397,12 @@ func TestHandler_Endpoints(t *testing.T) { for _, nodeUsage := range nodeUsages { numIndexes := len(nodeUsage.Disk.IndexUsage) + if nodeUsage.Disk.TotalUse < 75000 || nodeUsage.Disk.TotalUse > 300000 { + // Usage measurements are not consistent between machines, or + // over time, as features and implementations change, so checking + // for a range of sizes may be most useful way to test the details of this. + t.Fatalf("expected 75k < total < 300k, got %d", nodeUsage.Disk.TotalUse) + } if numIndexes != 2 { t.Fatalf("wrong length index usage list: expected %d, got %d", 2, numIndexes) } @@ -1574,139 +1577,6 @@ func TestQueryHistory(t *testing.T) { } } -func Test_DiskUsage_Roaring(t *testing.T) { - - // roaring-usage-1: keys, existence - // roaring-usage-1/f1: set, no keys - // roaring-usage-1/f2: set, keys - // roaring-usage-2: no keys, no existence - // roaring-usage-2/g1: no keys, no existence - // roaring-usage-3: keys, no existence, no fields - schemaString := `{"indexes": [{"fields": [{"options": {"keys": false,"cacheSize": 50000,"cacheType": "ranked","type": "set"},"name": "f1"},{"options": {"keys": true,"cacheSize": 50000,"cacheType": "ranked","type": "set"},"name": "f2"}],"options": {"trackExistence": true,"keys": true},"name": "roaring-usage-1"},{"fields": [{"options": {"keys": true,"cacheSize": 50000,"cacheType": "ranked","type": "set"},"name": "g1"}],"options": {"trackExistence": false,"keys": false},"name": "roaring-usage-2"},{"fields": [],"options": {"trackExistence": false,"keys": true},"name": "roaring-usage-3"}]} -` - - txsrc := []string{"roaring", "rbf"} - - exp0f := []string{`{"Test_DiskUsage_Roaring__0":{"diskUsage":{"capacity":%d,"totalInUse":16562,"indexes":{}}},"Test_DiskUsage_Roaring__1":{"diskUsage":{"capacity":%[1]d,"totalInUse":16562,"indexes":{}}},"Test_DiskUsage_Roaring__2":{"diskUsage":{"capacity":%[1]d,"totalInUse":16562,"indexes":{}}}} -`, `{"Test_DiskUsage_Roaring__0":{"diskUsage":{"capacity":%d,"totalInUse":16562,"indexes":{}}},"Test_DiskUsage_Roaring__1":{"diskUsage":{"capacity":%[1]d,"totalInUse":16562,"indexes":{}}},"Test_DiskUsage_Roaring__2":{"diskUsage":{"capacity":%[1]d,"totalInUse":16562,"indexes":{}}}} -`} - - exp1f := []string{`{"Test_DiskUsage_Roaring__0":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}},"Test_DiskUsage_Roaring__1":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}},"Test_DiskUsage_Roaring__2":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}}} -`, `{"Test_DiskUsage_Roaring__0":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}},"Test_DiskUsage_Roaring__1":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}},"Test_DiskUsage_Roaring__2":{"diskUsage":{"capacity":%[1]d,"totalInUse":115896,"indexes":{"roaring-usage-1":{"total":33156,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"_exists":{"total":32791,"fragments":0,"keys":0},"f1":{"total":32791,"fragments":0,"keys":0},"f2":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-2":{"total":32896,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{"g1":{"total":32793,"fragments":0,"keys":0}}},"roaring-usage-3":{"total":32770,"indexKeys":0,"fieldKeysTotal":0,"fragments":0,"fields":{}}}}}} -`} - - for n := 0; n < 2; n++ { - cluster := test.MustRunCluster(t, 3, []server.CommandOption{test.OptTxSrc(txsrc[n])}) - defer cluster.Close() - cmd := cluster.GetNode(0) - h := cmd.Handler.(*http.Handler).Handler - holder := cmd.Server.Holder() - - sysInfo := gopsutil.NewSystemInfo() - capacity, err := sysInfo.DiskCapacity(holder.Path()) - if err != nil { - t.Fatalf("unable to check disk capacity: %s", err) - } - - // check usage for empty cluster - dumpFile(holder.Path() + "/.id") - dumpFile(holder.Path() + "/.startup.log") - dumpFile(holder.Path() + "/.topology") - dumpFile(holder.Path() + "/idalloc.db") - dumpDir(holder.Path()) - exp0 := fmt.Sprintf(exp0f[n], capacity) - - 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) - } - body := w.Body.String() - if body != exp0 { - t.Fatalf("unexpected body:\n %s\nexpected:\n %s", body, exp0) - } - - // create schema - w = httptest.NewRecorder() - h.ServeHTTP(w, test.MustNewHTTPRequest("POST", "/schema", strings.NewReader(schemaString))) - - if w.Code != gohttp.StatusNoContent { - bod, err := ioutil.ReadAll(w.Result().Body) - if err != nil { - t.Errorf("reading body: %v", err) - } - t.Fatalf("unexpected code: %v, bod: %s", w.Code, bod) - } - - idx, err := cmd.API.Index(context.Background(), "roaring-usage-1") - if err != nil { - t.Fatalf("getting index: %v", err) - } - if idx.Name() != "roaring-usage-1" { - t.Fatalf("index did not get set, got %v", idx.Name()) - } - - // set some bits - test.MustNewHTTPRequest("POST", "/index/roaring-usage-1/query", strings.NewReader(`Set("col1", f2=10)`)) - test.MustNewHTTPRequest("POST", "/index/roaring-usage-1/query", strings.NewReader(`Set("col1", f2="row1")`)) - test.MustNewHTTPRequest("POST", "/index/roaring-usage-2/query", strings.NewReader(`Set(21, g1=42)`)) - // check usage with data - exp1 := fmt.Sprintf(exp1f[n], capacity) - - dumpDir(holder.Path()) - - 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) - } - body = w.Body.String() - if body != exp1 { - t.Fatalf("unexpected body:\n %s\nexpected:\n %s", body, exp1) - } - } -} - -func dumpFile(pth string) { - file, err := os.Open(pth) - defer file.Close() - msg := pth - if err != nil { - fmt.Printf("\n", err) - return - } - b, err := ioutil.ReadAll(file) - if err != nil { - fmt.Printf("\n", err) - } - msg += fmt.Sprintf(" (%d bytes):", len(b)) - if len(b) <= 1000 { - msg += fmt.Sprintf(string(b)) - } - fmt.Printf(msg) - fmt.Printf("\n") -} - -func dumpDir(pth string) { - var files []string - - err := filepath.Walk(pth, func(path string, info os.FileInfo, err error) error { - files = append(files, path) - return nil - }) - if err != nil { - panic(err) - } - for _, pth2 := range files { - file, _ := os.Open(pth2) - b, _ := ioutil.ReadAll(file) - fmt.Printf("%10d %s\n", len(b), pth2) - file.Close() - } -} - func mustJSONDecode(t *testing.T, r io.Reader) (ret map[string]interface{}) { dec := json.NewDecoder(r) err := dec.Decode(&ret) From 1e7c4d7e8e0a939ae6ab08640d454344c4a47ae9 Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Thu, 21 Jan 2021 10:12:09 -0600 Subject: [PATCH 15/25] Include metadata AKA 'other' in response --- api.go | 2 ++ txfactory.go | 14 ++++++++------ 2 files changed, 10 insertions(+), 6 deletions(-) diff --git a/api.go b/api.go index 913433807..4b27e0766 100644 --- a/api.go +++ b/api.go @@ -832,6 +832,7 @@ type IndexUsage struct { IndexKeys uint64 `json:"indexKeys"` FieldKeysTotal uint64 `json:"fieldKeysTotal"` Fragments uint64 `json:"fragments"` + Metadata uint64 `json:"metadata"` Fields map[string]FieldUsage `json:"fields"` } @@ -840,6 +841,7 @@ type FieldUsage struct { Total uint64 `json:"total"` Fragments uint64 `json:"fragments"` Keys uint64 `json:"keys"` + Metadata uint64 `json:"metadata"` } // Usage gets the disk usage, in a map[nodeID]NodeUsage. diff --git a/txfactory.go b/txfactory.go index c1801fb5e..811c33570 100644 --- a/txfactory.go +++ b/txfactory.go @@ -614,6 +614,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { fieldUsages := make(map[string]FieldUsage) fragmentsTotal := uint64(0) fieldKeysTotal := uint64(0) + fieldMetaBytesTotal := uint64(0) fieldsTotal := uint64(0) flds := idx.Fields() for _, fld := range flds { @@ -653,6 +654,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { fUsage.Total += fragmentUsage // add to running total + fieldMetaBytesTotal += fUsage.Metadata fieldKeysTotal += fUsage.Keys fragmentsTotal += fUsage.Fragments fieldsTotal += fUsage.Total @@ -675,6 +677,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { indexUsage[index] = IndexUsage{ Total: indexMetaBytes + indexKeysBytes + fieldsTotal, + Metadata: indexMetaBytes + fieldMetaBytesTotal, IndexKeys: indexKeysBytes, FieldKeysTotal: fieldKeysTotal, Fragments: fragmentsTotal, @@ -701,12 +704,10 @@ func (f *TxFactory) fieldUsage(indexPath string, fld *Field) (FieldUsage, error) // row keys keysBytes := int64(0) var err error - if fld.usesKeys { - keysBytes, err = fileSize(fld.TranslateStorePath()) - if err != nil { - // if file doesn't exist, size = 0 - keysBytes = 0 - } + keysBytes, err = fileSize(fld.TranslateStorePath()) + if err != nil { + // if file doesn't exist, size = 0 + keysBytes = 0 } // field metadata, e.g. rowAttrs @@ -728,6 +729,7 @@ func (f *TxFactory) fieldUsage(indexPath string, fld *Field) (FieldUsage, error) fieldUsage = FieldUsage{ Total: metaBytes + fragmentBytes, // metaBytes includes keys + Metadata: metaBytes - uint64(keysBytes), Fragments: fragmentBytes, Keys: uint64(keysBytes), } From e981b162f00555deab68600b746da306757b742e Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Fri, 29 Jan 2021 17:47:24 -0600 Subject: [PATCH 16/25] Upgrade lattice --- lattice | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lattice b/lattice index 28c2313ec..2f0302c1d 160000 --- a/lattice +++ b/lattice @@ -1 +1 @@ -Subproject commit 28c2313ecfcd7e083d42d4e409483e968b4c421b +Subproject commit 2f0302c1d124433f0e1af5ae6c3bb7e4a64ca520 From 367425bba155fb8fa98701d0408c5c69b5b35506 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Fri, 29 Jan 2021 13:12:59 -0600 Subject: [PATCH 17/25] Add duration header to all gRPC query results --- server/grpc.go | 24 ++++++++++++-- server/grpc_test.go | 76 +++++++++++++++++++++++++++++++++++++++++++-- 2 files changed, 96 insertions(+), 4 deletions(-) diff --git a/server/grpc.go b/server/grpc.go index c2e68bb1f..8ffc49ad3 100644 --- a/server/grpc.go +++ b/server/grpc.go @@ -20,6 +20,7 @@ import ( "fmt" "net" "net/http" + "strconv" "strings" "sync" "time" @@ -34,6 +35,7 @@ import ( "google.golang.org/grpc" "google.golang.org/grpc/codes" "google.golang.org/grpc/credentials" + "google.golang.org/grpc/metadata" "google.golang.org/grpc/reflection" "google.golang.org/grpc/status" ) @@ -146,6 +148,10 @@ func (h *GRPCHandler) QuerySQL(req *pb.QuerySQLRequest, stream pb.Pilosa_QuerySQ return err } + stream.SendHeader(metadata.New(map[string]string{ + "duration": strconv.Itoa(int(duration)), + })) + err = newDurationRowser(results, duration).ToRows(stream.Send) if err != nil { return errors.Wrap(err, "streaming result") @@ -183,7 +189,12 @@ func (h *GRPCHandler) QuerySQLUnary(ctx context.Context, req *pb.QuerySQLRequest if err != nil { return nil, err } - table.Duration = int64(time.Since(start)) + duration := time.Since(start) + table.Duration = int64(duration) + grpc.SendHeader(ctx, metadata.New(map[string]string{ + "duration": strconv.Itoa(int(duration)), + })) + return table, nil } @@ -197,6 +208,11 @@ func (h *GRPCHandler) QueryPQL(req *pb.QueryPQLRequest, stream pb.Pilosa_QueryPQ t := time.Now() resp, err := h.api.Query(stream.Context(), &query) durQuery := time.Since(t) + + stream.SendHeader(metadata.New(map[string]string{ + "duration": strconv.Itoa(int(durQuery)), + })) + // TODO: what about resp.CollumnAttrSets? if err != nil { return errToStatusError(err) @@ -262,7 +278,11 @@ func (h *GRPCHandler) QueryPQLUnary(ctx context.Context, req *pb.QueryPQLRequest } durFormat := time.Since(t) - table.Duration = int64(durQuery + durFormat) + duration := durQuery + durFormat + table.Duration = int64(duration) + grpc.SendHeader(ctx, metadata.New(map[string]string{ + "duration": strconv.Itoa(int(duration)), + })) h.stats.Timing(pilosa.MetricGRPCUnaryQueryDurationSeconds, durQuery, 0.1) h.stats.Timing(pilosa.MetricGRPCUnaryFormatDurationSeconds, durFormat, 0.1) diff --git a/server/grpc_test.go b/server/grpc_test.go index 3cc808bb2..89ae9451f 100644 --- a/server/grpc_test.go +++ b/server/grpc_test.go @@ -29,7 +29,9 @@ import ( "github.com/pilosa/pilosa/v2/sql" "github.com/pilosa/pilosa/v2/test" "github.com/pkg/errors" + "google.golang.org/grpc" "google.golang.org/grpc/codes" + "google.golang.org/grpc/metadata" "google.golang.org/grpc/status" ) @@ -354,9 +356,11 @@ func TestQueryPQLUnary(t *testing.T) { i := m.MustCreateIndex(t, "i", pilosa.IndexOptions{}) m.MustCreateField(t, i.Name(), "f", pilosa.OptFieldKeys()) - ctx := context.Background() gh := server.NewGRPCHandler(m.API) + stream := &MockServerTransportStream{} + ctx := grpc.NewContextWithServerTransportStream(context.Background(), stream) + resp, err := gh.QueryPQLUnary(ctx, &pb.QueryPQLRequest{ Index: i.Name(), Pql: `Set(0, f="zero")`, @@ -369,6 +373,10 @@ func TestQueryPQLUnary(t *testing.T) { if resp.Duration == 0 { t.Fatal("duration not recorded") } + duration, err := stream.GetDuration() + if duration == 0 || err != nil { + t.Fatal("duration header not recorded") + } _, err = gh.QueryPQLUnary(ctx, &pb.QueryPQLRequest{ Index: i.Name(), @@ -400,6 +408,11 @@ func TestQueryPQL(t *testing.T) { t.Fatal(err) } + duration, err := mock.GetDuration() + if duration == 0 || err != nil { + t.Fatal("duration header not recorded") + } + if len(mock.Results) != 1 { t.Fatal("expecting one result") } @@ -481,7 +494,9 @@ type ( func TestQuerySQL(t *testing.T) { - ctx := context.Background() + stream := &MockServerTransportStream{} + ctx := grpc.NewContextWithServerTransportStream(context.Background(), stream) + gh, tearDownFunc := setUpTestQuerySQLUnary(ctx, t) defer tearDownFunc() @@ -924,6 +939,11 @@ func TestQuerySQL(t *testing.T) { if resp.Duration == 0 { t.Fatal("duration not recorded") } + duration, err := stream.GetDuration() + if duration == 0 || err != nil { + t.Fatal("duration header not recorded") + } + stream.ClearMD() tr := toTableResponse(resp) if err := test.eq(test.exp, tr); err != nil { t.Fatalf("sql: %s, error: %+v", test.sql, err) @@ -942,6 +962,10 @@ func TestQuerySQL(t *testing.T) { if mock.Results[0].Duration == 0 { t.Fatal("duration not recorded") } + duration, err := mock.GetDuration() + if duration == 0 || err != nil { + t.Fatal("duration header not recorded") + } if len(mock.Results) > 1 && mock.Results[1].Duration != 0 { t.Fatal("duration on second result expected to be zero") } @@ -1383,7 +1407,43 @@ func equalUnordered(exp tableResponse, got tableResponse) error { return nil } +type MockServerTransportStream struct { + header metadata.MD +} + +func (stream *MockServerTransportStream) Method() string { + return "" +} + +func (stream *MockServerTransportStream) SetHeader(md metadata.MD) error { + // Should probably merge md with value of stream.header, but this works since we have only one metadata value + stream.header = md + return nil +} + +func (stream *MockServerTransportStream) SendHeader(md metadata.MD) error { + stream.header = md + return nil +} + +func (stream *MockServerTransportStream) SetTrailer(md metadata.MD) error { + return nil +} + +func (stream *MockServerTransportStream) GetDuration() (int, error) { + duration, ok := stream.header["duration"] + if ok { + return strconv.Atoi(duration[0]) + } + return 0, errors.New("duration not recorded") +} + +func (stream *MockServerTransportStream) ClearMD() { + stream.header = metadata.New(map[string]string{}) +} + type mockPilosa_QuerySQLServer struct { + MockServerTransportStream pb.Pilosa_QuerySQLServer Results []*pb.RowResponse } @@ -1393,6 +1453,18 @@ func (m *mockPilosa_QuerySQLServer) Send(result *pb.RowResponse) error { return nil } +func (m *mockPilosa_QuerySQLServer) SendHeader(md metadata.MD) error { + return m.MockServerTransportStream.SendHeader(md) +} + +func (m *mockPilosa_QuerySQLServer) SetHeader(md metadata.MD) error { + return m.MockServerTransportStream.SetHeader(md) +} + +func (m *mockPilosa_QuerySQLServer) SetTrailer(md metadata.MD) { + m.MockServerTransportStream.SetTrailer(md) +} + func (m *mockPilosa_QuerySQLServer) Context() context.Context { return context.Background() } From e2331372d89e5d741e8e5d70d0baef513fe49edc Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Mon, 1 Feb 2021 09:23:30 -0600 Subject: [PATCH 18/25] Handle errors --- server/grpc.go | 20 ++++++++++++++++---- server/grpc_test.go | 2 +- 2 files changed, 17 insertions(+), 5 deletions(-) diff --git a/server/grpc.go b/server/grpc.go index 8ffc49ad3..f50c9ac80 100644 --- a/server/grpc.go +++ b/server/grpc.go @@ -148,9 +148,12 @@ func (h *GRPCHandler) QuerySQL(req *pb.QuerySQLRequest, stream pb.Pilosa_QuerySQ return err } - stream.SendHeader(metadata.New(map[string]string{ + err = stream.SendHeader(metadata.New(map[string]string{ "duration": strconv.Itoa(int(duration)), })) + if err != nil { + return errors.Wrap(err, "sending header") + } err = newDurationRowser(results, duration).ToRows(stream.Send) if err != nil { @@ -191,9 +194,12 @@ func (h *GRPCHandler) QuerySQLUnary(ctx context.Context, req *pb.QuerySQLRequest } duration := time.Since(start) table.Duration = int64(duration) - grpc.SendHeader(ctx, metadata.New(map[string]string{ + err = grpc.SendHeader(ctx, metadata.New(map[string]string{ "duration": strconv.Itoa(int(duration)), })) + if err != nil { + return nil, errors.Wrap(err, "sending header") + } return table, nil } @@ -209,9 +215,12 @@ func (h *GRPCHandler) QueryPQL(req *pb.QueryPQLRequest, stream pb.Pilosa_QueryPQ resp, err := h.api.Query(stream.Context(), &query) durQuery := time.Since(t) - stream.SendHeader(metadata.New(map[string]string{ + err = stream.SendHeader(metadata.New(map[string]string{ "duration": strconv.Itoa(int(durQuery)), })) + if err != nil { + return errors.Wrap(err, "sending header") + } // TODO: what about resp.CollumnAttrSets? if err != nil { @@ -280,9 +289,12 @@ func (h *GRPCHandler) QueryPQLUnary(ctx context.Context, req *pb.QueryPQLRequest duration := durQuery + durFormat table.Duration = int64(duration) - grpc.SendHeader(ctx, metadata.New(map[string]string{ + err = grpc.SendHeader(ctx, metadata.New(map[string]string{ "duration": strconv.Itoa(int(duration)), })) + if err != nil { + return nil, errors.Wrap(err, "sending header") + } h.stats.Timing(pilosa.MetricGRPCUnaryQueryDurationSeconds, durQuery, 0.1) h.stats.Timing(pilosa.MetricGRPCUnaryFormatDurationSeconds, durFormat, 0.1) diff --git a/server/grpc_test.go b/server/grpc_test.go index 89ae9451f..843fe2fe1 100644 --- a/server/grpc_test.go +++ b/server/grpc_test.go @@ -1462,7 +1462,7 @@ func (m *mockPilosa_QuerySQLServer) SetHeader(md metadata.MD) error { } func (m *mockPilosa_QuerySQLServer) SetTrailer(md metadata.MD) { - m.MockServerTransportStream.SetTrailer(md) + _ = m.MockServerTransportStream.SetTrailer(md) } func (m *mockPilosa_QuerySQLServer) Context() context.Context { From c30e3f0c2d320fe6e5aae81ac120a82f0b918b9d Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Mon, 1 Feb 2021 09:30:24 -0600 Subject: [PATCH 19/25] Move duration header to fix error handling --- server/grpc.go | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/server/grpc.go b/server/grpc.go index f50c9ac80..f7fe1192c 100644 --- a/server/grpc.go +++ b/server/grpc.go @@ -215,13 +215,6 @@ func (h *GRPCHandler) QueryPQL(req *pb.QueryPQLRequest, stream pb.Pilosa_QueryPQ resp, err := h.api.Query(stream.Context(), &query) durQuery := time.Since(t) - err = stream.SendHeader(metadata.New(map[string]string{ - "duration": strconv.Itoa(int(durQuery)), - })) - if err != nil { - return errors.Wrap(err, "sending header") - } - // TODO: what about resp.CollumnAttrSets? if err != nil { return errToStatusError(err) @@ -240,6 +233,13 @@ func (h *GRPCHandler) QueryPQL(req *pb.QueryPQLRequest, stream pb.Pilosa_QueryPQ return errors.Wrap(err, "wrapping as type ToRowser") } + err = stream.SendHeader(metadata.New(map[string]string{ + "duration": strconv.Itoa(int(durQuery)), + })) + if err != nil { + return errors.Wrap(err, "sending header") + } + t = time.Now() if err := newDurationRowser(toRowser, durQuery).ToRows(stream.Send); err != nil { return errToStatusError(err) From 0e22ed71ccd3ce533820ce3ff09784000fa9d125 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Mon, 1 Feb 2021 10:18:38 -0600 Subject: [PATCH 20/25] Fix instances of context.Background that need mocked context --- server/grpc_test.go | 6 ++++-- server/handler_test.go | 5 ++++- 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/server/grpc_test.go b/server/grpc_test.go index 843fe2fe1..d24aa0cc3 100644 --- a/server/grpc_test.go +++ b/server/grpc_test.go @@ -977,7 +977,8 @@ func TestQuerySQL(t *testing.T) { func TestQuerySQLUnaryWithError(t *testing.T) { - ctx := context.Background() + stream := &MockServerTransportStream{} + ctx := grpc.NewContextWithServerTransportStream(context.Background(), stream) gh, tearDownFunc := setUpTestQuerySQLUnary(ctx, t) defer tearDownFunc() @@ -1023,7 +1024,8 @@ func TestCRUDIndexes(t *testing.T) { m := test.RunCommand(t) defer m.Close() - ctx := context.Background() + stream := &MockServerTransportStream{} + ctx := grpc.NewContextWithServerTransportStream(context.Background(), stream) gh := server.NewGRPCHandler(m.API) t.Run("CreateIndex", func(t *testing.T) { diff --git a/server/handler_test.go b/server/handler_test.go index 0ffaf3d0e..834399587 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -40,6 +40,7 @@ import ( pb "github.com/pilosa/pilosa/v2/proto" "github.com/pilosa/pilosa/v2/server" "github.com/pilosa/pilosa/v2/test" + "google.golang.org/grpc" ) func TestHandler_PostSchemaCluster(t *testing.T) { @@ -1517,7 +1518,9 @@ func TestQueryHistory(t *testing.T) { test.Do(t, "POST", cmd.URL()+"/index/i0/field/f0", "") gh := server.NewGRPCHandler(cmd.API) - _, err = gh.QuerySQLUnary(context.Background(), &pb.QuerySQLRequest{ + stream := &MockServerTransportStream{} + ctx := grpc.NewContextWithServerTransportStream(context.Background(), stream) + _, err = gh.QuerySQLUnary(ctx, &pb.QuerySQLRequest{ Sql: `select * from i0`, }) From 92426a9d1b054fafadb3ccc9e0fdd6eb1000dda2 Mon Sep 17 00:00:00 2001 From: nagamocha3000 Date: Mon, 1 Feb 2021 20:23:34 +0300 Subject: [PATCH 21/25] Close process on fragment.openStorage error When *fragment.openStorage is invoked in both f.importValue and f.importValueSmallWrite and it returns an error, this means there's some underlying error with the storage device and at the point of this commit, the sane thing to do is to close the process, otherwise the operation of Pilosa might proceed in an inconsistent state thus precipiatting other silent but hairy errors along the way such as dereferencing *fragment.gen later on which is set to nil once openStorage fails. --- fragment.go | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/fragment.go b/fragment.go index 8592182f4..bd29a0246 100644 --- a/fragment.go +++ b/fragment.go @@ -2651,7 +2651,13 @@ func (f *fragment) importValueSmallWrite(tx Tx, columnIDs []uint64, values []int } return nil }(); err != nil { - _ = f.openStorage(true) + errOpenStorage := f.openStorage(true) + if errOpenStorage != nil { + f.Logger.Printf("failed to import data into fragment: %v", err) + f.Logger.Printf("recovery with openStorage failed for fragment: %v", errOpenStorage) + f.Logger.Debugf("%s", debug.Stack()) + os.Exit(1) + } return err } rowSet := make(map[uint64]struct{}, bitDepth+1) @@ -2705,7 +2711,13 @@ func (f *fragment) importValue(tx Tx, columnIDs []uint64, values []int64, bitDep // Flush changes in bulk back to the transaction. return txb.Flush() }(); err != nil { - _ = f.openStorage(true) + errOpenStorage := f.openStorage(true) + if errOpenStorage != nil { + f.Logger.Printf("failed to import data into fragment: %v", err) + f.Logger.Printf("recovery with openStorage failed for fragment: %v", errOpenStorage) + f.Logger.Debugf("%s", debug.Stack()) + os.Exit(1) + } return err } // Keep stats accurate. We don't call incrementOpN here because it may From 6c139935f749c75ece9e0a2a62befb8dee091ca5 Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Wed, 3 Feb 2021 10:05:19 -0600 Subject: [PATCH 22/25] Add memory info to /ui/usage response --- api.go | 32 +++++++++++++++++++++++++++----- 1 file changed, 27 insertions(+), 5 deletions(-) diff --git a/api.go b/api.go index 4b27e0766..d2c09fe85 100644 --- a/api.go +++ b/api.go @@ -816,7 +816,8 @@ func (api *API) Node() *Node { // NodeUsage represents all usage measurements for one node. type NodeUsage struct { - Disk DiskUsage `json:"diskUsage"` + Disk DiskUsage `json:"diskUsage"` + Memory MemoryUsage `json:"memoryUsage"` } // DiskUsage represents the storage space used on disk by one node. @@ -844,7 +845,13 @@ type FieldUsage struct { Metadata uint64 `json:"metadata"` } -// Usage gets the disk usage, in a map[nodeID]NodeUsage. +// MemoryUsage represents the memory used by one node. +type MemoryUsage struct { + Capacity uint64 `json:"capacity"` + TotalUse uint64 `json:"totalInUse"` +} + +// Usage gets the resource 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() @@ -860,22 +867,37 @@ func (api *API) Usage(ctx context.Context, remote bool) (map[string]NodeUsage, e totalSize += s.Total } - capacity, err := api.server.systemInfo.DiskCapacity(api.holder.path) + // NOTE: these errors are ignored in api.Info(), but checked here + si := api.server.systemInfo + diskCapacity, err := si.DiskCapacity(api.holder.path) if err != nil { api.server.logger.Printf("couldn't read disk capacity: %s", err) } + memoryCapacity, err := si.MemTotal() + if err != nil { + api.server.logger.Printf("couldn't read memory capacity: %s", err) + } + memoryUse, err := si.MemUsed() + if err != nil { + api.server.logger.Printf("couldn't read memory usage: %s", err) + } + // Insert into result. nodeUsage := NodeUsage{ Disk: DiskUsage{ - Capacity: capacity, + Capacity: diskCapacity, TotalUse: totalSize, IndexUsage: indexDetails, }, + Memory: MemoryUsage{ + Capacity: memoryCapacity, + TotalUse: memoryUse, + }, } nodeUsages[api.server.nodeID] = nodeUsage - // Collect diskUsage from remote nodes + // Collect usage from remote nodes if !remote { nodes := api.cluster.Nodes() for _, node := range nodes { From 0227d8306fd7f42214f4f9fb35f2629c0734a58c Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Tue, 9 Feb 2021 15:00:11 -0600 Subject: [PATCH 23/25] Upgrade lattice --- lattice | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lattice b/lattice index 2f0302c1d..fa773628a 160000 --- a/lattice +++ b/lattice @@ -1 +1 @@ -Subproject commit 2f0302c1d124433f0e1af5ae6c3bb7e4a64ca520 +Subproject commit fa773628a276e2590785a87fbc236c7e88ea6284 From e7f272f37f707f2462655182e5c12a9a8aac7886 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Wed, 10 Feb 2021 12:43:27 -0600 Subject: [PATCH 24/25] updates existence field on importroaring fixes issue (1411) --- api.go | 13 ++++++++++++- executor.go | 2 -- http/client_test.go | 46 +++++++++++++++++++++++++++++++++++++++++++++ http/handler.go | 1 - index.go | 1 - 5 files changed, 58 insertions(+), 5 deletions(-) diff --git a/api.go b/api.go index d2c09fe85..79246c271 100644 --- a/api.go +++ b/api.go @@ -416,6 +416,12 @@ func importWorker(importWork chan importJob) { case RequestActionSet: fileMagic := uint32(binary.LittleEndian.Uint16(viewData[0:2])) if fileMagic == roaring.MagicNumber { // if pilosa roaring format + if ef := j.field.idx.existenceField(); ef != nil { + err = ef.importRoaring(j.ctx, tx, viewData, j.shard, "standard", false) + if err != nil { + return errors.Wrap(err, "importing pilosa roaring existence") + } + } err := j.field.importRoaring(j.ctx, tx, viewData, j.shard, viewName, doClear) if err != nil { return errors.Wrap(err, "importing pilosa roaring") @@ -425,6 +431,12 @@ func importWorker(importWork chan importJob) { // field.importRoaring changes the standard roaring run format to pilosa roaring data := make([]byte, len(viewData)) copy(data, viewData) + if ef := j.field.idx.existenceField(); ef != nil { + err = ef.importRoaring(j.ctx, tx, data, j.shard, "standard", false) + if err != nil { + return errors.Wrap(err, "importing pilosa roaring existence") + } + } err := j.field.importRoaring(j.ctx, tx, data, j.shard, viewName, doClear) if err != nil { @@ -468,7 +480,6 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, span, ctx := tracing.StartSpanFromContext(ctx, "API.ImportRoaring") span.LogKV("index", indexName, "field", fieldName) defer span.Finish() - if err := api.validate(apiField); err != nil { return errors.Wrap(err, "validating api method") } diff --git a/executor.go b/executor.go index 6a50d4c2c..4b88d4270 100644 --- a/executor.go +++ b/executor.go @@ -4569,7 +4569,6 @@ func (e *executor) executeUnionRows(ctx context.Context, qcx *Qcx, index string, // executeAllCallShard executes an All() call for a local shard. func (e *executor) executeAllCallShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (res *Row, err0 error) { - span, _ := tracing.StartSpanFromContext(ctx, "Executor.executeAllCallShard") defer span.Finish() @@ -4596,7 +4595,6 @@ func (e *executor) executeAllCallShard(ctx context.Context, qcx *Qcx, index stri } defer finisher(&err0) - if existenceRow, err = existenceFrag.row(tx, 0); err != nil { return nil, err } diff --git a/http/client_test.go b/http/client_test.go index b1b154647..4d2c7f315 100644 --- a/http/client_test.go +++ b/http/client_test.go @@ -1435,3 +1435,49 @@ func TestClient_ServerInfoHasTxSrc(t *testing.T) { } pilosa.MustTxsrcToTxtype(si.TxSrc) // panics if invalid } +func TestClient_ImportRoaringExists(t *testing.T) { + cluster := test.MustNewCluster(t, 1) + err := cluster.Start() + if err != nil { + t.Fatalf("starting cluster: %v", err) + } + defer cluster.Close() + + node := cluster.GetNode(0) + _, err = node.API.CreateIndex(context.Background(), "i", pilosa.IndexOptions{TrackExistence: true}) + if err != nil { + t.Fatalf("creating index: %v", err) + } + _, err = node.API.CreateField(context.Background(), "i", "f", pilosa.OptFieldTypeSet(pilosa.CacheTypeRanked, 100)) + if err != nil { + t.Fatalf("creating field: %v", err) + } + // Send import request. + host := node.URL() + c := MustNewClient(host, http.GetHTTPClient(nil)) + // [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 65537] + roaringReq := makeImportRoaringRequest(false, "3B3001000100000900010000000100010009000100") + + if err := c.ImportRoaring(context.Background(), &cluster.GetNode(0).API.Node().URI, "i", "f", 0, false, roaringReq); err != nil { + t.Fatal(err) + } + expected := []uint64{1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 65537} + var qr pilosa.QueryResponse + qr, err = node.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: "Row(f=0)"}) + if err != nil { + t.Fatalf("%v", err) + } + got := qr.Results[0].(*pilosa.Row).Columns() + if !reflect.DeepEqual(got, expected) { + t.Fatalf(" Row unexpected columns: got %+v expected: %+v", got, expected) + } + qr, err = node.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: "All()"}) + if err != nil { + t.Fatalf("%v", err) + } + got = qr.Results[0].(*pilosa.Row).Columns() + if !reflect.DeepEqual(got, expected) { + t.Fatalf("All unexpected columns: got %+v expected: %+v", got, expected) + } + +} diff --git a/http/handler.go b/http/handler.go index a6dfd5828..b7e44f67e 100644 --- a/http/handler.go +++ b/http/handler.go @@ -2486,7 +2486,6 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request http.Error(w, error, code) return } - // Get index and field type to determine how to handle the // import data. indexName := mux.Vars(r)["index"] diff --git a/index.go b/index.go index 6129289f5..94f47c6c7 100644 --- a/index.go +++ b/index.go @@ -482,7 +482,6 @@ func (i *Index) Fields() []*Field { func (i *Index) existenceField() *Field { i.mu.RLock() defer i.mu.RUnlock() - return i.existenceFld } From a61ed011fc1c4fd0ee3553e0350d7abb08ad9c53 Mon Sep 17 00:00:00 2001 From: tgruben Date: Wed, 10 Feb 2021 14:51:11 -0600 Subject: [PATCH 25/25] Revert "Update existence field on import-roaring requests" --- api.go | 13 +------------ executor.go | 2 ++ http/client_test.go | 46 --------------------------------------------- http/handler.go | 1 + index.go | 1 + 5 files changed, 5 insertions(+), 58 deletions(-) diff --git a/api.go b/api.go index 79246c271..d2c09fe85 100644 --- a/api.go +++ b/api.go @@ -416,12 +416,6 @@ func importWorker(importWork chan importJob) { case RequestActionSet: fileMagic := uint32(binary.LittleEndian.Uint16(viewData[0:2])) if fileMagic == roaring.MagicNumber { // if pilosa roaring format - if ef := j.field.idx.existenceField(); ef != nil { - err = ef.importRoaring(j.ctx, tx, viewData, j.shard, "standard", false) - if err != nil { - return errors.Wrap(err, "importing pilosa roaring existence") - } - } err := j.field.importRoaring(j.ctx, tx, viewData, j.shard, viewName, doClear) if err != nil { return errors.Wrap(err, "importing pilosa roaring") @@ -431,12 +425,6 @@ func importWorker(importWork chan importJob) { // field.importRoaring changes the standard roaring run format to pilosa roaring data := make([]byte, len(viewData)) copy(data, viewData) - if ef := j.field.idx.existenceField(); ef != nil { - err = ef.importRoaring(j.ctx, tx, data, j.shard, "standard", false) - if err != nil { - return errors.Wrap(err, "importing pilosa roaring existence") - } - } err := j.field.importRoaring(j.ctx, tx, data, j.shard, viewName, doClear) if err != nil { @@ -480,6 +468,7 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, span, ctx := tracing.StartSpanFromContext(ctx, "API.ImportRoaring") span.LogKV("index", indexName, "field", fieldName) defer span.Finish() + if err := api.validate(apiField); err != nil { return errors.Wrap(err, "validating api method") } diff --git a/executor.go b/executor.go index 4b88d4270..6a50d4c2c 100644 --- a/executor.go +++ b/executor.go @@ -4569,6 +4569,7 @@ func (e *executor) executeUnionRows(ctx context.Context, qcx *Qcx, index string, // executeAllCallShard executes an All() call for a local shard. func (e *executor) executeAllCallShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (res *Row, err0 error) { + span, _ := tracing.StartSpanFromContext(ctx, "Executor.executeAllCallShard") defer span.Finish() @@ -4595,6 +4596,7 @@ func (e *executor) executeAllCallShard(ctx context.Context, qcx *Qcx, index stri } defer finisher(&err0) + if existenceRow, err = existenceFrag.row(tx, 0); err != nil { return nil, err } diff --git a/http/client_test.go b/http/client_test.go index 4d2c7f315..b1b154647 100644 --- a/http/client_test.go +++ b/http/client_test.go @@ -1435,49 +1435,3 @@ func TestClient_ServerInfoHasTxSrc(t *testing.T) { } pilosa.MustTxsrcToTxtype(si.TxSrc) // panics if invalid } -func TestClient_ImportRoaringExists(t *testing.T) { - cluster := test.MustNewCluster(t, 1) - err := cluster.Start() - if err != nil { - t.Fatalf("starting cluster: %v", err) - } - defer cluster.Close() - - node := cluster.GetNode(0) - _, err = node.API.CreateIndex(context.Background(), "i", pilosa.IndexOptions{TrackExistence: true}) - if err != nil { - t.Fatalf("creating index: %v", err) - } - _, err = node.API.CreateField(context.Background(), "i", "f", pilosa.OptFieldTypeSet(pilosa.CacheTypeRanked, 100)) - if err != nil { - t.Fatalf("creating field: %v", err) - } - // Send import request. - host := node.URL() - c := MustNewClient(host, http.GetHTTPClient(nil)) - // [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 65537] - roaringReq := makeImportRoaringRequest(false, "3B3001000100000900010000000100010009000100") - - if err := c.ImportRoaring(context.Background(), &cluster.GetNode(0).API.Node().URI, "i", "f", 0, false, roaringReq); err != nil { - t.Fatal(err) - } - expected := []uint64{1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 65537} - var qr pilosa.QueryResponse - qr, err = node.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: "Row(f=0)"}) - if err != nil { - t.Fatalf("%v", err) - } - got := qr.Results[0].(*pilosa.Row).Columns() - if !reflect.DeepEqual(got, expected) { - t.Fatalf(" Row unexpected columns: got %+v expected: %+v", got, expected) - } - qr, err = node.API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: "All()"}) - if err != nil { - t.Fatalf("%v", err) - } - got = qr.Results[0].(*pilosa.Row).Columns() - if !reflect.DeepEqual(got, expected) { - t.Fatalf("All unexpected columns: got %+v expected: %+v", got, expected) - } - -} diff --git a/http/handler.go b/http/handler.go index b7e44f67e..a6dfd5828 100644 --- a/http/handler.go +++ b/http/handler.go @@ -2486,6 +2486,7 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request http.Error(w, error, code) return } + // Get index and field type to determine how to handle the // import data. indexName := mux.Vars(r)["index"] diff --git a/index.go b/index.go index 94f47c6c7..6129289f5 100644 --- a/index.go +++ b/index.go @@ -482,6 +482,7 @@ func (i *Index) Fields() []*Field { func (i *Index) existenceField() *Field { i.mu.RLock() defer i.mu.RUnlock() + return i.existenceFld }