syncthing/Godeps/_workspace/src/github.com/syndtr/goleveldb/leveldb/db_write.go

295 lines
6.2 KiB
Go
Raw Normal View History

2014-07-06 14:46:48 +02:00
// Copyright (c) 2012, Suryandaru Triandana <syndtr@gmail.com>
// All rights reserved.
//
// Use of this source code is governed by a BSD-style license that can be
// found in the LICENSE file.
package leveldb
import (
"time"
"github.com/syndtr/goleveldb/leveldb/memdb"
"github.com/syndtr/goleveldb/leveldb/opt"
"github.com/syndtr/goleveldb/leveldb/util"
)
2014-07-23 08:31:36 +02:00
func (db *DB) writeJournal(b *Batch) error {
w, err := db.journal.Next()
2014-07-06 14:46:48 +02:00
if err != nil {
return err
}
if _, err := w.Write(b.encode()); err != nil {
return err
}
2014-07-23 08:31:36 +02:00
if err := db.journal.Flush(); err != nil {
2014-07-06 14:46:48 +02:00
return err
}
if b.sync {
2014-07-23 08:31:36 +02:00
return db.journalWriter.Sync()
2014-07-06 14:46:48 +02:00
}
return nil
}
2014-07-23 08:31:36 +02:00
func (db *DB) jWriter() {
defer db.closeW.Done()
2014-07-06 14:46:48 +02:00
for {
select {
2014-07-23 08:31:36 +02:00
case b := <-db.journalC:
2014-07-06 14:46:48 +02:00
if b != nil {
2014-07-23 08:31:36 +02:00
db.journalAckC <- db.writeJournal(b)
2014-07-06 14:46:48 +02:00
}
2014-07-23 08:31:36 +02:00
case _, _ = <-db.closeC:
2014-07-06 14:46:48 +02:00
return
}
}
}
2014-08-15 09:16:30 +02:00
func (db *DB) rotateMem(n int) (mem *memDB, err error) {
2014-07-06 14:46:48 +02:00
// Wait for pending memdb compaction.
2014-07-23 08:31:36 +02:00
err = db.compSendIdle(db.mcompCmdC)
2014-07-06 14:46:48 +02:00
if err != nil {
return
}
// Create new memdb and journal.
2014-07-23 08:31:36 +02:00
mem, err = db.newMem(n)
2014-07-06 14:46:48 +02:00
if err != nil {
return
}
// Schedule memdb compaction.
2014-07-23 08:31:36 +02:00
db.compTrigger(db.mcompTriggerC)
2014-07-06 14:46:48 +02:00
return
}
2014-08-15 09:16:30 +02:00
func (db *DB) flush(n int) (mem *memDB, nn int, err error) {
2014-07-06 14:46:48 +02:00
delayed := false
2014-08-15 09:16:30 +02:00
flush := func() (retry bool) {
2014-07-23 08:31:36 +02:00
v := db.s.version()
2014-07-06 14:46:48 +02:00
defer v.release()
2014-07-23 08:31:36 +02:00
mem = db.getEffectiveMem()
2014-08-15 09:16:30 +02:00
defer func() {
if retry {
mem.decref()
mem = nil
}
}()
2014-09-02 09:43:42 +02:00
nn = mem.mdb.Free()
2014-07-06 14:46:48 +02:00
switch {
case v.tLen(0) >= kL0_SlowdownWritesTrigger && !delayed:
delayed = true
time.Sleep(time.Millisecond)
case nn >= n:
return false
case v.tLen(0) >= kL0_StopWritesTrigger:
delayed = true
2014-07-23 08:31:36 +02:00
err = db.compSendIdle(db.tcompCmdC)
2014-07-06 14:46:48 +02:00
if err != nil {
return false
}
default:
// Allow memdb to grow if it has no entry.
2014-09-02 09:43:42 +02:00
if mem.mdb.Len() == 0 {
2014-07-06 14:46:48 +02:00
nn = n
2014-08-15 09:16:30 +02:00
} else {
mem.decref()
mem, err = db.rotateMem(n)
if err == nil {
2014-09-02 09:43:42 +02:00
nn = mem.mdb.Free()
2014-08-15 09:16:30 +02:00
} else {
nn = 0
}
2014-07-06 14:46:48 +02:00
}
return false
}
return true
}
start := time.Now()
for flush() {
}
if delayed {
2014-07-23 08:31:36 +02:00
db.logf("db@write delayed T·%v", time.Since(start))
2014-07-06 14:46:48 +02:00
}
return
}
// Write apply the given batch to the DB. The batch will be applied
// sequentially.
//
// It is safe to modify the contents of the arguments after Write returns.
2014-07-23 08:31:36 +02:00
func (db *DB) Write(b *Batch, wo *opt.WriteOptions) (err error) {
err = db.ok()
2014-07-06 14:46:48 +02:00
if err != nil || b == nil || b.len() == 0 {
return
}
b.init(wo.GetSync())
// The write happen synchronously.
select {
2014-07-23 08:31:36 +02:00
case db.writeC <- b:
if <-db.writeMergedC {
return <-db.writeAckC
2014-07-06 14:46:48 +02:00
}
2014-07-23 08:31:36 +02:00
case db.writeLockC <- struct{}{}:
case _, _ = <-db.closeC:
2014-07-06 14:46:48 +02:00
return ErrClosed
}
merged := 0
2014-11-04 05:00:11 +01:00
danglingMerge := false
2014-07-06 14:46:48 +02:00
defer func() {
2014-11-04 05:00:11 +01:00
if danglingMerge {
db.writeMergedC <- false
} else {
<-db.writeLockC
}
2014-07-06 14:46:48 +02:00
for i := 0; i < merged; i++ {
2014-07-23 08:31:36 +02:00
db.writeAckC <- err
2014-07-06 14:46:48 +02:00
}
}()
2014-07-23 08:31:36 +02:00
mem, memFree, err := db.flush(b.size())
2014-07-06 14:46:48 +02:00
if err != nil {
return
}
2014-08-15 09:16:30 +02:00
defer mem.decref()
2014-07-06 14:46:48 +02:00
// Calculate maximum size of the batch.
m := 1 << 20
if x := b.size(); x <= 128<<10 {
m = x + (128 << 10)
}
m = minInt(m, memFree)
// Merge with other batch.
drain:
for b.size() < m && !b.sync {
select {
2014-07-23 08:31:36 +02:00
case nb := <-db.writeC:
2014-07-06 14:46:48 +02:00
if b.size()+nb.size() <= m {
b.append(nb)
2014-07-23 08:31:36 +02:00
db.writeMergedC <- true
2014-07-06 14:46:48 +02:00
merged++
} else {
2014-11-04 05:00:11 +01:00
danglingMerge = true
2014-07-06 14:46:48 +02:00
break drain
}
default:
break drain
}
}
// Set batch first seq number relative from last seq.
2014-07-23 08:31:36 +02:00
b.seq = db.seq + 1
2014-07-06 14:46:48 +02:00
// Write journal concurrently if it is large enough.
if b.size() >= (128 << 10) {
// Push the write batch to the journal writer
select {
2014-07-23 08:31:36 +02:00
case _, _ = <-db.closeC:
2014-07-06 14:46:48 +02:00
err = ErrClosed
return
2014-07-23 08:31:36 +02:00
case db.journalC <- b:
2014-07-06 14:46:48 +02:00
// Write into memdb
2014-09-02 09:43:42 +02:00
b.memReplay(mem.mdb)
2014-07-06 14:46:48 +02:00
}
// Wait for journal writer
select {
2014-07-23 08:31:36 +02:00
case _, _ = <-db.closeC:
2014-07-06 14:46:48 +02:00
err = ErrClosed
return
2014-07-23 08:31:36 +02:00
case err = <-db.journalAckC:
2014-07-06 14:46:48 +02:00
if err != nil {
// Revert memdb if error detected
2014-09-02 09:43:42 +02:00
b.revertMemReplay(mem.mdb)
2014-07-06 14:46:48 +02:00
return
}
}
} else {
2014-07-23 08:31:36 +02:00
err = db.writeJournal(b)
2014-07-06 14:46:48 +02:00
if err != nil {
return
}
2014-09-02 09:43:42 +02:00
b.memReplay(mem.mdb)
2014-07-06 14:46:48 +02:00
}
// Set last seq number.
2014-07-23 08:31:36 +02:00
db.addSeq(uint64(b.len()))
2014-07-06 14:46:48 +02:00
if b.size() >= memFree {
2014-07-23 08:31:36 +02:00
db.rotateMem(0)
2014-07-06 14:46:48 +02:00
}
return
}
// Put sets the value for the given key. It overwrites any previous value
// for that key; a DB is not a multi-map.
//
// It is safe to modify the contents of the arguments after Put returns.
2014-07-23 08:31:36 +02:00
func (db *DB) Put(key, value []byte, wo *opt.WriteOptions) error {
2014-07-06 14:46:48 +02:00
b := new(Batch)
b.Put(key, value)
2014-07-23 08:31:36 +02:00
return db.Write(b, wo)
2014-07-06 14:46:48 +02:00
}
// Delete deletes the value for the given key. It returns ErrNotFound if
// the DB does not contain the key.
//
// It is safe to modify the contents of the arguments after Delete returns.
2014-07-23 08:31:36 +02:00
func (db *DB) Delete(key []byte, wo *opt.WriteOptions) error {
2014-07-06 14:46:48 +02:00
b := new(Batch)
b.Delete(key)
2014-07-23 08:31:36 +02:00
return db.Write(b, wo)
2014-07-06 14:46:48 +02:00
}
2014-07-06 23:13:10 +02:00
func isMemOverlaps(icmp *iComparer, mem *memdb.DB, min, max []byte) bool {
2014-07-06 14:46:48 +02:00
iter := mem.NewIterator(nil)
defer iter.Release()
2014-07-06 23:13:10 +02:00
return (max == nil || (iter.First() && icmp.uCompare(max, iKey(iter.Key()).ukey()) >= 0)) &&
(min == nil || (iter.Last() && icmp.uCompare(min, iKey(iter.Key()).ukey()) <= 0))
2014-07-06 14:46:48 +02:00
}
// CompactRange compacts the underlying DB for the given key range.
// In particular, deleted and overwritten versions are discarded,
// and the data is rearranged to reduce the cost of operations
// needed to access the data. This operation should typically only
// be invoked by users who understand the underlying implementation.
//
// A nil Range.Start is treated as a key before all keys in the DB.
// And a nil Range.Limit is treated as a key after all keys in the DB.
// Therefore if both is nil then it will compact entire DB.
2014-07-23 08:31:36 +02:00
func (db *DB) CompactRange(r util.Range) error {
if err := db.ok(); err != nil {
2014-07-06 14:46:48 +02:00
return err
}
2014-11-04 05:00:11 +01:00
// Lock writer.
2014-07-06 14:46:48 +02:00
select {
2014-07-23 08:31:36 +02:00
case db.writeLockC <- struct{}{}:
case _, _ = <-db.closeC:
2014-07-06 14:46:48 +02:00
return ErrClosed
}
// Check for overlaps in memdb.
2014-07-23 08:31:36 +02:00
mem := db.getEffectiveMem()
2014-08-15 09:16:30 +02:00
defer mem.decref()
2014-09-02 09:43:42 +02:00
if isMemOverlaps(db.s.icmp, mem.mdb, r.Start, r.Limit) {
2014-07-06 14:46:48 +02:00
// Memdb compaction.
2014-07-23 08:31:36 +02:00
if _, err := db.rotateMem(0); err != nil {
<-db.writeLockC
2014-07-06 14:46:48 +02:00
return err
}
2014-07-23 08:31:36 +02:00
<-db.writeLockC
if err := db.compSendIdle(db.mcompCmdC); err != nil {
2014-07-06 14:46:48 +02:00
return err
}
} else {
2014-07-23 08:31:36 +02:00
<-db.writeLockC
2014-07-06 14:46:48 +02:00
}
// Table compaction.
2014-07-23 08:31:36 +02:00
return db.compSendRange(db.tcompCmdC, -1, r.Start, r.Limit)
2014-07-06 14:46:48 +02:00
}