From 0d9179f10c0132e91ec5b0ebfbb1e016d9e7de16 Mon Sep 17 00:00:00 2001 From: Nia Weiss Date: Thu, 21 Jan 2021 11:11:20 -0500 Subject: [PATCH 01/32] roaring: fix intersectionAnyRunBitmap when processing single-word runs When a run started and ended within a single word, the entirety of the word would be checked. This would cause small runs to be processed incorrectly, and caused Distinct-on-sets to select rows that did not match the specified filter. --- roaring/roaring.go | 26 +++++++++++++------------- roaring/roaring_container_test.go | 12 ++++++++++++ 2 files changed, 25 insertions(+), 13 deletions(-) diff --git a/roaring/roaring.go b/roaring/roaring.go index 2705a6ba3..9c3a97aa9 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -4144,26 +4144,26 @@ func intersectionAnyRunBitmap(a, b *Container) bool { bb := b.bitmap()[:1024] runs := a.runs() for _, r := range runs { - loWord, loBit := r.Start/64, r.Start%64 - hiWord, hiBit := r.Last/64, r.Last%64 - if loBit != 0 { - w := bb[loWord] - mask := (uint64(1) << loBit) - 1 - if w&^mask != 0 { + if r.Start/64 == r.Last/64 { + mask := (^uint64(0) << (r.Start % 64)) &^ + (^uint64(0) << ((r.Last % 64) + 1)) + if mask&bb[r.Start/64] != 0 { return true } + continue } - for i := loWord; i < hiWord; i++ { + + firstWord, lastWord := r.Start/64, r.Last/64 + for i := firstWord + 1; i < lastWord; i++ { if bb[i] != 0 { return true } } - if hiBit != 0 { - w := bb[hiWord] - mask := (uint64(1) << hiBit) - 1 - if w&mask != 0 { - return true - } + + firstMask := ^uint64(0) << (r.Start % 64) + lastMask := ^(^uint64(0) << ((r.Last % 64) + 1)) + if (firstMask&bb[firstWord])|(lastMask&bb[lastWord]) != 0 { + return true } } return false diff --git a/roaring/roaring_container_test.go b/roaring/roaring_container_test.go index 3b747b5c9..9fbf4b214 100644 --- a/roaring/roaring_container_test.go +++ b/roaring/roaring_container_test.go @@ -97,3 +97,15 @@ func TestIntersectVariants(t *testing.T) { } } } + +func TestIntersectionAnyRunBitmapSingleWordRegression(t *testing.T) { + // In a previous version, single-word runs would match any bit within the word. + // Verify that this no longer happens. + any := intersectionAnyRunBitmap( + NewContainerRun([]Interval16{{1, 2}}), + NewContainerBitmapN([]uint64{0b1001}, 2), + ) + if any { + t.Errorf("matched an exclusive single-word run") + } +} From acdff02fedb8917c33d9777dcbbe6833147d9a9d Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Fri, 22 Jan 2021 11:40:10 -0700 Subject: [PATCH 02/32] Add ingest time & latency stats to query benchmarks --- cmd/pilosa-bench/main.go | 45 ++++++++++++++++++++++++++ scripts/bench_read.sh | 11 +++++-- scripts/etc/gloat/query.count.yml | 1 + scripts/etc/gloat/query.difference.yml | 1 + scripts/etc/gloat/query.groupby.yml | 1 + scripts/etc/gloat/query.intersect.yml | 1 + scripts/etc/gloat/query.row-bsi.yml | 1 + scripts/etc/gloat/query.row-range.yml | 1 + scripts/etc/gloat/query.row.yml | 1 + scripts/etc/gloat/query.topk.yml | 1 + scripts/etc/gloat/query.union.yml | 1 + scripts/etc/gloat/query.xor.yml | 1 + 12 files changed, 64 insertions(+), 2 deletions(-) diff --git a/cmd/pilosa-bench/main.go b/cmd/pilosa-bench/main.go index b33bf0fcf..1502dae23 100644 --- a/cmd/pilosa-bench/main.go +++ b/cmd/pilosa-bench/main.go @@ -16,12 +16,14 @@ package main import ( "context" + "expvar" "flag" "fmt" "io/ioutil" "log" "math/rand" "net/http" + _ "net/http/pprof" "os" "sort" "strings" @@ -32,6 +34,14 @@ import ( "golang.org/x/sync/errgroup" ) +var ( + requestCountVar = expvar.NewInt("request_count") + requestCurrentLatencyVar = expvar.NewFloat("request_current_latency") // seconds + requestAvgLatencyVar = expvar.NewFloat("request_avg_latency") // seconds + requestTotalLatencyVar = expvar.NewFloat("request_total_latency") // seconds + requestPerSecVar = expvar.NewFloat("request_per_sec") +) + func main() { if err := run(context.Background(), os.Args[1:]); err == flag.ErrHelp { os.Exit(1) @@ -86,6 +96,13 @@ func run(ctx context.Context, args []string) (err error) { return err } + // Set up HTTP endpoint to provide /debug endpoints. + fmt.Println("Serving debug endpoint at http://localhost:7070/debug") + go func() { _ = http.ListenAndServe(":7070", nil) }() + + // Run separate goroutine to calculate the current req/sec & latency. + go monitor() + // Load all id/keys for each field. log.Printf("loading field identifiers") fieldIDMap, err := loadFields(ctx, client) @@ -145,10 +162,15 @@ func run(ctx context.Context, args []string) (err error) { log.Printf("[query] %s", q) g.Go(func() error { + t := time.Now() _, err = client.Query(ctx, key.index, &pilosa.QueryRequest{Index: key.index, Query: q}) if err != nil { return err } + elapsed := time.Since(t).Seconds() + requestCountVar.Add(1) + requestTotalLatencyVar.Add(elapsed) + requestAvgLatencyVar.Set(requestTotalLatencyVar.Value() / float64(requestCountVar.Value())) return nil }) } @@ -156,6 +178,29 @@ func run(ctx context.Context, args []string) (err error) { return g.Wait() } +// monitor runs in a separate goroutine and updates metrics. +func monitor() { + ticker := time.NewTicker(1 * time.Second) + defer ticker.Stop() + + var lastTime time.Time + var lastN int64 + var lastLatency float64 + for range ticker.C { + now, n := time.Now(), requestCountVar.Value() + latency := requestTotalLatencyVar.Value() + + if !lastTime.IsZero() { + elapsed := lastTime.Sub(now).Seconds() + if n > 0 { + requestCurrentLatencyVar.Set((lastLatency - latency) / float64(n)) + } + requestPerSecVar.Set(float64(lastN-n) / elapsed) + } + lastTime, lastN, lastLatency = now, n, latency + } +} + func generateQuery(typ, index, field string, info *pilosa.FieldInfo, identifiers *pilosa.RowIdentifiers, opt queryOptions) (string, error) { switch typ { case "row": diff --git a/scripts/bench_read.sh b/scripts/bench_read.sh index bfbf6baac..a070bfb94 100755 --- a/scripts/bench_read.sh +++ b/scripts/bench_read.sh @@ -25,14 +25,21 @@ for TYPE in row row-bsi row-range count intersect union difference xor groupby t do WORKFLOW_PATH="${BASH_SOURCE%/*}/etc/gloat/query.${TYPE}.yml" WORKFLOW_NAME="$(gloat workflow name $WORKFLOW_PATH)" - TITLE="$WORKFLOW_NAME, $DATE ($SHA)" - + # Execute RBF/Roaring benchmark. + STARTTIME=$(date +%s) RBF_PATH=gloat/data/query/${TYPE}/rbf/${DATE}.tar.gz TXSRC=rbf gloat run -v -o "$RBF_PATH" $WORKFLOW_PATH + RBF_ELAPSED=$(($(date +%s) - $STARTTIME)) + RBF_LATENCY=$(gloat metric -n -name request_avg_latency "$RBF_PATH") + STARTTIME=$(date +%s) ROARING_PATH=gloat/data/query/${TYPE}/roaring/${DATE}.tar.gz TXSRC=roaring gloat run -v -o "$ROARING_PATH" $WORKFLOW_PATH + ROARING_ELAPSED=$(($(date +%s) - $STARTTIME)) + ROARING_LATENCY=$(gloat metric -n -name request_avg_latency "$ROARING_PATH") + + TITLE="$WORKFLOW_NAME, $DATE ($SHA) elapsed rbf=$RBF_ELAPSEDroaring=$ROARING_ELAPSED> latency rbf=$RBF_LATENCY roaring=$ROARING_LATENCY" # 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 diff --git a/scripts/etc/gloat/query.count.yml b/scripts/etc/gloat/query.count.yml index 5690c13a8..b0215d62c 100644 --- a/scripts/etc/gloat/query.count.yml +++ b/scripts/etc/gloat/query.count.yml @@ -8,3 +8,4 @@ health_regexp: "NORMAL" vars_urls: - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/query.difference.yml b/scripts/etc/gloat/query.difference.yml index 3a73a8e46..d3b17694f 100644 --- a/scripts/etc/gloat/query.difference.yml +++ b/scripts/etc/gloat/query.difference.yml @@ -8,3 +8,4 @@ health_regexp: "NORMAL" vars_urls: - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/query.groupby.yml b/scripts/etc/gloat/query.groupby.yml index a5c4c5d00..25b97fce4 100644 --- a/scripts/etc/gloat/query.groupby.yml +++ b/scripts/etc/gloat/query.groupby.yml @@ -8,3 +8,4 @@ health_regexp: "NORMAL" vars_urls: - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/query.intersect.yml b/scripts/etc/gloat/query.intersect.yml index 00bbfbd2e..4cb0993e9 100644 --- a/scripts/etc/gloat/query.intersect.yml +++ b/scripts/etc/gloat/query.intersect.yml @@ -8,3 +8,4 @@ health_regexp: "NORMAL" vars_urls: - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/query.row-bsi.yml b/scripts/etc/gloat/query.row-bsi.yml index 04fe472eb..7c86b5a41 100644 --- a/scripts/etc/gloat/query.row-bsi.yml +++ b/scripts/etc/gloat/query.row-bsi.yml @@ -8,3 +8,4 @@ health_regexp: "NORMAL" vars_urls: - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/query.row-range.yml b/scripts/etc/gloat/query.row-range.yml index f05b5cb4c..75ab2c6d6 100644 --- a/scripts/etc/gloat/query.row-range.yml +++ b/scripts/etc/gloat/query.row-range.yml @@ -8,3 +8,4 @@ health_regexp: "NORMAL" vars_urls: - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/query.row.yml b/scripts/etc/gloat/query.row.yml index 142f12004..7f6997ea3 100644 --- a/scripts/etc/gloat/query.row.yml +++ b/scripts/etc/gloat/query.row.yml @@ -8,3 +8,4 @@ health_regexp: "NORMAL" vars_urls: - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/query.topk.yml b/scripts/etc/gloat/query.topk.yml index 10da63427..25ef7fe50 100644 --- a/scripts/etc/gloat/query.topk.yml +++ b/scripts/etc/gloat/query.topk.yml @@ -8,3 +8,4 @@ health_regexp: "NORMAL" vars_urls: - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/query.union.yml b/scripts/etc/gloat/query.union.yml index a2270d974..1b2777b90 100644 --- a/scripts/etc/gloat/query.union.yml +++ b/scripts/etc/gloat/query.union.yml @@ -8,3 +8,4 @@ health_regexp: "NORMAL" vars_urls: - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars diff --git a/scripts/etc/gloat/query.xor.yml b/scripts/etc/gloat/query.xor.yml index f2eb857b8..3328fb0b1 100644 --- a/scripts/etc/gloat/query.xor.yml +++ b/scripts/etc/gloat/query.xor.yml @@ -8,3 +8,4 @@ health_regexp: "NORMAL" vars_urls: - http://localhost:10101/debug/vars + - http://localhost:7070/debug/vars From 649202b0816bcebd05d17f43ec435bd69e7a7ac9 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Mon, 21 Dec 2020 13:25:15 -0600 Subject: [PATCH 03/32] Fix cluster size setter: expose failures --- executor_test.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/executor_test.go b/executor_test.go index c0532b4f9..6a202ca3a 100644 --- a/executor_test.go +++ b/executor_test.go @@ -5517,7 +5517,7 @@ func TestExecutor_Execute_DistinctFailure(t *testing.T) { func TestExecutor_Execute_GroupBy(t *testing.T) { groupByTest := func(t *testing.T, clusterSize int) { - c := test.MustRunCluster(t, 1) + c := test.MustRunCluster(t, clusterSize) defer c.Close() c.CreateField(t, "i", pilosa.IndexOptions{}, "general") c.CreateField(t, "i", pilosa.IndexOptions{}, "sub") @@ -5924,7 +5924,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) { }) } - for size := range []int{1, 3} { + for _, size := range []int{1, 3} { t.Run(fmt.Sprintf("%d_nodes", size), func(t *testing.T) { groupByTest(t, size) }) From dde318ac8cb929a2c2fb7b12c431d99e8fd648b6 Mon Sep 17 00:00:00 2001 From: Nia Weiss Date: Mon, 25 Jan 2021 10:15:45 -0500 Subject: [PATCH 04/32] Move globally computed GroupBy rows calls into EmbeddedData This fixes a bug where a globally computed Rows call would be computed with a subset of the shards. --- executor.go | 17 +++++++++++++++++ pql/ast.go | 1 + row.go | 4 ++++ 3 files changed, 22 insertions(+) diff --git a/executor.go b/executor.go index 0c02fd6c8..b1a0ae1eb 100644 --- a/executor.go +++ b/executor.go @@ -2797,6 +2797,12 @@ func (e *executor) executeGroupBy(ctx context.Context, qcx *Qcx, index string, c } if hasLimit || hasCol { // we need to perform this query cluster-wide ahead of executeGroupByShard + if idx, ok := child.Args["valueidx"].(int64); ok { + // The rows query was already completed on the initiating node. + childRows[i] = opt.EmbeddedData[idx].Columns() + continue + } + childRows[i], err = e.executeRows(ctx, qcx, index, child, shards, opt) if err != nil { return nil, errors.Wrap(err, "getting rows for ") @@ -2804,6 +2810,13 @@ func (e *executor) executeGroupBy(ctx context.Context, qcx *Qcx, index string, c if len(childRows[i]) == 0 { // there are no results because this field has no values. return &GroupCounts{}, nil } + + // Stuff the result into opt.EmbeddedData so that it gets sent to other nodes in the map-reduce. + // This is flagged as "NoSplit" to ensure that the entire row gets sent out. + rowsRow := NewRow(childRows[i]...) + rowsRow.NoSplit = true + child.Args["valueidx"] = int64(len(opt.EmbeddedData)) + opt.EmbeddedData = append(opt.EmbeddedData, rowsRow) } } @@ -5527,6 +5540,10 @@ func makeEmbeddedDataForShards(allRows []*Row, shards []uint64) []*Row { if row == nil || len(row.segments) == 0 { continue } + if row.NoSplit { + newRows[i] = row + continue + } segments := row.segments segmentIndex := 0 newRows[i] = &Row{ diff --git a/pql/ast.go b/pql/ast.go index 40e2acaba..54ae2dee7 100644 --- a/pql/ast.go +++ b/pql/ast.go @@ -388,6 +388,7 @@ var callInfoByFunc = map[string]callInfo{ "from": nil, "to": nil, "like": "", + "valueidx": int64(0), }, }, "Shift": {allowUnknown: false, diff --git a/row.go b/row.go index 46a0ec478..3ff47cba5 100644 --- a/row.go +++ b/row.go @@ -42,6 +42,10 @@ type Row struct { // query. Knowing the index and field, we can figure out how to // interpret the row data. Field string + + // NoSplit indicates that this row may not be split. + // This is used for `Rows` calls in a GroupBy. + NoSplit bool } // NewRow returns a new instance of Row. From d1a9c91a96df1a208e8a7f1450743e73ebc75a81 Mon Sep 17 00:00:00 2001 From: Seebs Date: Mon, 25 Jan 2021 14:41:05 -0600 Subject: [PATCH 05/32] Test CountRange on non-container ranges. This test was supposed to check against all the container types, but especially bitmaps, but turns out not to work because the containers turn into RLE containers. Oops. Now, we start with every-other-bit for the first 8k, then start filling in the holes, so we get some bitmap containers and then start generating RLE containers. --- tx_internal_test.go | 27 ++++++++++++++++++++------- 1 file changed, 20 insertions(+), 7 deletions(-) diff --git a/tx_internal_test.go b/tx_internal_test.go index a5da9cccd..3ff4ddd63 100644 --- a/tx_internal_test.go +++ b/tx_internal_test.go @@ -37,20 +37,29 @@ func requireCountRangeSampleData(tb testing.TB) (*fragment, Tx) { // request that each container get its own copy of the bitmap. var bitmapSample [1025]uint64 for i := range arraySample { - arraySample[i] = uint16(i) + arraySample[i] = uint16(i * 2) } - for i := 0; i < 4096/64; i++ { - bitmapSample[i] = ^uint64(0) + // Put corresponding bits in the bitmap... + for i := 0; i < 4096/32; i++ { + // bit 0 is 0x1, bit 2 is 0x4, so even-numbered bits + // are 0x5555.... + bitmapSample[i] = 0x5555555555555555 } bm := roaring.NewSliceBitmap() for n := 0; n < 4096 && n < countRangeMaxN; n++ { c := roaring.NewContainerArray(arraySample[:n]) bm.Put(uint64(n), c) } - for n := 4096; n < countRangeMaxN; n++ { + // Start filling in the missing bits. This starts us out with + // bitmap containers, but then eventually converts to things + // that are more likely to be run containers. At the end of this, + // we should have exactly the first 8,192 bits set, for a single + // run of 8k. + for n := 4096; n < 8192; n++ { c := roaring.NewContainerBitmapN(bitmapSample[:], int32(n)) bm.Put(uint64(n), c) - bitmapSample[n/64] |= 1 << (n % 64) + w := n - 4096 + bitmapSample[w/32] |= 1 << (((n % 32) * 2) + 1) } var asBytes bytes.Buffer n, err := bm.WriteTo(&asBytes) @@ -90,11 +99,14 @@ func TestTx_CountRange(t *testing.T) { expected := uint64(0) j := uint64(0) for i := uint64(0); i < countRangeMaxN; i += 7 { + expected += i if i%4 == 3 { expected -= (j * 7) + 21 j += 7 } - got, err := tx.CountRange("i", "f", viewStandard, 0, uint64(j)<<16, uint64(i)<<16) + // Every other bit gets set, for a total of i bits in container + // i, so they're all in the first (i*2) bits of the container. + got, err := tx.CountRange("i", "f", viewStandard, 0, uint64(j)<<16, (uint64(i)<<16)+(i*2)) if err != nil { t.Fatalf("counting range: %v", err) } @@ -102,7 +114,8 @@ func TestTx_CountRange(t *testing.T) { t.Fatalf("counting from container %d to %d, expected %d, got %d", j, i, expected, got) } - expected += (i * 7) + 21 + // The -i here undoes the +i at the top of this loop. + expected += (i * 7) + 21 - i } } From a238afb21a75a0ed584d8c82c3164cba034d0792 Mon Sep 17 00:00:00 2001 From: Seebs Date: Mon, 25 Jan 2021 14:50:37 -0600 Subject: [PATCH 06/32] Handle BitmapPtr cells in countRange We need to be able to count bits in BitmapPtr containers. This only comes up if you have a non-container-aligned range count, which we never do in real production yet, but the API allows it so it should work. In order to do this, we need to provide the tx to countRange so it can grab pages as needed. Arguably, we should be able to avoid actually creating/copying that page since we're only using it internally, never returning it, but this is a pretty rare case and probably not performance-critical. --- rbf/rbf.go | 6 +++++- rbf/tx.go | 7 +++---- 2 files changed, 8 insertions(+), 5 deletions(-) diff --git a/rbf/rbf.go b/rbf/rbf.go index f13a5a7fc..365907121 100644 --- a/rbf/rbf.go +++ b/rbf/rbf.go @@ -462,7 +462,7 @@ func (c *leafCell) lastValue(tx *Tx) uint16 { // We have to take int32 rather than uint16 because the interval is [start, end), // and otherwise we have no way to ask to count the entire container (the // high bit will be missed). -func (c *leafCell) countRange(start, end int32) (n int) { +func (c *leafCell) countRange(tx *Tx, start, end int32) (n int) { // If the full range is being queried, simply use the precalculated count. if start == 0 && end > math.MaxUint16 { return c.BitN @@ -475,6 +475,10 @@ func (c *leafCell) countRange(start, end int32) (n int) { return int(roaring.RunCountRange(toInterval16(c.Data), start, end)) case ContainerTypeBitmap: return int(roaring.BitmapCountRange(toArray64(c.Data), start, end)) + case ContainerTypeBitmapPtr: + _, a, err := tx.leafCellBitmap(toPgno(c.Data)) + panicOn(err) + return int(roaring.BitmapCountRange(a, start, end)) default: panic(fmt.Sprintf("invalid container type: %d", c.Type)) } diff --git a/rbf/tx.go b/rbf/tx.go index 705edfc5b..a952104ec 100644 --- a/rbf/tx.go +++ b/rbf/tx.go @@ -1318,7 +1318,6 @@ func (tx *Tx) CountRange(name string, start, end uint64) (uint64, error) { } else if err != nil { return 0, err } - var n uint64 for { if err := csr.Next(); err == io.EOF { @@ -1341,7 +1340,7 @@ func (tx *Tx) CountRange(name string, start, end uint64) (uint64, error) { // If range is entirely in one container then just count that range. if skey == ekey { - return uint64(c.countRange(int32(lowbits(start)), ebits)), nil + return uint64(c.countRange(tx, int32(lowbits(start)), ebits)), nil } // INVAR: skey < ekey @@ -1351,7 +1350,7 @@ func (tx *Tx) CountRange(name string, start, end uint64) (uint64, error) { break } if k == skey { - n += uint64(c.countRange(int32(lowbits(start)), roaring.MaxContainerVal+1)) + n += uint64(c.countRange(tx, int32(lowbits(start)), roaring.MaxContainerVal+1)) continue } if k < ekey { @@ -1359,7 +1358,7 @@ func (tx *Tx) CountRange(name string, start, end uint64) (uint64, error) { continue } if k == ekey && ebits > 0 { - n += uint64(c.countRange(0, ebits)) + n += uint64(c.countRange(tx, 0, ebits)) break } } From 55a9952b9278226698d18579fb28681773b74565 Mon Sep 17 00:00:00 2001 From: Maxton Huff Date: Mon, 25 Jan 2021 16:19:12 -0600 Subject: [PATCH 07/32] mmap limit comparison error message --- server/server.go | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/server/server.go b/server/server.go index 291c986e1..92de59e3e 100644 --- a/server/server.go +++ b/server/server.go @@ -23,11 +23,13 @@ import ( "bytes" "context" "crypto/tls" + "encoding/binary" "io" "log" "math/rand" "net" "os" + "os/exec" "os/signal" "runtime" "strconv" @@ -152,6 +154,12 @@ func (m *Command) Start() (err error) { return errors.Wrap(err, "setting up server") } + cmd, err := exec.Command("sysctl vm.max_map_count").Output() + data := binary.BigEndian.Uint64(cmd) + if m.Config.MaxMapCount >= data { + m.logger.Printf("grpc server error: Config max map limit (%v) is greater than current system limits (%v)", m.Config.MaxMapCount, data) + } + // Set up networking (i.e. gossip) err = m.setupNetworking() if err != nil { From 3809fe673440df5c5bbb590b26f19fcd59d2130c Mon Sep 17 00:00:00 2001 From: Maxton Huff Date: Tue, 26 Jan 2021 09:35:51 -0600 Subject: [PATCH 08/32] add error check --- server/server.go | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/server/server.go b/server/server.go index 92de59e3e..8ab8402b9 100644 --- a/server/server.go +++ b/server/server.go @@ -24,6 +24,7 @@ import ( "context" "crypto/tls" "encoding/binary" + "fmt" "io" "log" "math/rand" @@ -154,10 +155,14 @@ func (m *Command) Start() (err error) { return errors.Wrap(err, "setting up server") } - cmd, err := exec.Command("sysctl vm.max_map_count").Output() + cmd, err := exec.Command("sysctl", "vm.max_map_count").Output() data := binary.BigEndian.Uint64(cmd) - if m.Config.MaxMapCount >= data { - m.logger.Printf("grpc server error: Config max map limit (%v) is greater than current system limits (%v)", m.Config.MaxMapCount, data) + if err != nil { + fmt.Println("Error: ", err) + } else { + if m.Config.MaxMapCount >= data { + m.logger.Printf("grpc server error: Config max map limit (%v) is greater than current system limits (%v)", m.Config.MaxMapCount, data) + } } // Set up networking (i.e. gossip) From 3982a8e970fd07edd60dc756b57fa74a50382473 Mon Sep 17 00:00:00 2001 From: Maxton Huff Date: Tue, 26 Jan 2021 11:26:36 -0600 Subject: [PATCH 09/32] format messages and change mmap comparison logic --- server/server.go | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/server/server.go b/server/server.go index 8ab8402b9..ac06bc0fa 100644 --- a/server/server.go +++ b/server/server.go @@ -24,7 +24,6 @@ import ( "context" "crypto/tls" "encoding/binary" - "fmt" "io" "log" "math/rand" @@ -155,13 +154,13 @@ func (m *Command) Start() (err error) { return errors.Wrap(err, "setting up server") } - cmd, err := exec.Command("sysctl", "vm.max_map_count").Output() - data := binary.BigEndian.Uint64(cmd) + result, err := exec.Command("sysctl", "vm.max_map_count").Output() if err != nil { - fmt.Println("Error: ", err) + m.logger.Printf("Tried unsuccessfully to check system mmap limit: %v", err) } else { - if m.Config.MaxMapCount >= data { - m.logger.Printf("grpc server error: Config max map limit (%v) is greater than current system limits (%v)", m.Config.MaxMapCount, data) + sysMmapLimit := binary.BigEndian.Uint64(result) + if m.Config.MaxMapCount > sysMmapLimit { + m.logger.Printf("WARNING: Config max map limit (%v) is greater than current system limits (%v)", m.Config.MaxMapCount, sysMmapLimit) } } From ac09c11badb630a375d65b930e606f17e3e14662 Mon Sep 17 00:00:00 2001 From: Nia Weiss Date: Tue, 26 Jan 2021 12:30:13 -0500 Subject: [PATCH 10/32] 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 11/32] 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 12/32] 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 927db378b403923a138df1b0fdd4b0ac2b66317e Mon Sep 17 00:00:00 2001 From: Maxton Huff Date: Wed, 27 Jan 2021 09:52:13 -0600 Subject: [PATCH 13/32] add linux OS check and the way mmap limit is read --- lattice | 2 +- server/server.go | 22 +++++++++++++--------- 2 files changed, 14 insertions(+), 10 deletions(-) diff --git a/lattice b/lattice index 36f453c1e..28c2313ec 160000 --- a/lattice +++ b/lattice @@ -1 +1 @@ -Subproject commit 36f453c1ea3bf86c546a8ad4a88f2a926724d683 +Subproject commit 28c2313ecfcd7e083d42d4e409483e968b4c421b diff --git a/server/server.go b/server/server.go index ac06bc0fa..530e4102a 100644 --- a/server/server.go +++ b/server/server.go @@ -23,16 +23,16 @@ import ( "bytes" "context" "crypto/tls" - "encoding/binary" "io" + "io/ioutil" "log" "math/rand" "net" "os" - "os/exec" "os/signal" "runtime" "strconv" + "strings" "sync" "syscall" "time" @@ -154,13 +154,17 @@ func (m *Command) Start() (err error) { return errors.Wrap(err, "setting up server") } - result, err := exec.Command("sysctl", "vm.max_map_count").Output() - if err != nil { - m.logger.Printf("Tried unsuccessfully to check system mmap limit: %v", err) - } else { - sysMmapLimit := binary.BigEndian.Uint64(result) - if m.Config.MaxMapCount > sysMmapLimit { - m.logger.Printf("WARNING: Config max map limit (%v) is greater than current system limits (%v)", m.Config.MaxMapCount, sysMmapLimit) + if runtime.GOOS == "linux" { + result, err := ioutil.ReadFile("/proc/sys/vm/max_map_count") + if err != nil { + m.logger.Printf("Tried unsuccessfully to check system mmap limit: %v", err) + } else { + sysMmapLimit, err := strconv.ParseUint(strings.TrimSuffix(string(result), "\n"), 10, 64) + if err != nil { + m.logger.Printf("Tried unsuccessfully to check system mmap limit: %v", err) + } else if m.Config.MaxMapCount > sysMmapLimit { + m.logger.Printf("WARNING: Config max map limit (%v) is greater than current system limits (%v)", m.Config.MaxMapCount, sysMmapLimit) + } } } From 60b6eecc7a81ad7438c3de013ee7a117c1eb7dce Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Thu, 28 Jan 2021 16:36:45 -0700 Subject: [PATCH 14/32] 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 15/32] 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 16/32] 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 17/32] 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 18/32] 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 19/32] 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 20/32] 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 21/32] 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 22/32] 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 23/32] 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 24/32] 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 25/32] 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 26/32] 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 27/32] 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 28/32] 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 29/32] 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 30/32] 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 31/32] 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 32/32] 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 {