From e397d35ed5c723d6a6539ccff582234b292e7f69 Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Wed, 13 Jan 2021 17:53:05 -0600 Subject: [PATCH 01/14] 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 02/14] 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 03/14] 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 04/14] 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 05/14] 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 06/14] 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 07/14] 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 08/14] 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 09/14] 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 10/14] 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 11/14] 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 12/14] 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 13/14] 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 14/14] 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`, })