2021-07-28 09:47:38 +07:00
|
|
|
package olddb
|
2020-08-17 13:45:52 +07:00
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"fmt"
|
|
|
|
|
2021-07-29 18:53:13 +07:00
|
|
|
"github.com/ledgerwatch/erigon-lib/kv"
|
2021-06-19 15:21:53 +07:00
|
|
|
"github.com/ledgerwatch/erigon/ethdb"
|
2021-07-29 17:23:23 +07:00
|
|
|
"github.com/ledgerwatch/log/v3"
|
2020-08-17 13:45:52 +07:00
|
|
|
)
|
|
|
|
|
|
|
|
// TxDb - provides Database interface around ethdb.Tx
|
|
|
|
// It's not thread-safe!
|
2020-08-24 18:07:59 +07:00
|
|
|
// TxDb not usable after .Commit()/.Rollback() call, but usable after .CommitAndBegin() call
|
2021-03-22 13:47:01 +07:00
|
|
|
// you can put unlimited amount of data into this class
|
2020-08-17 13:45:52 +07:00
|
|
|
// Walk and MultiWalk methods - work outside of Tx object yet, will implement it later
|
2022-08-10 19:04:13 +07:00
|
|
|
// Deprecated
|
|
|
|
// nolint
|
2020-08-17 13:45:52 +07:00
|
|
|
type TxDb struct {
|
2021-06-19 15:21:53 +07:00
|
|
|
db ethdb.Database
|
2021-07-28 09:47:38 +07:00
|
|
|
tx kv.Tx
|
|
|
|
cursors map[string]kv.Cursor
|
2021-06-19 15:21:53 +07:00
|
|
|
txFlags ethdb.TxFlags
|
2021-02-10 20:04:22 +03:00
|
|
|
len uint64
|
2020-08-17 13:45:52 +07:00
|
|
|
}
|
|
|
|
|
2022-08-10 19:04:13 +07:00
|
|
|
// nolint
|
2021-07-28 09:47:38 +07:00
|
|
|
func WrapIntoTxDB(tx kv.RwTx) *TxDb {
|
|
|
|
return &TxDb{tx: tx, cursors: map[string]kv.Cursor{}}
|
2021-05-01 14:42:23 +07:00
|
|
|
}
|
|
|
|
|
2020-08-17 13:45:52 +07:00
|
|
|
func (m *TxDb) Close() {
|
|
|
|
panic("don't call me")
|
|
|
|
}
|
|
|
|
|
2021-06-19 15:21:53 +07:00
|
|
|
func (m *TxDb) Begin(ctx context.Context, flags ethdb.TxFlags) (ethdb.DbWithPendingMutations, error) {
|
2020-08-26 13:02:10 +07:00
|
|
|
batch := m
|
|
|
|
if m.tx != nil {
|
2021-02-10 20:04:22 +03:00
|
|
|
panic("nested transactions not supported")
|
2020-08-26 13:02:10 +07:00
|
|
|
}
|
|
|
|
|
2021-02-10 20:04:22 +03:00
|
|
|
if err := batch.begin(ctx, flags); err != nil {
|
2020-08-17 13:45:52 +07:00
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
return batch, nil
|
|
|
|
}
|
|
|
|
|
2021-07-28 09:47:38 +07:00
|
|
|
func (m *TxDb) cursor(bucket string) (kv.Cursor, error) {
|
2021-01-15 16:38:09 +07:00
|
|
|
c, ok := m.cursors[bucket]
|
|
|
|
if !ok {
|
2021-04-02 13:36:49 +07:00
|
|
|
var err error
|
|
|
|
c, err = m.tx.Cursor(bucket)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2021-01-15 16:38:09 +07:00
|
|
|
m.cursors[bucket] = c
|
|
|
|
}
|
2021-04-02 13:36:49 +07:00
|
|
|
return c, nil
|
2021-01-15 16:38:09 +07:00
|
|
|
}
|
|
|
|
|
2021-03-20 17:12:54 +03:00
|
|
|
func (m *TxDb) IncrementSequence(bucket string, amount uint64) (res uint64, err error) {
|
2021-07-28 09:47:38 +07:00
|
|
|
return m.tx.(kv.RwTx).IncrementSequence(bucket, amount)
|
2021-03-20 17:12:54 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) ReadSequence(bucket string) (res uint64, err error) {
|
|
|
|
return m.tx.ReadSequence(bucket)
|
2020-11-14 20:48:29 +07:00
|
|
|
}
|
|
|
|
|
2022-07-26 12:47:05 +07:00
|
|
|
func (m *TxDb) Put(table string, k, v []byte) error {
|
|
|
|
m.len += uint64(len(k) + len(v))
|
|
|
|
c, err := m.cursor(table)
|
2021-04-02 13:36:49 +07:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2022-07-26 12:47:05 +07:00
|
|
|
return c.(kv.RwCursor).Put(k, v)
|
2020-08-17 13:45:52 +07:00
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) Append(bucket string, key []byte, value []byte) error {
|
|
|
|
m.len += uint64(len(key) + len(value))
|
2021-04-02 13:36:49 +07:00
|
|
|
c, err := m.cursor(bucket)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2021-07-28 09:47:38 +07:00
|
|
|
return c.(kv.RwCursor).Append(key, value)
|
2020-11-28 21:24:47 +07:00
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) AppendDup(bucket string, key []byte, value []byte) error {
|
|
|
|
m.len += uint64(len(key) + len(value))
|
2021-04-02 13:36:49 +07:00
|
|
|
c, err := m.cursor(bucket)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2021-07-28 09:47:38 +07:00
|
|
|
return c.(kv.RwCursorDupSort).AppendDup(key, value)
|
2020-08-17 13:45:52 +07:00
|
|
|
}
|
|
|
|
|
2022-07-26 12:47:05 +07:00
|
|
|
func (m *TxDb) Delete(table string, k []byte) error {
|
2020-10-29 20:19:31 +07:00
|
|
|
m.len += uint64(len(k))
|
2022-07-26 12:47:05 +07:00
|
|
|
c, err := m.cursor(table)
|
2021-04-02 13:36:49 +07:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2022-07-26 12:47:05 +07:00
|
|
|
return c.(kv.RwCursor).Delete(k)
|
2020-08-17 13:45:52 +07:00
|
|
|
}
|
|
|
|
|
2021-06-19 15:21:53 +07:00
|
|
|
func (m *TxDb) begin(ctx context.Context, flags ethdb.TxFlags) error {
|
2021-07-28 09:47:38 +07:00
|
|
|
db := m.db.(ethdb.HasRwKV).RwKV()
|
2021-03-21 16:15:25 +03:00
|
|
|
|
2021-07-28 09:47:38 +07:00
|
|
|
var tx kv.Tx
|
2021-03-21 16:15:25 +03:00
|
|
|
var err error
|
2021-06-19 15:21:53 +07:00
|
|
|
if flagsðdb.RO != 0 {
|
2021-07-28 09:47:38 +07:00
|
|
|
tx, err = db.BeginRo(ctx)
|
2021-03-21 16:15:25 +03:00
|
|
|
} else {
|
2021-07-28 09:47:38 +07:00
|
|
|
tx, err = db.BeginRw(ctx)
|
2021-03-21 16:15:25 +03:00
|
|
|
}
|
2020-08-17 13:45:52 +07:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2020-08-26 13:02:10 +07:00
|
|
|
m.tx = tx
|
2021-07-28 09:47:38 +07:00
|
|
|
m.cursors = make(map[string]kv.Cursor, 16)
|
2020-08-17 13:45:52 +07:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-07-28 09:47:38 +07:00
|
|
|
func (m *TxDb) RwKV() kv.RwDB {
|
2020-08-26 13:02:10 +07:00
|
|
|
panic("not allowed to get KV interface because you will loose transaction, please use .Tx() method")
|
2020-08-17 13:45:52 +07:00
|
|
|
}
|
|
|
|
|
2021-02-28 11:11:28 +07:00
|
|
|
// Last can only be called from the transaction thread
|
2020-08-17 13:45:52 +07:00
|
|
|
func (m *TxDb) Last(bucket string) ([]byte, []byte, error) {
|
2021-04-02 13:36:49 +07:00
|
|
|
c, err := m.cursor(bucket)
|
|
|
|
if err != nil {
|
|
|
|
return []byte{}, nil, err
|
|
|
|
}
|
|
|
|
return c.Last()
|
2020-08-17 13:45:52 +07:00
|
|
|
}
|
|
|
|
|
2021-04-05 16:04:58 +03:00
|
|
|
func (m *TxDb) GetOne(bucket string, key []byte) ([]byte, error) {
|
2021-04-02 13:36:49 +07:00
|
|
|
c, err := m.cursor(bucket)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
_, v, err := c.SeekExact(key)
|
2021-04-05 16:04:58 +03:00
|
|
|
return v, err
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) Get(bucket string, key []byte) ([]byte, error) {
|
|
|
|
dat, err := m.GetOne(bucket, key)
|
2021-06-19 15:21:53 +07:00
|
|
|
return ethdb.GetOneWrapper(dat, err)
|
2020-08-17 13:45:52 +07:00
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) Has(bucket string, key []byte) (bool, error) {
|
|
|
|
v, err := m.Get(bucket, key)
|
|
|
|
if err != nil {
|
|
|
|
return false, err
|
|
|
|
}
|
|
|
|
return v != nil, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) BatchSize() int {
|
|
|
|
return int(m.len)
|
|
|
|
}
|
|
|
|
|
2021-06-05 22:17:04 +07:00
|
|
|
func (m *TxDb) ForEach(bucket string, fromPrefix []byte, walker func(k, v []byte) error) error {
|
|
|
|
return m.tx.ForEach(bucket, fromPrefix, walker)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) ForPrefix(bucket string, prefix []byte, walker func(k, v []byte) error) error {
|
|
|
|
return m.tx.ForPrefix(bucket, prefix, walker)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) ForAmount(bucket string, prefix []byte, amount uint32, walker func(k, v []byte) error) error {
|
|
|
|
return m.tx.ForAmount(bucket, prefix, amount, walker)
|
|
|
|
}
|
2020-08-17 13:45:52 +07:00
|
|
|
|
2021-03-22 19:41:52 +07:00
|
|
|
func (m *TxDb) Commit() error {
|
2020-08-26 13:02:10 +07:00
|
|
|
if m.tx == nil {
|
2021-03-22 19:41:52 +07:00
|
|
|
return fmt.Errorf("second call .Commit() on same transaction")
|
2020-08-17 13:45:52 +07:00
|
|
|
}
|
2021-04-03 09:26:00 +03:00
|
|
|
if err := m.tx.Commit(); err != nil {
|
2021-03-22 19:41:52 +07:00
|
|
|
return err
|
2020-08-17 13:45:52 +07:00
|
|
|
}
|
2020-08-26 13:02:10 +07:00
|
|
|
m.tx = nil
|
2020-08-17 13:45:52 +07:00
|
|
|
m.cursors = nil
|
|
|
|
m.len = 0
|
2021-03-22 19:41:52 +07:00
|
|
|
return nil
|
2020-08-17 13:45:52 +07:00
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) Rollback() {
|
2020-08-26 13:02:10 +07:00
|
|
|
if m.tx == nil {
|
2020-08-17 13:45:52 +07:00
|
|
|
return
|
|
|
|
}
|
2020-08-26 13:02:10 +07:00
|
|
|
m.tx.Rollback()
|
2020-08-17 13:45:52 +07:00
|
|
|
m.cursors = nil
|
2020-08-26 13:02:10 +07:00
|
|
|
m.tx = nil
|
2020-08-17 13:45:52 +07:00
|
|
|
m.len = 0
|
|
|
|
}
|
|
|
|
|
2021-07-28 09:47:38 +07:00
|
|
|
func (m *TxDb) Tx() kv.Tx {
|
2020-08-26 13:02:10 +07:00
|
|
|
return m.tx
|
|
|
|
}
|
|
|
|
|
2020-09-09 02:39:43 +07:00
|
|
|
func (m *TxDb) BucketExists(name string) (bool, error) {
|
2021-07-28 09:47:38 +07:00
|
|
|
migrator, ok := m.tx.(kv.BucketMigrator)
|
2020-09-09 02:39:43 +07:00
|
|
|
if !ok {
|
|
|
|
return false, fmt.Errorf("%T doesn't implement ethdb.TxMigrator interface", m.tx)
|
|
|
|
}
|
2021-07-24 11:28:05 +07:00
|
|
|
return migrator.ExistsBucket(name)
|
2020-09-09 02:39:43 +07:00
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) ClearBuckets(buckets ...string) error {
|
|
|
|
for i := range buckets {
|
|
|
|
name := buckets[i]
|
|
|
|
|
2021-07-28 09:47:38 +07:00
|
|
|
migrator, ok := m.tx.(kv.BucketMigrator)
|
2020-09-09 02:39:43 +07:00
|
|
|
if !ok {
|
|
|
|
return fmt.Errorf("%T doesn't implement ethdb.TxMigrator interface", m.tx)
|
|
|
|
}
|
|
|
|
if err := migrator.ClearBucket(name); err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *TxDb) DropBuckets(buckets ...string) error {
|
|
|
|
for i := range buckets {
|
|
|
|
name := buckets[i]
|
|
|
|
log.Info("Dropping bucket", "name", name)
|
2021-07-28 09:47:38 +07:00
|
|
|
migrator, ok := m.tx.(kv.BucketMigrator)
|
2020-09-09 02:39:43 +07:00
|
|
|
if !ok {
|
|
|
|
return fmt.Errorf("%T doesn't implement ethdb.TxMigrator interface", m.tx)
|
|
|
|
}
|
|
|
|
if err := migrator.DropBucket(name); err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|