mirror of
https://github.com/luxfi/zapdb.git
synced 2026-07-27 06:54:45 +00:00
Introduce SSTable sha256 checksums (#689)
Add SHA256 checksums for SSTables in MANIFEST. If a table no longer matches this checksum, that table would be skipped over with an error. Tested that it works with previous badger directories. As new tables get created, Badger would store their checksums in MANIFEST. Modified `badger info` to show the checksums stored in MANIFEST, so user can manually compare the output from `sha256sum <filename>` if needed. Fixes #680 .
This commit is contained in:
+5
-4
@@ -169,8 +169,9 @@ func printInfo(dir, valueDir string) error {
|
||||
})
|
||||
for _, tableID := range tableIDs {
|
||||
tableFile := table.IDToFilename(tableID)
|
||||
file, ok := fileinfoByName[tableFile]
|
||||
if ok {
|
||||
tm, ok1 := manifest.Tables[tableID]
|
||||
file, ok2 := fileinfoByName[tableFile]
|
||||
if ok1 && ok2 {
|
||||
fileinfoMarked[tableFile] = true
|
||||
emptyString := ""
|
||||
fileSize := file.Size()
|
||||
@@ -180,8 +181,8 @@ func printInfo(dir, valueDir string) error {
|
||||
}
|
||||
levelSizes[level] += fileSize
|
||||
// (Put level on every line to make easier to process with sed/perl.)
|
||||
fmt.Printf("[%25s] %-12s %6s L%d%s\n", dur(baseTime, file.ModTime()),
|
||||
tableFile, hbytes(fileSize), level, emptyString)
|
||||
fmt.Printf("[%25s] %-12s %6s L%d %x%s\n", dur(baseTime, file.ModTime()),
|
||||
tableFile, hbytes(fileSize), level, tm.Checksum, emptyString)
|
||||
} else {
|
||||
fmt.Printf("%s [MISSING]\n", tableFile)
|
||||
numMissing++
|
||||
|
||||
@@ -829,7 +829,7 @@ func (db *DB) handleFlushTask(ft flushTask) error {
|
||||
db.elog.Errorf("ERROR while syncing level directory: %v", dirSyncErr)
|
||||
}
|
||||
|
||||
tbl, err := table.OpenTable(fd, db.opt.TableLoadingMode)
|
||||
tbl, err := table.OpenTable(fd, db.opt.TableLoadingMode, nil)
|
||||
if err != nil {
|
||||
db.elog.Printf("ERROR while opening table: %v", err)
|
||||
return err
|
||||
|
||||
@@ -22,6 +22,7 @@ import (
|
||||
"math/rand"
|
||||
"os"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
@@ -122,7 +123,7 @@ func newLevelsController(db *DB, mf *Manifest) (*levelsController, error) {
|
||||
tick := time.NewTicker(3 * time.Second)
|
||||
defer tick.Stop()
|
||||
|
||||
for fileID, tableManifest := range mf.Tables {
|
||||
for fileID, tf := range mf.Tables {
|
||||
fname := table.NewFilename(fileID, db.opt.Dir)
|
||||
select {
|
||||
case <-tick.C:
|
||||
@@ -137,7 +138,7 @@ func newLevelsController(db *DB, mf *Manifest) (*levelsController, error) {
|
||||
if fileID > maxFileID {
|
||||
maxFileID = fileID
|
||||
}
|
||||
go func(fname string, level int) {
|
||||
go func(fname string, tf TableManifest) {
|
||||
var rerr error
|
||||
defer func() {
|
||||
throttle.Done(rerr)
|
||||
@@ -149,16 +150,22 @@ func newLevelsController(db *DB, mf *Manifest) (*levelsController, error) {
|
||||
return
|
||||
}
|
||||
|
||||
t, err := table.OpenTable(fd, db.opt.TableLoadingMode)
|
||||
t, err := table.OpenTable(fd, db.opt.TableLoadingMode, tf.Checksum)
|
||||
if err != nil {
|
||||
rerr = errors.Wrapf(err, "Opening table: %q", fname)
|
||||
if strings.HasPrefix(err.Error(), "CHECKSUM_MISMATCH:") {
|
||||
db.opt.Errorf(err.Error())
|
||||
db.opt.Errorf("Ignoring table %s", fd.Name())
|
||||
// Do not set rerr. We will continue without this table.
|
||||
} else {
|
||||
rerr = errors.Wrapf(err, "Opening table: %q", fname)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
mu.Lock()
|
||||
tables[level] = append(tables[level], t)
|
||||
tables[tf.Level] = append(tables[tf.Level], t)
|
||||
mu.Unlock()
|
||||
}(fname, int(tableManifest.Level))
|
||||
}(fname, tf)
|
||||
}
|
||||
if err := throttle.Finish(); err != nil {
|
||||
closeAllTables(tables)
|
||||
@@ -226,7 +233,7 @@ func (s *levelsController) deleteLSMTree() (int, error) {
|
||||
// Generate the manifest changes.
|
||||
changes := []*pb.ManifestChange{}
|
||||
for _, table := range all {
|
||||
changes = append(changes, makeTableDeleteChange(table.ID()))
|
||||
changes = append(changes, newDeleteChange(table.ID()))
|
||||
}
|
||||
changeSet := pb.ManifestChangeSet{Changes: changes}
|
||||
if err := s.kv.manifest.addChanges(changeSet.Changes); err != nil {
|
||||
@@ -486,7 +493,7 @@ func (s *levelsController) compactBuildTables(
|
||||
return
|
||||
}
|
||||
|
||||
tbl, err := table.OpenTable(fd, s.kv.opt.TableLoadingMode)
|
||||
tbl, err := table.OpenTable(fd, s.kv.opt.TableLoadingMode, nil)
|
||||
// decrRef is added below.
|
||||
resultCh <- newTableResult{tbl, errors.Wrapf(err, "Unable to open table: %q", fd.Name())}
|
||||
}(builder)
|
||||
@@ -534,13 +541,14 @@ func (s *levelsController) compactBuildTables(
|
||||
func buildChangeSet(cd *compactDef, newTables []*table.Table) pb.ManifestChangeSet {
|
||||
changes := []*pb.ManifestChange{}
|
||||
for _, table := range newTables {
|
||||
changes = append(changes, makeTableCreateChange(table.ID(), cd.nextLevel.level))
|
||||
changes = append(changes,
|
||||
newCreateChange(table.ID(), cd.nextLevel.level, table.Checksum))
|
||||
}
|
||||
for _, table := range cd.top {
|
||||
changes = append(changes, makeTableDeleteChange(table.ID()))
|
||||
changes = append(changes, newDeleteChange(table.ID()))
|
||||
}
|
||||
for _, table := range cd.bot {
|
||||
changes = append(changes, makeTableDeleteChange(table.ID()))
|
||||
changes = append(changes, newDeleteChange(table.ID()))
|
||||
}
|
||||
return pb.ManifestChangeSet{Changes: changes}
|
||||
}
|
||||
@@ -748,7 +756,7 @@ func (s *levelsController) addLevel0Table(t *table.Table) error {
|
||||
// the proper order. (That means this update happens before that of some compaction which
|
||||
// deletes the table.)
|
||||
err := s.kv.manifest.addChanges([]*pb.ManifestChange{
|
||||
makeTableCreateChange(t.ID(), 0),
|
||||
newCreateChange(t.ID(), 0, t.Checksum),
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
+16
-13
@@ -42,7 +42,7 @@ import (
|
||||
// reconstruct the manifest at startup.
|
||||
type Manifest struct {
|
||||
Levels []levelManifest
|
||||
Tables map[uint64]tableManifest
|
||||
Tables map[uint64]TableManifest
|
||||
|
||||
// Contains total number of creation and deletion changes in the manifest -- used to compute
|
||||
// whether it'd be useful to rewrite the manifest.
|
||||
@@ -54,7 +54,7 @@ func createManifest() Manifest {
|
||||
levels := make([]levelManifest, 0)
|
||||
return Manifest{
|
||||
Levels: levels,
|
||||
Tables: make(map[uint64]tableManifest),
|
||||
Tables: make(map[uint64]TableManifest),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -64,10 +64,11 @@ type levelManifest struct {
|
||||
Tables map[uint64]struct{} // Set of table id's
|
||||
}
|
||||
|
||||
// tableManifest contains information about a specific level
|
||||
// TableManifest contains information about a specific level
|
||||
// in the LSM tree.
|
||||
type tableManifest struct {
|
||||
Level uint8
|
||||
type TableManifest struct {
|
||||
Level uint8
|
||||
Checksum []byte
|
||||
}
|
||||
|
||||
// manifestFile holds the file pointer (and other info) about the manifest file, which is a log
|
||||
@@ -98,7 +99,7 @@ const (
|
||||
func (m *Manifest) asChanges() []*pb.ManifestChange {
|
||||
changes := make([]*pb.ManifestChange, 0, len(m.Tables))
|
||||
for id, tm := range m.Tables {
|
||||
changes = append(changes, makeTableCreateChange(id, int(tm.Level)))
|
||||
changes = append(changes, newCreateChange(id, int(tm.Level), tm.Checksum))
|
||||
}
|
||||
return changes
|
||||
}
|
||||
@@ -384,8 +385,9 @@ func applyManifestChange(build *Manifest, tc *pb.ManifestChange) error {
|
||||
if _, ok := build.Tables[tc.Id]; ok {
|
||||
return fmt.Errorf("MANIFEST invalid, table %d exists", tc.Id)
|
||||
}
|
||||
build.Tables[tc.Id] = tableManifest{
|
||||
Level: uint8(tc.Level),
|
||||
build.Tables[tc.Id] = TableManifest{
|
||||
Level: uint8(tc.Level),
|
||||
Checksum: append([]byte{}, tc.Checksum...),
|
||||
}
|
||||
for len(build.Levels) <= int(tc.Level) {
|
||||
build.Levels = append(build.Levels, levelManifest{make(map[uint64]struct{})})
|
||||
@@ -417,15 +419,16 @@ func applyChangeSet(build *Manifest, changeSet *pb.ManifestChangeSet) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func makeTableCreateChange(id uint64, level int) *pb.ManifestChange {
|
||||
func newCreateChange(id uint64, level int, checksum []byte) *pb.ManifestChange {
|
||||
return &pb.ManifestChange{
|
||||
Id: id,
|
||||
Op: pb.ManifestChange_CREATE,
|
||||
Level: uint32(level),
|
||||
Id: id,
|
||||
Op: pb.ManifestChange_CREATE,
|
||||
Level: uint32(level),
|
||||
Checksum: checksum,
|
||||
}
|
||||
}
|
||||
|
||||
func makeTableDeleteChange(id uint64) *pb.ManifestChange {
|
||||
func newDeleteChange(id uint64) *pb.ManifestChange {
|
||||
return &pb.ManifestChange{
|
||||
Id: id,
|
||||
Op: pb.ManifestChange_DELETE,
|
||||
|
||||
+7
-7
@@ -169,7 +169,7 @@ func TestOverlappingKeyRangeError(t *testing.T) {
|
||||
lh0 := newLevelHandler(kv, 0)
|
||||
lh1 := newLevelHandler(kv, 1)
|
||||
f := buildTestTable(t, "k", 2)
|
||||
t1, err := table.OpenTable(f, options.MemoryMap)
|
||||
t1, err := table.OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer t1.DecrRef()
|
||||
|
||||
@@ -190,7 +190,7 @@ func TestOverlappingKeyRangeError(t *testing.T) {
|
||||
lc.runCompactDef(0, cd)
|
||||
|
||||
f = buildTestTable(t, "l", 2)
|
||||
t2, err := table.OpenTable(f, options.MemoryMap)
|
||||
t2, err := table.OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer t2.DecrRef()
|
||||
done = lh0.tryAddLevel0Table(t2)
|
||||
@@ -221,14 +221,14 @@ func TestManifestRewrite(t *testing.T) {
|
||||
require.Equal(t, 0, m.Deletions)
|
||||
|
||||
err = mf.addChanges([]*pb.ManifestChange{
|
||||
makeTableCreateChange(0, 0),
|
||||
newCreateChange(0, 0, nil),
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
for i := uint64(0); i < uint64(deletionsThreshold*3); i++ {
|
||||
ch := []*pb.ManifestChange{
|
||||
makeTableCreateChange(i+1, 0),
|
||||
makeTableDeleteChange(i),
|
||||
newCreateChange(i+1, 0, nil),
|
||||
newDeleteChange(i),
|
||||
}
|
||||
err := mf.addChanges(ch)
|
||||
require.NoError(t, err)
|
||||
@@ -238,7 +238,7 @@ func TestManifestRewrite(t *testing.T) {
|
||||
mf = nil
|
||||
mf, m, err = helpOpenOrCreateManifestFile(dir, false, deletionsThreshold)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, map[uint64]tableManifest{
|
||||
uint64(deletionsThreshold * 3): {Level: 0},
|
||||
require.Equal(t, map[uint64]TableManifest{
|
||||
uint64(deletionsThreshold * 3): {Level: 0, Checksum: []byte{}},
|
||||
}, m.Tables)
|
||||
}
|
||||
|
||||
+102
-48
@@ -3,11 +3,12 @@
|
||||
|
||||
package pb
|
||||
|
||||
import proto "github.com/golang/protobuf/proto"
|
||||
import fmt "fmt"
|
||||
import math "math"
|
||||
|
||||
import io "io"
|
||||
import (
|
||||
fmt "fmt"
|
||||
proto "github.com/golang/protobuf/proto"
|
||||
io "io"
|
||||
math "math"
|
||||
)
|
||||
|
||||
// Reference imports to suppress errors if they are not otherwise used.
|
||||
var _ = proto.Marshal
|
||||
@@ -31,6 +32,7 @@ var ManifestChange_Operation_name = map[int32]string{
|
||||
0: "CREATE",
|
||||
1: "DELETE",
|
||||
}
|
||||
|
||||
var ManifestChange_Operation_value = map[string]int32{
|
||||
"CREATE": 0,
|
||||
"DELETE": 1,
|
||||
@@ -39,8 +41,9 @@ var ManifestChange_Operation_value = map[string]int32{
|
||||
func (x ManifestChange_Operation) String() string {
|
||||
return proto.EnumName(ManifestChange_Operation_name, int32(x))
|
||||
}
|
||||
|
||||
func (ManifestChange_Operation) EnumDescriptor() ([]byte, []int) {
|
||||
return fileDescriptor_pb_7d6edc481ce0a5cd, []int{3, 0}
|
||||
return fileDescriptor_f80abaa17e25ccc8, []int{3, 0}
|
||||
}
|
||||
|
||||
type KV struct {
|
||||
@@ -59,7 +62,7 @@ func (m *KV) Reset() { *m = KV{} }
|
||||
func (m *KV) String() string { return proto.CompactTextString(m) }
|
||||
func (*KV) ProtoMessage() {}
|
||||
func (*KV) Descriptor() ([]byte, []int) {
|
||||
return fileDescriptor_pb_7d6edc481ce0a5cd, []int{0}
|
||||
return fileDescriptor_f80abaa17e25ccc8, []int{0}
|
||||
}
|
||||
func (m *KV) XXX_Unmarshal(b []byte) error {
|
||||
return m.Unmarshal(b)
|
||||
@@ -76,8 +79,8 @@ func (m *KV) XXX_Marshal(b []byte, deterministic bool) ([]byte, error) {
|
||||
return b[:n], nil
|
||||
}
|
||||
}
|
||||
func (dst *KV) XXX_Merge(src proto.Message) {
|
||||
xxx_messageInfo_KV.Merge(dst, src)
|
||||
func (m *KV) XXX_Merge(src proto.Message) {
|
||||
xxx_messageInfo_KV.Merge(m, src)
|
||||
}
|
||||
func (m *KV) XXX_Size() int {
|
||||
return m.Size()
|
||||
@@ -131,7 +134,7 @@ func (m *KV) GetMeta() []byte {
|
||||
}
|
||||
|
||||
type KVList struct {
|
||||
Kv []*KV `protobuf:"bytes,1,rep,name=kv" json:"kv,omitempty"`
|
||||
Kv []*KV `protobuf:"bytes,1,rep,name=kv,proto3" json:"kv,omitempty"`
|
||||
XXX_NoUnkeyedLiteral struct{} `json:"-"`
|
||||
XXX_unrecognized []byte `json:"-"`
|
||||
XXX_sizecache int32 `json:"-"`
|
||||
@@ -141,7 +144,7 @@ func (m *KVList) Reset() { *m = KVList{} }
|
||||
func (m *KVList) String() string { return proto.CompactTextString(m) }
|
||||
func (*KVList) ProtoMessage() {}
|
||||
func (*KVList) Descriptor() ([]byte, []int) {
|
||||
return fileDescriptor_pb_7d6edc481ce0a5cd, []int{1}
|
||||
return fileDescriptor_f80abaa17e25ccc8, []int{1}
|
||||
}
|
||||
func (m *KVList) XXX_Unmarshal(b []byte) error {
|
||||
return m.Unmarshal(b)
|
||||
@@ -158,8 +161,8 @@ func (m *KVList) XXX_Marshal(b []byte, deterministic bool) ([]byte, error) {
|
||||
return b[:n], nil
|
||||
}
|
||||
}
|
||||
func (dst *KVList) XXX_Merge(src proto.Message) {
|
||||
xxx_messageInfo_KVList.Merge(dst, src)
|
||||
func (m *KVList) XXX_Merge(src proto.Message) {
|
||||
xxx_messageInfo_KVList.Merge(m, src)
|
||||
}
|
||||
func (m *KVList) XXX_Size() int {
|
||||
return m.Size()
|
||||
@@ -179,7 +182,7 @@ func (m *KVList) GetKv() []*KV {
|
||||
|
||||
type ManifestChangeSet struct {
|
||||
// A set of changes that are applied atomically.
|
||||
Changes []*ManifestChange `protobuf:"bytes,1,rep,name=changes" json:"changes,omitempty"`
|
||||
Changes []*ManifestChange `protobuf:"bytes,1,rep,name=changes,proto3" json:"changes,omitempty"`
|
||||
XXX_NoUnkeyedLiteral struct{} `json:"-"`
|
||||
XXX_unrecognized []byte `json:"-"`
|
||||
XXX_sizecache int32 `json:"-"`
|
||||
@@ -189,7 +192,7 @@ func (m *ManifestChangeSet) Reset() { *m = ManifestChangeSet{} }
|
||||
func (m *ManifestChangeSet) String() string { return proto.CompactTextString(m) }
|
||||
func (*ManifestChangeSet) ProtoMessage() {}
|
||||
func (*ManifestChangeSet) Descriptor() ([]byte, []int) {
|
||||
return fileDescriptor_pb_7d6edc481ce0a5cd, []int{2}
|
||||
return fileDescriptor_f80abaa17e25ccc8, []int{2}
|
||||
}
|
||||
func (m *ManifestChangeSet) XXX_Unmarshal(b []byte) error {
|
||||
return m.Unmarshal(b)
|
||||
@@ -206,8 +209,8 @@ func (m *ManifestChangeSet) XXX_Marshal(b []byte, deterministic bool) ([]byte, e
|
||||
return b[:n], nil
|
||||
}
|
||||
}
|
||||
func (dst *ManifestChangeSet) XXX_Merge(src proto.Message) {
|
||||
xxx_messageInfo_ManifestChangeSet.Merge(dst, src)
|
||||
func (m *ManifestChangeSet) XXX_Merge(src proto.Message) {
|
||||
xxx_messageInfo_ManifestChangeSet.Merge(m, src)
|
||||
}
|
||||
func (m *ManifestChangeSet) XXX_Size() int {
|
||||
return m.Size()
|
||||
@@ -229,6 +232,7 @@ type ManifestChange struct {
|
||||
Id uint64 `protobuf:"varint,1,opt,name=Id,proto3" json:"Id,omitempty"`
|
||||
Op ManifestChange_Operation `protobuf:"varint,2,opt,name=Op,proto3,enum=pb.ManifestChange_Operation" json:"Op,omitempty"`
|
||||
Level uint32 `protobuf:"varint,3,opt,name=Level,proto3" json:"Level,omitempty"`
|
||||
Checksum []byte `protobuf:"bytes,4,opt,name=Checksum,proto3" json:"Checksum,omitempty"`
|
||||
XXX_NoUnkeyedLiteral struct{} `json:"-"`
|
||||
XXX_unrecognized []byte `json:"-"`
|
||||
XXX_sizecache int32 `json:"-"`
|
||||
@@ -238,7 +242,7 @@ func (m *ManifestChange) Reset() { *m = ManifestChange{} }
|
||||
func (m *ManifestChange) String() string { return proto.CompactTextString(m) }
|
||||
func (*ManifestChange) ProtoMessage() {}
|
||||
func (*ManifestChange) Descriptor() ([]byte, []int) {
|
||||
return fileDescriptor_pb_7d6edc481ce0a5cd, []int{3}
|
||||
return fileDescriptor_f80abaa17e25ccc8, []int{3}
|
||||
}
|
||||
func (m *ManifestChange) XXX_Unmarshal(b []byte) error {
|
||||
return m.Unmarshal(b)
|
||||
@@ -255,8 +259,8 @@ func (m *ManifestChange) XXX_Marshal(b []byte, deterministic bool) ([]byte, erro
|
||||
return b[:n], nil
|
||||
}
|
||||
}
|
||||
func (dst *ManifestChange) XXX_Merge(src proto.Message) {
|
||||
xxx_messageInfo_ManifestChange.Merge(dst, src)
|
||||
func (m *ManifestChange) XXX_Merge(src proto.Message) {
|
||||
xxx_messageInfo_ManifestChange.Merge(m, src)
|
||||
}
|
||||
func (m *ManifestChange) XXX_Size() int {
|
||||
return m.Size()
|
||||
@@ -288,13 +292,49 @@ func (m *ManifestChange) GetLevel() uint32 {
|
||||
return 0
|
||||
}
|
||||
|
||||
func (m *ManifestChange) GetChecksum() []byte {
|
||||
if m != nil {
|
||||
return m.Checksum
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func init() {
|
||||
proto.RegisterEnum("pb.ManifestChange_Operation", ManifestChange_Operation_name, ManifestChange_Operation_value)
|
||||
proto.RegisterType((*KV)(nil), "pb.KV")
|
||||
proto.RegisterType((*KVList)(nil), "pb.KVList")
|
||||
proto.RegisterType((*ManifestChangeSet)(nil), "pb.ManifestChangeSet")
|
||||
proto.RegisterType((*ManifestChange)(nil), "pb.ManifestChange")
|
||||
proto.RegisterEnum("pb.ManifestChange_Operation", ManifestChange_Operation_name, ManifestChange_Operation_value)
|
||||
}
|
||||
|
||||
func init() { proto.RegisterFile("pb.proto", fileDescriptor_f80abaa17e25ccc8) }
|
||||
|
||||
var fileDescriptor_f80abaa17e25ccc8 = []byte{
|
||||
// 342 bytes of a gzipped FileDescriptorProto
|
||||
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x64, 0x91, 0x4d, 0x6a, 0xf2, 0x40,
|
||||
0x18, 0xc7, 0x9d, 0x31, 0x46, 0x7d, 0x5e, 0x5f, 0x49, 0x87, 0x52, 0x42, 0x3f, 0x42, 0x48, 0x37,
|
||||
0x2e, 0x24, 0x0b, 0x7b, 0x02, 0x6b, 0xb3, 0x10, 0x15, 0x61, 0x2a, 0x6e, 0x25, 0xd1, 0xa7, 0x35,
|
||||
0x44, 0x93, 0x21, 0x19, 0x43, 0x7b, 0x91, 0xd2, 0x0b, 0xf4, 0x2e, 0x5d, 0xf6, 0x08, 0xc5, 0x5e,
|
||||
0xa4, 0x64, 0xfc, 0x00, 0xe9, 0xee, 0xff, 0x31, 0xcf, 0x7f, 0xf1, 0x1b, 0xa8, 0x89, 0xc0, 0x15,
|
||||
0x69, 0x22, 0x13, 0x46, 0x45, 0xe0, 0xbc, 0x11, 0xa0, 0x83, 0x29, 0x33, 0xa0, 0x1c, 0xe1, 0xab,
|
||||
0x49, 0x6c, 0xd2, 0x6a, 0xf0, 0x42, 0xb2, 0x73, 0xa8, 0xe4, 0xfe, 0x6a, 0x83, 0x26, 0x55, 0xd9,
|
||||
0xce, 0xb0, 0x2b, 0xa8, 0x6f, 0x32, 0x4c, 0x67, 0x6b, 0x94, 0xbe, 0x59, 0x56, 0x4d, 0xad, 0x08,
|
||||
0x46, 0x28, 0x7d, 0x66, 0x42, 0x35, 0xc7, 0x34, 0x0b, 0x93, 0xd8, 0xd4, 0x6c, 0xd2, 0xd2, 0xf8,
|
||||
0xc1, 0xb2, 0x1b, 0x00, 0x7c, 0x11, 0x61, 0x8a, 0xd9, 0xcc, 0x97, 0x66, 0x45, 0x95, 0xf5, 0x7d,
|
||||
0xd2, 0x95, 0x8c, 0x81, 0xa6, 0x06, 0x75, 0x35, 0xa8, 0xb4, 0x63, 0x83, 0x3e, 0x98, 0x0e, 0xc3,
|
||||
0x4c, 0xb2, 0x0b, 0xa0, 0x51, 0x6e, 0x12, 0xbb, 0xdc, 0xfa, 0xd7, 0xd1, 0x5d, 0x11, 0xb8, 0x83,
|
||||
0x29, 0xa7, 0x51, 0xee, 0x74, 0xe1, 0x6c, 0xe4, 0xc7, 0xe1, 0x13, 0x66, 0xb2, 0xb7, 0xf4, 0xe3,
|
||||
0x67, 0x7c, 0x44, 0xc9, 0xda, 0x50, 0x9d, 0x2b, 0x93, 0xed, 0x2f, 0x58, 0x71, 0x71, 0xfa, 0x8e,
|
||||
0x1f, 0x9e, 0x38, 0x1f, 0x04, 0x9a, 0xa7, 0x1d, 0x6b, 0x02, 0xed, 0x2f, 0x14, 0x08, 0x8d, 0xd3,
|
||||
0xfe, 0x82, 0xb5, 0x81, 0x8e, 0x85, 0x82, 0xd0, 0xec, 0x5c, 0xff, 0xdd, 0x72, 0xc7, 0x02, 0x53,
|
||||
0x5f, 0x86, 0x49, 0xcc, 0xe9, 0x58, 0x14, 0xd4, 0x86, 0x98, 0xe3, 0x4a, 0xb1, 0xf9, 0xcf, 0x77,
|
||||
0x86, 0x5d, 0x42, 0xad, 0xb7, 0xc4, 0x79, 0x94, 0x6d, 0xd6, 0x8a, 0x4c, 0x83, 0x1f, 0xbd, 0x73,
|
||||
0x0b, 0xf5, 0xe3, 0x04, 0x03, 0xd0, 0x7b, 0xdc, 0xeb, 0x4e, 0x3c, 0xa3, 0x54, 0xe8, 0x07, 0x6f,
|
||||
0xe8, 0x4d, 0x3c, 0x83, 0xdc, 0x1b, 0x9f, 0x5b, 0x8b, 0x7c, 0x6d, 0x2d, 0xf2, 0xbd, 0xb5, 0xc8,
|
||||
0xfb, 0x8f, 0x55, 0x0a, 0x74, 0xf5, 0x85, 0x77, 0xbf, 0x01, 0x00, 0x00, 0xff, 0xff, 0x50, 0xdf,
|
||||
0x4a, 0x84, 0xce, 0x01, 0x00, 0x00,
|
||||
}
|
||||
|
||||
func (m *KV) Marshal() (dAtA []byte, err error) {
|
||||
size := m.Size()
|
||||
dAtA = make([]byte, size)
|
||||
@@ -446,6 +486,12 @@ func (m *ManifestChange) MarshalTo(dAtA []byte) (int, error) {
|
||||
i++
|
||||
i = encodeVarintPb(dAtA, i, uint64(m.Level))
|
||||
}
|
||||
if len(m.Checksum) > 0 {
|
||||
dAtA[i] = 0x22
|
||||
i++
|
||||
i = encodeVarintPb(dAtA, i, uint64(len(m.Checksum)))
|
||||
i += copy(dAtA[i:], m.Checksum)
|
||||
}
|
||||
if m.XXX_unrecognized != nil {
|
||||
i += copy(dAtA[i:], m.XXX_unrecognized)
|
||||
}
|
||||
@@ -546,6 +592,10 @@ func (m *ManifestChange) Size() (n int) {
|
||||
if m.Level != 0 {
|
||||
n += 1 + sovPb(uint64(m.Level))
|
||||
}
|
||||
l = len(m.Checksum)
|
||||
if l > 0 {
|
||||
n += 1 + l + sovPb(uint64(l))
|
||||
}
|
||||
if m.XXX_unrecognized != nil {
|
||||
n += len(m.XXX_unrecognized)
|
||||
}
|
||||
@@ -1028,6 +1078,37 @@ func (m *ManifestChange) Unmarshal(dAtA []byte) error {
|
||||
break
|
||||
}
|
||||
}
|
||||
case 4:
|
||||
if wireType != 2 {
|
||||
return fmt.Errorf("proto: wrong wireType = %d for field Checksum", wireType)
|
||||
}
|
||||
var byteLen int
|
||||
for shift := uint(0); ; shift += 7 {
|
||||
if shift >= 64 {
|
||||
return ErrIntOverflowPb
|
||||
}
|
||||
if iNdEx >= l {
|
||||
return io.ErrUnexpectedEOF
|
||||
}
|
||||
b := dAtA[iNdEx]
|
||||
iNdEx++
|
||||
byteLen |= (int(b) & 0x7F) << shift
|
||||
if b < 0x80 {
|
||||
break
|
||||
}
|
||||
}
|
||||
if byteLen < 0 {
|
||||
return ErrInvalidLengthPb
|
||||
}
|
||||
postIndex := iNdEx + byteLen
|
||||
if postIndex > l {
|
||||
return io.ErrUnexpectedEOF
|
||||
}
|
||||
m.Checksum = append(m.Checksum[:0], dAtA[iNdEx:postIndex]...)
|
||||
if m.Checksum == nil {
|
||||
m.Checksum = []byte{}
|
||||
}
|
||||
iNdEx = postIndex
|
||||
default:
|
||||
iNdEx = preIndex
|
||||
skippy, err := skipPb(dAtA[iNdEx:])
|
||||
@@ -1154,30 +1235,3 @@ var (
|
||||
ErrInvalidLengthPb = fmt.Errorf("proto: negative length found during unmarshaling")
|
||||
ErrIntOverflowPb = fmt.Errorf("proto: integer overflow")
|
||||
)
|
||||
|
||||
func init() { proto.RegisterFile("pb.proto", fileDescriptor_pb_7d6edc481ce0a5cd) }
|
||||
|
||||
var fileDescriptor_pb_7d6edc481ce0a5cd = []byte{
|
||||
// 325 bytes of a gzipped FileDescriptorProto
|
||||
0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x64, 0x51, 0x4f, 0x4e, 0xfa, 0x40,
|
||||
0x14, 0x66, 0x86, 0x52, 0xe0, 0xfd, 0x7e, 0x92, 0xfa, 0x62, 0x4c, 0x13, 0xb5, 0x69, 0xea, 0x86,
|
||||
0x05, 0xe9, 0x02, 0x4f, 0x80, 0xd8, 0x05, 0x01, 0x42, 0x32, 0x12, 0xb6, 0xa4, 0x95, 0xa7, 0x36,
|
||||
0x60, 0x3b, 0x69, 0x87, 0x46, 0x8f, 0xe0, 0x05, 0x8c, 0x47, 0x72, 0xe9, 0x11, 0x0c, 0x5e, 0xc4,
|
||||
0x74, 0x00, 0x13, 0xe2, 0xee, 0xfb, 0xf7, 0xbe, 0xc5, 0xf7, 0xa0, 0x21, 0x23, 0x5f, 0x66, 0xa9,
|
||||
0x4a, 0x91, 0xcb, 0xc8, 0x7b, 0x63, 0xc0, 0x87, 0x33, 0xb4, 0xa0, 0xba, 0xa4, 0x17, 0x9b, 0xb9,
|
||||
0xac, 0xfd, 0x5f, 0x94, 0x10, 0x4f, 0xa0, 0x56, 0x84, 0xab, 0x35, 0xd9, 0x5c, 0x6b, 0x5b, 0x82,
|
||||
0x67, 0xd0, 0x5c, 0xe7, 0x94, 0xcd, 0x9f, 0x48, 0x85, 0x76, 0x55, 0x3b, 0x8d, 0x52, 0x18, 0x93,
|
||||
0x0a, 0xd1, 0x86, 0x7a, 0x41, 0x59, 0x1e, 0xa7, 0x89, 0x6d, 0xb8, 0xac, 0x6d, 0x88, 0x3d, 0xc5,
|
||||
0x0b, 0x00, 0x7a, 0x96, 0x71, 0x46, 0xf9, 0x3c, 0x54, 0x76, 0x4d, 0x9b, 0xcd, 0x9d, 0xd2, 0x53,
|
||||
0x88, 0x60, 0xe8, 0x42, 0x53, 0x17, 0x6a, 0xec, 0xb9, 0x60, 0x0e, 0x67, 0xa3, 0x38, 0x57, 0x78,
|
||||
0x0a, 0x7c, 0x59, 0xd8, 0xcc, 0xad, 0xb6, 0xff, 0x75, 0x4d, 0x5f, 0x46, 0xfe, 0x70, 0x26, 0xf8,
|
||||
0xb2, 0xf0, 0x7a, 0x70, 0x3c, 0x0e, 0x93, 0xf8, 0x9e, 0x72, 0xd5, 0x7f, 0x0c, 0x93, 0x07, 0xba,
|
||||
0x25, 0x85, 0x1d, 0xa8, 0xdf, 0x69, 0x92, 0xef, 0x2e, 0xb0, 0xbc, 0x38, 0xcc, 0x89, 0x7d, 0xc4,
|
||||
0x7b, 0x65, 0xd0, 0x3a, 0xf4, 0xb0, 0x05, 0x7c, 0xb0, 0xd0, 0x43, 0x18, 0x82, 0x0f, 0x16, 0xd8,
|
||||
0x01, 0x3e, 0x91, 0x7a, 0x84, 0x56, 0xf7, 0xfc, 0x6f, 0x97, 0x3f, 0x91, 0x94, 0x85, 0x2a, 0x4e,
|
||||
0x13, 0xc1, 0x27, 0xb2, 0x5c, 0x6d, 0x44, 0x05, 0xad, 0xf4, 0x36, 0x47, 0x62, 0x4b, 0xbc, 0x4b,
|
||||
0x68, 0xfe, 0xc6, 0x10, 0xc0, 0xec, 0x8b, 0xa0, 0x37, 0x0d, 0xac, 0x4a, 0x89, 0x6f, 0x82, 0x51,
|
||||
0x30, 0x0d, 0x2c, 0x76, 0x6d, 0x7d, 0x6c, 0x1c, 0xf6, 0xb9, 0x71, 0xd8, 0xd7, 0xc6, 0x61, 0xef,
|
||||
0xdf, 0x4e, 0x25, 0x32, 0xf5, 0x9b, 0xae, 0x7e, 0x02, 0x00, 0x00, 0xff, 0xff, 0x8c, 0x0e, 0xe6,
|
||||
0x7c, 0xb2, 0x01, 0x00, 0x00,
|
||||
}
|
||||
|
||||
+3
-2
@@ -43,6 +43,7 @@ message ManifestChange {
|
||||
CREATE = 0;
|
||||
DELETE = 1;
|
||||
}
|
||||
Operation Op = 2;
|
||||
uint32 Level = 3; // Only used for CREATE
|
||||
Operation Op = 2;
|
||||
uint32 Level = 3; // Only used for CREATE
|
||||
bytes Checksum = 4; // Only used for CREATE
|
||||
}
|
||||
|
||||
+50
-47
@@ -17,8 +17,8 @@
|
||||
package table
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"crypto/sha256"
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -33,6 +33,7 @@ import (
|
||||
"github.com/AndreasBriese/bbloom"
|
||||
"github.com/dgraph-io/badger/options"
|
||||
"github.com/dgraph-io/badger/y"
|
||||
"github.com/dgraph-io/dgraph/x"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
@@ -69,6 +70,8 @@ type Table struct {
|
||||
id uint64 // file id, part of filename
|
||||
|
||||
bf bbloom.Bloom
|
||||
|
||||
Checksum []byte
|
||||
}
|
||||
|
||||
// IncrRef increments the refcount (having to do with whether the file should be deleted)
|
||||
@@ -115,7 +118,7 @@ func (b block) NewIterator() *blockIterator {
|
||||
// entry. Returns a table with one reference count on it (decrementing which may delete the file!
|
||||
// -- consider t.Close() instead). The fd has to writeable because we call Truncate on it before
|
||||
// deleting.
|
||||
func OpenTable(fd *os.File, loadingMode options.FileLoadingMode) (*Table, error) {
|
||||
func OpenTable(fd *os.File, mode options.FileLoadingMode, cksum []byte) (*Table, error) {
|
||||
fileInfo, err := fd.Stat()
|
||||
if err != nil {
|
||||
// It's OK to ignore fd.Close() errs in this function because we have only read
|
||||
@@ -134,26 +137,24 @@ func OpenTable(fd *os.File, loadingMode options.FileLoadingMode) (*Table, error)
|
||||
fd: fd,
|
||||
ref: 1, // Caller is given one reference.
|
||||
id: id,
|
||||
loadingMode: loadingMode,
|
||||
loadingMode: mode,
|
||||
}
|
||||
|
||||
t.tableSize = int(fileInfo.Size())
|
||||
|
||||
if loadingMode == options.MemoryMap {
|
||||
t.mmap, err = y.Mmap(fd, false, fileInfo.Size())
|
||||
if err != nil {
|
||||
_ = fd.Close()
|
||||
return nil, y.Wrapf(err, "Unable to map file")
|
||||
}
|
||||
} else if loadingMode == options.LoadToRAM {
|
||||
err = t.loadToRAM()
|
||||
if err != nil {
|
||||
_ = fd.Close()
|
||||
return nil, y.Wrap(err)
|
||||
}
|
||||
// We first load to RAM, so we can read the index and do checksum.
|
||||
if err := t.loadToRAM(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := t.readIndex(loadingMode); err != nil {
|
||||
// Enforce checksum before we read index. Otherwise, if the file was
|
||||
// truncated, we'd end up with panics in readIndex.
|
||||
if len(cksum) > 0 && !bytes.Equal(t.Checksum, cksum) {
|
||||
return nil, x.Errorf(
|
||||
"CHECKSUM_MISMATCH: Table checksum does not match checksum in MANIFEST."+
|
||||
" NOT including table %s. This would lead to missing data."+
|
||||
"\n sha256 %x Expected\n sha256 %x Found\n", filename, cksum, t.Checksum)
|
||||
}
|
||||
if err := t.readIndex(); err != nil {
|
||||
return nil, y.Wrap(err)
|
||||
}
|
||||
|
||||
@@ -170,6 +171,21 @@ func OpenTable(fd *os.File, loadingMode options.FileLoadingMode) (*Table, error)
|
||||
if it2.Valid() {
|
||||
t.biggest = it2.Key()
|
||||
}
|
||||
|
||||
switch mode {
|
||||
case options.LoadToRAM:
|
||||
// No need to do anything. t.mmap is already filled.
|
||||
case options.MemoryMap:
|
||||
t.mmap, err = y.Mmap(fd, false, fileInfo.Size())
|
||||
if err != nil {
|
||||
_ = fd.Close()
|
||||
return nil, y.Wrapf(err, "Unable to map file")
|
||||
}
|
||||
case options.FileIO:
|
||||
t.mmap = nil
|
||||
default:
|
||||
panic(fmt.Sprintf("Invalid loading mode: %v", mode))
|
||||
}
|
||||
return t, nil
|
||||
}
|
||||
|
||||
@@ -203,7 +219,10 @@ func (t *Table) readNoFail(off int, sz int) []byte {
|
||||
return res
|
||||
}
|
||||
|
||||
func (t *Table) readIndex(loadingMode options.FileLoadingMode) error {
|
||||
func (t *Table) readIndex() error {
|
||||
if len(t.mmap) != t.tableSize {
|
||||
panic("Table size does not match the read bytes")
|
||||
}
|
||||
readPos := t.tableSize
|
||||
|
||||
// Read bloom filter.
|
||||
@@ -243,39 +262,17 @@ func (t *Table) readIndex(loadingMode options.FileLoadingMode) error {
|
||||
t.blockIndex = append(t.blockIndex, ko)
|
||||
}
|
||||
|
||||
// Execute this index read serially, because all disks are orders of magnitude faster when read
|
||||
// serially compared to executing random reads.
|
||||
// Execute this index read serially, because we already have table data in memory.
|
||||
var h header
|
||||
var offset int
|
||||
var r *bufio.Reader
|
||||
if loadingMode == options.LoadToRAM {
|
||||
// We already read the table to put it into t.mmap. So, no point reading it again from disk.
|
||||
// Instead use the read buffer.
|
||||
r = bufio.NewReader(bytes.NewReader(t.mmap))
|
||||
} else {
|
||||
if _, err := t.fd.Seek(0, io.SeekStart); err != nil {
|
||||
return err
|
||||
}
|
||||
r = bufio.NewReader(t.fd)
|
||||
}
|
||||
hbuf := make([]byte, h.Size())
|
||||
for idx := range t.blockIndex {
|
||||
ko := &t.blockIndex[idx]
|
||||
if _, err := r.Discard(ko.offset - offset); err != nil {
|
||||
return err
|
||||
}
|
||||
offset = ko.offset
|
||||
if _, err := io.ReadFull(r, hbuf); err != nil {
|
||||
return err
|
||||
}
|
||||
offset += len(hbuf)
|
||||
|
||||
hbuf := t.readNoFail(ko.offset, h.Size())
|
||||
h.Decode(hbuf)
|
||||
y.AssertTrue(h.plen == 0)
|
||||
ko.key = make([]byte, h.klen)
|
||||
if _, err := io.ReadFull(r, ko.key); err != nil {
|
||||
return err
|
||||
}
|
||||
offset += len(ko.key)
|
||||
|
||||
key := t.readNoFail(ko.offset+len(hbuf), int(h.klen))
|
||||
ko.key = append([]byte{}, key...)
|
||||
}
|
||||
|
||||
return nil
|
||||
@@ -343,11 +340,17 @@ func NewFilename(id uint64, dir string) string {
|
||||
}
|
||||
|
||||
func (t *Table) loadToRAM() error {
|
||||
if _, err := t.fd.Seek(0, io.SeekStart); err != nil {
|
||||
return err
|
||||
}
|
||||
t.mmap = make([]byte, t.tableSize)
|
||||
read, err := t.fd.ReadAt(t.mmap, 0)
|
||||
sum := sha256.New()
|
||||
tee := io.TeeReader(t.fd, sum)
|
||||
read, err := tee.Read(t.mmap)
|
||||
if err != nil || read != t.tableSize {
|
||||
return y.Wrapf(err, "Unable to load file in memory. Table file: %s", t.Filename())
|
||||
}
|
||||
t.Checksum = sum.Sum(nil)
|
||||
y.NumReads.Add(1)
|
||||
y.NumBytesRead.Add(int64(read))
|
||||
return nil
|
||||
|
||||
+25
-25
@@ -79,7 +79,7 @@ func TestTableIterator(t *testing.T) {
|
||||
for _, n := range []int{99, 100, 101} {
|
||||
t.Run(fmt.Sprintf("n=%d", n), func(t *testing.T) {
|
||||
f := buildTestTable(t, "key", n)
|
||||
table, err := OpenTable(f, options.MemoryMap)
|
||||
table, err := OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer table.DecrRef()
|
||||
it := table.NewIterator(false)
|
||||
@@ -101,7 +101,7 @@ func TestSeekToFirst(t *testing.T) {
|
||||
for _, n := range []int{99, 100, 101, 199, 200, 250, 9999, 10000} {
|
||||
t.Run(fmt.Sprintf("n=%d", n), func(t *testing.T) {
|
||||
f := buildTestTable(t, "key", n)
|
||||
table, err := OpenTable(f, options.MemoryMap)
|
||||
table, err := OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer table.DecrRef()
|
||||
it := table.NewIterator(false)
|
||||
@@ -119,7 +119,7 @@ func TestSeekToLast(t *testing.T) {
|
||||
for _, n := range []int{99, 100, 101, 199, 200, 250, 9999, 10000} {
|
||||
t.Run(fmt.Sprintf("n=%d", n), func(t *testing.T) {
|
||||
f := buildTestTable(t, "key", n)
|
||||
table, err := OpenTable(f, options.MemoryMap)
|
||||
table, err := OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer table.DecrRef()
|
||||
it := table.NewIterator(false)
|
||||
@@ -140,7 +140,7 @@ func TestSeekToLast(t *testing.T) {
|
||||
|
||||
func TestSeek(t *testing.T) {
|
||||
f := buildTestTable(t, "k", 10000)
|
||||
table, err := OpenTable(f, options.MemoryMap)
|
||||
table, err := OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer table.DecrRef()
|
||||
|
||||
@@ -175,7 +175,7 @@ func TestSeek(t *testing.T) {
|
||||
|
||||
func TestSeekForPrev(t *testing.T) {
|
||||
f := buildTestTable(t, "k", 10000)
|
||||
table, err := OpenTable(f, options.MemoryMap)
|
||||
table, err := OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer table.DecrRef()
|
||||
|
||||
@@ -213,7 +213,7 @@ func TestIterateFromStart(t *testing.T) {
|
||||
for _, n := range []int{99, 100, 101, 199, 200, 250, 9999, 10000} {
|
||||
t.Run(fmt.Sprintf("n=%d", n), func(t *testing.T) {
|
||||
f := buildTestTable(t, "key", n)
|
||||
table, err := OpenTable(f, options.MemoryMap)
|
||||
table, err := OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer table.DecrRef()
|
||||
ti := table.NewIterator(false)
|
||||
@@ -240,7 +240,7 @@ func TestIterateFromEnd(t *testing.T) {
|
||||
for _, n := range []int{99, 100, 101, 199, 200, 250, 9999, 10000} {
|
||||
t.Run(fmt.Sprintf("n=%d", n), func(t *testing.T) {
|
||||
f := buildTestTable(t, "key", n)
|
||||
table, err := OpenTable(f, options.FileIO)
|
||||
table, err := OpenTable(f, options.FileIO, nil)
|
||||
require.NoError(t, err)
|
||||
defer table.DecrRef()
|
||||
ti := table.NewIterator(false)
|
||||
@@ -263,7 +263,7 @@ func TestIterateFromEnd(t *testing.T) {
|
||||
|
||||
func TestTable(t *testing.T) {
|
||||
f := buildTestTable(t, "key", 10000)
|
||||
table, err := OpenTable(f, options.FileIO)
|
||||
table, err := OpenTable(f, options.FileIO, nil)
|
||||
require.NoError(t, err)
|
||||
defer table.DecrRef()
|
||||
ti := table.NewIterator(false)
|
||||
@@ -290,7 +290,7 @@ func TestTable(t *testing.T) {
|
||||
|
||||
func TestIterateBackAndForth(t *testing.T) {
|
||||
f := buildTestTable(t, "key", 10000)
|
||||
table, err := OpenTable(f, options.MemoryMap)
|
||||
table, err := OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer table.DecrRef()
|
||||
|
||||
@@ -331,7 +331,7 @@ func TestIterateBackAndForth(t *testing.T) {
|
||||
|
||||
func TestUniIterator(t *testing.T) {
|
||||
f := buildTestTable(t, "key", 10000)
|
||||
table, err := OpenTable(f, options.MemoryMap)
|
||||
table, err := OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer table.DecrRef()
|
||||
{
|
||||
@@ -367,7 +367,7 @@ func TestConcatIteratorOneTable(t *testing.T) {
|
||||
{"k2", "a2"},
|
||||
})
|
||||
|
||||
tbl, err := OpenTable(f, options.MemoryMap)
|
||||
tbl, err := OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer tbl.DecrRef()
|
||||
|
||||
@@ -387,13 +387,13 @@ func TestConcatIterator(t *testing.T) {
|
||||
f := buildTestTable(t, "keya", 10000)
|
||||
f2 := buildTestTable(t, "keyb", 10000)
|
||||
f3 := buildTestTable(t, "keyc", 10000)
|
||||
tbl, err := OpenTable(f, options.MemoryMap)
|
||||
tbl, err := OpenTable(f, options.MemoryMap, nil)
|
||||
require.NoError(t, err)
|
||||
defer tbl.DecrRef()
|
||||
tbl2, err := OpenTable(f2, options.LoadToRAM)
|
||||
tbl2, err := OpenTable(f2, options.LoadToRAM, nil)
|
||||
require.NoError(t, err)
|
||||
defer tbl2.DecrRef()
|
||||
tbl3, err := OpenTable(f3, options.LoadToRAM)
|
||||
tbl3, err := OpenTable(f3, options.LoadToRAM, nil)
|
||||
require.NoError(t, err)
|
||||
defer tbl3.DecrRef()
|
||||
|
||||
@@ -472,10 +472,10 @@ func TestMergingIterator(t *testing.T) {
|
||||
{"k1", "b1"},
|
||||
{"k2", "b2"},
|
||||
})
|
||||
tbl1, err := OpenTable(f1, options.LoadToRAM)
|
||||
tbl1, err := OpenTable(f1, options.LoadToRAM, nil)
|
||||
require.NoError(t, err)
|
||||
defer tbl1.DecrRef()
|
||||
tbl2, err := OpenTable(f2, options.LoadToRAM)
|
||||
tbl2, err := OpenTable(f2, options.LoadToRAM, nil)
|
||||
require.NoError(t, err)
|
||||
defer tbl2.DecrRef()
|
||||
it1 := tbl1.NewIterator(false)
|
||||
@@ -512,10 +512,10 @@ func TestMergingIteratorReversed(t *testing.T) {
|
||||
{"k1", "b1"},
|
||||
{"k2", "b2"},
|
||||
})
|
||||
tbl1, err := OpenTable(f1, options.LoadToRAM)
|
||||
tbl1, err := OpenTable(f1, options.LoadToRAM, nil)
|
||||
require.NoError(t, err)
|
||||
defer tbl1.DecrRef()
|
||||
tbl2, err := OpenTable(f2, options.LoadToRAM)
|
||||
tbl2, err := OpenTable(f2, options.LoadToRAM, nil)
|
||||
require.NoError(t, err)
|
||||
defer tbl2.DecrRef()
|
||||
it1 := tbl1.NewIterator(true)
|
||||
@@ -551,10 +551,10 @@ func TestMergingIteratorTakeOne(t *testing.T) {
|
||||
})
|
||||
f2 := buildTable(t, [][]string{})
|
||||
|
||||
t1, err := OpenTable(f1, options.LoadToRAM)
|
||||
t1, err := OpenTable(f1, options.LoadToRAM, nil)
|
||||
require.NoError(t, err)
|
||||
defer t1.DecrRef()
|
||||
t2, err := OpenTable(f2, options.LoadToRAM)
|
||||
t2, err := OpenTable(f2, options.LoadToRAM, nil)
|
||||
require.NoError(t, err)
|
||||
defer t2.DecrRef()
|
||||
|
||||
@@ -591,10 +591,10 @@ func TestMergingIteratorTakeTwo(t *testing.T) {
|
||||
{"k2", "a2"},
|
||||
})
|
||||
|
||||
t1, err := OpenTable(f1, options.LoadToRAM)
|
||||
t1, err := OpenTable(f1, options.LoadToRAM, nil)
|
||||
require.NoError(t, err)
|
||||
defer t1.DecrRef()
|
||||
t2, err := OpenTable(f2, options.LoadToRAM)
|
||||
t2, err := OpenTable(f2, options.LoadToRAM, nil)
|
||||
require.NoError(t, err)
|
||||
defer t2.DecrRef()
|
||||
|
||||
@@ -636,7 +636,7 @@ func BenchmarkRead(b *testing.B) {
|
||||
}
|
||||
|
||||
f.Write(builder.Finish())
|
||||
tbl, err := OpenTable(f, options.MemoryMap)
|
||||
tbl, err := OpenTable(f, options.MemoryMap, nil)
|
||||
y.Check(err)
|
||||
defer tbl.DecrRef()
|
||||
|
||||
@@ -666,7 +666,7 @@ func BenchmarkReadAndBuild(b *testing.B) {
|
||||
}
|
||||
|
||||
f.Write(builder.Finish())
|
||||
tbl, err := OpenTable(f, options.MemoryMap)
|
||||
tbl, err := OpenTable(f, options.MemoryMap, nil)
|
||||
y.Check(err)
|
||||
defer tbl.DecrRef()
|
||||
|
||||
@@ -706,7 +706,7 @@ func BenchmarkReadMerged(b *testing.B) {
|
||||
y.Check(builder.Add([]byte(k), y.ValueStruct{Value: []byte(v), Meta: 123, UserMeta: 0}))
|
||||
}
|
||||
f.Write(builder.Finish())
|
||||
tbl, err := OpenTable(f, options.MemoryMap)
|
||||
tbl, err := OpenTable(f, options.MemoryMap, nil)
|
||||
y.Check(err)
|
||||
tables = append(tables, tbl)
|
||||
defer tbl.DecrRef()
|
||||
|
||||
Reference in New Issue
Block a user