2019-05-27 13:51:49 +00:00
|
|
|
package ethdb
|
|
|
|
|
|
|
|
import (
|
2020-10-27 22:30:18 +00:00
|
|
|
"bytes"
|
2020-06-04 09:35:42 +00:00
|
|
|
"context"
|
2020-10-27 22:30:18 +00:00
|
|
|
"strings"
|
|
|
|
"sync"
|
|
|
|
"sync/atomic"
|
|
|
|
"unsafe"
|
|
|
|
|
2020-09-01 06:48:25 +00:00
|
|
|
"github.com/c2h5oh/datasize"
|
2020-10-27 22:30:18 +00:00
|
|
|
"github.com/google/btree"
|
2019-05-27 13:51:49 +00:00
|
|
|
"github.com/ledgerwatch/turbo-geth/common"
|
2020-06-15 12:30:54 +00:00
|
|
|
"github.com/ledgerwatch/turbo-geth/metrics"
|
2019-05-27 13:51:49 +00:00
|
|
|
)
|
|
|
|
|
2020-07-21 08:33:03 +00:00
|
|
|
var (
|
|
|
|
dbCommitBigBatchTimer = metrics.NewRegisteredTimer("db/commit/big_batch", nil)
|
|
|
|
dbCommitSmallBatchTimer = metrics.NewRegisteredTimer("db/commit/small_batch", nil)
|
|
|
|
)
|
2020-06-15 12:30:54 +00:00
|
|
|
|
2019-05-27 13:51:49 +00:00
|
|
|
type mutation struct {
|
2020-10-27 22:30:18 +00:00
|
|
|
puts *btree.BTree
|
|
|
|
mu sync.RWMutex
|
|
|
|
searchItem MutationItem
|
|
|
|
size int
|
|
|
|
db Database
|
|
|
|
}
|
|
|
|
|
|
|
|
type MutationItem struct {
|
|
|
|
table string
|
|
|
|
key []byte
|
|
|
|
value []byte
|
|
|
|
}
|
|
|
|
|
|
|
|
func (mi *MutationItem) Less(than btree.Item) bool {
|
|
|
|
i := than.(*MutationItem)
|
|
|
|
c := strings.Compare(mi.table, i.table)
|
|
|
|
if c != 0 {
|
|
|
|
return c < 0
|
|
|
|
}
|
|
|
|
return bytes.Compare(mi.key, i.key) < 0
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
2020-06-05 09:25:33 +00:00
|
|
|
func (m *mutation) KV() KV {
|
|
|
|
if casted, ok := m.db.(HasKV); ok {
|
2020-03-20 11:30:14 +00:00
|
|
|
return casted.KV()
|
|
|
|
}
|
2020-05-15 08:58:36 +00:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) getMem(table string, key []byte) ([]byte, bool) {
|
2019-05-27 13:51:49 +00:00
|
|
|
m.mu.RLock()
|
|
|
|
defer m.mu.RUnlock()
|
2020-10-27 22:30:18 +00:00
|
|
|
m.searchItem.table = table
|
|
|
|
m.searchItem.key = key
|
|
|
|
i := m.puts.Get(&m.searchItem)
|
|
|
|
if i == nil {
|
|
|
|
return nil, false
|
|
|
|
}
|
|
|
|
return i.(*MutationItem).value, true
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// Can only be called from the worker thread
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) Get(table string, key []byte) ([]byte, error) {
|
|
|
|
if value, ok := m.getMem(table, key); ok {
|
2019-05-27 13:51:49 +00:00
|
|
|
if value == nil {
|
|
|
|
return nil, ErrKeyNotFound
|
|
|
|
}
|
|
|
|
return value, nil
|
|
|
|
}
|
|
|
|
if m.db != nil {
|
2020-10-27 22:30:18 +00:00
|
|
|
return m.db.Get(table, key)
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
return nil, ErrKeyNotFound
|
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) Last(table string) ([]byte, []byte, error) {
|
|
|
|
return m.db.Last(table)
|
2020-08-12 03:49:52 +00:00
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) Reserve(table string, key []byte, i int) ([]byte, error) {
|
|
|
|
return m.db.(DbWithPendingMutations).Reserve(table, key, i)
|
2020-09-28 17:18:36 +00:00
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) GetIndexChunk(table string, key []byte, timestamp uint64) ([]byte, error) {
|
2020-04-20 10:35:33 +00:00
|
|
|
if m.db != nil {
|
2020-10-27 22:30:18 +00:00
|
|
|
return m.db.GetIndexChunk(table, key, timestamp)
|
2020-04-20 10:35:33 +00:00
|
|
|
}
|
2020-04-21 08:15:40 +00:00
|
|
|
return nil, ErrKeyNotFound
|
2020-04-20 10:35:33 +00:00
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) hasMem(table string, key []byte) bool {
|
2020-04-20 10:35:33 +00:00
|
|
|
m.mu.RLock()
|
|
|
|
defer m.mu.RUnlock()
|
2020-10-27 22:30:18 +00:00
|
|
|
m.searchItem.table = table
|
|
|
|
m.searchItem.key = key
|
|
|
|
return m.puts.Has(&m.searchItem)
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) Has(table string, key []byte) (bool, error) {
|
|
|
|
if m.hasMem(table, key) {
|
2019-05-27 13:51:49 +00:00
|
|
|
return true, nil
|
|
|
|
}
|
|
|
|
if m.db != nil {
|
2020-10-27 22:30:18 +00:00
|
|
|
return m.db.Has(table, key)
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
return false, nil
|
|
|
|
}
|
|
|
|
|
2020-06-04 09:35:42 +00:00
|
|
|
func (m *mutation) DiskSize(ctx context.Context) (common.StorageSize, error) {
|
2019-05-27 13:51:49 +00:00
|
|
|
if m.db == nil {
|
2020-06-04 09:35:42 +00:00
|
|
|
return 0, nil
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
2020-07-29 04:31:46 +00:00
|
|
|
sz, err := m.db.(HasStats).DiskSize(ctx)
|
|
|
|
if err != nil {
|
|
|
|
return 0, err
|
2020-06-04 09:35:42 +00:00
|
|
|
}
|
2020-07-29 04:31:46 +00:00
|
|
|
return common.StorageSize(sz), nil
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) Put(table string, key []byte, value []byte) error {
|
2019-05-27 13:51:49 +00:00
|
|
|
m.mu.Lock()
|
|
|
|
defer m.mu.Unlock()
|
2019-11-21 18:38:00 +00:00
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
newMi := &MutationItem{table: table, key: key, value: value}
|
|
|
|
i := m.puts.ReplaceOrInsert(newMi)
|
|
|
|
m.size += int(unsafe.Sizeof(newMi)) + len(key) + len(value)
|
|
|
|
if i != nil {
|
|
|
|
oldMi := i.(*MutationItem)
|
|
|
|
m.size -= (int(unsafe.Sizeof(oldMi)) + len(oldMi.key) + len(oldMi.value))
|
|
|
|
}
|
2019-05-27 13:51:49 +00:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) Append(table string, key []byte, value []byte) error {
|
|
|
|
return m.Put(table, key, value)
|
2020-09-01 06:48:25 +00:00
|
|
|
}
|
|
|
|
|
2019-05-27 13:51:49 +00:00
|
|
|
func (m *mutation) MultiPut(tuples ...[]byte) (uint64, error) {
|
|
|
|
m.mu.Lock()
|
|
|
|
defer m.mu.Unlock()
|
|
|
|
l := len(tuples)
|
|
|
|
for i := 0; i < l; i += 3 {
|
2020-10-27 22:30:18 +00:00
|
|
|
newMi := &MutationItem{table: string(tuples[i]), key: tuples[i+1], value: tuples[i+2]}
|
|
|
|
i := m.puts.ReplaceOrInsert(newMi)
|
|
|
|
m.size += int(unsafe.Sizeof(newMi)) + len(newMi.key) + len(newMi.value)
|
|
|
|
if i != nil {
|
|
|
|
oldMi := i.(*MutationItem)
|
|
|
|
m.size -= (int(unsafe.Sizeof(oldMi)) + len(oldMi.key) + len(oldMi.value))
|
|
|
|
}
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
return 0, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *mutation) BatchSize() int {
|
|
|
|
m.mu.RLock()
|
|
|
|
defer m.mu.RUnlock()
|
2020-10-27 22:30:18 +00:00
|
|
|
return m.size
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
2019-11-21 15:12:38 +00:00
|
|
|
// IdealBatchSize defines the size of the data batches should ideally add in one write.
|
|
|
|
func (m *mutation) IdealBatchSize() int {
|
2020-08-26 06:03:50 +00:00
|
|
|
return int(512 * datasize.MB)
|
2019-11-21 15:12:38 +00:00
|
|
|
}
|
|
|
|
|
2019-10-30 17:33:01 +00:00
|
|
|
// WARNING: Merged mem/DB walk is not implemented
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) Walk(table string, startkey []byte, fixedbits int, walker func([]byte, []byte) (bool, error)) error {
|
2019-12-20 12:25:40 +00:00
|
|
|
m.panicOnEmptyDB()
|
2020-10-27 22:30:18 +00:00
|
|
|
return m.db.Walk(table, startkey, fixedbits, walker)
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
2019-10-30 17:33:01 +00:00
|
|
|
// WARNING: Merged mem/DB walk is not implemented
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) MultiWalk(table string, startkeys [][]byte, fixedbits []int, walker func(int, []byte, []byte) error) error {
|
2019-12-20 12:25:40 +00:00
|
|
|
m.panicOnEmptyDB()
|
2020-10-27 22:30:18 +00:00
|
|
|
return m.db.MultiWalk(table, startkeys, fixedbits, walker)
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) Delete(table string, key []byte) error {
|
|
|
|
return m.Put(table, key, nil)
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
2020-09-28 17:18:36 +00:00
|
|
|
func (m *mutation) CommitAndBegin(ctx context.Context) error {
|
2020-08-24 11:07:59 +00:00
|
|
|
_, err := m.Commit()
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
2020-09-28 17:18:36 +00:00
|
|
|
func (m *mutation) RollbackAndBegin(ctx context.Context) error {
|
|
|
|
m.Rollback()
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
func (m *mutation) doCommit(tx Tx) error {
|
|
|
|
var prevTable string
|
|
|
|
var c Cursor
|
|
|
|
var innerErr error
|
|
|
|
var isEndOfBucket bool
|
|
|
|
m.puts.Ascend(func(i btree.Item) bool {
|
|
|
|
mi := i.(*MutationItem)
|
|
|
|
if mi.table != prevTable {
|
|
|
|
if c != nil {
|
|
|
|
c.Close()
|
|
|
|
}
|
|
|
|
c = tx.Cursor(mi.table)
|
|
|
|
prevTable = mi.table
|
|
|
|
firstKey, _, err := c.Seek(mi.key)
|
|
|
|
if err != nil {
|
|
|
|
innerErr = err
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
isEndOfBucket = firstKey == nil
|
|
|
|
}
|
|
|
|
if isEndOfBucket {
|
|
|
|
if len(mi.value) > 0 {
|
|
|
|
if err := c.Append(mi.key, mi.value); err != nil {
|
|
|
|
innerErr = err
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
}
|
|
|
|
} else if len(mi.value) == 0 {
|
|
|
|
if err := c.Delete(mi.key); err != nil {
|
|
|
|
innerErr = err
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
if err := c.Put(mi.key, mi.value); err != nil {
|
|
|
|
innerErr = err
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return true
|
|
|
|
})
|
|
|
|
return innerErr
|
|
|
|
}
|
|
|
|
|
2019-05-27 13:51:49 +00:00
|
|
|
func (m *mutation) Commit() (uint64, error) {
|
|
|
|
if m.db == nil {
|
|
|
|
return 0, nil
|
|
|
|
}
|
|
|
|
m.mu.Lock()
|
|
|
|
defer m.mu.Unlock()
|
2020-10-27 22:30:18 +00:00
|
|
|
if tx, ok := m.db.(HasTx); ok {
|
|
|
|
if err := m.doCommit(tx.Tx()); err != nil {
|
|
|
|
return 0, err
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
if err := m.db.(HasKV).KV().Update(context.Background(), func(tx Tx) error {
|
|
|
|
return m.doCommit(tx)
|
|
|
|
}); err != nil {
|
|
|
|
return 0, err
|
2019-11-21 18:38:00 +00:00
|
|
|
}
|
2020-10-27 14:35:25 +00:00
|
|
|
}
|
|
|
|
|
2020-10-27 22:30:18 +00:00
|
|
|
m.puts.Clear(false /* addNodesToFreelist */)
|
|
|
|
m.size = 0
|
|
|
|
return 0, nil
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func (m *mutation) Rollback() {
|
|
|
|
m.mu.Lock()
|
|
|
|
defer m.mu.Unlock()
|
2020-10-27 22:30:18 +00:00
|
|
|
m.puts.Clear(false /* addNodesToFreelist */)
|
|
|
|
m.size = 0
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func (m *mutation) Keys() ([][]byte, error) {
|
|
|
|
m.mu.RLock()
|
|
|
|
defer m.mu.RUnlock()
|
2020-05-23 10:27:05 +00:00
|
|
|
tuples := common.NewTuples(m.puts.Len(), 2, 1)
|
2020-10-27 22:30:18 +00:00
|
|
|
var innerErr error
|
|
|
|
m.puts.Ascend(func(i btree.Item) bool {
|
|
|
|
mi := i.(*MutationItem)
|
|
|
|
if err := tuples.Append([]byte(mi.table), mi.key); err != nil {
|
|
|
|
innerErr = err
|
|
|
|
return false
|
2019-11-21 18:38:00 +00:00
|
|
|
}
|
2020-10-27 22:30:18 +00:00
|
|
|
return true
|
|
|
|
})
|
|
|
|
return tuples.Values, innerErr
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func (m *mutation) Close() {
|
|
|
|
m.Rollback()
|
|
|
|
}
|
|
|
|
|
2019-10-30 17:33:01 +00:00
|
|
|
func (m *mutation) NewBatch() DbWithPendingMutations {
|
2019-05-27 13:51:49 +00:00
|
|
|
mm := &mutation{
|
2020-04-15 09:33:22 +00:00
|
|
|
db: m,
|
2020-10-27 22:30:18 +00:00
|
|
|
puts: btree.New(32),
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
return mm
|
|
|
|
}
|
|
|
|
|
2020-10-25 08:38:55 +00:00
|
|
|
func (m *mutation) Begin(ctx context.Context, flags TxFlags) (DbWithPendingMutations, error) {
|
|
|
|
return m.db.Begin(ctx, flags)
|
2020-08-17 06:45:52 +00:00
|
|
|
}
|
|
|
|
|
2019-12-20 12:25:40 +00:00
|
|
|
func (m *mutation) panicOnEmptyDB() {
|
|
|
|
if m.db == nil {
|
|
|
|
panic("Not implemented")
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-05-27 13:51:49 +00:00
|
|
|
func (m *mutation) MemCopy() Database {
|
2019-12-20 12:25:40 +00:00
|
|
|
m.panicOnEmptyDB()
|
|
|
|
return m.db
|
2019-05-27 13:51:49 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// [TURBO-GETH] Freezer support (not implemented yet)
|
|
|
|
// Ancients returns an error as we don't have a backing chain freezer.
|
|
|
|
func (m *mutation) Ancients() (uint64, error) {
|
|
|
|
return 0, errNotSupported
|
|
|
|
}
|
|
|
|
|
|
|
|
// TruncateAncients returns an error as we don't have a backing chain freezer.
|
|
|
|
func (m *mutation) TruncateAncients(items uint64) error {
|
|
|
|
return errNotSupported
|
|
|
|
}
|
2019-12-20 12:25:40 +00:00
|
|
|
|
|
|
|
func NewRWDecorator(db Database) *RWCounterDecorator {
|
|
|
|
return &RWCounterDecorator{
|
|
|
|
db,
|
|
|
|
DBCounterStats{},
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
type RWCounterDecorator struct {
|
|
|
|
Database
|
|
|
|
DBCounterStats
|
|
|
|
}
|
|
|
|
|
|
|
|
type DBCounterStats struct {
|
2020-04-15 09:33:22 +00:00
|
|
|
Put uint64
|
|
|
|
Get uint64
|
|
|
|
GetS uint64
|
|
|
|
GetAsOf uint64
|
|
|
|
Has uint64
|
|
|
|
Walk uint64
|
|
|
|
WalkAsOf uint64
|
|
|
|
MultiWalk uint64
|
|
|
|
MultiWalkAsOf uint64
|
|
|
|
Delete uint64
|
|
|
|
MultiPut uint64
|
2019-12-20 12:25:40 +00:00
|
|
|
}
|
|
|
|
|
2020-08-10 23:55:32 +00:00
|
|
|
func (d *RWCounterDecorator) Put(bucket string, key, value []byte) error {
|
2019-12-20 12:25:40 +00:00
|
|
|
atomic.AddUint64(&d.DBCounterStats.Put, 1)
|
|
|
|
return d.Database.Put(bucket, key, value)
|
|
|
|
}
|
|
|
|
|
2020-08-10 23:55:32 +00:00
|
|
|
func (d *RWCounterDecorator) Get(bucket string, key []byte) ([]byte, error) {
|
2019-12-20 12:25:40 +00:00
|
|
|
atomic.AddUint64(&d.DBCounterStats.Get, 1)
|
|
|
|
return d.Database.Get(bucket, key)
|
|
|
|
}
|
2020-01-07 10:41:33 +00:00
|
|
|
|
2020-08-10 23:55:32 +00:00
|
|
|
func (d *RWCounterDecorator) Has(bucket string, key []byte) (bool, error) {
|
2019-12-20 12:25:40 +00:00
|
|
|
atomic.AddUint64(&d.DBCounterStats.Has, 1)
|
|
|
|
return d.Database.Has(bucket, key)
|
|
|
|
}
|
2020-08-10 23:55:32 +00:00
|
|
|
func (d *RWCounterDecorator) Walk(bucket string, startkey []byte, fixedbits int, walker func([]byte, []byte) (bool, error)) error {
|
2019-12-20 12:25:40 +00:00
|
|
|
atomic.AddUint64(&d.DBCounterStats.Walk, 1)
|
|
|
|
return d.Database.Walk(bucket, startkey, fixedbits, walker)
|
|
|
|
}
|
2020-08-10 23:55:32 +00:00
|
|
|
func (d *RWCounterDecorator) MultiWalk(bucket string, startkeys [][]byte, fixedbits []int, walker func(int, []byte, []byte) error) error {
|
2019-12-20 12:25:40 +00:00
|
|
|
atomic.AddUint64(&d.DBCounterStats.MultiWalk, 1)
|
|
|
|
return d.Database.MultiWalk(bucket, startkeys, fixedbits, walker)
|
|
|
|
}
|
2020-08-10 23:55:32 +00:00
|
|
|
func (d *RWCounterDecorator) Delete(bucket string, key []byte) error {
|
2019-12-20 12:25:40 +00:00
|
|
|
atomic.AddUint64(&d.DBCounterStats.Delete, 1)
|
|
|
|
return d.Database.Delete(bucket, key)
|
|
|
|
}
|
|
|
|
func (d *RWCounterDecorator) MultiPut(tuples ...[]byte) (uint64, error) {
|
|
|
|
atomic.AddUint64(&d.DBCounterStats.MultiPut, 1)
|
|
|
|
return d.Database.MultiPut(tuples...)
|
|
|
|
}
|
|
|
|
func (d *RWCounterDecorator) NewBatch() DbWithPendingMutations {
|
|
|
|
mm := &mutation{
|
2020-04-15 09:33:22 +00:00
|
|
|
db: d,
|
2020-10-27 22:30:18 +00:00
|
|
|
puts: btree.New(32),
|
2019-12-20 12:25:40 +00:00
|
|
|
}
|
|
|
|
return mm
|
|
|
|
}
|