mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 11:27:50 +00:00
Merge pull request #1512 from travisturner/datadir-take-2
Core-316 Reorganize Pilosa datadir
This commit is contained in:
commit
ab50403c35
12 changed files with 134 additions and 184 deletions
4
bolt.go
4
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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
17
dbshard.go
17
dbshard.go
|
|
@ -33,6 +33,15 @@ import (
|
|||
|
||||
var _ = sort.Sort
|
||||
|
||||
const (
|
||||
// backendsDir is the default backends directory used to store the
|
||||
// data for each backend.
|
||||
backendsDir = "backends"
|
||||
|
||||
// backendDirPrefix is the default prefix of each backend directory.
|
||||
backendDirPrefix = "backend"
|
||||
)
|
||||
|
||||
// types to support a database file per shard
|
||||
|
||||
type DBHolder struct {
|
||||
|
|
@ -124,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()
|
||||
|
|
@ -561,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 + ".index.txstores@@@" + sep + "store" + 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.
|
||||
|
|
@ -574,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 + ".index.txstores@@@" + sep + "store" + ty.FileSuffix() + "@" + sep
|
||||
return per.HolderDir + sep + idx.name + sep + backendsDir + sep + backendDirPrefix + ty.FileSuffix() + sep
|
||||
}
|
||||
|
||||
var ErrNoData = fmt.Errorf("no data")
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -90,13 +89,13 @@ 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)
|
||||
}
|
||||
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,59 +172,61 @@ 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
|
||||
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
|
||||
|
|
@ -235,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)
|
||||
|
|
@ -253,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
|
||||
|
|
|
|||
101
holder.go
101
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"
|
||||
)
|
||||
|
||||
|
|
@ -54,8 +51,20 @@ 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"
|
||||
|
||||
// IndexesDir is the default indexes directory used by the holder.
|
||||
IndexesDir = "indexes"
|
||||
|
||||
// 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"
|
||||
|
||||
// RowAttrsFileName is the name of the file used for the row attributes store.
|
||||
RowAttrsFileName = "row-attributes"
|
||||
)
|
||||
|
||||
func init() {
|
||||
|
|
@ -296,7 +305,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 +314,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, IndexesDir)
|
||||
}
|
||||
|
||||
type HolderInfo struct {
|
||||
FragmentInfo map[string]FragmentInfo
|
||||
FragmentNames []string
|
||||
|
|
@ -590,7 +604,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 +619,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 +643,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")
|
||||
}
|
||||
|
|
@ -645,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()]
|
||||
|
|
@ -850,13 +860,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")
|
||||
}
|
||||
|
|
@ -871,16 +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 {
|
||||
continue
|
||||
}
|
||||
|
||||
return true, nil
|
||||
}
|
||||
return false, nil
|
||||
|
|
@ -992,18 +992,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.
|
||||
|
|
@ -1307,7 +1296,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
|
||||
|
|
@ -1483,38 +1472,13 @@ 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
|
||||
// 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")
|
||||
}
|
||||
|
|
@ -1793,7 +1757,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),
|
||||
|
|
@ -2270,14 +2234,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)
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
||||
|
|
@ -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)
|
||||
|
|
|
|||
21
index.go
21
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, FieldsDir)
|
||||
}
|
||||
|
||||
// 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()
|
||||
|
||||
|
|
@ -309,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
|
||||
|
|
@ -508,7 +509,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 {
|
||||
|
|
@ -790,7 +791,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
|
||||
}
|
||||
|
|
|
|||
8
rbf.go
8
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
|
||||
|
|
|
|||
16
rrtx.go
16
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,12 +453,8 @@ 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) {
|
||||
continue
|
||||
}
|
||||
viewsDir := filepath.Join(fieldPath, "views")
|
||||
file, err := os.Open(viewsDir)
|
||||
if os.IsNotExist(err) {
|
||||
|
|
@ -506,7 +502,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,12 +525,8 @@ 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) {
|
||||
continue
|
||||
}
|
||||
viewsDir := filepath.Join(fieldPath, "views")
|
||||
file, err := os.Open(viewsDir)
|
||||
if os.IsNotExist(err) {
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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/%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/views", m.Config.DataDir, indexName, 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)
|
||||
|
|
|
|||
46
txfactory.go
46
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 {
|
||||
|
|
@ -594,6 +577,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 +588,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)
|
||||
|
|
@ -704,7 +691,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, 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)
|
||||
|
|
@ -763,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
|
||||
}
|
||||
|
||||
|
|
@ -1001,19 +987,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
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue