From 7789e24965279c46e8a0dea0099f8798d0af4579 Mon Sep 17 00:00:00 2001 From: Travis Date: Mon, 1 Mar 2021 22:54:25 -0600 Subject: [PATCH 1/7] introduce "indexes" directory between datadir and index --- holder.go | 66 +++++++++++-------------------------------- holder_test.go | 5 ++-- server/server_test.go | 4 +-- txfactory.go | 6 +++- 4 files changed, 27 insertions(+), 54 deletions(-) diff --git a/holder.go b/holder.go index 53ba12b63..4ff2a2f03 100644 --- a/holder.go +++ b/holder.go @@ -17,9 +17,7 @@ package pilosa import ( "context" "fmt" - "io/ioutil" "os" - "path" "path/filepath" "regexp" "runtime" @@ -40,7 +38,6 @@ import ( "github.com/pilosa/pilosa/v2/topology" "github.com/pilosa/pilosa/v2/tracing" "github.com/pkg/errors" - uuid "github.com/satori/go.uuid" "golang.org/x/sync/errgroup" ) @@ -56,6 +53,9 @@ const ( // DefaultDiscoDir is the default data directory used by the disco implementation. DefaultDiscoDir = ".disco" + + // DefaultIndexesDir is the default indexes directory used by the holder. + DefaultIndexesDir = "indexes" ) func init() { @@ -296,7 +296,7 @@ func NewHolder(path string, cfg *HolderConfig) *Holder { storage.SetRowCacheOn(cfg.RowcacheOn) - txf, err := NewTxFactory(cfg.StorageConfig.Backend, path, h) + txf, err := NewTxFactory(cfg.StorageConfig.Backend, h.IndexesPath(), h) panicOn(err) h.txf = txf h.txf.blueGreenOffIfRunningBlueGreen() @@ -305,11 +305,16 @@ func NewHolder(path string, cfg *HolderConfig) *Holder { return h } -// Path() returns the path directory the holder was created with. +// Path returns the path directory the holder was created with. func (h *Holder) Path() string { return h.path } +// IndexesPath returns the path of the indexes directory. +func (h *Holder) IndexesPath() string { + return filepath.Join(h.path, DefaultIndexesDir) +} + type HolderInfo struct { FragmentInfo map[string]FragmentInfo FragmentNames []string @@ -590,7 +595,7 @@ func (h *Holder) Open() error { defer func() { h.opening = false }() if h.txf == nil { - txf, err := NewTxFactory(h.cfg.StorageConfig.Backend, h.path, h) + txf, err := NewTxFactory(h.cfg.StorageConfig.Backend, h.IndexesPath(), h) if err != nil { return errors.Wrap(err, "Holder.Open NewTxFactory()") } @@ -605,7 +610,7 @@ func (h *Holder) Open() error { h.setFileLimit() h.Logger.Printf("open holder path: %s", h.path) - if err := os.MkdirAll(h.path, 0777); err != nil { + if err := os.MkdirAll(h.IndexesPath(), 0777); err != nil { return errors.Wrap(err, "creating directory") } @@ -629,7 +634,7 @@ func (h *Holder) Open() error { } // Open path to read all index directories. - f, err := os.Open(h.path) + f, err := os.Open(h.IndexesPath()) if err != nil { return errors.Wrap(err, "opening directory") } @@ -850,13 +855,13 @@ func (h *Holder) HasData() (bool, error) { return true, nil } // Open path to read all index directories. - if _, err := os.Stat(h.path); os.IsNotExist(err) { + if _, err := os.Stat(h.IndexesPath()); os.IsNotExist(err) { return false, nil } else if err != nil { return false, errors.Wrap(err, "statting data dir") } - f, err := os.Open(h.path) + f, err := os.Open(h.IndexesPath()) if err != nil { return false, errors.Wrap(err, "opening data dir") } @@ -992,18 +997,7 @@ func (h *Holder) applySchema(schema *Schema) error { // IndexPath returns the path where a given index is stored. func (h *Holder) IndexPath(name string) string { - return filepath.Join(h.path, name) -} - -// HolderPathFromIndexPath is -// used by test/index.go:71 in test.Index.Reopen() to get the right -// path into a test Holder that doesn't know its own proper path. -// If the Holder changes index paths to being something other than -// holderPath + "/" + indexName, this will need adjusting too. -func (h *Holder) HolderPathFromIndexPath(indexPath, indexName string) string { - n := len(indexPath) - hpath2 := indexPath[:n-(len(indexName)+1)] - return hpath2 + return filepath.Join(h.IndexesPath(), name) } // Index returns the index by name. @@ -1483,31 +1477,6 @@ func (h *Holder) setFileLimit() { } } -func (h *Holder) LoadNodeID() (string, error) { - idPath := path.Join(h.path, ".id") - h.Logger.Printf("load NodeID: %s", idPath) - if err := os.MkdirAll(h.path, 0777); err != nil { - return "", errors.Wrap(err, "creating directory") - } - - nodeIDBytes, err := ioutil.ReadFile(idPath) - if err == nil { - nodeid := strings.TrimSpace(string(nodeIDBytes)) - h.Logger.Printf("I am NodeID: %s", nodeid) - return nodeid, nil - } - if !os.IsNotExist(err) { - return "", errors.Wrap(err, "reading file") - } - nodeID := uuid.NewV4().String() - err = ioutil.WriteFile(idPath, []byte(nodeID), 0600) - if err != nil { - return "", errors.Wrap(err, "writing file") - } - h.Logger.Printf("I am NodeID: %s", nodeID) - return nodeID, nil -} - // Log startup time and version to $DATA_DIR/.startup.log func (h *Holder) logStartup() error { RFC3339NanoFixedWidth := "2006-01-02T15:04:05.000000 07:00" @@ -2270,14 +2239,13 @@ func (h *Holder) Txf() *TxFactory { return h.txf } -// Begin starts a transaction on the holder. The index and shard +// BeginTx starts a transaction on the holder. The index and shard // must be specified. func (h *Holder) BeginTx(writable bool, idx *Index, shard uint64) (Tx, error) { return h.txf.NewTx(Txo{Write: writable, Index: idx, Shard: shard}), nil } func (h *Holder) HasRoaringData() (has bool, err error) { - idxs := h.Indexes() for _, idx := range idxs { paths, err := listFilesUnderDir(idx.path, false, "", true) diff --git a/holder_test.go b/holder_test.go index 6904cef88..5f5aca69d 100644 --- a/holder_test.go +++ b/holder_test.go @@ -240,7 +240,7 @@ func TestHolder_Open(t *testing.T) { t.Fatal(err) } else if err := h.Holder.Close(); err != nil { t.Fatal(err) - } else if err := os.Truncate(filepath.Join(h.Path(), "foo", "bar", "views", "standard", "fragments", "0"), 20); err != nil { + } else if err := os.Truncate(filepath.Join(h.IndexesPath(), "foo", "bar", "views", "standard", "fragments", "0"), 20); err != nil { t.Fatal(err) } @@ -353,7 +353,8 @@ func TestHolder_HasData(t *testing.T) { }) t.Run("Peek", func(t *testing.T) { - h := test.NewHolder(t) + h := test.MustOpenHolder(t) + defer h.Close() if ok, err := h.HasData(); ok || err != nil { t.Fatal("expected HasData to return false, no err, but", ok, err) diff --git a/server/server_test.go b/server/server_test.go index 578e4a9d8..f14aa75e2 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -839,7 +839,7 @@ func TestMain_ImportTimestamp(t *testing.T) { t.Fatal(err) } // Ensure the correct views were created. - dir := fmt.Sprintf("%s/%s/%s/views", m.Config.DataDir, indexName, fieldName) + dir := fmt.Sprintf("%s/%s/%s/%s/views", m.Config.DataDir, pilosa.DefaultIndexesDir, indexName, fieldName) files, err := ioutil.ReadDir(dir) if err != nil { t.Fatal(err) @@ -895,7 +895,7 @@ func TestMain_ImportTimestampNoStandardView(t *testing.T) { } // Ensure the correct views were created. - dir := fmt.Sprintf("%s/%s/%s/views", m.Config.DataDir, indexName, fieldName) + dir := fmt.Sprintf("%s/%s/%s/%s/views", m.Config.DataDir, pilosa.DefaultIndexesDir, indexName, fieldName) files, err := ioutil.ReadDir(dir) if err != nil { t.Fatal(err) diff --git a/txfactory.go b/txfactory.go index ed03bf125..1cee1d286 100644 --- a/txfactory.go +++ b/txfactory.go @@ -594,6 +594,10 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { if err != nil { return indexUsage, 0, errors.Wrap(err, "expanding data directory") } + indexesPath, err := expandDirName(f.holder.IndexesPath()) + if err != nil { + return indexUsage, 0, errors.Wrap(err, "expanding indexes directory") + } idxs := f.holder.Indexes() @@ -601,7 +605,7 @@ func (f *TxFactory) IndexUsageDetails() (map[string]IndexUsage, uint64, error) { defer qcx.Abort() for _, idx := range idxs { index := idx.name - indexPath := path.Join(holderPath, index) + indexPath := path.Join(indexesPath, index) // field usage fieldUsages := make(map[string]FieldUsage) From 4e7ec34c6d524d763010177baca6f6c8c8dfefff Mon Sep 17 00:00:00 2001 From: Travis Date: Tue, 2 Mar 2021 13:47:16 -0600 Subject: [PATCH 2/7] change directory ".disco" to "disco" --- holder.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/holder.go b/holder.go index 4ff2a2f03..91c778361 100644 --- a/holder.go +++ b/holder.go @@ -52,7 +52,7 @@ const ( existenceFieldName = "_exists" // DefaultDiscoDir is the default data directory used by the disco implementation. - DefaultDiscoDir = ".disco" + DefaultDiscoDir = "disco" // DefaultIndexesDir is the default indexes directory used by the holder. DefaultIndexesDir = "indexes" From 345d076fdf7d36652fa58a29ea0200bc39a8a3ca Mon Sep 17 00:00:00 2001 From: Travis Date: Thu, 4 Mar 2021 22:42:52 -0600 Subject: [PATCH 3/7] introduce "fields" directory between index and field --- bsi_test.go | 3 --- dbshard_internal_test.go | 45 ++++++++++++++++++++-------------------- holder.go | 5 ++++- index.go | 15 +++++++++----- rrtx.go | 8 +++---- server/server_test.go | 4 ++-- txfactory.go | 16 +++++++------- 7 files changed, 50 insertions(+), 46 deletions(-) diff --git a/bsi_test.go b/bsi_test.go index 57803a7f9..091daf409 100644 --- a/bsi_test.go +++ b/bsi_test.go @@ -64,7 +64,6 @@ func TestBSIAdd(t *testing.T) { if max < i { max = i } - t.Log("num: ", i) break } idToIndex[id] = int(i) @@ -99,8 +98,6 @@ func TestBSIAdd(t *testing.T) { } }) } - t.Log("min", min) - t.Log("max", max) } type bsiAddCase struct { diff --git a/dbshard_internal_test.go b/dbshard_internal_test.go index a5ae25348..a35e3cce1 100644 --- a/dbshard_internal_test.go +++ b/dbshard_internal_test.go @@ -25,7 +25,6 @@ import ( "github.com/pilosa/pilosa/v2/rbf" "github.com/pilosa/pilosa/v2/shardwidth" txkey "github.com/pilosa/pilosa/v2/short_txkey" - //txkey "github.com/pilosa/pilosa/v2/txkey" ) // Shard per db evaluation @@ -95,8 +94,8 @@ func Test_DBPerShard_GetShardsForIndex_LocalOnly(t *testing.T) { idx, err = NewIndex(holder, filepath.Join(tmpdir, index), index) panicOn(err) } - estd := "rick/_exists/views/standard" - std := "rick/f/views/standard" + estd := "rick/fields/_exists/views/standard" + std := "rick/fields/f/views/standard" shards, err := holder.txf.GetShardsForIndex(idx, tmpdir+sep+std, false) panicOn(err) @@ -128,7 +127,7 @@ func Test_DBPerShard_GetShardsForIndex_LocalOnly(t *testing.T) { expect0 := txkey.FieldView{Field: "_exists", View: "standard"} expect1 := txkey.FieldView{Field: "f", View: "standard"} if len(fvs) != 2 { - panic(fmt.Sprintf("fvs should be len 2, got '%#v'", fvs)) + panic(fmt.Sprintf("fvs should be len 2, got '%#v' (%s)", fvs, src)) } if fvs[0] != expect0 { panic(fmt.Sprintf("expected fvs[0]='%#v', but got '%#v'", expect0, fvs[0])) @@ -155,7 +154,7 @@ func Test_DBPerShard_GetShardsForIndex_LocalOnly(t *testing.T) { expect0 := txkey.FieldView{Field: "_exists", View: "standard"} expect1 := txkey.FieldView{Field: "f", View: "standard"} if len(fvs) != 2 { - panic(fmt.Sprintf("fvs should be len 2, got '%#v'", fvs)) + panic(fmt.Sprintf("fvs should be len 2, got '%#v' (%s)", fvs, src)) } if fvs[0] != expect0 { panic(fmt.Sprintf("expected fvs[0]='%#v', but got '%#v'", expect0, fvs[0])) @@ -173,24 +172,24 @@ func Test_DBPerShard_GetShardsForIndex_LocalOnly(t *testing.T) { // data for Test_DBPerShard_GetShardsForIndex // var sampleRoaringDirList = map[string]string{"roaring": ` -rick/f/views/standard/fragments/215.cache -rick/f/views/standard/fragments/221.cache -rick/f/views/standard/fragments/223.cache -rick/f/views/standard/fragments/93.cache -rick/f/views/standard/fragments/217.cache -rick/f/views/standard/fragments/219.cache -rick/f/views/standard/fragments/217 -rick/f/views/standard/fragments/219 -rick/f/views/standard/fragments/215 -rick/f/views/standard/fragments/221 -rick/f/views/standard/fragments/223 -rick/f/views/standard/fragments/93 -rick/_exists/views/standard/fragments/221 -rick/_exists/views/standard/fragments/215 -rick/_exists/views/standard/fragments/217 -rick/_exists/views/standard/fragments/93 -rick/_exists/views/standard/fragments/219 -rick/_exists/views/standard/fragments/223 +rick/fields/f/views/standard/fragments/215.cache +rick/fields/f/views/standard/fragments/221.cache +rick/fields/f/views/standard/fragments/223.cache +rick/fields/f/views/standard/fragments/93.cache +rick/fields/f/views/standard/fragments/217.cache +rick/fields/f/views/standard/fragments/219.cache +rick/fields/f/views/standard/fragments/217 +rick/fields/f/views/standard/fragments/219 +rick/fields/f/views/standard/fragments/215 +rick/fields/f/views/standard/fragments/221 +rick/fields/f/views/standard/fragments/223 +rick/fields/f/views/standard/fragments/93 +rick/fields/_exists/views/standard/fragments/221 +rick/fields/_exists/views/standard/fragments/215 +rick/fields/_exists/views/standard/fragments/217 +rick/fields/_exists/views/standard/fragments/93 +rick/fields/_exists/views/standard/fragments/219 +rick/fields/_exists/views/standard/fragments/223 `, "bolt": ` rick.index.txstores@@@/store-boltdb@@/shard.0093-boltdb@/bolt.db diff --git a/holder.go b/holder.go index 91c778361..f385f0b89 100644 --- a/holder.go +++ b/holder.go @@ -56,6 +56,9 @@ const ( // DefaultIndexesDir is the default indexes directory used by the holder. DefaultIndexesDir = "indexes" + + // DefaultFieldsDir is the default fields directory used by each index. + DefaultFieldsDir = "fields" ) func init() { @@ -1762,7 +1765,7 @@ func (s *holderSyncer) resetTranslationSync() error { //////////////////////////////////////////////////////////// -// translationSyncer provides an interface allowing a function +// TranslationSyncer provides an interface allowing a function // to notify the server that an action has occurred which requires // the translation sync process to be reset. In general, this // includes anything which modifies schema (add/remove index, etc), diff --git a/index.go b/index.go index 080fa6ba1..0f0725dc9 100644 --- a/index.go +++ b/index.go @@ -140,6 +140,11 @@ func (i *Index) Path() string { return i.path } +// FieldsPath returns the path of the fields directory. +func (i *Index) FieldsPath() string { + return filepath.Join(i.path, DefaultFieldsDir) +} + // TranslateStorePath returns the translation database path for a partition. func (i *Index) TranslateStorePath(partitionID int) string { return filepath.Join(i.path, translateStoreDir, strconv.Itoa(partitionID)) @@ -204,8 +209,8 @@ func (i *Index) OpenWithSchema(idx *disco.Index) error { // not validated against the schema as they are opened. func (i *Index) open(idx *disco.Index) (err error) { // Ensure the path exists. - i.holder.Logger.Debugf("ensure index path exists: %s", i.path) - if err := os.MkdirAll(i.path, 0777); err != nil { + i.holder.Logger.Debugf("ensure index path exists: %s", i.FieldsPath()) + if err := os.MkdirAll(i.FieldsPath(), 0777); err != nil { return errors.Wrap(err, "creating directory") } @@ -286,9 +291,9 @@ var indexQueue = make(chan struct{}, 8) // openFields opens and initializes the fields inside the index. func (i *Index) openFields(idx *disco.Index) error { - f, err := os.Open(i.path) + f, err := os.Open(i.FieldsPath()) if err != nil { - return errors.Wrap(err, "opening directory") + return errors.Wrap(err, "opening fields directory") } defer f.Close() @@ -508,7 +513,7 @@ func (i *Index) BeginTx(writable bool, shard uint64) (Tx, error) { } // fieldPath returns the path to a field in the index. -func (i *Index) fieldPath(name string) string { return filepath.Join(i.path, name) } +func (i *Index) fieldPath(name string) string { return filepath.Join(i.FieldsPath(), name) } // Field returns a field in the index by name. func (i *Index) Field(name string) *Field { diff --git a/rrtx.go b/rrtx.go index 2a16807ae..64715fb50 100644 --- a/rrtx.go +++ b/rrtx.go @@ -430,7 +430,7 @@ func roaringGetFieldView2Shards(idx *Index) (vs *FieldView2Shards, err error) { vs = NewFieldView2Shards() // A) open the index directory - f, err := os.Open(idx.path) + f, err := os.Open(idx.FieldsPath()) if err != nil { return nil, errors.Wrap(err, "opening directory") } @@ -453,7 +453,7 @@ func roaringGetFieldView2Shards(idx *Index) (vs *FieldView2Shards, err error) { //vv("roaringGetFieldView2Shards B) on field '%v'", field) - fieldPath := filepath.Join(idx.path, field) + fieldPath := filepath.Join(idx.FieldsPath(), field) // Skip embedded db files too. if idx.holder.txf.IsTxDatabasePath(field) { @@ -506,7 +506,7 @@ func roaringGetFieldView2Shards(idx *Index) (vs *FieldView2Shards, err error) { func (tx *RoaringTx) GetSortedFieldViewList(idx *Index, shard uint64) (fvs []txkey.FieldView, err error) { // A) open the index directory - f, err := os.Open(idx.path) + f, err := os.Open(idx.FieldsPath()) if err != nil { return nil, errors.Wrap(err, "opening directory") } @@ -529,7 +529,7 @@ func (tx *RoaringTx) GetSortedFieldViewList(idx *Index, shard uint64) (fvs []txk //vv("B) on field '%v'", field) - fieldPath := filepath.Join(idx.path, field) + fieldPath := filepath.Join(idx.FieldsPath(), field) // Skip embedded db files too. if idx.holder.txf.IsTxDatabasePath(field) { diff --git a/server/server_test.go b/server/server_test.go index f14aa75e2..43e67f5ab 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -839,7 +839,7 @@ func TestMain_ImportTimestamp(t *testing.T) { t.Fatal(err) } // Ensure the correct views were created. - dir := fmt.Sprintf("%s/%s/%s/%s/views", m.Config.DataDir, pilosa.DefaultIndexesDir, indexName, fieldName) + dir := fmt.Sprintf("%s/%s/%s/%s/%s/views", m.Config.DataDir, pilosa.DefaultIndexesDir, indexName, pilosa.DefaultFieldsDir, fieldName) files, err := ioutil.ReadDir(dir) if err != nil { t.Fatal(err) @@ -895,7 +895,7 @@ func TestMain_ImportTimestampNoStandardView(t *testing.T) { } // Ensure the correct views were created. - dir := fmt.Sprintf("%s/%s/%s/%s/views", m.Config.DataDir, pilosa.DefaultIndexesDir, indexName, fieldName) + dir := fmt.Sprintf("%s/%s/%s/%s/%s/views", m.Config.DataDir, pilosa.DefaultIndexesDir, indexName, pilosa.DefaultFieldsDir, fieldName) files, err := ioutil.ReadDir(dir) if err != nil { t.Fatal(err) diff --git a/txfactory.go b/txfactory.go index 1cee1d286..b86566d71 100644 --- a/txfactory.go +++ b/txfactory.go @@ -708,7 +708,7 @@ func (f *TxFactory) fieldUsage(indexPath string, fld *Field) (FieldUsage, error) } // field metadata, e.g. rowAttrs - fieldPath := path.Join(indexPath, field) + fieldPath := path.Join(indexPath, DefaultFieldsDir, 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) @@ -1005,19 +1005,19 @@ func fragmentSpecFromRoaringPath(path string) (field, view string, shard uint64, } // sample path: - // field view shard - // myfield/views/standard/fragments/0 + // field view shard + // fields/myfield/views/standard/fragments/0 s := strings.Split(path, "/") n := len(s) - if n != 5 { + if n != 6 { err = fmt.Errorf("len(s)=%v, but expected 5. path='%v'", n, path) return } - field = s[0] - view = s[2] - shard, err = strconv.ParseUint(s[4], 10, 64) + field = s[1] + view = s[3] + shard, err = strconv.ParseUint(s[5], 10, 64) if err != nil { - err = fmt.Errorf("fragmentSpecFromRoaringPath(path='%v') could not parse shard '%v' as uint: '%v'", path, s[4], err) + err = fmt.Errorf("fragmentSpecFromRoaringPath(path='%v') could not parse shard '%v' as uint: '%v'", path, s[5], err) } return } From 4738a818d2bf339b8cc5d6da25edb02e7c7692d5 Mon Sep 17 00:00:00 2001 From: Travis Date: Thu, 4 Mar 2021 23:12:42 -0600 Subject: [PATCH 4/7] change attributes file ".data" to "column/row-attributes" --- holder.go | 8 +++++++- holder_test.go | 4 ++-- index.go | 2 +- 3 files changed, 10 insertions(+), 4 deletions(-) diff --git a/holder.go b/holder.go index f385f0b89..065ac4ded 100644 --- a/holder.go +++ b/holder.go @@ -59,6 +59,12 @@ const ( // DefaultFieldsDir is the default fields directory used by each index. DefaultFieldsDir = "fields" + + // ColumnAttrsFileName is the name of the file used for the column attributes store. + ColumnAttrsFileName = "column-attributes" + + // RowAttrsFileName is the name of the file used for the row attributes store. + RowAttrsFileName = "row-attributes" ) func init() { @@ -1304,7 +1310,7 @@ func (h *Holder) newIndex(path, name string) (*Index, error) { index.serializer = h.serializer index.Schemator = h.schemator index.newAttrStore = h.NewAttrStore - index.columnAttrs = h.NewAttrStore(filepath.Join(index.path, ".data")) + index.columnAttrs = h.NewAttrStore(filepath.Join(index.path, ColumnAttrsFileName)) index.OpenTranslateStore = h.OpenTranslateStore index.translationSyncer = h.translationSyncer return index, nil diff --git a/holder_test.go b/holder_test.go index 5f5aca69d..f3ba426af 100644 --- a/holder_test.go +++ b/holder_test.go @@ -63,7 +63,7 @@ func TestHolder_Open(t *testing.T) { t.Fatal(err) } else if err := h.Holder.Close(); err != nil { t.Fatal(err) - } else if err := os.Truncate(filepath.Join(h.IndexPath("test"), ".data"), 2); err != nil { + } else if err := os.Truncate(filepath.Join(h.IndexPath("test"), pilosa.ColumnAttrsFileName), 2); err != nil { t.Fatal(err) } @@ -135,7 +135,7 @@ func TestHolder_Open(t *testing.T) { t.Fatal(err) } else if err := h.Holder.Close(); err != nil { t.Fatal(err) - } else if err := os.Truncate(filepath.Join(h.Path(), "foo", "bar", ".data"), 2); err != nil { + } else if err := os.Truncate(filepath.Join(h.Path(), "foo", "bar", pilosa.RowAttrsFileName), 2); err != nil { t.Fatal(err) } diff --git a/index.go b/index.go index 0f0725dc9..ce7c235dd 100644 --- a/index.go +++ b/index.go @@ -795,7 +795,7 @@ func (i *Index) newField(path, name string) (*Field, error) { f.broadcaster = i.broadcaster f.schemator = i.Schemator f.serializer = i.serializer - f.rowAttrStore = i.newAttrStore(filepath.Join(f.path, ".data")) + f.rowAttrStore = i.newAttrStore(filepath.Join(f.path, RowAttrsFileName)) f.OpenTranslateStore = i.OpenTranslateStore return f, nil } From 4661f3ac0875e03d432d07a604315c7eb935849f Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 5 Mar 2021 12:16:10 -0600 Subject: [PATCH 5/7] reorganize the storage backends directory --- bolt.go | 4 ++-- dbshard.go | 13 +++++++++++-- dbshard_internal_test.go | 42 +++++++++++++++++++++++----------------- holder.go | 12 ++---------- index.go | 4 ---- rbf.go | 8 ++++---- rrtx.go | 8 -------- txfactory.go | 21 ++------------------ 8 files changed, 45 insertions(+), 67 deletions(-) diff --git a/bolt.go b/bolt.go index 9a40ef53a..7e654d76d 100644 --- a/bolt.go +++ b/bolt.go @@ -115,8 +115,8 @@ func DumpAllBolt() { // boltPath is a helper for determining the full directory // in which the bolt database will be stored. func boltPath(path string) string { - if !strings.HasSuffix(path, "-boltdb@") { - return path + "-boltdb@" + if !strings.HasSuffix(path, "-boltdb") { + return path + "-boltdb" } return path } diff --git a/dbshard.go b/dbshard.go index b2062abe2..1cdfaa912 100644 --- a/dbshard.go +++ b/dbshard.go @@ -33,6 +33,15 @@ import ( var _ = sort.Sort +const ( + // DefaultBackendsDir is the default backends directory used to store the + // data for each backend. + DefaultBackendsDir = "backends" + + // DefaultBackendDirPrefix is the default prefix of each backend directory. + DefaultBackendDirPrefix = "backend" +) + // types to support a database file per shard type DBHolder struct { @@ -561,7 +570,7 @@ func (dbs *DBShard) pathForType(ty txtype) string { // what here for roaring? well, roaringRegistrar.OpenDBWrapper() // is a no-op anyhow. so doesn't need to be correct atm. - path := dbs.HolderPath + sep + dbs.Index + ".index.txstores@@@" + sep + "store" + ty.FileSuffix() + "@" + sep + fmt.Sprintf("shard.%04v%v", dbs.Shard, ty.FileSuffix()) + path := dbs.HolderPath + sep + dbs.Index + sep + DefaultBackendsDir + sep + DefaultBackendDirPrefix + ty.FileSuffix() + sep + fmt.Sprintf("shard.%04v%v", dbs.Shard, ty.FileSuffix()) if ty == boltTxn { // special case: // bolt doesn't use a directory like the others, just a direct path. @@ -574,7 +583,7 @@ func (dbs *DBShard) pathForType(ty txtype) string { // prefixForType and pathForType must be kept in sync! func (per *DBPerShard) prefixForType(idx *Index, ty txtype) string { // top level paths will end in "@@" - return per.HolderDir + sep + idx.name + ".index.txstores@@@" + sep + "store" + ty.FileSuffix() + "@" + sep + return per.HolderDir + sep + idx.name + sep + DefaultBackendsDir + sep + DefaultBackendDirPrefix + ty.FileSuffix() + sep } var ErrNoData = fmt.Errorf("no data") diff --git a/dbshard_internal_test.go b/dbshard_internal_test.go index a35e3cce1..398ed7569 100644 --- a/dbshard_internal_test.go +++ b/dbshard_internal_test.go @@ -89,7 +89,7 @@ func Test_DBPerShard_GetShardsForIndex_LocalOnly(t *testing.T) { holder := NewHolder(tmpdir, cfg) index := "rick" - idx := makeSampleRoaringDir(tmpdir, index, src, 1, holder, v2s) + idx := makeSampleRoaringDir(t, tmpdir, index, src, 1, holder, v2s) if idx == nil { idx, err = NewIndex(holder, filepath.Join(tmpdir, index), index) panicOn(err) @@ -192,39 +192,41 @@ rick/fields/_exists/views/standard/fragments/219 rick/fields/_exists/views/standard/fragments/223 `, "bolt": ` -rick.index.txstores@@@/store-boltdb@@/shard.0093-boltdb@/bolt.db -rick.index.txstores@@@/store-boltdb@@/shard.0215-boltdb@/bolt.db -rick.index.txstores@@@/store-boltdb@@/shard.0217-boltdb@/bolt.db -rick.index.txstores@@@/store-boltdb@@/shard.0219-boltdb@/bolt.db -rick.index.txstores@@@/store-boltdb@@/shard.0221-boltdb@/bolt.db -rick.index.txstores@@@/store-boltdb@@/shard.0223-boltdb@/bolt.db +rick/backends/backend-boltdb/shard.0093-bolt/bolt.db +rick/backends/backend-boltdb/shard.0215-bolt/bolt.db +rick/backends/backend-boltdb/shard.0217-bolt/bolt.db +rick/backends/backend-boltdb/shard.0219-bolt/bolt.db +rick/backends/backend-boltdb/shard.0221-bolt/bolt.db +rick/backends/backend-boltdb/shard.0223-bolt/bolt.db `, "rbf": ` -rick.index.txstores@@@/store-rbfdb@@/shard.0093-rbfdb@ -rick.index.txstores@@@/store-rbfdb@@/shard.0215-rbfdb@ -rick.index.txstores@@@/store-rbfdb@@/shard.0217-rbfdb@ -rick.index.txstores@@@/store-rbfdb@@/shard.0219-rbfdb@ -rick.index.txstores@@@/store-rbfdb@@/shard.0221-rbfdb@ -rick.index.txstores@@@/store-rbfdb@@/shard.0223-rbfdb@ +rick/backends/backend-rbf/shard.0093-rbf +rick/backends/backend-rbf/shard.0215-rbf +rick/backends/backend-rbf/shard.0217-rbf +rick/backends/backend-rbf/shard.0219-rbf +rick/backends/backend-rbf/shard.0221-rbf +rick/backends/backend-rbf/shard.0223-rbf `, } -func makeSampleRoaringDir(root, index, backend string, minBytes int, h *Holder, view2shards *FieldView2Shards) (idx *Index) { +func makeSampleRoaringDir(t *testing.T, root, index, backend string, minBytes int, h *Holder, view2shards *FieldView2Shards) (idx *Index) { shards := []uint64{0, 93, 215, 217, 219, 221, 223} fns := strings.Split(sampleRoaringDirList[backend], "\n") firstDone := false for i, fn := range fns { + // This check is here because in sampleRoaringDirList, the first entry + // of each map value is a line feed, so the strings.Split() above + // results in a blank entry for the first item. This means that the + // slice of shards above has an initial entry "0" which is not used. if fn == "" { continue } var shard uint64 - if backend != "roaring" { - // only have shards for the non-roaring - shard = shards[i] - } switch backend { case "bolt", "rbf": + shard = shards[i] + idx = helperCreateDBShard(h, index, shard) // first time only, we'll actually make all the shards at this point because @@ -234,6 +236,9 @@ func makeSampleRoaringDir(root, index, backend string, minBytes int, h *Holder, makeTxTestDBWithViewsShards(h, idx, view2shards) } continue + case "roaring": + default: + t.Fatalf("invalid backend: %s", backend) } path := root + sep + filepath.Dir(fn) @@ -252,6 +257,7 @@ func makeSampleRoaringDir(root, index, backend string, minBytes int, h *Holder, func helperCreateDBShard(h *Holder, index string, shard uint64) *Index { idx, err := h.CreateIndexIfNotExists(index, IndexOptions{}) panicOn(err) + // TODO: It's not clear that this is actually doing anything. dbs, err := h.txf.dbPerShard.GetDBShard(index, shard, idx) panicOn(err) _ = dbs diff --git a/holder.go b/holder.go index 065ac4ded..93fb36188 100644 --- a/holder.go +++ b/holder.go @@ -659,10 +659,6 @@ func (h *Holder) Open() error { if !fi.IsDir() || strings.HasPrefix(fi.Name(), ".") { continue } - // Skip embedded db files too. - if h.txf.IsTxDatabasePath(fi.Name()) { - continue - } // Only continue with indexes which are present in schema. idx, ok := schema[fi.Name()] @@ -885,10 +881,6 @@ func (h *Holder) HasData() (bool, error) { if !fi.IsDir() { continue } - // Skip embedded db files too. - if h.txf.IsTxDatabasePath(fi.Name()) { - continue - } // Skip DisCo data directory. if fi.Name() == DefaultDiscoDir { @@ -1486,13 +1478,13 @@ func (h *Holder) setFileLimit() { } } -// Log startup time and version to $DATA_DIR/.startup.log +// Log startup time and version to $DATA_DIR/startup.log func (h *Holder) logStartup() error { 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) + f, err := os.OpenFile(h.path+"/startup.log", os.O_APPEND|os.O_WRONLY|os.O_CREATE, 0600) if err != nil { return errors.Wrap(err, "opening startup log") } diff --git a/index.go b/index.go index ce7c235dd..411f37035 100644 --- a/index.go +++ b/index.go @@ -314,10 +314,6 @@ fileLoop: if !fi.IsDir() { continue } - // Skip embedded db files too. - if i.holder.txf.IsTxDatabasePath(fi.Name()) { - continue - } var cfm *CreateFieldMessage = &CreateFieldMessage{} var err error diff --git a/rbf.go b/rbf.go index 874d996f4..e49bd6251 100644 --- a/rbf.go +++ b/rbf.go @@ -144,15 +144,15 @@ func (r *rbfDBRegistrar) unregister(w *RbfDBWrapper) { // rbfPath is a helper for determining the full directory // in which the RBF database will be stored. func rbfPath(path string) string { - if !strings.HasSuffix(path, "-rbfdb@") { - return path + "-rbfdb@" + if !strings.HasSuffix(path, "-rbf") { + return path + "-rbf" } return path } -// OpenDBWrapper opens the database in the path directoy +// OpenDBWrapper opens the database in the path directory // without deleting any prior content. Any -// database directory will have the "-rbfdb@" suffix. +// database directory will have the "-rbf" suffix. // // OpenDBWrapper will check the registry and make a new instance only // if one does not exist for its path. Otherwise it returns diff --git a/rrtx.go b/rrtx.go index 64715fb50..c2ee55748 100644 --- a/rrtx.go +++ b/rrtx.go @@ -455,10 +455,6 @@ func roaringGetFieldView2Shards(idx *Index) (vs *FieldView2Shards, err error) { fieldPath := filepath.Join(idx.FieldsPath(), field) - // Skip embedded db files too. - if idx.holder.txf.IsTxDatabasePath(field) { - continue - } viewsDir := filepath.Join(fieldPath, "views") file, err := os.Open(viewsDir) if os.IsNotExist(err) { @@ -531,10 +527,6 @@ func (tx *RoaringTx) GetSortedFieldViewList(idx *Index, shard uint64) (fvs []txk fieldPath := filepath.Join(idx.FieldsPath(), field) - // Skip embedded db files too. - if idx.holder.txf.IsTxDatabasePath(field) { - continue - } viewsDir := filepath.Join(fieldPath, "views") file, err := os.Open(viewsDir) if os.IsNotExist(err) { diff --git a/txfactory.go b/txfactory.go index b86566d71..b9f51c601 100644 --- a/txfactory.go +++ b/txfactory.go @@ -423,10 +423,6 @@ const ( boltTxn txtype = 4 ) -// these need to be skipped by the holder.go field scanner that -// calls IsTxDatabasePath -var allTypesWithSuffixes = []txtype{rbfTxn, boltTxn} - // FileSuffix is used to determine backend directory names. // We append '@' to be sure we never collide with a field name // inside the index directory. In the future for different @@ -437,26 +433,13 @@ func (ty txtype) FileSuffix() string { case roaringTxn: return "" case rbfTxn: - return "-rbfdb@" + return "-rbf" case boltTxn: - return "-boltdb@" + return "-boltdb" } panic(fmt.Sprintf("unkown txtype %v", int(ty))) } -func (txf *TxFactory) IsTxDatabasePath(path string) bool { - if strings.HasSuffix(filepath.Base(path), ".txstores@@@") { - // top level dir - return true - } - for _, ty := range allTypesWithSuffixes { - if strings.HasSuffix(path, ty.FileSuffix()) { - return true - } - } - return false -} - func (txf *TxFactory) NeedsSnapshot() (b bool) { for _, ty := range txf.types { switch ty { From 5191d0c3f4a87bfd9a1fdeb1f3b97037c5038bac Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 5 Mar 2021 12:22:02 -0600 Subject: [PATCH 6/7] remove unnecessary skip check --- holder.go | 6 ------ 1 file changed, 6 deletions(-) diff --git a/holder.go b/holder.go index 93fb36188..6fe1abbd5 100644 --- a/holder.go +++ b/holder.go @@ -881,12 +881,6 @@ func (h *Holder) HasData() (bool, error) { if !fi.IsDir() { continue } - - // Skip DisCo data directory. - if fi.Name() == DefaultDiscoDir { - continue - } - return true, nil } return false, nil From cad78f1d34a9fbf5b4ae95ebaca63ee057382f0b Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 5 Mar 2021 16:38:27 -0600 Subject: [PATCH 7/7] remove "Default*" from const names --- dbshard.go | 16 ++++++---------- holder.go | 14 +++++++------- index.go | 2 +- server/server.go | 2 +- server/server_test.go | 4 ++-- txfactory.go | 5 ++--- 6 files changed, 19 insertions(+), 24 deletions(-) diff --git a/dbshard.go b/dbshard.go index 1cdfaa912..d381c9f85 100644 --- a/dbshard.go +++ b/dbshard.go @@ -34,12 +34,12 @@ import ( var _ = sort.Sort const ( - // DefaultBackendsDir is the default backends directory used to store the + // backendsDir is the default backends directory used to store the // data for each backend. - DefaultBackendsDir = "backends" + backendsDir = "backends" - // DefaultBackendDirPrefix is the default prefix of each backend directory. - DefaultBackendDirPrefix = "backend" + // backendDirPrefix is the default prefix of each backend directory. + backendDirPrefix = "backend" ) // types to support a database file per shard @@ -133,10 +133,6 @@ func (dbs *DBShard) Close() (err error) { return } -func (dbs *DBShard) HolderString() string { - return dbs.HolderPath -} - // Cleanup must be called at every commit/rollback of a Tx, in // order to release the read-write mutex that guarantees a single // writer at a time. Each tx must take care to call cleanup() @@ -570,7 +566,7 @@ func (dbs *DBShard) pathForType(ty txtype) string { // what here for roaring? well, roaringRegistrar.OpenDBWrapper() // is a no-op anyhow. so doesn't need to be correct atm. - path := dbs.HolderPath + sep + dbs.Index + sep + DefaultBackendsDir + sep + DefaultBackendDirPrefix + ty.FileSuffix() + sep + fmt.Sprintf("shard.%04v%v", dbs.Shard, ty.FileSuffix()) + path := dbs.HolderPath + sep + dbs.Index + sep + backendsDir + sep + backendDirPrefix + ty.FileSuffix() + sep + fmt.Sprintf("shard.%04v%v", dbs.Shard, ty.FileSuffix()) if ty == boltTxn { // special case: // bolt doesn't use a directory like the others, just a direct path. @@ -583,7 +579,7 @@ func (dbs *DBShard) pathForType(ty txtype) string { // prefixForType and pathForType must be kept in sync! func (per *DBPerShard) prefixForType(idx *Index, ty txtype) string { // top level paths will end in "@@" - return per.HolderDir + sep + idx.name + sep + DefaultBackendsDir + sep + DefaultBackendDirPrefix + ty.FileSuffix() + sep + return per.HolderDir + sep + idx.name + sep + backendsDir + sep + backendDirPrefix + ty.FileSuffix() + sep } var ErrNoData = fmt.Errorf("no data") diff --git a/holder.go b/holder.go index 6fe1abbd5..38a1df3aa 100644 --- a/holder.go +++ b/holder.go @@ -51,14 +51,14 @@ const ( // existenceFieldName is the name of the internal field used to store existence values. existenceFieldName = "_exists" - // DefaultDiscoDir is the default data directory used by the disco implementation. - DefaultDiscoDir = "disco" + // DiscoDir is the default data directory used by the disco implementation. + DiscoDir = "disco" - // DefaultIndexesDir is the default indexes directory used by the holder. - DefaultIndexesDir = "indexes" + // IndexesDir is the default indexes directory used by the holder. + IndexesDir = "indexes" - // DefaultFieldsDir is the default fields directory used by each index. - DefaultFieldsDir = "fields" + // FieldsDir is the default fields directory used by each index. + FieldsDir = "fields" // ColumnAttrsFileName is the name of the file used for the column attributes store. ColumnAttrsFileName = "column-attributes" @@ -321,7 +321,7 @@ func (h *Holder) Path() string { // IndexesPath returns the path of the indexes directory. func (h *Holder) IndexesPath() string { - return filepath.Join(h.path, DefaultIndexesDir) + return filepath.Join(h.path, IndexesDir) } type HolderInfo struct { diff --git a/index.go b/index.go index 411f37035..6bec27db8 100644 --- a/index.go +++ b/index.go @@ -142,7 +142,7 @@ func (i *Index) Path() string { // FieldsPath returns the path of the fields directory. func (i *Index) FieldsPath() string { - return filepath.Join(i.path, DefaultFieldsDir) + return filepath.Join(i.path, FieldsDir) } // TranslateStorePath returns the translation database path for a partition. diff --git a/server/server.go b/server/server.go index ab94b40f5..40ff3a7ec 100644 --- a/server/server.go +++ b/server/server.go @@ -388,7 +388,7 @@ func (m *Command) SetupServer() error { if err != nil { return errors.Wrapf(err, "expanding directory name: %s", m.Config.DataDir) } - m.Config.Etcd.Dir = filepath.Join(path, pilosa.DefaultDiscoDir) + m.Config.Etcd.Dir = filepath.Join(path, pilosa.DiscoDir) } e := petcd.NewEtcdWithCache(m.Config.Etcd, m.Config.Cluster.ReplicaN) diff --git a/server/server_test.go b/server/server_test.go index 43e67f5ab..44aee924e 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -839,7 +839,7 @@ func TestMain_ImportTimestamp(t *testing.T) { t.Fatal(err) } // Ensure the correct views were created. - dir := fmt.Sprintf("%s/%s/%s/%s/%s/views", m.Config.DataDir, pilosa.DefaultIndexesDir, indexName, pilosa.DefaultFieldsDir, fieldName) + dir := fmt.Sprintf("%s/%s/%s/%s/%s/views", m.Config.DataDir, pilosa.IndexesDir, indexName, pilosa.FieldsDir, fieldName) files, err := ioutil.ReadDir(dir) if err != nil { t.Fatal(err) @@ -895,7 +895,7 @@ func TestMain_ImportTimestampNoStandardView(t *testing.T) { } // Ensure the correct views were created. - dir := fmt.Sprintf("%s/%s/%s/%s/%s/views", m.Config.DataDir, pilosa.DefaultIndexesDir, indexName, pilosa.DefaultFieldsDir, fieldName) + dir := fmt.Sprintf("%s/%s/%s/%s/%s/views", m.Config.DataDir, pilosa.IndexesDir, indexName, pilosa.FieldsDir, fieldName) files, err := ioutil.ReadDir(dir) if err != nil { t.Fatal(err) diff --git a/txfactory.go b/txfactory.go index b9f51c601..bb60299f8 100644 --- a/txfactory.go +++ b/txfactory.go @@ -691,7 +691,7 @@ func (f *TxFactory) fieldUsage(indexPath string, fld *Field) (FieldUsage, error) } // field metadata, e.g. rowAttrs - fieldPath := path.Join(indexPath, DefaultFieldsDir, field) + fieldPath := path.Join(indexPath, FieldsDir, 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) @@ -750,10 +750,9 @@ func directoryUsage(fname string, recursive bool) (uint64, error) { return size, nil } +// CloseIndex is a no-op. This seems to be in place for debugging purposes. func (f *TxFactory) CloseIndex(idx *Index) error { - // under roaring and all the new databases, this is a no-op. //idx.Dump("CloseIndex") - return nil }