diff --git a/.travis.yml b/.travis.yml index 79a7bb6c..e22ede6b 100644 --- a/.travis.yml +++ b/.travis.yml @@ -1,3 +1,4 @@ +dist: xenial language: go go: - 1.11 diff --git a/backup.go b/backup.go index a6673ff8..4b37379c 100644 --- a/backup.go +++ b/backup.go @@ -89,7 +89,7 @@ func OpenBackupEngine(opts *Options, path string) (*BackupEngine, error) { be := C.rocksdb_backup_engine_open(opts.c, cpath, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return &BackupEngine{ @@ -110,7 +110,7 @@ func (b *BackupEngine) CreateNewBackup(db *DB) error { C.rocksdb_backup_engine_create_new_backup(b.c, db.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } @@ -138,7 +138,7 @@ func (b *BackupEngine) RestoreDBFromLatestBackup(dbDir, walDir string, ro *Resto C.rocksdb_backup_engine_restore_db_from_latest_backup(b.c, cDbDir, cWalDir, ro.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil diff --git a/cache.go b/cache.go index ed708d7f..866326dc 100644 --- a/cache.go +++ b/cache.go @@ -9,7 +9,7 @@ type Cache struct { } // NewLRUCache creates a new LRU Cache object with the capacity given. -func NewLRUCache(capacity int) *Cache { +func NewLRUCache(capacity uint64) *Cache { return NewNativeCache(C.rocksdb_cache_create_lru(C.size_t(capacity))) } @@ -19,13 +19,13 @@ func NewNativeCache(c *C.rocksdb_cache_t) *Cache { } // GetUsage returns the Cache memory usage. -func (c *Cache) GetUsage() int { - return int(C.rocksdb_cache_get_usage(c.c)) +func (c *Cache) GetUsage() uint64 { + return uint64(C.rocksdb_cache_get_usage(c.c)) } // GetPinnedUsage returns the Cache pinned memory usage. -func (c *Cache) GetPinnedUsage() int { - return int(C.rocksdb_cache_get_pinned_usage(c.c)) +func (c *Cache) GetPinnedUsage() uint64 { + return uint64(C.rocksdb_cache_get_pinned_usage(c.c)) } // Destroy deallocates the Cache object. diff --git a/checkpoint.go b/checkpoint.go index a7d2bf40..4a6436d2 100644 --- a/checkpoint.go +++ b/checkpoint.go @@ -43,7 +43,7 @@ func (checkpoint *Checkpoint) CreateCheckpoint(checkpoint_dir string, log_size_f C.rocksdb_checkpoint_create(checkpoint.c, cDir, C.uint64_t(log_size_for_flush), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil diff --git a/checkpoint_test.go b/checkpoint_test.go index 9505740d..1ea10fdb 100644 --- a/checkpoint_test.go +++ b/checkpoint_test.go @@ -1,10 +1,11 @@ package gorocksdb import ( - "github.com/facebookgo/ensure" "io/ioutil" "os" "testing" + + "github.com/facebookgo/ensure" ) func TestCheckpoint(t *testing.T) { diff --git a/db.go b/db.go old mode 100644 new mode 100755 index e3c128ce..64735c61 --- a/db.go +++ b/db.go @@ -32,7 +32,7 @@ func OpenDb(opts *Options, name string) (*DB, error) { defer C.free(unsafe.Pointer(cName)) db := C.rocksdb_open(opts.c, cName, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return &DB{ @@ -51,7 +51,7 @@ func OpenDbWithTTL(opts *Options, name string, ttl int) (*DB, error) { defer C.free(unsafe.Pointer(cName)) db := C.rocksdb_open_with_ttl(opts.c, cName, C.int(ttl), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return &DB{ @@ -70,7 +70,7 @@ func OpenDbForReadOnly(opts *Options, name string, errorIfLogFileExist bool) (*D defer C.free(unsafe.Pointer(cName)) db := C.rocksdb_open_for_read_only(opts.c, cName, boolToChar(errorIfLogFileExist), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return &DB{ @@ -123,7 +123,7 @@ func OpenDbColumnFamilies( &cErr, ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, nil, errors.New(C.GoString(cErr)) } @@ -185,7 +185,7 @@ func OpenDbForReadOnlyColumnFamilies( &cErr, ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, nil, errors.New(C.GoString(cErr)) } @@ -211,12 +211,15 @@ func ListColumnFamilies(opts *Options, name string) ([]string, error) { defer C.free(unsafe.Pointer(cName)) cNames := C.rocksdb_list_column_families(opts.c, cName, &cLen, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } namesLen := int(cLen) names := make([]string, namesLen) - cNamesArr := (*[1 << 30]*C.char)(unsafe.Pointer(cNames))[:namesLen:namesLen] + // The maximum capacity of the following two slices is limited to (2^29)-1 to remain compatible + // with 32-bit platforms. The size of a `*C.char` (a pointer) is 4 Byte on a 32-bit system + // and (2^29)*4 == math.MaxInt32 + 1. -- See issue golang/go#13656 + cNamesArr := (*[(1 << 29) - 1]*C.char)(unsafe.Pointer(cNames))[:namesLen:namesLen] for i, n := range cNamesArr { names[i] = C.GoString(n) } @@ -243,7 +246,7 @@ func (db *DB) Get(opts *ReadOptions, key []byte) (*Slice, error) { ) cValue := C.rocksdb_get(db.c, opts.c, cKey, C.size_t(len(key)), &cValLen, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return NewSlice(cValue, cValLen), nil @@ -258,13 +261,13 @@ func (db *DB) GetBytes(opts *ReadOptions, key []byte) ([]byte, error) { ) cValue := C.rocksdb_get(db.c, opts.c, cKey, C.size_t(len(key)), &cValLen, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } if cValue == nil { return nil, nil } - defer C.free(unsafe.Pointer(cValue)) + defer C.rocksdb_free(unsafe.Pointer(cValue)) return C.GoBytes(unsafe.Pointer(cValue), C.int(cValLen)), nil } @@ -277,7 +280,7 @@ func (db *DB) GetCF(opts *ReadOptions, cf *ColumnFamilyHandle, key []byte) (*Sli ) cValue := C.rocksdb_get_cf(db.c, opts.c, cf.c, cKey, C.size_t(len(key)), &cValLen, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return NewSlice(cValue, cValLen), nil @@ -291,7 +294,7 @@ func (db *DB) GetPinned(opts *ReadOptions, key []byte) (*PinnableSliceHandle, er ) cHandle := C.rocksdb_get_pinned(db.c, opts.c, cKey, C.size_t(len(key)), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return NewNativePinnableSliceHandle(cHandle), nil @@ -320,7 +323,7 @@ func (db *DB) MultiGet(opts *ReadOptions, keys ...[]byte) (Slices, error) { for i, rocksErr := range rocksErrs { if rocksErr != nil { - defer C.free(unsafe.Pointer(rocksErr)) + defer C.rocksdb_free(unsafe.Pointer(rocksErr)) err := fmt.Errorf("getting %q failed: %v", string(keys[i]), C.GoString(rocksErr)) errs = append(errs, err) } @@ -372,7 +375,7 @@ func (db *DB) MultiGetCFMultiCF(opts *ReadOptions, cfs ColumnFamilyHandles, keys for i, rocksErr := range rocksErrs { if rocksErr != nil { - defer C.free(unsafe.Pointer(rocksErr)) + defer C.rocksdb_free(unsafe.Pointer(rocksErr)) err := fmt.Errorf("getting %q failed: %v", string(keys[i]), C.GoString(rocksErr)) errs = append(errs, err) } @@ -399,7 +402,7 @@ func (db *DB) Put(opts *WriteOptions, key, value []byte) error { ) C.rocksdb_put(db.c, opts.c, cKey, C.size_t(len(key)), cValue, C.size_t(len(value)), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -414,7 +417,7 @@ func (db *DB) PutCF(opts *WriteOptions, cf *ColumnFamilyHandle, key, value []byt ) C.rocksdb_put_cf(db.c, opts.c, cf.c, cKey, C.size_t(len(key)), cValue, C.size_t(len(value)), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -428,7 +431,7 @@ func (db *DB) Delete(opts *WriteOptions, key []byte) error { ) C.rocksdb_delete(db.c, opts.c, cKey, C.size_t(len(key)), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -442,7 +445,7 @@ func (db *DB) DeleteCF(opts *WriteOptions, cf *ColumnFamilyHandle, key []byte) e ) C.rocksdb_delete_cf(db.c, opts.c, cf.c, cKey, C.size_t(len(key)), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -457,7 +460,7 @@ func (db *DB) Merge(opts *WriteOptions, key []byte, value []byte) error { ) C.rocksdb_merge(db.c, opts.c, cKey, C.size_t(len(key)), cValue, C.size_t(len(value)), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -473,7 +476,7 @@ func (db *DB) MergeCF(opts *WriteOptions, cf *ColumnFamilyHandle, key []byte, va ) C.rocksdb_merge_cf(db.c, opts.c, cf.c, cKey, C.size_t(len(key)), cValue, C.size_t(len(value)), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -484,7 +487,7 @@ func (db *DB) Write(opts *WriteOptions, batch *WriteBatch) error { var cErr *C.char C.rocksdb_write(db.c, opts.c, batch.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -504,6 +507,20 @@ func (db *DB) NewIteratorCF(opts *ReadOptions, cf *ColumnFamilyHandle) *Iterator return NewNativeIterator(unsafe.Pointer(cIter)) } +func (db *DB) GetUpdatesSince(seqNumber uint64) (*WalIterator, error) { + var cErr *C.char + cIter := C.rocksdb_get_updates_since(db.c, C.uint64_t(seqNumber), nil, &cErr) + if cErr != nil { + defer C.rocksdb_free(unsafe.Pointer(cErr)) + return nil, errors.New(C.GoString(cErr)) + } + return NewNativeWalIterator(unsafe.Pointer(cIter)), nil +} + +func (db *DB) GetLatestSequenceNumber() uint64 { + return uint64(C.rocksdb_get_latest_sequence_number(db.c)) +} + // NewSnapshot creates a new snapshot of the database. func (db *DB) NewSnapshot() *Snapshot { cSnap := C.rocksdb_create_snapshot(db.c) @@ -521,7 +538,7 @@ func (db *DB) GetProperty(propName string) string { cprop := C.CString(propName) defer C.free(unsafe.Pointer(cprop)) cValue := C.rocksdb_property_value(db.c, cprop) - defer C.free(unsafe.Pointer(cValue)) + defer C.rocksdb_free(unsafe.Pointer(cValue)) return C.GoString(cValue) } @@ -530,7 +547,7 @@ func (db *DB) GetPropertyCF(propName string, cf *ColumnFamilyHandle) string { cProp := C.CString(propName) defer C.free(unsafe.Pointer(cProp)) cValue := C.rocksdb_property_value_cf(db.c, cf.c, cProp) - defer C.free(unsafe.Pointer(cValue)) + defer C.rocksdb_free(unsafe.Pointer(cValue)) return C.GoString(cValue) } @@ -543,7 +560,7 @@ func (db *DB) CreateColumnFamily(opts *Options, name string) (*ColumnFamilyHandl defer C.free(unsafe.Pointer(cName)) cHandle := C.rocksdb_create_column_family(db.c, opts.c, cName, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return NewNativeColumnFamilyHandle(cHandle), nil @@ -554,7 +571,7 @@ func (db *DB) DropColumnFamily(c *ColumnFamilyHandle) error { var cErr *C.char C.rocksdb_drop_column_family(db.c, c.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -668,7 +685,7 @@ func (db *DB) SetOptions(keys, values []string) error { &cErr, ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -729,7 +746,7 @@ func (db *DB) Flush(opts *FlushOptions) error { var cErr *C.char C.rocksdb_flush(db.c, opts.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -740,7 +757,7 @@ func (db *DB) DisableFileDeletions() error { var cErr *C.char C.rocksdb_disable_file_deletions(db.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -751,7 +768,7 @@ func (db *DB) EnableFileDeletions(force bool) error { var cErr *C.char C.rocksdb_enable_file_deletions(db.c, boolToChar(force), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -766,6 +783,50 @@ func (db *DB) DeleteFile(name string) { C.rocksdb_delete_file(db.c, cName) } +// DeleteFileInRange deletes SST files that contain keys between the Range, [r.Start, r.Limit] +func (db *DB) DeleteFileInRange(r Range) error { + cStartKey := byteToChar(r.Start) + cLimitKey := byteToChar(r.Limit) + + var cErr *C.char + + C.rocksdb_delete_file_in_range( + db.c, + cStartKey, C.size_t(len(r.Start)), + cLimitKey, C.size_t(len(r.Limit)), + &cErr, + ) + + if cErr != nil { + defer C.rocksdb_free(unsafe.Pointer(cErr)) + return errors.New(C.GoString(cErr)) + } + return nil +} + +// DeleteFileInRangeCF deletes SST files that contain keys between the Range, [r.Start, r.Limit], and +// belong to a given column family +func (db *DB) DeleteFileInRangeCF(cf *ColumnFamilyHandle, r Range) error { + cStartKey := byteToChar(r.Start) + cLimitKey := byteToChar(r.Limit) + + var cErr *C.char + + C.rocksdb_delete_file_in_range_cf( + db.c, + cf.c, + cStartKey, C.size_t(len(r.Start)), + cLimitKey, C.size_t(len(r.Limit)), + &cErr, + ) + + if cErr != nil { + defer C.rocksdb_free(unsafe.Pointer(cErr)) + return errors.New(C.GoString(cErr)) + } + return nil +} + // IngestExternalFile loads a list of external SST files. func (db *DB) IngestExternalFile(filePaths []string, opts *IngestExternalFileOptions) error { cFilePaths := make([]*C.char, len(filePaths)) @@ -789,7 +850,7 @@ func (db *DB) IngestExternalFile(filePaths []string, opts *IngestExternalFileOpt ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -819,7 +880,7 @@ func (db *DB) IngestExternalFileCF(handle *ColumnFamilyHandle, filePaths []strin ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -834,7 +895,7 @@ func (db *DB) NewCheckpoint() (*Checkpoint, error) { db.c, &cErr, ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } @@ -856,7 +917,7 @@ func DestroyDb(name string, opts *Options) error { defer C.free(unsafe.Pointer(cName)) C.rocksdb_destroy_db(opts.c, cName, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -871,7 +932,7 @@ func RepairDb(name string, opts *Options) error { defer C.free(unsafe.Pointer(cName)) C.rocksdb_repair_db(opts.c, cName, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil diff --git a/db_test.go b/db_test.go old mode 100644 new mode 100755 index 3c0df0a1..4ccc7aa8 --- a/db_test.go +++ b/db_test.go @@ -19,9 +19,8 @@ func TestDBCRUD(t *testing.T) { var ( givenKey = []byte("hello") - givenVal1 = []byte("world1") - givenVal2 = []byte("world2") - givenVal3 = []byte("world2") + givenVal1 = []byte("") + givenVal2 = []byte("world1") wo = NewDefaultWriteOptions() ro = NewDefaultReadOptions() ) @@ -46,7 +45,7 @@ func TestDBCRUD(t *testing.T) { v3, err := db.GetPinned(ro, givenKey) defer v3.Destroy() ensure.Nil(t, err) - ensure.DeepEqual(t, v3.Data(), givenVal3) + ensure.DeepEqual(t, v3.Data(), givenVal2) // delete ensure.Nil(t, db.Delete(wo, givenKey)) @@ -58,7 +57,7 @@ func TestDBCRUD(t *testing.T) { v5, err := db.GetPinned(ro, givenKey) defer v5.Destroy() ensure.Nil(t, err) - ensure.Nil(t, v5.Data()) + ensure.True(t, v5.Data() == nil) } func TestDBCRUDDBPaths(t *testing.T) { @@ -75,12 +74,20 @@ func TestDBCRUDDBPaths(t *testing.T) { var ( givenKey = []byte("hello") - givenVal1 = []byte("world1") - givenVal2 = []byte("world2") + givenVal1 = []byte("") + givenVal2 = []byte("world1") + givenVal3 = []byte("world2") wo = NewDefaultWriteOptions() ro = NewDefaultReadOptions() ) + // retrieve before create + noexist, err := db.Get(ro, givenKey) + defer noexist.Free() + ensure.Nil(t, err) + ensure.False(t, noexist.Exists()) + ensure.DeepEqual(t, noexist.Data(), []byte(nil)) + // create ensure.Nil(t, db.Put(wo, givenKey, givenVal1)) @@ -88,6 +95,7 @@ func TestDBCRUDDBPaths(t *testing.T) { v1, err := db.Get(ro, givenKey) defer v1.Free() ensure.Nil(t, err) + ensure.True(t, v1.Exists()) ensure.DeepEqual(t, v1.Data(), givenVal1) // update @@ -95,13 +103,24 @@ func TestDBCRUDDBPaths(t *testing.T) { v2, err := db.Get(ro, givenKey) defer v2.Free() ensure.Nil(t, err) + ensure.True(t, v2.Exists()) ensure.DeepEqual(t, v2.Data(), givenVal2) + // update + ensure.Nil(t, db.Put(wo, givenKey, givenVal3)) + v3, err := db.Get(ro, givenKey) + defer v3.Free() + ensure.Nil(t, err) + ensure.True(t, v3.Exists()) + ensure.DeepEqual(t, v3.Data(), givenVal3) + // delete ensure.Nil(t, db.Delete(wo, givenKey)) - v3, err := db.Get(ro, givenKey) + v4, err := db.Get(ro, givenKey) + defer v4.Free() ensure.Nil(t, err) - ensure.True(t, v3.Data() == nil) + ensure.False(t, v4.Exists()) + ensure.DeepEqual(t, v4.Data(), []byte(nil)) } func newTestDB(t *testing.T, name string, applyOpts func(opts *Options)) *DB { diff --git a/dynflag.go b/dynflag.go index 92b39934..322867f1 100644 --- a/dynflag.go +++ b/dynflag.go @@ -2,5 +2,5 @@ package gorocksdb -// #cgo LDFLAGS: -lrocksdb -lstdc++ -lm -lz -lbz2 -lsnappy -llz4 -lzstd +// #cgo LDFLAGS: -lrocksdb -lstdc++ -lm -lz -lbz2 -lsnappy -llz4 -lzstd -ldl import "C" diff --git a/iterator.go b/iterator.go index 4a280e2c..fefb82f1 100644 --- a/iterator.go +++ b/iterator.go @@ -113,7 +113,7 @@ func (iter *Iterator) Err() error { var cErr *C.char C.rocksdb_iter_get_error(iter.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil diff --git a/memory_usage.go b/memory_usage.go index 740b877d..7b9a6ad6 100644 --- a/memory_usage.go +++ b/memory_usage.go @@ -42,7 +42,7 @@ func GetApproximateMemoryUsageByType(dbs []*DB, caches []*Cache) (*MemoryUsage, var cErr *C.char memoryUsage := C.rocksdb_approximate_memory_usage_create(consumers, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } @@ -55,4 +55,4 @@ func GetApproximateMemoryUsageByType(dbs []*DB, caches []*Cache) (*MemoryUsage, CacheTotal: uint64(C.rocksdb_approximate_memory_usage_get_cache_total(memoryUsage)), } return result, nil -} \ No newline at end of file +} diff --git a/options.go b/options.go index 0f1376ad..5f80d1af 100644 --- a/options.go +++ b/options.go @@ -60,6 +60,15 @@ const ( FatalInfoLogLevel = InfoLogLevel(4) ) +type WALRecoveryMode int + +const ( + TolerateCorruptedTailRecordsRecovery = WALRecoveryMode(0) + AbsoluteConsistencyRecovery = WALRecoveryMode(1) + PointInTimeRecovery = WALRecoveryMode(2) + SkipAnyCorruptedRecordsRecovery = WALRecoveryMode(3) +) + // Options represent all of the available options when opening a database with Open. type Options struct { c *C.rocksdb_options_t @@ -102,7 +111,7 @@ func GetOptionsFromString(base *Options, optStr string) (*Options, error) { newOpt := NewDefaultOptions() C.rocksdb_get_options_from_string(base.c, cOptStr, newOpt.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } @@ -333,7 +342,7 @@ func (opts *Options) OptimizeUniversalStyleCompaction(memtable_memory_budget uin // so you may wish to adjust this parameter to control memory usage. // Also, a larger write buffer will result in a longer recovery time // the next time the database is opened. -// Default: 4MB +// Default: 64MB func (opts *Options) SetWriteBufferSize(value int) { C.rocksdb_options_set_write_buffer_size(opts.c, C.size_t(value)) } @@ -801,6 +810,14 @@ func (opts *Options) SetDisableAutoCompactions(value bool) { C.rocksdb_options_set_disable_auto_compactions(opts.c, C.int(btoi(value))) } +// SetWALRecoveryMode sets the recovery mode +// +// Recovery mode to control the consistency while replaying WAL +// Default: TolerateCorruptedTailRecordsRecovery +func (opts *Options) SetWALRecoveryMode(mode WALRecoveryMode) { + C.rocksdb_options_set_wal_recovery_mode(opts.c, C.int(mode)) +} + // SetWALTtlSeconds sets the WAL ttl in seconds. // // The following two options affect how archived logs will be deleted. @@ -829,6 +846,13 @@ func (opts *Options) SetWalSizeLimitMb(value uint64) { C.rocksdb_options_set_WAL_size_limit_MB(opts.c, C.uint64_t(value)) } +// SetEnablePipelinedWrite enables pipelined write +// +// Default: false +func (opts *Options) SetEnablePipelinedWrite(value bool) { + C.rocksdb_options_set_enable_pipelined_write(opts.c, boolToChar(value)) +} + // SetManifestPreallocationSize sets the number of bytes // to preallocate (via fallocate) the manifest files. // @@ -967,7 +991,7 @@ func (opts *Options) SetFIFOCompactionOptions(value *FIFOCompactionOptions) { // GetStatisticsString returns the statistics as a string. func (opts *Options) GetStatisticsString() string { sString := C.rocksdb_options_statistics_get_string(opts.c) - defer C.free(unsafe.Pointer(sString)) + defer C.rocksdb_free(unsafe.Pointer(sString)) return C.GoString(sString) } @@ -1142,6 +1166,36 @@ func (opts *Options) SetAllowIngestBehind(value bool) { C.rocksdb_options_set_allow_ingest_behind(opts.c, boolToChar(value)) } +// SetMemTablePrefixBloomSizeRatio sets memtable_prefix_bloom_size_ratio +// if prefix_extractor is set and memtable_prefix_bloom_size_ratio is not 0, +// create prefix bloom for memtable with the size of +// write_buffer_size * memtable_prefix_bloom_size_ratio. +// If it is larger than 0.25, it is sanitized to 0.25. +// +// Default: 0 (disable) +func (opts *Options) SetMemTablePrefixBloomSizeRatio(value float64) { + C.rocksdb_options_set_memtable_prefix_bloom_size_ratio(opts.c, C.double(value)) +} + +// SetOptimizeFiltersForHits sets optimize_filters_for_hits +// This flag specifies that the implementation should optimize the filters +// mainly for cases where keys are found rather than also optimize for keys +// missed. This would be used in cases where the application knows that +// there are very few misses or the performance in the case of misses is not +// important. +// +// For now, this flag allows us to not store filters for the last level i.e +// the largest level which contains data of the LSM store. For keys which +// are hits, the filters in this level are not useful because we will search +// for the data anyway. NOTE: the filters in other levels are still useful +// even for key hit because they tell us whether to look in that level or go +// to the higher level. +// +// Default: false +func (opts *Options) SetOptimizeFiltersForHits(value bool) { + C.rocksdb_options_set_optimize_filters_for_hits(opts.c, C.int(btoi(value))) +} + // Destroy deallocates the Options object. func (opts *Options) Destroy() { C.rocksdb_options_destroy(opts.c) diff --git a/slice.go b/slice.go index 223ef321..01ecfd3c 100644 --- a/slice.go +++ b/slice.go @@ -32,7 +32,8 @@ func StringToSlice(data string) *Slice { return NewSlice(C.CString(data), C.size_t(len(data))) } -// Data returns the data of the slice. +// Data returns the data of the slice. If the key doesn't exist this will be a +// nil slice. func (s *Slice) Data() []byte { return charToByte(s.data, s.size) } @@ -42,10 +43,15 @@ func (s *Slice) Size() int { return int(s.size) } +// Exists returns if the key exists +func (s *Slice) Exists() bool { + return s.data != nil +} + // Free frees the slice data. func (s *Slice) Free() { if !s.freed { - C.free(unsafe.Pointer(s.data)) + C.rocksdb_free(unsafe.Pointer(s.data)) s.freed = true } } diff --git a/slice_transform.go b/slice_transform.go index e66e4d84..8b9b2362 100644 --- a/slice_transform.go +++ b/slice_transform.go @@ -23,6 +23,11 @@ func NewFixedPrefixTransform(prefixLen int) SliceTransform { return NewNativeSliceTransform(C.rocksdb_slicetransform_create_fixed_prefix(C.size_t(prefixLen))) } +// NewNoopPrefixTransform creates a new no-op prefix transform. +func NewNoopPrefixTransform() SliceTransform { + return NewNativeSliceTransform(C.rocksdb_slicetransform_create_noop()) +} + // NewNativeSliceTransform creates a SliceTransform object. func NewNativeSliceTransform(c *C.rocksdb_slicetransform_t) SliceTransform { return nativeSliceTransform{c} diff --git a/slice_transform_test.go b/slice_transform_test.go index 1c551183..d60c7326 100644 --- a/slice_transform_test.go +++ b/slice_transform_test.go @@ -35,6 +35,13 @@ func TestFixedPrefixTransformOpen(t *testing.T) { defer db.Close() } +func TestNewNoopPrefixTransform(t *testing.T) { + db := newTestDB(t, "TestNewNoopPrefixTransform", func(opts *Options) { + opts.SetPrefixExtractor(NewNoopPrefixTransform()) + }) + defer db.Close() +} + type testSliceTransform struct { initiated bool } diff --git a/sst_file_writer.go b/sst_file_writer.go index 0f4689c2..54f2c139 100644 --- a/sst_file_writer.go +++ b/sst_file_writer.go @@ -30,7 +30,7 @@ func (w *SSTFileWriter) Open(path string) error { defer C.free(unsafe.Pointer(cPath)) C.rocksdb_sstfilewriter_open(w.c, cPath, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -44,7 +44,7 @@ func (w *SSTFileWriter) Add(key, value []byte) error { var cErr *C.char C.rocksdb_sstfilewriter_add(w.c, cKey, C.size_t(len(key)), cValue, C.size_t(len(value)), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -55,7 +55,7 @@ func (w *SSTFileWriter) Finish() error { var cErr *C.char C.rocksdb_sstfilewriter_finish(w.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil diff --git a/transaction.go b/transaction.go index 49a04bd7..67c9ef09 100644 --- a/transaction.go +++ b/transaction.go @@ -26,7 +26,7 @@ func (transaction *Transaction) Commit() error { ) C.rocksdb_transaction_commit(transaction.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -40,7 +40,7 @@ func (transaction *Transaction) Rollback() error { C.rocksdb_transaction_rollback(transaction.c, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -57,7 +57,7 @@ func (transaction *Transaction) Get(opts *ReadOptions, key []byte) (*Slice, erro transaction.c, opts.c, cKey, C.size_t(len(key)), &cValLen, &cErr, ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return NewSlice(cValue, cValLen), nil @@ -74,7 +74,7 @@ func (transaction *Transaction) GetForUpdate(opts *ReadOptions, key []byte) (*Sl transaction.c, opts.c, cKey, C.size_t(len(key)), &cValLen, C.uchar(byte(1)) /*exclusive*/, &cErr, ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return NewSlice(cValue, cValLen), nil @@ -91,7 +91,7 @@ func (transaction *Transaction) Put(key, value []byte) error { transaction.c, cKey, C.size_t(len(key)), cValue, C.size_t(len(value)), &cErr, ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -105,7 +105,7 @@ func (transaction *Transaction) Delete(key []byte) error { ) C.rocksdb_transaction_delete(transaction.c, cKey, C.size_t(len(key)), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil diff --git a/transactiondb.go b/transactiondb.go index f5d2fd70..cfdeac9c 100644 --- a/transactiondb.go +++ b/transactiondb.go @@ -30,7 +30,7 @@ func OpenTransactionDb( db := C.rocksdb_transactiondb_open( opts.c, transactionDBOpts.c, cName, &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return &TransactionDB{ @@ -83,7 +83,7 @@ func (db *TransactionDB) Get(opts *ReadOptions, key []byte) (*Slice, error) { db.c, opts.c, cKey, C.size_t(len(key)), &cValLen, &cErr, ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } return NewSlice(cValue, cValLen), nil @@ -100,7 +100,7 @@ func (db *TransactionDB) Put(opts *WriteOptions, key, value []byte) error { db.c, opts.c, cKey, C.size_t(len(key)), cValue, C.size_t(len(value)), &cErr, ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -114,7 +114,7 @@ func (db *TransactionDB) Delete(opts *WriteOptions, key []byte) error { ) C.rocksdb_transactiondb_delete(db.c, opts.c, cKey, C.size_t(len(key)), &cErr) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return errors.New(C.GoString(cErr)) } return nil @@ -129,7 +129,7 @@ func (db *TransactionDB) NewCheckpoint() (*Checkpoint, error) { db.c, &cErr, ) if cErr != nil { - defer C.free(unsafe.Pointer(cErr)) + defer C.rocksdb_free(unsafe.Pointer(cErr)) return nil, errors.New(C.GoString(cErr)) } diff --git a/util.go b/util.go index ae099cd3..b6637ec3 100644 --- a/util.go +++ b/util.go @@ -1,4 +1,5 @@ package gorocksdb + // #include import "C" diff --git a/wal_iterator.go b/wal_iterator.go new file mode 100755 index 00000000..7805d7c9 --- /dev/null +++ b/wal_iterator.go @@ -0,0 +1,49 @@ +package gorocksdb + +// #include +// #include "rocksdb/c.h" +import "C" +import ( + "errors" + "unsafe" +) + +type WalIterator struct { + c *C.rocksdb_wal_iterator_t +} + +func NewNativeWalIterator(c unsafe.Pointer) *WalIterator { + return &WalIterator{(*C.rocksdb_wal_iterator_t)(c)} +} + +func (iter *WalIterator) Valid() bool { + return C.rocksdb_wal_iter_valid(iter.c) != 0 +} + +func (iter *WalIterator) Next() { + C.rocksdb_wal_iter_next(iter.c) +} + +func (iter *WalIterator) Err() error { + var cErr *C.char + C.rocksdb_wal_iter_status(iter.c, &cErr) + if cErr != nil { + defer C.rocksdb_free(unsafe.Pointer(cErr)) + return errors.New(C.GoString(cErr)) + } + return nil +} + +func (iter *WalIterator) Destroy() { + C.rocksdb_wal_iter_destroy(iter.c) + iter.c = nil +} + +// C.rocksdb_wal_iter_get_batch in the official rocksdb c wrapper has memory leak +// see https://github.com/facebook/rocksdb/pull/5515 +// https://github.com/facebook/rocksdb/issues/5536 +func (iter *WalIterator) GetBatch() (*WriteBatch, uint64) { + var cSeq C.uint64_t + cB := C.rocksdb_wal_iter_get_batch(iter.c, &cSeq) + return NewNativeWriteBatch(cB), uint64(cSeq) +} diff --git a/write_batch.go b/write_batch.go index 88b5aac8..f894427b 100644 --- a/write_batch.go +++ b/write_batch.go @@ -41,6 +41,12 @@ func (wb *WriteBatch) PutCF(cf *ColumnFamilyHandle, key, value []byte) { C.rocksdb_writebatch_put_cf(wb.c, cf.c, cKey, C.size_t(len(key)), cValue, C.size_t(len(value))) } +// Append a blob of arbitrary size to the records in this batch. +func (wb *WriteBatch) PutLogData(blob []byte) { + cBlob := byteToChar(blob) + C.rocksdb_writebatch_put_log_data(wb.c, cBlob, C.size_t(len(blob))) +} + // Merge queues a merge of "value" with the existing value of "key". func (wb *WriteBatch) Merge(key, value []byte) { cKey := byteToChar(key) @@ -68,6 +74,21 @@ func (wb *WriteBatch) DeleteCF(cf *ColumnFamilyHandle, key []byte) { C.rocksdb_writebatch_delete_cf(wb.c, cf.c, cKey, C.size_t(len(key))) } +// DeleteRange deletes keys that are between [startKey, endKey) +func (wb *WriteBatch) DeleteRange(startKey []byte, endKey []byte) { + cStartKey := byteToChar(startKey) + cEndKey := byteToChar(endKey) + C.rocksdb_writebatch_delete_range(wb.c, cStartKey, C.size_t(len(startKey)), cEndKey, C.size_t(len(endKey))) +} + +// DeleteRangeCF deletes keys that are between [startKey, endKey) and +// belong to a given column family +func (wb *WriteBatch) DeleteRangeCF(cf *ColumnFamilyHandle, startKey []byte, endKey []byte) { + cStartKey := byteToChar(startKey) + cEndKey := byteToChar(endKey) + C.rocksdb_writebatch_delete_range_cf(wb.c, cf.c, cStartKey, C.size_t(len(startKey)), cEndKey, C.size_t(len(endKey))) +} + // Data returns the serialized version of this batch. func (wb *WriteBatch) Data() []byte { var cSize C.size_t diff --git a/write_batch_test.go b/write_batch_test.go index d6000e57..72eeb36e 100644 --- a/write_batch_test.go +++ b/write_batch_test.go @@ -39,6 +39,18 @@ func TestWriteBatch(t *testing.T) { defer v2.Free() ensure.Nil(t, err) ensure.True(t, v2.Data() == nil) + + // DeleteRange test + wb.Clear() + wb.DeleteRange(givenKey1, givenKey2) + + // perform the batch + ensure.Nil(t, db.Write(wo, wb)) + + v1, err = db.Get(ro, givenKey1) + defer v1.Free() + ensure.Nil(t, err) + ensure.True(t, v1.Data() == nil) } func TestWriteBatchIterator(t *testing.T) {