From 560a9eb741ef66aa7986fd7eee1ee7713c6a7434 Mon Sep 17 00:00:00 2001 From: Victor Grishchenko Date: Fri, 22 Dec 2017 19:22:36 +0500 Subject: [PATCH 1/3] extend the write batch parser --- write_batch.go | 122 +++++++++++++++++++++++++++++++++---------------- 1 file changed, 83 insertions(+), 39 deletions(-) diff --git a/write_batch.go b/write_batch.go index a4a5b8ba..0a45a9b2 100644 --- a/write_batch.go +++ b/write_batch.go @@ -2,7 +2,10 @@ package gorocksdb // #include "rocksdb/c.h" import "C" -import "io" +import ( + "errors" + "io" +) // WriteBatch is a batching of Puts, Merges and Deletes. type WriteBatch struct { @@ -102,14 +105,26 @@ type WriteBatchRecordType byte // Types of batch records. const ( - WriteBatchRecordTypeDeletion WriteBatchRecordType = 0x0 - WriteBatchRecordTypeValue WriteBatchRecordType = 0x1 - WriteBatchRecordTypeMerge WriteBatchRecordType = 0x2 - WriteBatchRecordTypeLogData WriteBatchRecordType = 0x3 + WriteBatchRecordTypeDeletion WriteBatchRecordType = 0x0 + WriteBatchRecordTypeValue WriteBatchRecordType = 0x1 + WriteBatchRecordTypeMerge WriteBatchRecordType = 0x2 + WriteBatchRecordTypeLogData WriteBatchRecordType = 0x3 + WriteBatchRecordTypeCFDeletion WriteBatchRecordType = 0x4 + WriteBatchRecordTypeCFValue WriteBatchRecordType = 0x5 + WriteBatchRecordTypeCFMerge WriteBatchRecordType = 0x6 + WriteBatchRecordTypeSingleDeletion WriteBatchRecordType = 0x7 + WriteBatchRecordTypeCFSingleDeletion WriteBatchRecordType = 0x8 + WriteBatchRecordTypeNoop WriteBatchRecordType = 0xD + WriteBatchRecordTypeBeginPrepareXID WriteBatchRecordType = 0x9 + WriteBatchRecordTypeEndPrepareXID WriteBatchRecordType = 0xA + WriteBatchRecordTypeCommitXID WriteBatchRecordType = 0xB + WriteBatchRecordTypeRollbackXID WriteBatchRecordType = 0xC + WriteBatchRecordTypeNotUsed WriteBatchRecordType = 0x7F ) // WriteBatchRecord represents a record inside a WriteBatch. type WriteBatchRecord struct { + CF int Key []byte Value []byte Type WriteBatchRecordType @@ -133,32 +148,35 @@ func (iter *WriteBatchIterator) Next() bool { iter.record.Value = nil // parse the record type - recordType := WriteBatchRecordType(iter.data[0]) - iter.record.Type = recordType - iter.data = iter.data[1:] - - // parse the key - x, n := iter.decodeVarint(iter.data) - if n == 0 { - iter.err = io.ErrShortBuffer - return false - } - k := n + int(x) - iter.record.Key = iter.data[n:k] - iter.data = iter.data[k:] - - // parse the data - if recordType == WriteBatchRecordTypeValue || recordType == WriteBatchRecordTypeMerge { - x, n := iter.decodeVarint(iter.data) - if n == 0 { - iter.err = io.ErrShortBuffer - return false + iter.record.Type = iter.decodeRecType() + + switch iter.record.Type { + case WriteBatchRecordTypeDeletion, WriteBatchRecordTypeSingleDeletion, + WriteBatchRecordTypeBeginPrepareXID, WriteBatchRecordTypeCommitXID, + WriteBatchRecordTypeRollbackXID: + iter.record.Key = iter.decodeSlice() + case WriteBatchRecordTypeValue, WriteBatchRecordTypeMerge: + iter.record.Key = iter.decodeSlice() + if iter.err == nil { + iter.record.Value = iter.decodeSlice() } - k := n + int(x) - iter.record.Value = iter.data[n:k] - iter.data = iter.data[k:] + case WriteBatchRecordTypeCFDeletion, WriteBatchRecordTypeCFValue, + WriteBatchRecordTypeCFMerge, WriteBatchRecordTypeCFSingleDeletion: + iter.record.CF = int(iter.decodeVarint()) + if iter.err == nil { + iter.record.Key = iter.decodeSlice() + } + if iter.err == nil { + iter.record.Value = iter.decodeSlice() + } + case WriteBatchRecordTypeEndPrepareXID, WriteBatchRecordTypeNoop, + WriteBatchRecordTypeNotUsed: + default: + iter.err = errors.New("unsupported wal record type") } - return true + + return iter.err == nil + } // Record returns the current record. @@ -171,19 +189,45 @@ func (iter *WriteBatchIterator) Error() error { return iter.err } -func (iter *WriteBatchIterator) decodeVarint(buf []byte) (x uint64, n int) { - // x, n already 0 - for shift := uint(0); shift < 64; shift += 7 { - if n >= len(buf) { - return 0, 0 - } - b := uint64(buf[n]) +func (iter *WriteBatchIterator) decodeSlice() []byte { + l := int(iter.decodeVarint()) + if l > len(iter.data) { + iter.err = io.ErrShortBuffer + } + if iter.err != nil { + return []byte{} + } + ret := iter.data[:l] + iter.data = iter.data[l:] + return ret +} + +func (iter *WriteBatchIterator) decodeRecType() WriteBatchRecordType { + if len(iter.data) == 0 { + iter.err = io.ErrShortBuffer + return WriteBatchRecordTypeNotUsed + } + t := iter.data[0] + iter.data = iter.data[1:] + return WriteBatchRecordType(t) +} + +func (iter *WriteBatchIterator) decodeVarint() uint64 { + var n int + var x uint64 + for shift := uint(0); shift < 64 && n < len(iter.data); shift += 7 { + b := uint64(iter.data[n]) n++ x |= (b & 0x7F) << shift if (b & 0x80) == 0 { - return x, n + iter.data = iter.data[n:] + return x } } - // The number is too large to represent in a 64-bit value. - return 0, 0 + if n == len(iter.data) { + iter.err = io.ErrShortBuffer + } else { + iter.err = errors.New("malformed varint") + } + return 0 } From 0052ea713881858d03f8e801712586b8b1deca36 Mon Sep 17 00:00:00 2001 From: Victor Grishchenko Date: Mon, 25 Dec 2017 11:48:14 +0500 Subject: [PATCH 2/3] write batch record types, according to db/write_batch.cc ReadRecordFromWriteBatch() --- write_batch.go | 79 +++++++++++++++++++++++++++++++-------------- write_batch_test.go | 4 +-- 2 files changed, 57 insertions(+), 26 deletions(-) diff --git a/write_batch.go b/write_batch.go index 0a45a9b2..803c7a0a 100644 --- a/write_batch.go +++ b/write_batch.go @@ -105,21 +105,26 @@ type WriteBatchRecordType byte // Types of batch records. const ( - WriteBatchRecordTypeDeletion WriteBatchRecordType = 0x0 - WriteBatchRecordTypeValue WriteBatchRecordType = 0x1 - WriteBatchRecordTypeMerge WriteBatchRecordType = 0x2 - WriteBatchRecordTypeLogData WriteBatchRecordType = 0x3 - WriteBatchRecordTypeCFDeletion WriteBatchRecordType = 0x4 - WriteBatchRecordTypeCFValue WriteBatchRecordType = 0x5 - WriteBatchRecordTypeCFMerge WriteBatchRecordType = 0x6 - WriteBatchRecordTypeSingleDeletion WriteBatchRecordType = 0x7 - WriteBatchRecordTypeCFSingleDeletion WriteBatchRecordType = 0x8 - WriteBatchRecordTypeNoop WriteBatchRecordType = 0xD - WriteBatchRecordTypeBeginPrepareXID WriteBatchRecordType = 0x9 - WriteBatchRecordTypeEndPrepareXID WriteBatchRecordType = 0xA - WriteBatchRecordTypeCommitXID WriteBatchRecordType = 0xB - WriteBatchRecordTypeRollbackXID WriteBatchRecordType = 0xC - WriteBatchRecordTypeNotUsed WriteBatchRecordType = 0x7F + WriteBatchDeletionRecord WriteBatchRecordType = 0x0 + WriteBatchValueRecord WriteBatchRecordType = 0x1 + WriteBatchMergeRecord WriteBatchRecordType = 0x2 + WriteBatchLogDataRecord WriteBatchRecordType = 0x3 + WriteBatchCFDeletionRecord WriteBatchRecordType = 0x4 + WriteBatchCFValueRecord WriteBatchRecordType = 0x5 + WriteBatchCFMergeRecord WriteBatchRecordType = 0x6 + WriteBatchSingleDeletionRecord WriteBatchRecordType = 0x7 + WriteBatchCFSingleDeletionRecord WriteBatchRecordType = 0x8 + WriteBatchBeginPrepareXIDRecord WriteBatchRecordType = 0x9 + WriteBatchEndPrepareXIDRecord WriteBatchRecordType = 0xA + WriteBatchCommitXIDRecord WriteBatchRecordType = 0xB + WriteBatchRollbackXIDRecord WriteBatchRecordType = 0xC + WriteBatchNoopRecord WriteBatchRecordType = 0xD + WriteBatchRangeDeletion WriteBatchRecordType = 0xF + WriteBatchCFRangeDeletion WriteBatchRecordType = 0xE + WriteBatchCFBlobIndex WriteBatchRecordType = 0x10 + WriteBatchBlobIndex WriteBatchRecordType = 0x11 + WriteBatchBeginPersistedPrepareXIDRecord WriteBatchRecordType = 0x12 + WriteBatchNotUsedRecord WriteBatchRecordType = 0x7F ) // WriteBatchRecord represents a record inside a WriteBatch. @@ -127,6 +132,7 @@ type WriteBatchRecord struct { CF int Key []byte Value []byte + Blob []byte Type WriteBatchRecordType } @@ -144,24 +150,40 @@ func (iter *WriteBatchIterator) Next() bool { return false } // reset the current record + iter.record.CF = 0 iter.record.Key = nil iter.record.Value = nil + iter.record.Blob = nil // parse the record type iter.record.Type = iter.decodeRecType() switch iter.record.Type { - case WriteBatchRecordTypeDeletion, WriteBatchRecordTypeSingleDeletion, - WriteBatchRecordTypeBeginPrepareXID, WriteBatchRecordTypeCommitXID, - WriteBatchRecordTypeRollbackXID: + case + WriteBatchDeletionRecord, + WriteBatchSingleDeletionRecord: iter.record.Key = iter.decodeSlice() - case WriteBatchRecordTypeValue, WriteBatchRecordTypeMerge: + case + WriteBatchCFDeletionRecord, + WriteBatchCFSingleDeletionRecord: + iter.record.CF = int(iter.decodeVarint()) + if iter.err == nil { + iter.record.Key = iter.decodeSlice() + } + case + WriteBatchValueRecord, + WriteBatchMergeRecord, + WriteBatchRangeDeletion, + WriteBatchBlobIndex: iter.record.Key = iter.decodeSlice() if iter.err == nil { iter.record.Value = iter.decodeSlice() } - case WriteBatchRecordTypeCFDeletion, WriteBatchRecordTypeCFValue, - WriteBatchRecordTypeCFMerge, WriteBatchRecordTypeCFSingleDeletion: + case + WriteBatchCFValueRecord, + WriteBatchCFRangeDeletion, + WriteBatchCFMergeRecord, + WriteBatchCFBlobIndex: iter.record.CF = int(iter.decodeVarint()) if iter.err == nil { iter.record.Key = iter.decodeSlice() @@ -169,8 +191,17 @@ func (iter *WriteBatchIterator) Next() bool { if iter.err == nil { iter.record.Value = iter.decodeSlice() } - case WriteBatchRecordTypeEndPrepareXID, WriteBatchRecordTypeNoop, - WriteBatchRecordTypeNotUsed: + case WriteBatchLogDataRecord: + iter.record.Blob = iter.decodeSlice() + case + WriteBatchNoopRecord, + WriteBatchBeginPrepareXIDRecord, + WriteBatchBeginPersistedPrepareXIDRecord: + case + WriteBatchEndPrepareXIDRecord, + WriteBatchCommitXIDRecord, + WriteBatchRollbackXIDRecord: + iter.record.Blob = iter.decodeSlice() default: iter.err = errors.New("unsupported wal record type") } @@ -205,7 +236,7 @@ func (iter *WriteBatchIterator) decodeSlice() []byte { func (iter *WriteBatchIterator) decodeRecType() WriteBatchRecordType { if len(iter.data) == 0 { iter.err = io.ErrShortBuffer - return WriteBatchRecordTypeNotUsed + return WriteBatchNotUsedRecord } t := iter.data[0] iter.data = iter.data[1:] diff --git a/write_batch_test.go b/write_batch_test.go index 8913d3cc..d6000e57 100644 --- a/write_batch_test.go +++ b/write_batch_test.go @@ -61,13 +61,13 @@ func TestWriteBatchIterator(t *testing.T) { iter := wb.NewIterator() ensure.True(t, iter.Next()) record := iter.Record() - ensure.DeepEqual(t, record.Type, WriteBatchRecordTypeValue) + ensure.DeepEqual(t, record.Type, WriteBatchValueRecord) ensure.DeepEqual(t, record.Key, givenKey1) ensure.DeepEqual(t, record.Value, givenVal1) ensure.True(t, iter.Next()) record = iter.Record() - ensure.DeepEqual(t, record.Type, WriteBatchRecordTypeDeletion) + ensure.DeepEqual(t, record.Type, WriteBatchDeletionRecord) ensure.DeepEqual(t, record.Key, givenKey2) // there shouldn't be any left From 4ab001a8f8655a0c190bd33d8fbb3fc321b6e346 Mon Sep 17 00:00:00 2001 From: Victor Grishchenko Date: Mon, 25 Dec 2017 12:02:47 +0500 Subject: [PATCH 3/3] rm unnecessary value/blob/xid distinction --- write_batch.go | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/write_batch.go b/write_batch.go index 803c7a0a..88b5aac8 100644 --- a/write_batch.go +++ b/write_batch.go @@ -132,7 +132,6 @@ type WriteBatchRecord struct { CF int Key []byte Value []byte - Blob []byte Type WriteBatchRecordType } @@ -153,7 +152,6 @@ func (iter *WriteBatchIterator) Next() bool { iter.record.CF = 0 iter.record.Key = nil iter.record.Value = nil - iter.record.Blob = nil // parse the record type iter.record.Type = iter.decodeRecType() @@ -192,7 +190,7 @@ func (iter *WriteBatchIterator) Next() bool { iter.record.Value = iter.decodeSlice() } case WriteBatchLogDataRecord: - iter.record.Blob = iter.decodeSlice() + iter.record.Value = iter.decodeSlice() case WriteBatchNoopRecord, WriteBatchBeginPrepareXIDRecord, @@ -201,7 +199,7 @@ func (iter *WriteBatchIterator) Next() bool { WriteBatchEndPrepareXIDRecord, WriteBatchCommitXIDRecord, WriteBatchRollbackXIDRecord: - iter.record.Blob = iter.decodeSlice() + iter.record.Value = iter.decodeSlice() default: iter.err = errors.New("unsupported wal record type") }