diff --git a/Makefile b/Makefile index 65840bf4d..ba3d994a8 100644 --- a/Makefile +++ b/Makefile @@ -72,6 +72,9 @@ testvsub-race: cd ..; \ done +tour: + ./tournament.sh + bench: go test ./... -bench=. -run=NoneZ -timeout=127m $(TESTFLAGS) diff --git a/badger.go b/badger.go index 39500f262..f4d16facc 100644 --- a/badger.go +++ b/badger.go @@ -290,6 +290,10 @@ func (r *badgerRegistrar) OpenDBWrapper(bpath string, doAllocZero bool) (DBWrapp // have the -badgerdb suffix. bpath = badgerPath(bpath) + err := os.MkdirAll(bpath, 0755) + if err != nil { + return nil, err + } r.mu.Lock() defer r.mu.Unlock() w, ok := r.path2db[bpath] @@ -1734,7 +1738,7 @@ func (tx *BadgerTx) countBitsSet(bkey []byte) (n int) { return } -func (tx *BadgerTx) Dump(short bool) { +func (tx *BadgerTx) Dump(short bool, shard uint64) { fmt.Printf("BadgerTx %p Dump: %v\n", tx, stringifiedBadgerKeysTx(tx, short)) } diff --git a/bluegreentx.go b/bluegreentx.go index 033f2809c..53f6ef6d5 100644 --- a/bluegreentx.go +++ b/bluegreentx.go @@ -18,6 +18,7 @@ import ( "bytes" "fmt" "io" + "math" "reflect" "sort" "sync" @@ -152,15 +153,15 @@ func (c *blueGreenTx) Type() string { var blueGreenTxDumpMut sync.Mutex -func (c *blueGreenTx) Dump(short bool) { +func (c *blueGreenTx) Dump(short bool, shard uint64) { blueGreenTxDumpMut.Lock() defer blueGreenTxDumpMut.Unlock() fmt.Printf("%v blueGreenTx.Dump ============== \n", FileLine(2)) fmt.Printf("A(%v) Dump:\n", c.as) - c.a.Dump(short) + c.a.Dump(short, shard) fmt.Printf("B(%v) Dump:\n", c.bs) - c.b.Dump(short) + c.b.Dump(short, shard) if !short { fmt.Printf("dbPerShard.DumpAll(): idx=%p\n", c.idx) @@ -211,21 +212,21 @@ func (c *blueGreenTx) compareTxState(index, field, view string, shard uint64) { } if aFound != bFound { - c.Dump(c.short) + c.Dump(c.short, shard) panic(fmt.Sprintf("compareTxState[%v]: A(%v) ContainerIterator had aFound=%v, but B(%v) had bFound=%v; at '%v'", here, c.as, aFound, c.bs, bFound, stack())) } if aErr != nil || bErr != nil { if aErr != nil && bErr != nil { - c.Dump(c.short) + c.Dump(c.short, shard) panic(fmt.Sprintf("compareTxState[%v]: A(%v) reported err '%v'; B(%v) reported err '%v' at %v", here, c.as, aErr, c.bs, bErr, stack())) } if aErr != nil { - c.Dump(c.short) + c.Dump(c.short, shard) panic(fmt.Sprintf("compareTxState[%v]: A(%v) reported err %v at %v; but B(%v) did not", here, c.as, aErr, c.bs, stack())) } if bErr != nil { - c.Dump(c.short) + c.Dump(c.short, shard) panic(fmt.Sprintf("compareTxState[%v]: B(%v) reported err %v at %v; but A(%v) did not", here, c.bs, bErr, c.as, stack())) } } @@ -234,17 +235,17 @@ func (c *blueGreenTx) compareTxState(index, field, view string, shard uint64) { if !bIter.Next() { AlwaysPrintf("compareTxState[%v]: A(%v) found key %v, B(%v) didn't, dump to follow, stack=\n %v\n\n and here is dump:", here, c.as, aKey, c.bs, stack()) - c.Dump(c.short) + c.Dump(c.short, shard) panic(fmt.Sprintf("compareTxState[%v]: A(%v) found key %v, B(%v) didn't, at %v", here, c.as, aKey, c.bs, stack())) } bKey, bValue := bIter.Value() if bKey != aKey { AlwaysPrintf("problem in caller %v", Caller(2)) - c.Dump(c.short) + c.Dump(c.short, shard) panic(fmt.Sprintf("compareTxState[%v]: A(%v) found key %v, B(%v) found %v, at %v", here, c.as, aKey, c.bs, bKey, stack())) } if err := aValue.BitwiseCompare(bValue); err != nil { - c.Dump(c.short) + c.Dump(c.short, shard) //vv("compareTxState[%v]: key %v differs: %v; A=%v; B=%v; at stack=%v", here, aKey, err, c.as, c.bs, stack()) panic(fmt.Sprintf("compareTxState[%v]: key %v differs: %v; A=%v; B=%v; at stack=%v", here, aKey, err, c.as, c.bs, stack())) } @@ -253,7 +254,7 @@ func (c *blueGreenTx) compareTxState(index, field, view string, shard uint64) { // end checking everything in A, but does B have more? if bIter.Next() { AlwaysPrintf("bIter has more than it should. problem in caller %v. _sn_ %v", Caller(2), c.Sn()) - c.Dump(c.short) + c.Dump(c.short, shard) bKey, _ := bIter.Value() panic(fmt.Sprintf("compareTxState[%v]: B(%v) found key %v, A(%v) didn't, (a.sn=%v) (b.sn=%v) at %v", here, c.bs, bKey, c.as, c.a.Sn(), c.b.Sn(), stack())) } @@ -430,7 +431,7 @@ func (c *blueGreenTx) ImportRoaringBits(index, field, view string, shard uint64, // ================== begin save comments. //c.checkDatabase() ////vv("got past database check at TOP of ImportRoaringBits") - //c.Dump(c.short) + //c.Dump(c.short, shard) ////vv("done with top dump; clear=%v", clear) // ================== end save comments. defer func() { @@ -789,7 +790,7 @@ func (c *blueGreenTx) CountRange(index, field, view string, shard uint64, start, c.checker.see(index, field, view, shard) defer func() { if r := recover(); r != nil { - c.Dump(c.short) + c.Dump(c.short, shard) AlwaysPrintf("see CountRange() panic '%v' at '%v'", r, stack()) panic(r) } @@ -822,7 +823,7 @@ func (c *blueGreenTx) OffsetRange(index, field, view string, shard, offset, star err = roaringBitmapDiff(a, b) if err != nil { - c.Dump(false) + c.Dump(false, shard) panicOn(fmt.Errorf("on _sn_ %v OffsetRange(index='%v', field='%v', view='%v', shard='%v', offset: %v start: %v, end: %v) err: %v", c.Sn(), index, field, view, int(shard), offset, start, end, err)) } compareErrors(errA, errB) @@ -835,7 +836,7 @@ func (c *blueGreenTx) RoaringBitmapReader(index, field, view string, shard uint6 c.checker.see(index, field, view, shard) defer func() { if r := recover(); r != nil { - c.Dump(c.short) + c.Dump(c.short, shard) AlwaysPrintf("see RoaringBitmapReader() panic '%v' at '%v'", r, stack()) panic(r) } @@ -895,7 +896,7 @@ func (c *blueGreenTx) SliceOfShards(index, field, view, optionalViewPath string) //c.checker.see(index, field, view, shard) // don't have shard. defer func() { if r := recover(); r != nil { - c.Dump(c.short) + c.Dump(c.short, math.MaxUint64) AlwaysPrintf("see SliceOfShards() panic '%v' at '%v'", r, stack()) panic(r) } @@ -921,7 +922,7 @@ func (c *blueGreenTx) SliceOfShards(index, field, view, optionalViewPath string) for _, kb := range slcB { if !ma[kb] { //vv("blueGreenTx SliceOfShards diference! B(%v) had shard %v, but A(%v) did not. cpa='%#v'; cpb='%#v'; in the SliceOfShards returned slice.", c.bs, kb, c.as, cpa, cpb) - c.Dump(c.short) + c.Dump(c.short, math.MaxUint64) panic(fmt.Sprintf("blueGreenTx SliceOfShards diference! B(%v) had shard %v, but A(%v) did not. cpa='%#v'; cpb='%#v'; in the SliceOfShards returned slice.", c.bs, kb, c.as, cpa, cpb)) } delete(ma, kb) diff --git a/catcher.go b/catcher.go index e2b4766e5..be1095783 100644 --- a/catcher.go +++ b/catcher.go @@ -58,8 +58,8 @@ func (c *catcherTx) ImportRoaringBits(index, field, view string, shard uint64, r return c.b.ImportRoaringBits(index, field, view, shard, rit, clear, log, rowSize, data) } -func (c *catcherTx) Dump(short bool) { - c.b.Dump(short) +func (c *catcherTx) Dump(short bool, shard uint64) { + c.b.Dump(short, shard) } func (c *catcherTx) Readonly() bool { diff --git a/dbshard.go b/dbshard.go index 90a2ee201..5584ba09e 100644 --- a/dbshard.go +++ b/dbshard.go @@ -195,7 +195,7 @@ func (dbs *DBShard) DeleteDBPath() (err error) { type DBPerShard struct { Mu sync.Mutex - Dir string // holder dir + HolderDir string dbh *DBHolder @@ -260,7 +260,7 @@ func (txf *TxFactory) NewDBPerShard(types []txtype, holderDir string) (d *DBPerS d = &DBPerShard{ types: types, - Dir: holderDir, + HolderDir: holderDir, dbh: NewDBHolder(), Flatmap: make(map[*DBShard]struct{}), txf: txf, @@ -340,7 +340,7 @@ func (dbs *DBShard) DumpAll() { panicOn(err) defer tx.Rollback() fmt.Printf("\n============= dumping dbs.W[%v] %v ========\n", i, ty) - tx.Dump(short) + tx.Dump(short, dbs.Shard) switch ty { case roaringTxn: @@ -367,15 +367,30 @@ func (per *DBPerShard) DumpAll() { } } if !found1 { - AlwaysPrintf("DBPerShard.DumpAll() sees no databases. dir='%v'", per.Dir) + AlwaysPrintf("DBPerShard.DumpAll() sees no databases. dir='%v'", per.HolderDir) } } +// if you know the shard, you can use this +// pathForType and prefixForType must be kept in sync! func (dbs *DBShard) pathForType(ty txtype) string { // top level paths will end in "@@" - return dbs.HolderPath + sep + dbs.Index + ".index.txstores@@@" + sep + "store" + ty.FileSuffix() + "@" + sep + fmt.Sprintf("shard.%04v", dbs.Shard) + + // what here for roaring? well, roaringRegistrar.OpenDBWrapper() + // is a no-op anyhow. so doesn't need to be correct atm. + + return dbs.HolderPath + sep + dbs.Index + ".index.txstores@@@" + sep + "store" + ty.FileSuffix() + "@" + sep + fmt.Sprintf("shard.%04v%v", dbs.Shard, ty.FileSuffix()) } +// if you don't know the shard, you have to use this. +// 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 +} + +var ErrNoData = fmt.Errorf("no data") + func (per *DBPerShard) GetDBShard(index string, shard uint64, idx *Index) (dbs *DBShard, err error) { per.Mu.Lock() @@ -398,12 +413,11 @@ func (per *DBPerShard) GetDBShard(index string, shard uint64, idx *Index) (dbs * ParentDBIndex: dbi, Index: index, Shard: shard, - HolderPath: per.Dir, - //Path: per.Path(index, shard), - idx: idx, - per: per, - useOpenList: per.useOpenList, - hasRoaring: per.hasRoaring, + HolderPath: per.HolderDir, + idx: idx, + per: per, + useOpenList: per.useOpenList, + hasRoaring: per.hasRoaring, } dbi.Shard[shard] = dbs } @@ -422,8 +436,8 @@ func (per *DBPerShard) GetDBShard(index string, shard uint64, idx *Index) (dbs * default: panic(fmt.Sprintf("unknown txtyp: '%v'", ty)) } - - w, err := registry.OpenDBWrapper(dbs.pathForType(ty), DetectMemAccessPastTx) + path := dbs.pathForType(ty) + w, err := registry.OpenDBWrapper(path, DetectMemAccessPastTx) panicOn(err) h := idx.Holder() w.SetHolder(h) @@ -466,20 +480,48 @@ func (per *DBPerShard) Close() (err error) { return } -// DBPerShardGetShardsForIndex returns the indexes from the B (green) database if -// blue-green comparison is in use, rather than from the A (blue) database. -func DBPerShardGetShardsForIndex(idx *Index, roaringViewPath string) (sliceOfShards []uint64, err error) { +// DBPerShardGetShardsForIndex returns the shards for idx. +func (f *TxFactory) GetShardsForIndex(idx *Index, roaringViewPath string) (sliceOfShards []uint64, err error) { - // follow the blueGreen convention of returning the answer for 'B' or - // the last wrapper type. - types := idx.holder.txf.Types() - ty := types[len(types)-1] - return TypedDBPerShardGetLocalShardsForIndex(ty, idx, roaringViewPath) + var shards [][]uint64 + for _, ty := range f.types { + var slc []uint64 + slc, err = f.dbPerShard.TypedDBPerShardGetShardsForIndex(ty, idx, roaringViewPath) + if err != nil { + return + } + shards = append(shards, slc) + } + n := len(f.types) + if n != 1 && n != 2 { + panic(fmt.Sprintf("internal error. only green or blue/green supported. we see types len %v", n)) + } + if !f.blueGreenOff && n == 2 { + // this is a blue green check which cannot live inside Tx because we don't know the shard yet. + // Therefore it has to be above Tx, since we are getting all the shards to choose from here. + // + // But, we still want to check for blue-green consistency. In fact, this was written + // in response to an issue with balancing/re-balancing shards being different + // over the cluster of nodes between blue and green. + + b := sliceToMap(shards[0]) + g := sliceToMap(shards[1]) + blueMinusGreenDiff := mapDiff(b, g) + greenMinusBlueDiff := mapDiff(g, b) + if len(blueMinusGreenDiff) == 0 && len(greenMinusBlueDiff) == 0 { + // ok + } else { + vv("blue[%v] and green[%v] have different shards for index '%v': blueMinusGreenDiff: %v, greenMinusBlueDiff: %v; blueShards='%v', greenShards='%v'; idx.path='%v'", f.types[0].String(), f.types[1].String(), idx.name, blueMinusGreenDiff, greenMinusBlueDiff, asInts(shards[0]), asInts(shards[1]), idx.path) + panic(fmt.Sprintf("blue[%v] and green[%v] have different shards for index '%v': blueMinusGreenDiff: %v, greenMinusBlueDiff: %v; blueShards='%v', greenShards='%v'; idx.path='%v'", f.types[0].String(), f.types[1].String(), idx.name, blueMinusGreenDiff, greenMinusBlueDiff, asInts(shards[0]), asInts(shards[1]), idx.path)) + } + } + // If we are populating blue from green, it does matter that we return green. + return shards[n-1], nil } // if roaringViewPath is "" then for ty == roaringTxn we go to disk to discover // all the view paths under idx for type ty. -func TypedDBPerShardGetLocalShardsForIndex(ty txtype, idx *Index, roaringViewPath string) (sliceOfShards []uint64, err error) { +func (per *DBPerShard) TypedDBPerShardGetShardsForIndex(ty txtype, idx *Index, roaringViewPath string) (sliceOfShards []uint64, err error) { if ty == roaringTxn { rx := &RoaringTx{ @@ -503,7 +545,8 @@ func TypedDBPerShardGetLocalShardsForIndex(ty txtype, idx *Index, roaringViewPat return rx.SliceOfShards("", "", "", roaringViewPath) } requiredSuffix := ty.FileSuffix() - path := idx.Path() + //path := idx.Path() + path := per.prefixForType(idx, ty) ignoreEmpty := false includeRoot := true @@ -517,8 +560,13 @@ func TypedDBPerShardGetLocalShardsForIndex(ty txtype, idx *Index, roaringViewPat panic(fmt.Sprintf("should have 2 parts: nm='%v', base(nm)='%v'; requiredSuffix='%v'", nm, base, requiredSuffix)) } prefix := splt[0] + const shardPrefix = "shard." + const lenOfShardPrefix = len(shardPrefix) + if !strings.HasPrefix(prefix, shardPrefix) { + continue + } // Parse filename into integer. - shard, err := strconv.ParseUint(prefix, 10, 64) + shard, err := strconv.ParseUint(prefix[lenOfShardPrefix:], 10, 64) if err != nil { continue } @@ -528,8 +576,11 @@ func TypedDBPerShardGetLocalShardsForIndex(ty txtype, idx *Index, roaringViewPat } func listDirUnderDir(root string, includeRoot bool, requiredSuffix string, ignoreEmpty bool) (files []string, err error) { + //vv("listDirUnderDir(root ='%v', suffix='%v')", root, requiredSuffix) if !dirExists(root) { - return nil, fmt.Errorf("listFilesUnderDir error: root directory '%v' not found", root) + //vv("warning: listFilesUnderDir error: root directory '%v' not found", root) + //return nil, fmt.Errorf("listFilesUnderDir error: root directory '%v' not found", root) + return } n := len(root) + 1 if includeRoot { diff --git a/dbshard_internal_test.go b/dbshard_internal_test.go index 175c24f65..2405dbdb3 100644 --- a/dbshard_internal_test.go +++ b/dbshard_internal_test.go @@ -64,9 +64,9 @@ func TestShardPerDB_SetBit(t *testing.T) { } } -// test that we find all shards -func Test_DBPerShard_GetShardsForIndex(t *testing.T) { - tmpdir, err := ioutil.TempDir("", "TestDBPerShardGetShardsForIndex") +// test that we find all *local* shards +func Test_DBPerShard_GetShardsForIndex_LocalOnly(t *testing.T) { + tmpdir, err := ioutil.TempDir("", "Test_DBPerShard_GetShardsForIndex_LocalOnly") panicOn(err) orig := os.Getenv("PILOSA_TXSRC") @@ -84,7 +84,7 @@ func Test_DBPerShard_GetShardsForIndex(t *testing.T) { estd := "rick/_exists/views/standard" std := "rick/f/views/standard" - sos, err := DBPerShardGetShardsForIndex(idx, tmpdir+sep+std) + sos, err := holder.txf.GetShardsForIndex(idx, tmpdir+sep+std) panicOn(err) for _, shard := range []uint64{93, 223, 221, 215, 219, 217} { if !inSlice(sos, shard) { @@ -93,7 +93,7 @@ func Test_DBPerShard_GetShardsForIndex(t *testing.T) { } if src == "roaring" { // check estd too - sos, err = DBPerShardGetShardsForIndex(idx, tmpdir+sep+estd) + sos, err = holder.txf.GetShardsForIndex(idx, tmpdir+sep+estd) panicOn(err) for _, shard := range []uint64{93, 223, 221, 215, 219, 217} { if !inSlice(sos, shard) { @@ -137,58 +137,58 @@ rick/_exists/views/standard/fragments/219 rick/_exists/views/standard/fragments/223 `, "lmdb": ` -rick/0219-lmdb@/data.mdb -rick/0219-lmdb@/lock.mdb -rick/0093-lmdb@/data.mdb -rick/0093-lmdb@/lock.mdb -rick/0223-lmdb@/data.mdb -rick/0223-lmdb@/lock.mdb -rick/0215-lmdb@/data.mdb -rick/0215-lmdb@/lock.mdb -rick/0217-lmdb@/data.mdb -rick/0217-lmdb@/lock.mdb -rick/0221-lmdb@/data.mdb -rick/0221-lmdb@/lock.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0219-lmdb@/data.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0219-lmdb@/lock.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0093-lmdb@/data.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0093-lmdb@/lock.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0223-lmdb@/data.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0223-lmdb@/lock.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0215-lmdb@/data.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0215-lmdb@/lock.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0217-lmdb@/data.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0217-lmdb@/lock.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0221-lmdb@/data.mdb +rick.index.txstores@@@/store-lmdb@@/shard.0221-lmdb@/lock.mdb `, "badger": ` -rick/0219-badgerdb@/000000.vlog -rick/0219-badgerdb@/KEYREGISTRY -rick/0219-badgerdb@/MANIFEST -rick/0219-badgerdb@/LOCK -rick/0221-badgerdb@/000000.vlog -rick/0221-badgerdb@/KEYREGISTRY -rick/0221-badgerdb@/MANIFEST -rick/0221-badgerdb@/LOCK -rick/0223-badgerdb@/000000.vlog -rick/0223-badgerdb@/KEYREGISTRY -rick/0223-badgerdb@/MANIFEST -rick/0223-badgerdb@/LOCK -rick/0093-badgerdb@/000000.vlog -rick/0093-badgerdb@/KEYREGISTRY -rick/0093-badgerdb@/MANIFEST -rick/0093-badgerdb@/LOCK -rick/0217-badgerdb@/000000.vlog -rick/0217-badgerdb@/KEYREGISTRY -rick/0217-badgerdb@/MANIFEST -rick/0217-badgerdb@/LOCK -rick/0215-badgerdb@/000000.vlog -rick/0215-badgerdb@/KEYREGISTRY -rick/0215-badgerdb@/MANIFEST -rick/0215-badgerdb@/LOCK +rick.index.txstores@@@/store-badgerdb@@/shard.0219-badgerdb@/000000.vlog +rick.index.txstores@@@/store-badgerdb@@/shard.0219-badgerdb@/KEYREGISTRY +rick.index.txstores@@@/store-badgerdb@@/shard.0219-badgerdb@/MANIFEST +rick.index.txstores@@@/store-badgerdb@@/shard.0219-badgerdb@/LOCK +rick.index.txstores@@@/store-badgerdb@@/shard.0221-badgerdb@/000000.vlog +rick.index.txstores@@@/store-badgerdb@@/shard.0221-badgerdb@/KEYREGISTRY +rick.index.txstores@@@/store-badgerdb@@/shard.0221-badgerdb@/MANIFEST +rick.index.txstores@@@/store-badgerdb@@/shard.0221-badgerdb@/LOCK +rick.index.txstores@@@/store-badgerdb@@/shard.0223-badgerdb@/000000.vlog +rick.index.txstores@@@/store-badgerdb@@/shard.0223-badgerdb@/KEYREGISTRY +rick.index.txstores@@@/store-badgerdb@@/shard.0223-badgerdb@/MANIFEST +rick.index.txstores@@@/store-badgerdb@@/shard.0223-badgerdb@/LOCK +rick.index.txstores@@@/store-badgerdb@@/shard.0093-badgerdb@/000000.vlog +rick.index.txstores@@@/store-badgerdb@@/shard.0093-badgerdb@/KEYREGISTRY +rick.index.txstores@@@/store-badgerdb@@/shard.0093-badgerdb@/MANIFEST +rick.index.txstores@@@/store-badgerdb@@/shard.0093-badgerdb@/LOCK +rick.index.txstores@@@/store-badgerdb@@/shard.0217-badgerdb@/000000.vlog +rick.index.txstores@@@/store-badgerdb@@/shard.0217-badgerdb@/KEYREGISTRY +rick.index.txstores@@@/store-badgerdb@@/shard.0217-badgerdb@/MANIFEST +rick.index.txstores@@@/store-badgerdb@@/shard.0217-badgerdb@/LOCK +rick.index.txstores@@@/store-badgerdb@@/shard.0215-badgerdb@/000000.vlog +rick.index.txstores@@@/store-badgerdb@@/shard.0215-badgerdb@/KEYREGISTRY +rick.index.txstores@@@/store-badgerdb@@/shard.0215-badgerdb@/MANIFEST +rick.index.txstores@@@/store-badgerdb@@/shard.0215-badgerdb@/LOCK `, "rbf": ` -rick/0223-rbfdb@/wal/0000000000000001.wal -rick/0223-rbfdb@/data -rick/0093-rbfdb@/wal/0000000000000001.wal -rick/0093-rbfdb@/data -rick/0217-rbfdb@/wal/0000000000000001.wal -rick/0217-rbfdb@/data -rick/0215-rbfdb@/wal/0000000000000001.wal -rick/0215-rbfdb@/data -rick/0221-rbfdb@/wal/0000000000000001.wal -rick/0221-rbfdb@/data -rick/0219-rbfdb@/wal/0000000000000001.wal -rick/0219-rbfdb@/data +rick.index.txstores@@@/store-rbfdb@@/shard.0223-rbfdb@/wal/0000000000000001.wal +rick.index.txstores@@@/store-rbfdb@@/shard.0223-rbfdb@/data +rick.index.txstores@@@/store-rbfdb@@/shard.0093-rbfdb@/wal/0000000000000001.wal +rick.index.txstores@@@/store-rbfdb@@/shard.0093-rbfdb@/data +rick.index.txstores@@@/store-rbfdb@@/shard.0217-rbfdb@/wal/0000000000000001.wal +rick.index.txstores@@@/store-rbfdb@@/shard.0217-rbfdb@/data +rick.index.txstores@@@/store-rbfdb@@/shard.0215-rbfdb@/wal/0000000000000001.wal +rick.index.txstores@@@/store-rbfdb@@/shard.0215-rbfdb@/data +rick.index.txstores@@@/store-rbfdb@@/shard.0221-rbfdb@/wal/0000000000000001.wal +rick.index.txstores@@@/store-rbfdb@@/shard.0221-rbfdb@/data +rick.index.txstores@@@/store-rbfdb@@/shard.0219-rbfdb@/wal/0000000000000001.wal +rick.index.txstores@@@/store-rbfdb@@/shard.0219-rbfdb@/data `, } diff --git a/dbshard_test.go b/dbshard_test.go index d69aa7943..fb6d7e00f 100644 --- a/dbshard_test.go +++ b/dbshard_test.go @@ -28,8 +28,6 @@ import ( "github.com/pilosa/pilosa/v2/test" ) -var sep = string(os.PathSeparator) - func skipForNonLMDB(t *testing.T) { src := os.Getenv("PILOSA_TXSRC") if src != "lmdb" { @@ -62,13 +60,15 @@ func Test_DBPerShard_multiple_shards_used(t *testing.T) { hldr.SetBit(index, "general", 11, ShardWidth+2) types := pilosa.MustTxsrcToTxtype("lmdb") - tx_suffix := types[0].FileSuffix() - root := hldr.Path() + sep + index - shards := []string{"0000", "0001", "0002"} + idx := hldr.Index(index) + shardsU := []uint64{0, 1, 2} pathShard := []string{} + // check that 3 different shard databases/files were made for i := 0; i < 2; i++ { - path := root + sep + shards[i] + tx_suffix + + path, err := hldr.Txf().GetDBShardPath(index, shardsU[i], idx, types[0], !writable) + panicOn(err) pathShard = append(pathShard, path) if !DirExists(pathShard[i]) { diff --git a/executor.go b/executor.go index 4ef7e90a4..6268e0fd1 100644 --- a/executor.go +++ b/executor.go @@ -142,6 +142,7 @@ func (e *executor) Close() error { // Execute executes a PQL query. func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shards []uint64, opt *execOptions) (QueryResponse, error) { + span, ctx := tracing.StartSpanFromContext(ctx, "Executor.Execute") span.LogKV("pql", q.String()) defer span.Finish() @@ -3834,15 +3835,15 @@ func (e *executor) executeClearBitField(ctx context.Context, qcx *Qcx, index str shard := colID / ShardWidth - idx := e.Holder.Index(index) - - tx, finisher := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) - defer finisher(&err) - ret := false for _, node := range e.Cluster.shardNodes(index, shard) { // Update locally if host matches. if node.ID == e.Node.ID { + + idx := e.Holder.Index(index) + tx, finisher := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) + defer finisher(&err) + val, err := f.ClearBit(tx, rowID, colID) if err != nil { return false, err @@ -4109,7 +4110,7 @@ func (e *executor) executeSetRowShard(ctx context.Context, qcx *Qcx, index strin } // executeSet executes a Set() call. -func (e *executor) executeSet(ctx context.Context, qcx *Qcx, index string, c *pql.Call, opt *execOptions) (_ bool, err error) { +func (e *executor) executeSet(ctx context.Context, qcx *Qcx, index string, c *pql.Call, opt *execOptions) (_ bool, err0 error) { span, ctx := tracing.StartSpanFromContext(ctx, "Executor.executeSet") defer span.Finish() @@ -4128,9 +4129,6 @@ func (e *executor) executeSet(ctx context.Context, qcx *Qcx, index string, c *pq shard := colID / ShardWidth - tx, finisher := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) - defer finisher(&err) - // Read field name. fieldName, err := c.FieldArg() if err != nil { @@ -4145,11 +4143,15 @@ func (e *executor) executeSet(ctx context.Context, qcx *Qcx, index string, c *pq // Set column on existence field. if ef := idx.existenceField(); ef != nil { + // we create tx here, rather than just above, to avoid creating an extra empty shard. + tx, finisher := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) + defer finisher(&err0) + if _, err := ef.SetBit(tx, 0, colID, nil); err != nil { return false, errors.Wrap(err, "setting existence column") } + finisher(nil) // commit to free of the write lock needed inside executeSetBitField } - finisher(nil) // commit to free of the write lock needed inside executeSetBitField switch f.Type() { case FieldTypeInt, FieldTypeDecimal: @@ -4206,13 +4208,14 @@ func (e *executor) executeSetBitField(ctx context.Context, qcx *Qcx, index strin shard := colID / ShardWidth ret := false - idx := e.Holder.Index(index) - tx, finisher := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) - defer finisher(&err0) - for _, node := range e.Cluster.shardNodes(index, shard) { // Update locally if host matches. if node.ID == e.Node.ID { + + idx := e.Holder.Index(index) + tx, finisher := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) + defer finisher(&err0) + val, err := f.SetBit(tx, rowID, colID, timestamp) if err != nil { return false, err @@ -4245,13 +4248,14 @@ func (e *executor) executeSetValueField(ctx context.Context, qcx *Qcx, index str shard := colID / ShardWidth ret := false - idx := e.Holder.Index(index) - tx, finisher := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) - defer finisher(&err) - for _, node := range e.Cluster.shardNodes(index, shard) { // Update locally if host matches. if node.ID == e.Node.ID { + + idx := e.Holder.Index(index) + tx, finisher := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) + defer finisher(&err) + val, err := f.SetValue(tx, colID, value) if err != nil { return false, err @@ -4284,13 +4288,14 @@ func (e *executor) executeClearValueField(ctx context.Context, qcx *Qcx, index s shard := colID / ShardWidth ret := false - idx := e.Holder.Index(index) - tx, finisher := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) - defer finisher(&err) - for _, node := range e.Cluster.shardNodes(index, shard) { // Update locally if host matches. if node.ID == e.Node.ID { + + idx := e.Holder.Index(index) + tx, finisher := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) + defer finisher(&err) + val, err := f.ClearValue(tx, colID) if err != nil { return false, err diff --git a/holder.go b/holder.go index 74b648e45..41697d03d 100644 --- a/holder.go +++ b/holder.go @@ -675,7 +675,7 @@ func (h *Holder) Open() error { // under blue_green, we must sync blue from green before we turn on checking. if err := h.txf.green2blue(h); err != nil { - return errors.Wrap(err, "Holder.Open h.txf.UpdateBlueFromGreen(h)") + return errors.Wrap(err, "Holder.Open h.txf.green2blue(h)") } h.txf.blueGreenOnIfRunningBlueGreen() diff --git a/index.go b/index.go index eade152d5..4e7d3459a 100644 --- a/index.go +++ b/index.go @@ -691,8 +691,8 @@ func FormatQualifiedIndexName(index string) string { // Dump prints to stdout the contents of the roaring Containers // stored in idx. Mostly for debugging. func (idx *Index) Dump(label string) { - //fileline := FileLine(2) - fmt.Printf("\nDump: %v\n\n", label) + fileline := FileLine(2) + fmt.Printf("\n%v Dump: %v\n\n", fileline, label) idx.holder.txf.dbPerShard.DumpAll() } diff --git a/lmdb.go b/lmdb.go index e18907d55..4758c4b26 100644 --- a/lmdb.go +++ b/lmdb.go @@ -1575,7 +1575,7 @@ func (tx *LMDBTx) countBitsSet(bkey []byte) (n int) { return } -func (tx *LMDBTx) Dump(short bool) { +func (tx *LMDBTx) Dump(short bool, shard uint64) { fmt.Printf("%v\n", stringifiedLMDBKeysTx(tx, short)) } @@ -1650,11 +1650,11 @@ func (w *LMDBWrapper) DeleteField(index, field, fieldPath string) (err error) { // under blue-green roaring_lmdb, the directory will not be found, b/c roaring will have // already done the os.RemoveAll(). BUT, RemoveAll returns nil error in this case. Docs: // "If the path does not exist, RemoveAll returns nil (no error)" - //w.DeleteDBPath(&DBShard{Path: fieldPath}) - //if err != nil { - //return errors.Wrap(err, "removing directory") - //} + err = os.RemoveAll(fieldPath) + if err != nil { + return errors.Wrap(err, "removing directory") + } prefix := txkey.FieldPrefix(index, field) return w.DeletePrefix(prefix) } diff --git a/lmdb_other.go b/lmdb_other.go index 3a5c04995..955ea66d6 100644 --- a/lmdb_other.go +++ b/lmdb_other.go @@ -404,7 +404,7 @@ func (tx *LMDBTx) IsDone() (done bool) { panic("lmdb only available on 64-bit arch") } -func (tx *LMDBTx) Dump(short bool) { +func (tx *LMDBTx) Dump(short bool, shard uint64) { panic("lmdb only available on 64-bit arch") } diff --git a/rbf.go b/rbf.go index d88f13fd7..79284de3c 100644 --- a/rbf.go +++ b/rbf.go @@ -426,8 +426,8 @@ func (tx *RBFTx) Pointer() string { return fmt.Sprintf("%p", tx) } -func (tx *RBFTx) Dump(short bool) { - tx.tx.Dump(short) +func (tx *RBFTx) Dump(short bool, shard uint64) { + tx.tx.Dump(short, shard) } // Readonly is true if the transaction is not read-and-write, but only doing reads. diff --git a/rbf/tx.go b/rbf/tx.go index bf1093c7e..f559fb0b0 100644 --- a/rbf/tx.go +++ b/rbf/tx.go @@ -1334,10 +1334,10 @@ func (si *emptyContainerIterator) Value() (uint64, *roaring.Container) { panic("emptyContainerIterator never has any Values") } -func (tx *Tx) Dump(short bool) { - fmt.Println(tx.DumpString(short)) +func (tx *Tx) Dump(short bool, shard uint64) { + fmt.Println(tx.DumpString(short, shard)) } -func (tx *Tx) DumpString(short bool) (r string) { +func (tx *Tx) DumpString(short bool, shard uint64) (r string) { r = "allkeys:[\n" diff --git a/rbf/tx_test.go b/rbf/tx_test.go index fedd5f47d..5a86068f8 100644 --- a/rbf/tx_test.go +++ b/rbf/tx_test.go @@ -16,6 +16,7 @@ package rbf_test import ( "fmt" + "math" "math/rand" "testing" "time" @@ -483,7 +484,7 @@ func TestTx_Dump(t *testing.T) { } // test that we don't crash, and get *something* back - s := tx.DumpString(true) + s := tx.DumpString(true, math.MaxUint64) if s == "" { panic("should have had 3 containers!") } diff --git a/rrtx.go b/rrtx.go index 131afb26c..7ea06f659 100644 --- a/rrtx.go +++ b/rrtx.go @@ -55,8 +55,10 @@ func (tx *RoaringTx) Type() string { return RoaringTxn } -func (tx *RoaringTx) Dump(short bool) { - fmt.Printf("%v\n", tx.Index.StringifiedRoaringKeys(short, false, tx.o)) +func (tx *RoaringTx) Dump(short bool, shard uint64) { + o := tx.o + o.Shard = shard + fmt.Printf("%v\n", tx.Index.StringifiedRoaringKeys(short, false, o)) } func (tx *RoaringTx) UseRowCache() bool { @@ -188,6 +190,7 @@ func (tx *RoaringTx) RemoveContainer(index, field, view string, shard uint64, ke } func (tx *RoaringTx) Add(index, field, view string, shard uint64, batched bool, a ...uint64) (changeCount int, err error) { + //vv("RoaringTx.Add(index='%v', shard='%v') stack=\n%v", index, shard, stack()) b, err := tx.bitmap(index, field, view, shard) if err != nil { return 0, err diff --git a/stattx.go b/stattx.go index 20b30a57c..383ad0d78 100644 --- a/stattx.go +++ b/stattx.go @@ -310,8 +310,8 @@ func (c *statTx) ImportRoaringBits(index, field, view string, shard uint64, rit return c.b.ImportRoaringBits(index, field, view, shard, rit, clear, log, rowSize, data) } -func (c *statTx) Dump(short bool) { - c.b.Dump(short) +func (c *statTx) Dump(short bool, shard uint64) { + c.b.Dump(short, shard) } func (c *statTx) Readonly() bool { diff --git a/tx.go b/tx.go index 8546ac8d3..c6bc1066b 100644 --- a/tx.go +++ b/tx.go @@ -202,7 +202,7 @@ type Tx interface { Group() *TxGroup // Dump is for debugging, what does this Tx see as its database? - Dump(short bool) + Dump(short bool, shard uint64) // Options returns the options used to create this Tx. This // can be implementd by embedding Txo, and Txo provides the diff --git a/txfactory.go b/txfactory.go index ee625ca9f..033295935 100644 --- a/txfactory.go +++ b/txfactory.go @@ -827,6 +827,9 @@ func (idx *Index) StringifiedRoaringKeys(hashOnly, showOps bool, o Txo) (r strin if err != nil { continue // ignore .meta paths } + if shard != o.Shard { + continue // only print the shard the Txo is on. + } abspath := idx.path + sep + relpath s, _, err := stringifiedRawRoaringFragment(abspath, index, field, view, shard, showOps, hashOnly, os.Stdout) @@ -839,7 +842,7 @@ func (idx *Index) StringifiedRoaringKeys(hashOnly, showOps bool, o Txo) (r strin n++ } if n == 0 { - return "" // new convention that empty database => empty string returned. + return "" } // note that we can have a bitmap present, but it can be empty r += "]\n all-in-blake3:" + hash.Blake3sum16([]byte(r)) + "\n" @@ -1197,14 +1200,19 @@ func (f *TxFactory) green2blue(holder *Holder) (err error) { for _, idx := range idxs { - blueShards, err := TypedDBPerShardGetLocalShardsForIndex(blueDest, idx, "") + // scan directories + blueShards, err := f.dbPerShard.TypedDBPerShardGetShardsForIndex(blueDest, idx, "") if err != nil { return errors.Wrap(err, fmt.Sprintf("GetDBShard(index='%v') error fetching blueShards", idx.name)) } - greenShards, err := TypedDBPerShardGetLocalShardsForIndex(greenSrc, idx, "") + //vv("from blueDest='%v', blueShards = '%#v'", blueDest, blueShards) + + // scan directories + greenShards, err := f.dbPerShard.TypedDBPerShardGetShardsForIndex(greenSrc, idx, "") if err != nil { return errors.Wrap(err, fmt.Sprintf("GetDBShard(index='%v') error fetching greenShards", idx.name)) } + //vv("from greenSrc='%v', greenShards = '%#v'", greenSrc, greenShards) diff := f.shardSliceDiff(blueShards, greenShards) if diff != "" { @@ -1272,3 +1280,12 @@ func (f *TxFactory) shardSliceDiff(blueShards, greenShards []uint64) (diff strin diff += fmt.Sprintf("shard diff: blueMinusGreen shards: '%#v'; greenMinusBlue shards: '%#v'", bmg, gmb) return } + +func (f *TxFactory) GetDBShardPath(index string, shard uint64, idx *Index, ty txtype, write bool) (shardPath string, err error) { + dbs, err := f.dbPerShard.GetDBShard(index, shard, idx) + if err != nil { + return "", errors.Wrap(err, fmt.Sprintf("GetDBShardPath(index='%v', shard='%v', ty='%v')", index, shard, ty.String())) + } + shardPath = dbs.pathForType(ty) + return +} diff --git a/txfactory_internal_test.go b/txfactory_internal_test.go index f21414ee0..57b2fd133 100644 --- a/txfactory_internal_test.go +++ b/txfactory_internal_test.go @@ -112,7 +112,7 @@ func Test_TxFactory_Qcx_query_context(t *testing.T) { // and b) we have an easy migration mechanism, to go from one storage format to another. // func Test_TxFactory_UpdateBlueFromGreen_OnStartup(t *testing.T) { - t.Skip("TODO(jea) bring this back in. broken by the local vs remote shard determination for a cluster") + //t.Skip("TODO(jea) bring this back in. broken by the local vs remote shard determination for a cluster") orig := os.Getenv("PILOSA_TXSRC") defer os.Setenv("PILOSA_TXSRC", orig) // must restore or will mess up other tests! @@ -218,6 +218,8 @@ func Test_TxFactory_UpdateBlueFromGreen_OnStartup(t *testing.T) { // open a holder with path again, now looking at both blue and green. // The Holder.Open should do the migration from green, populating blue. h4 := NewHolder(path, nil) + + //vv("about to h4.Open we should populate blue from green") panicOn(h4.Open()) defer h4.Close() @@ -232,7 +234,7 @@ func Test_TxFactory_UpdateBlueFromGreen_OnStartup(t *testing.T) { // go to verify it but blue has more data than green. // That will also cause query divergence. func Test_TxFactory_verifyBlueEqualsGreen(t *testing.T) { - t.Skip("TODO(jea) bring this back in. broken by the local vs remote shard determination for a cluster") + //t.Skip("TODO(jea) bring this back in. broken by the local vs remote shard determination for a cluster") orig := os.Getenv("PILOSA_TXSRC") defer os.Setenv("PILOSA_TXSRC", orig) // must restore or will mess up other tests! diff --git a/view.go b/view.go index 52c940986..40ff7fb5a 100644 --- a/view.go +++ b/view.go @@ -173,7 +173,7 @@ var workQueue = make(chan struct{}, runtime.NumCPU()*2) // replaces v.openFragments() with Tx generic code. func (v *view) openFragmentsInTx() error { - shards, err := DBPerShardGetShardsForIndex(v.idx, v.path) + shards, err := v.holder.txf.GetShardsForIndex(v.idx, v.path) if err != nil { return errors.Wrap(err, "DBPerShardGetShardsForIndex()") }