Managed Transactions

- Allow a way to specify the read and commit timestamp for a transaction. This would be useful for Dgraph.
- Add tests for managed transactions.

Various cleanup:

- Don't return an error while creating a transaction.
- Make Entry struct private, now that BatchSet is no longer public.
- Put all errors in one var block.
This commit is contained in:
Manish R Jain
2017-10-02 18:32:11 +11:00
parent e0ad839873
commit c849ea647b
9 changed files with 241 additions and 180 deletions
+36 -36
View File
@@ -22,43 +22,31 @@ import (
"github.com/pkg/errors"
)
// ErrInvalidDir is returned when Badger cannot find the directory
// from where it is supposed to load the key-value store.
var ErrInvalidDir = errors.New("Invalid Dir, directory does not exist")
// ErrValueLogSize is returned when opt.ValueLogFileSize option is not within the valid
// range.
var ErrValueLogSize = errors.New("Invalid ValueLogFileSize, must be between 1MB and 2GB")
// ErrKeyNotFound is returned when key isn't found on a txn.Get.
var ErrKeyNotFound = errors.New("Key not found")
// ErrTxnTooBig is returned if too many writes are fit into a single transaction.
var ErrTxnTooBig = errors.New("Txn is too big to fit into one request.")
// ErrConflict is returned when a transaction conflicts with another transaction. This can happen if
// the read rows had been updated concurrently by another transaction.
var ErrConflict = errors.New("Transaction Conflict. Please retry.")
// ErrReadOnlyTxn is returned if an update function is called on a read-only transaction.
var ErrReadOnlyTxn = errors.New("No sets or deletes are allowed in a read-only transaction.")
// ErrEmptyKey is returned if an empty key is passed on an update function.
var ErrEmptyKey = errors.New("Key cannot be empty.")
const maxKeySize = 1 << 20
func exceedsMaxKeySizeError(key []byte) error {
return errors.Errorf("Key with size %d exceeded %dMB limit. Key:\n%s",
len(key), maxKeySize<<20, hex.Dump(key[:1<<10]))
}
func exceedsMaxValueSizeError(value []byte, maxValueSize int64) error {
return errors.Errorf("Value with size %d exceeded ValueLogFileSize (%dMB). Key:\n%s",
len(value), maxValueSize<<20, hex.Dump(value[:1<<10]))
}
var (
// ErrInvalidDir is returned when Badger cannot find the directory
// from where it is supposed to load the key-value store.
ErrInvalidDir = errors.New("Invalid Dir, directory does not exist")
// ErrValueLogSize is returned when opt.ValueLogFileSize option is not within the valid
// range.
ErrValueLogSize = errors.New("Invalid ValueLogFileSize, must be between 1MB and 2GB")
// ErrKeyNotFound is returned when key isn't found on a txn.Get.
ErrKeyNotFound = errors.New("Key not found")
// ErrTxnTooBig is returned if too many writes are fit into a single transaction.
ErrTxnTooBig = errors.New("Txn is too big to fit into one request.")
// ErrConflict is returned when a transaction conflicts with another transaction. This can happen if
// the read rows had been updated concurrently by another transaction.
ErrConflict = errors.New("Transaction Conflict. Please retry.")
// ErrReadOnlyTxn is returned if an update function is called on a read-only transaction.
ErrReadOnlyTxn = errors.New("No sets or deletes are allowed in a read-only transaction.")
// ErrEmptyKey is returned if an empty key is passed on an update function.
ErrEmptyKey = errors.New("Key cannot be empty.")
// ErrRetry is returned when a log file containing the value is not found.
// This usually indicates that it may have been garbage collected, and the
// operation needs to be retried.
@@ -80,3 +68,15 @@ var (
// ErrInvalidRequest is returned if the user request is invalid.
ErrInvalidRequest = errors.New("Invalid request")
)
const maxKeySize = 1 << 20
func exceedsMaxKeySizeError(key []byte) error {
return errors.Errorf("Key with size %d exceeded %dMB limit. Key:\n%s",
len(key), maxKeySize<<20, hex.Dump(key[:1<<10]))
}
func exceedsMaxValueSizeError(value []byte, maxValueSize int64) error {
return errors.Errorf("Value with size %d exceeded ValueLogFileSize (%dMB). Key:\n%s",
len(value), maxValueSize<<20, hex.Dump(value[:1<<10]))
}
+5 -5
View File
@@ -217,7 +217,7 @@ func NewKV(optParam *Options) (out *KV, err error) {
}
first := true
fn := func(e Entry, vp valuePointer) error { // Function for replaying.
fn := func(e entry, vp valuePointer) error { // Function for replaying.
if first {
out.elog.Printf("First key=%s\n", e.Key)
}
@@ -503,7 +503,7 @@ var requestPool = sync.Pool{
},
}
func (s *KV) shouldWriteValueToLSM(e Entry) bool {
func (s *KV) shouldWriteValueToLSM(e entry) bool {
return len(e.Value) < s.opt.ValueThreshold
}
@@ -646,7 +646,7 @@ func (s *KV) doWrites(lc *y.Closer) {
}
}
func (s *KV) sendToWriteCh(entries []*Entry) (*request, error) {
func (s *KV) sendToWriteCh(entries []*entry) (*request, error) {
var count, size int64
for _, e := range entries {
size += int64(s.opt.estimateSize(e))
@@ -671,7 +671,7 @@ func (s *KV) sendToWriteCh(entries []*Entry) (*request, error) {
// batchSet applies a list of badger.Entry. If a request level error occurs it
// will be returned.
// Check(kv.BatchSet(entries))
func (s *KV) batchSet(entries []*Entry) error {
func (s *KV) batchSet(entries []*entry) error {
req, err := s.sendToWriteCh(entries)
if err != nil {
return err
@@ -690,7 +690,7 @@ func (s *KV) batchSet(entries []*Entry) error {
// err := kv.BatchSetAsync(entries, func(err error)) {
// Check(err)
// }
func (s *KV) batchSetAsync(entries []*Entry, f func(error)) error {
func (s *KV) batchSetAsync(entries []*entry, f func(error)) error {
req, err := s.sendToWriteCh(entries)
if err != nil {
return err
+27 -49
View File
@@ -59,22 +59,19 @@ func getItemValue(t *testing.T, item *KVItem) (val []byte) {
}
func txnSet(t *testing.T, kv *KV, key []byte, val []byte, meta byte) {
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
require.NoError(t, txn.Set(key, val, meta))
require.NoError(t, txn.Commit(nil))
}
func txnDelete(t *testing.T, kv *KV, key []byte) {
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
require.NoError(t, txn.Delete(key))
require.NoError(t, txn.Commit(nil))
}
func txnGet(t *testing.T, kv *KV, key []byte) (KVItem, error) {
txn, err := kv.NewTransaction(false)
require.NoError(t, err)
txn := kv.NewTransaction(false)
return txn.Get(key)
}
@@ -121,7 +118,7 @@ func TestConcurrentWrite(t *testing.T) {
opt.PrefetchSize = 10
opt.PrefetchValues = true
txn, err := kv.NewTransaction(true)
txn := kv.NewTransaction(true)
it := txn.NewIterator(opt)
defer it.Close()
var i, j int
@@ -242,10 +239,8 @@ func TestGetMore(t *testing.T) {
// n := 500000
n := 10000
m := 49 // Increasing would cause ErrTxnTooBig
fmt.Println("writing")
for i := 0; i < n; i += m {
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
for j := i; j < i+m && j < n; j++ {
require.NoError(t, txn.Set(data(j), data(j), 0))
}
@@ -253,7 +248,6 @@ func TestGetMore(t *testing.T) {
}
require.NoError(t, kv.validate())
fmt.Println("retrieving")
for i := 0; i < n; i++ {
item, err := txnGet(t, kv, data(i))
if err != nil {
@@ -263,10 +257,8 @@ func TestGetMore(t *testing.T) {
}
// Overwrite
fmt.Println("overwriting")
for i := 0; i < n; i += m {
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
for j := i; j < i+m && j < n; j++ {
require.NoError(t, txn.Set(data(j),
// Use a long value that will certainly exceed value threshold.
@@ -277,7 +269,6 @@ func TestGetMore(t *testing.T) {
}
require.NoError(t, kv.validate())
fmt.Println("testing")
for i := 0; i < n; i++ {
k := []byte(fmt.Sprintf("%09d", i))
expectedValue := fmt.Sprintf("zzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzzz%09d", i)
@@ -294,8 +285,7 @@ func TestGetMore(t *testing.T) {
fmt.Printf("wanted=%q Item: %s\n", k, item.ToString())
fmt.Printf("on re-run, got version: %+v\n", vs)
txn, err := kv.NewTransaction(false)
require.NoError(t, err)
txn := kv.NewTransaction(false)
itr := txn.NewIterator(DefaultIteratorOptions)
for itr.Seek(k0); itr.Valid(); itr.Next() {
item := itr.Item()
@@ -314,8 +304,7 @@ func TestGetMore(t *testing.T) {
if (i % 10000) == 0 {
fmt.Printf("Deleting i=%d\n", i)
}
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
for j := i; j < i+m && j < n; j++ {
require.NoError(t, txn.Delete([]byte(fmt.Sprintf("%09d", j))))
}
@@ -331,7 +320,6 @@ func TestGetMore(t *testing.T) {
item, err := txnGet(t, kv, []byte(k))
require.Equal(t, ErrKeyNotFound, err, "wanted=%q item=%s\n", k, item.ToString())
}
fmt.Println("Done and closing")
}
// Put a lot of data to move some data to disk.
@@ -352,10 +340,9 @@ func TestExistsMore(t *testing.T) {
m := 49
for i := 0; i < n; i += m {
if (i % 1000) == 0 {
fmt.Printf("Putting i=%d\n", i)
t.Logf("Putting i=%d\n", i)
}
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
for j := i; j < i+m && j < n; j++ {
require.NoError(t, txn.Set([]byte(fmt.Sprintf("%09d", j)),
[]byte(fmt.Sprintf("%09d", j)),
@@ -381,8 +368,7 @@ func TestExistsMore(t *testing.T) {
if (i % 1000) == 0 {
fmt.Printf("Deleting i=%d\n", i)
}
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
for j := i; j < i+m && j < n; j++ {
require.NoError(t, txn.Delete([]byte(fmt.Sprintf("%09d", j))))
}
@@ -428,8 +414,7 @@ func TestIterate2Basic(t *testing.T) {
opt.PrefetchValues = true
opt.PrefetchSize = 10
txn, err := kv.NewTransaction(false)
require.NoError(t, err)
txn := kv.NewTransaction(false)
it := txn.NewIterator(opt)
{
var count int
@@ -535,14 +520,12 @@ func TestIterateDeleted(t *testing.T) {
iterOpt := DefaultIteratorOptions
iterOpt.PrefetchValues = false
txn, err := ps.NewTransaction(false)
require.NoError(t, err)
txn := ps.NewTransaction(false)
idxIt := txn.NewIterator(iterOpt)
defer idxIt.Close()
count := 0
txn2, err := ps.NewTransaction(true)
require.NoError(t, err)
txn2 := ps.NewTransaction(true)
prefix := []byte("Key")
for idxIt.Seek(prefix); idxIt.Valid(); idxIt.Next() {
key := idxIt.Item().Key()
@@ -559,8 +542,7 @@ func TestIterateDeleted(t *testing.T) {
for _, prefetch := range [...]bool{true, false} {
t.Run(fmt.Sprintf("Prefetch=%t", prefetch), func(t *testing.T) {
txn, err := ps.NewTransaction(false)
require.NoError(t, err)
txn := ps.NewTransaction(false)
iterOpt = DefaultIteratorOptions
iterOpt.PrefetchValues = prefetch
idxIt = txn.NewIterator(iterOpt)
@@ -646,13 +628,12 @@ func TestBigKeyValuePairs(t *testing.T) {
bigV := make([]byte, opt.ValueLogFileSize+1)
small := make([]byte, 10)
txn, err := kv.NewTransaction(true)
txn := kv.NewTransaction(true)
require.Regexp(t, regexp.MustCompile("Key.*exceeded"), txn.Set(bigK, small, 0))
txn, err = kv.NewTransaction(true)
txn = kv.NewTransaction(true)
require.Regexp(t, regexp.MustCompile("Value.*exceeded"), txn.Set(small, bigV, 0))
txn, err = kv.NewTransaction(true)
require.NoError(t, err)
txn = kv.NewTransaction(true)
require.NoError(t, txn.Set(small, small, 0x00))
require.Regexp(t, regexp.MustCompile("Key.*exceeded"), txn.Set(bigK, bigV, 0x00))
@@ -678,9 +659,9 @@ func TestIteratorPrefetchSize(t *testing.T) {
n := 100
for i := 0; i < n; i++ {
if (i % 10) == 0 {
t.Logf("Put i=%d\n", i)
}
// if (i % 10) == 0 {
// t.Logf("Put i=%d\n", i)
// }
txnSet(t, kv, bkey(i), bval(i), byte(i%127))
}
@@ -690,8 +671,7 @@ func TestIteratorPrefetchSize(t *testing.T) {
opt.PrefetchSize = prefetchSize
var count int
txn, err := kv.NewTransaction(false)
require.NoError(t, err)
txn := kv.NewTransaction(false)
it := txn.NewIterator(opt)
{
t.Log("Starting first basic iteration")
@@ -724,11 +704,10 @@ func TestSetIfAbsentAsync(t *testing.T) {
n := 1000
for i := 0; i < n; i++ {
if (i % 10) == 0 {
t.Logf("Put i=%d\n", i)
}
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
// if (i % 10) == 0 {
// t.Logf("Put i=%d\n", i)
// }
txn := kv.NewTransaction(true)
_, err = txn.Get(bkey(i))
require.Equal(t, ErrKeyNotFound, err)
require.NoError(t, txn.Set(bkey(i), nil, byte(i%127)))
@@ -740,8 +719,7 @@ func TestSetIfAbsentAsync(t *testing.T) {
require.NoError(t, err)
opt := DefaultIteratorOptions
txn, err := kv.NewTransaction(false)
require.NoError(t, err)
txn := kv.NewTransaction(false)
var count int
it := txn.NewIterator(opt)
{
+2 -2
View File
@@ -99,6 +99,6 @@ var DefaultOptions = Options{
ValueThreshold: 20,
}
func (opt *Options) estimateSize(entry *Entry) int {
return entry.estimateSize(opt.ValueThreshold)
func (opt *Options) estimateSize(e *entry) int {
return e.estimateSize(opt.ValueThreshold)
}
+5 -5
View File
@@ -73,10 +73,10 @@ func (h *header) Decode(buf []byte) {
h.userMeta = buf[9]
}
// Entry provides Key, Value and if required, CASCounterCheck to kv.BatchSet() API.
// entry provides Key, Value and if required, CASCounterCheck to kv.BatchSet() API.
// If CASCounterCheck is provided, it would be compared against the current casCounter
// assigned to this key-value. Set be done on this key only if the counters match.
type Entry struct {
type entry struct {
Key []byte
Value []byte
Meta byte
@@ -86,7 +86,7 @@ type Entry struct {
offset uint32
}
func (e *Entry) estimateSize(threshold int) int {
func (e *entry) estimateSize(threshold int) int {
if len(e.Value) < threshold {
return len(e.Key) + len(e.Value) + 2 // Meta, UserMeta
}
@@ -94,7 +94,7 @@ func (e *Entry) estimateSize(threshold int) int {
}
// Encodes e to buf. Returns number of bytes written.
func encodeEntry(e *Entry, buf *bytes.Buffer) (int, error) {
func encodeEntry(e *entry, buf *bytes.Buffer) (int, error) {
var h header
h.klen = uint32(len(e.Key))
h.vlen = uint32(len(e.Value))
@@ -122,7 +122,7 @@ func encodeEntry(e *Entry, buf *bytes.Buffer) (int, error) {
return len(headerEnc) + len(e.Key) + len(e.Value) + len(crcBuf), nil
}
func (e Entry) print(prefix string) {
func (e entry) print(prefix string) {
fmt.Printf("%s Key: %s Meta: %d UserMeta: %d Offset: %d len(val)=%d",
prefix, e.Key, e.Meta, e.UserMeta, e.offset, len(e.Value))
}
+44 -17
View File
@@ -87,7 +87,21 @@ func (gs *globalTxnState) newCommitTs(txn *Txn) uint64 {
if gs.hasConflict(txn) {
return 0
}
ts := gs.nextCommit
var ts uint64
if txn.commitTs == 0 {
// This is the general case, when user doesn't specify the read and commit ts.
ts = gs.nextCommit
gs.nextCommit++
} else {
// If commitTs is set, use it instead.
ts = txn.commitTs
if gs.nextCommit <= ts { // Update this to max+1 commit ts, so replay works.
gs.nextCommit = ts + 1
}
}
for _, w := range txn.writes {
gs.commits[w] = ts // Update the commitTs.
}
@@ -96,8 +110,6 @@ func (gs *globalTxnState) newCommitTs(txn *Txn) uint64 {
panic(fmt.Sprintf("We shouldn't have the commit ts: %d", ts))
}
gs.pendingCommits[ts] = struct{}{}
gs.nextCommit++
return ts
}
@@ -128,13 +140,14 @@ func (gs *globalTxnState) doneCommit(cts uint64) {
}
type Txn struct {
readTs uint64
readTs uint64
commitTs uint64
update bool // update is used to conditionally keep track of reads.
reads []uint64 // contains fingerprints of keys read.
writes []uint64 // contains fingerprints of keys written.
pendingWrites map[string]*Entry // cache stores any writes done by txn.
pendingWrites map[string]*entry // cache stores any writes done by txn.
gs *globalTxnState
kv *KV
@@ -161,7 +174,7 @@ func (txn *Txn) Set(key, val []byte, userMeta byte) error {
fp := farm.Fingerprint64(key) // Avoid dealing with byte arrays.
txn.writes = append(txn.writes, fp)
e := &Entry{
e := &entry{
Key: key,
Value: val,
UserMeta: userMeta,
@@ -185,7 +198,7 @@ func (txn *Txn) Delete(key []byte) error {
fp := farm.Fingerprint64(key) // Avoid dealing with byte arrays.
txn.writes = append(txn.writes, fp)
e := &Entry{
e := &entry{
Key: key,
Meta: BitDelete,
}
@@ -257,7 +270,7 @@ func (txn *Txn) Commit(callback func(error)) error {
}
defer txn.gs.doneCommit(commitTs)
entries := make([]*Entry, 0, len(txn.pendingWrites)+1)
entries := make([]*entry, 0, len(txn.pendingWrites)+1)
for _, e := range txn.pendingWrites {
// Suffix the keys with commit ts, so the key versions are sorted in
// descending order of commit timestamp.
@@ -265,12 +278,12 @@ func (txn *Txn) Commit(callback func(error)) error {
e.Meta |= BitTxn
entries = append(entries, e)
}
entry := &Entry{
e := &entry{
Key: y.KeyWithTs(txnKey, commitTs),
Value: []byte(strconv.FormatUint(commitTs, 10)),
Meta: BitFinTxn,
}
entries = append(entries, entry)
entries = append(entries, e)
if callback == nil {
// If batchSet failed, LSM would not have been updated. So, no need to rollback anything.
@@ -282,6 +295,11 @@ func (txn *Txn) Commit(callback func(error)) error {
return txn.kv.batchSetAsync(entries, callback)
}
func (txn *Txn) CommitAt(commitTs uint64, callback func(error)) error {
txn.commitTs = commitTs
return txn.Commit(callback)
}
// NewIterator returns a new iterator. Depending upon the options, either only keys, or both
// key-value pairs would be fetched. The keys are returned in lexicographically sorted order.
// Usage:
@@ -320,14 +338,23 @@ func (txn *Txn) NewIterator(opt IteratorOptions) *Iterator {
}
return res
}
func (kv *KV) NewTransaction(update bool) (*Txn, error) {
func (kv *KV) NewTransaction(update bool) *Txn {
txn := &Txn{
update: update,
gs: kv.txnState,
kv: kv,
readTs: kv.txnState.readTs(),
pendingWrites: make(map[string]*Entry),
update: update,
gs: kv.txnState,
kv: kv,
readTs: kv.txnState.readTs(),
}
if update {
txn.pendingWrites = make(map[string]*entry)
}
return txn, nil
return txn
}
func (kv *KV) NewTransactionAt(readTs uint64, update bool) *Txn {
txn := kv.NewTransaction(update)
txn.readTs = readTs
return txn
}
+89 -22
View File
@@ -34,8 +34,7 @@ func TestTxnSimple(t *testing.T) {
require.NoError(t, err)
defer kv.Close()
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
for i := 0; i < 10; i++ {
k := []byte(fmt.Sprintf("key=%d", i))
@@ -63,8 +62,7 @@ func TestTxnVersions(t *testing.T) {
k := []byte("key")
for i := 1; i < 10; i++ {
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
txn.Set(k, []byte(fmt.Sprintf("valversion=%d", i)), 0)
require.NoError(t, txn.Commit(nil))
@@ -89,7 +87,7 @@ func TestTxnVersions(t *testing.T) {
}
for i := 1; i < 10; i++ {
txn, err := kv.NewTransaction(true)
txn := kv.NewTransaction(true)
require.NoError(t, err)
txn.readTs = uint64(i) // Read version at i.
@@ -114,7 +112,7 @@ func TestTxnVersions(t *testing.T) {
itr = txn.NewIterator(opt)
checkIterator(itr, i)
}
txn, err := kv.NewTransaction(true)
txn := kv.NewTransaction(true)
require.NoError(t, err)
item, err := txn.Get(k)
require.NoError(t, err)
@@ -138,8 +136,7 @@ func TestTxnWriteSkew(t *testing.T) {
ay := []byte("y")
// Set balance to $100 in each account.
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
val := []byte(strconv.Itoa(100))
txn.Set(ax, val, 0)
txn.Set(ay, val, 0)
@@ -160,8 +157,7 @@ func TestTxnWriteSkew(t *testing.T) {
}
// Start two transactions, each would read both accounts and deduct from one account.
txn1, err := kv.NewTransaction(true)
require.NoError(t, err)
txn1 := kv.NewTransaction(true)
sum := getBal(txn1, ax)
sum += getBal(txn1, ay)
@@ -175,8 +171,7 @@ func TestTxnWriteSkew(t *testing.T) {
require.Equal(t, 100, sum)
// Don't commit yet.
txn2, err := kv.NewTransaction(true)
require.NoError(t, err)
txn2 := kv.NewTransaction(true)
sum = getBal(txn2, ax)
sum += getBal(txn2, ay)
@@ -214,31 +209,27 @@ func TestTxnIterationEdgeCase(t *testing.T) {
kc := []byte("c")
// c1
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
txn.Set(kc, []byte("c1"), 0)
require.NoError(t, txn.Commit(nil))
require.Equal(t, uint64(1), kv.txnState.readTs())
// a2, c2
txn, err = kv.NewTransaction(true)
require.NoError(t, err)
txn = kv.NewTransaction(true)
txn.Set(ka, []byte("a2"), 0)
txn.Set(kc, []byte("c2"), 0)
require.NoError(t, txn.Commit(nil))
require.Equal(t, uint64(2), kv.txnState.readTs())
// b3
txn, err = kv.NewTransaction(true)
require.NoError(t, err)
txn = kv.NewTransaction(true)
txn.Set(ka, []byte("a3"), 0)
txn.Set(kb, []byte("b3"), 0)
require.NoError(t, txn.Commit(nil))
require.Equal(t, uint64(3), kv.txnState.readTs())
// b4 (del)
txn, err = kv.NewTransaction(true)
require.NoError(t, err)
txn = kv.NewTransaction(true)
txn.Delete(kb)
require.NoError(t, txn.Commit(nil))
require.Equal(t, uint64(4), kv.txnState.readTs())
@@ -256,8 +247,7 @@ func TestTxnIterationEdgeCase(t *testing.T) {
}
require.Equal(t, len(expected), i)
}
txn, err = kv.NewTransaction(true)
require.NoError(t, err)
txn = kv.NewTransaction(true)
itr := txn.NewIterator(DefaultIteratorOptions)
checkIterator(itr, []string{"a3", "c2"})
@@ -284,3 +274,80 @@ func TestTxnIterationEdgeCase(t *testing.T) {
itr = txn.NewIterator(rev)
checkIterator(itr, []string{"c1"})
}
func TestTxnManaged(t *testing.T) {
dir, err := ioutil.TempDir("", "badger")
require.NoError(t, err)
defer os.RemoveAll(dir)
kv, err := NewKV(getTestOptions(dir))
require.NoError(t, err)
defer kv.Close()
key := func(i int) []byte {
return []byte(fmt.Sprintf("key-%02d", i))
}
val := func(i int) []byte {
return []byte(fmt.Sprintf("val-%d", i))
}
// Write data at t=3.
txn := kv.NewTransactionAt(3, true)
for i := 0; i <= 3; i++ {
require.NoError(t, txn.Set(key(i), val(i), 0))
}
require.NoError(t, txn.CommitAt(3, nil))
// Read data at t=2.
txn = kv.NewTransactionAt(2, false)
for i := 0; i <= 3; i++ {
_, err := txn.Get(key(i))
require.Equal(t, ErrKeyNotFound, err)
}
// Read data at t=3.
txn = kv.NewTransactionAt(3, false)
for i := 0; i <= 3; i++ {
item, err := txn.Get(key(i))
require.NoError(t, err)
require.Equal(t, uint64(3), item.Version())
require.NoError(t, item.Value(func(v []byte) error {
require.Equal(t, val(i), v)
return nil
}))
}
// Write data at t=7.
txn = kv.NewTransactionAt(6, true)
for i := 0; i <= 7; i++ {
_, err := txn.Get(key(i))
if err == nil {
continue // Don't overwrite existing keys.
}
require.NoError(t, txn.Set(key(i), val(i), 0))
}
require.NoError(t, txn.CommitAt(7, nil))
// Read data at t=9.
txn = kv.NewTransactionAt(9, false)
for i := 0; i <= 9; i++ {
item, err := txn.Get(key(i))
if i <= 7 {
require.NoError(t, err)
} else {
require.Equal(t, ErrKeyNotFound, err)
}
if i <= 3 {
require.Equal(t, uint64(3), item.Version())
} else if i <= 7 {
require.Equal(t, uint64(7), item.Version())
}
if i <= 7 {
require.NoError(t, item.Value(func(v []byte) error {
require.Equal(t, val(i), v)
return nil
}))
}
}
}
+10 -10
View File
@@ -152,7 +152,7 @@ func (lf *logFile) sync() error {
var errStop = errors.New("Stop iteration")
type logEntry func(e Entry, vp valuePointer) error
type logEntry func(e entry, vp valuePointer) error
// iterate iterates over log file. It doesn't not allocate new memory for every kv pair.
// Therefore, the kv pair is only valid for the duration of fn call.
@@ -185,7 +185,7 @@ func (vlog *valueLog) iterate(lf *logFile, offset uint32, fn logEntry) error {
return err
}
var e Entry
var e entry
e.offset = recordOffset
h.Decode(hbuf[:])
if h.klen > maxKeySize {
@@ -269,12 +269,12 @@ func (vlog *valueLog) rewrite(f *logFile) error {
defer elog.Finish()
elog.Printf("Rewriting fid: %d", f.fid)
wb := make([]*Entry, 0, 1000)
wb := make([]*entry, 0, 1000)
var size int64
y.AssertTrue(vlog.kv != nil)
var count int
fe := func(e Entry) error {
fe := func(e entry) error {
count++
if count%10000 == 0 {
elog.Printf("Processing entry %d", count)
@@ -303,7 +303,7 @@ func (vlog *valueLog) rewrite(f *logFile) error {
}
if vp.Fid == f.fid && vp.Offset == e.offset {
// This new entry only contains the key, and a pointer to the value.
ne := new(Entry)
ne := new(entry)
ne.Meta = 0 // Remove all bits.
ne.UserMeta = e.UserMeta
ne.Key = make([]byte, len(e.Key))
@@ -326,7 +326,7 @@ func (vlog *valueLog) rewrite(f *logFile) error {
return nil
}
err := vlog.iterate(f, 0, func(e Entry, vp valuePointer) error {
err := vlog.iterate(f, 0, func(e entry, vp valuePointer) error {
return fe(e)
})
if err != nil {
@@ -627,7 +627,7 @@ func (vlog *valueLog) Replay(ptr valuePointer, fn logEntry) error {
type request struct {
// Input values
Entries []*Entry
Entries []*entry
// Output values and wait group stuff below
Ptrs []valuePointer
Wg sync.WaitGroup
@@ -795,7 +795,7 @@ func (vlog *valueLog) readValueBytes(vp valuePointer, consumer func([]byte) erro
}
// Test helper
func valueBytesToEntry(buf []byte) (e Entry) {
func valueBytesToEntry(buf []byte) (e entry) {
var h header
h.Decode(buf)
n := uint32(headerBufSize)
@@ -823,7 +823,7 @@ func (vlog *valueLog) pickLog() *logFile {
return vlog.filesMap[fids[idx]]
}
func discardEntry(e Entry, vs y.ValueStruct) bool {
func discardEntry(e entry, vs y.ValueStruct) bool {
if vs.Version != y.ParseTs(e.Key) {
// Version not found. Discard.
return true
@@ -865,7 +865,7 @@ func (vlog *valueLog) doRunGC(gcThreshold float64) error {
start := time.Now()
y.AssertTrue(vlog.kv != nil)
err := vlog.iterate(lf, 0, func(e Entry, vp valuePointer) error {
err := vlog.iterate(lf, 0, func(e entry, vp valuePointer) error {
esz := float64(vp.Len) / (1 << 20) // in MBs. +4 for the CAS stuff.
skipped += esz
if skipped < skipFirstM {
+23 -34
View File
@@ -43,23 +43,23 @@ func TestValueBasic(t *testing.T) {
const val2 = "samplevalb012345678901234567890123"
require.True(t, len(val1) >= kv.opt.ValueThreshold)
entry := &Entry{
e := &entry{
Key: []byte("samplekey"),
Value: []byte(val1),
Meta: BitValuePointer,
}
entry2 := &Entry{
e2 := &entry{
Key: []byte("samplekeyb"),
Value: []byte(val2),
Meta: BitValuePointer,
}
b := new(request)
b.Entries = []*Entry{entry, entry2}
b.Entries = []*entry{e, e2}
log.write([]*request{b})
require.Len(t, b.Ptrs, 2)
fmt.Printf("Pointer written: %+v %+v\n", b.Ptrs[0], b.Ptrs[1])
t.Logf("Pointer written: %+v %+v\n", b.Ptrs[0], b.Ptrs[1])
var buf1, buf2 []byte
var err1, err2 error
@@ -74,8 +74,8 @@ func TestValueBasic(t *testing.T) {
require.NoError(t, err1)
require.NoError(t, err2)
readEntries := []Entry{valueBytesToEntry(buf1), valueBytesToEntry(buf2)}
require.EqualValues(t, []Entry{
readEntries := []entry{valueBytesToEntry(buf1), valueBytesToEntry(buf2)}
require.EqualValues(t, []entry{
{
Key: []byte("samplekey"),
Value: []byte(val1),
@@ -100,16 +100,14 @@ func TestValueGC(t *testing.T) {
defer kv.Close()
sz := 32 << 10
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
for i := 0; i < 100; i++ {
v := make([]byte, sz)
rand.Read(v[:rand.Intn(sz)])
require.NoError(t, txn.Set([]byte(fmt.Sprintf("key%d", i)), v, 0))
if i%20 == 0 {
require.NoError(t, txn.Commit(nil))
txn, err = kv.NewTransaction(true)
require.NoError(t, err)
txn = kv.NewTransaction(true)
}
}
require.NoError(t, txn.Commit(nil))
@@ -149,16 +147,14 @@ func TestValueGC2(t *testing.T) {
defer kv.Close()
sz := 32 << 10
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
for i := 0; i < 100; i++ {
v := make([]byte, sz)
rand.Read(v[:rand.Intn(sz)])
require.NoError(t, txn.Set([]byte(fmt.Sprintf("key%d", i)), v, 0))
if i%20 == 0 {
require.NoError(t, txn.Commit(nil))
txn, err = kv.NewTransaction(true)
require.NoError(t, err)
txn = kv.NewTransaction(true)
}
}
require.NoError(t, txn.Commit(nil))
@@ -221,8 +217,7 @@ func TestValueGC3(t *testing.T) {
valueSize := 32 << 10
var value3 []byte
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
for i := 0; i < 100; i++ {
v := make([]byte, valueSize) // 32K * 100 will take >=3'276'800 B.
if i == 3 {
@@ -233,8 +228,7 @@ func TestValueGC3(t *testing.T) {
require.NoError(t, txn.Set([]byte(fmt.Sprintf("key%03d", i)), v, 0))
if i%20 == 0 {
require.NoError(t, txn.Commit(nil))
txn, err = kv.NewTransaction(true)
require.NoError(t, err)
txn = kv.NewTransaction(true)
}
}
require.NoError(t, txn.Commit(nil))
@@ -246,8 +240,7 @@ func TestValueGC3(t *testing.T) {
Reverse: false,
}
txn, err = kv.NewTransaction(true)
require.NoError(t, err)
txn = kv.NewTransaction(true)
it := txn.NewIterator(itOpt)
defer it.Close()
// Walk a few keys
@@ -309,7 +302,7 @@ func TestChecksums(t *testing.T) {
require.True(t, len(v0) >= kv.opt.ValueThreshold)
// Use a vlog with K0=V0 and a (corrupted) second transaction(k1,k2)
buf := createVlog(t, []*Entry{
buf := createVlog(t, []*entry{
{Key: k0, Value: v0},
{Key: k1, Value: v1},
{Key: k2, Value: v2},
@@ -335,8 +328,7 @@ func TestChecksums(t *testing.T) {
// last due to checksum failure).
kv, err = NewKV(opts)
require.NoError(t, err)
txn, err := kv.NewTransaction(false)
require.NoError(t, err)
txn := kv.NewTransaction(false)
iter := txn.NewIterator(DefaultIteratorOptions)
iter.Seek(k0)
require.True(t, iter.Valid())
@@ -379,7 +371,7 @@ func TestPartialAppendToValueLog(t *testing.T) {
// Create truncated vlog to simulate a partial append.
// k0 - single transaction, k1 and k2 in another transaction
buf := createVlog(t, []*Entry{
buf := createVlog(t, []*entry{
{Key: k0, Value: v0},
{Key: k1, Value: v1},
{Key: k2, Value: v2},
@@ -419,16 +411,14 @@ func TestValueLogTrigger(t *testing.T) {
// Write a lot of data, so it creates some work for valug log GC.
sz := 32 << 10
txn, err := kv.NewTransaction(true)
require.NoError(t, err)
txn := kv.NewTransaction(true)
for i := 0; i < 100; i++ {
v := make([]byte, sz)
rand.Read(v[:rand.Intn(sz)])
require.NoError(t, txn.Set([]byte(fmt.Sprintf("key%d", i)), v, 0))
if i%20 == 0 {
require.NoError(t, txn.Commit(nil))
txn, err = kv.NewTransaction(true)
require.NoError(t, err)
txn = kv.NewTransaction(true)
}
}
require.NoError(t, txn.Commit(nil))
@@ -456,7 +446,7 @@ func TestValueLogTrigger(t *testing.T) {
require.Equal(t, ErrRejected, err, "Error should be returned after closing KV.")
}
func createVlog(t *testing.T, entries []*Entry) []byte {
func createVlog(t *testing.T, entries []*entry) []byte {
dir, err := ioutil.TempDir("", "badger")
require.NoError(t, err)
defer os.RemoveAll(dir)
@@ -467,7 +457,7 @@ func createVlog(t *testing.T, entries []*Entry) []byte {
require.NoError(t, err)
txnSet(t, kv, entries[0].Key, entries[0].Value, entries[0].Meta)
entries = entries[1:]
txn, err := kv.NewTransaction(true)
txn := kv.NewTransaction(true)
for _, entry := range entries {
require.NoError(t, txn.Set(entry.Key, entry.Value, entry.Meta))
}
@@ -482,8 +472,7 @@ func createVlog(t *testing.T, entries []*Entry) []byte {
func checkKeys(t *testing.T, kv *KV, keys [][]byte) {
i := 0
txn, err := kv.NewTransaction(false)
require.NoError(t, err)
txn := kv.NewTransaction(false)
iter := txn.NewIterator(IteratorOptions{})
for iter.Seek(keys[0]); iter.Valid(); iter.Next() {
require.Equal(t, iter.Item().Key(), keys[i])
@@ -513,11 +502,11 @@ func BenchmarkReadWrite(b *testing.B) {
b.ResetTimer()
for i := 0; i < b.N; i++ {
e := new(Entry)
e := new(entry)
e.Key = make([]byte, 16)
e.Value = make([]byte, vsz)
bl := new(request)
bl.Entries = []*Entry{e}
bl.Entries = []*entry{e}
var ptrs []valuePointer