multi: Migrate to UTXO database.
This migrates the UTXO set and state from the block database to the new UTXO database. The benefit of doing this is that it will allow the reading and writing to the UTXO set in the database to be optimized specifically for use with the UTXO set. An overview of the changes is as follows: - Introduce upgrade code to migrate UTXO data from the block database to the new UTXO database - Update all references to use the UTXO database rather than the block database where appropriate
This commit is contained in:
parent
6ce69b0992
commit
73d1b5f7de
@ -1162,14 +1162,6 @@ func (b *BlockChain) createChainState() error {
|
||||
return err
|
||||
}
|
||||
|
||||
// Create the bucket that houses the utxo set. Note that the
|
||||
// genesis block coinbase transaction is intentionally not
|
||||
// inserted here since it is not spendable by consensus rules.
|
||||
_, err = meta.CreateBucket(utxoSetBucketName)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Add the genesis block to the block index.
|
||||
err = dbPutBlockNode(dbTx, node)
|
||||
if err != nil {
|
||||
|
||||
@ -144,7 +144,7 @@ func chainSetup(dbName string, params *chaincfg.Params) (*BlockChain, func(), er
|
||||
TimeSource: NewMedianTime(),
|
||||
SigCache: sigCache,
|
||||
UtxoCache: NewUtxoCache(&UtxoCacheConfig{
|
||||
DB: db,
|
||||
DB: utxoDb,
|
||||
MaxSize: 100 * 1024 * 1024, // 100 MiB
|
||||
}),
|
||||
})
|
||||
@ -787,7 +787,7 @@ func (g *chaingenHarness) ExpectUtxoSetState(blockName string) {
|
||||
|
||||
// Fetch the utxo set state from the database.
|
||||
var gotState *utxoSetState
|
||||
err := g.chain.db.View(func(dbTx database.Tx) error {
|
||||
err := g.chain.utxoDb.View(func(dbTx database.Tx) error {
|
||||
var err error
|
||||
gotState, err = dbFetchUtxoSetState(dbTx)
|
||||
return err
|
||||
|
||||
@ -137,7 +137,7 @@ func chainSetup(dbName string, params *chaincfg.Params) (*blockchain.BlockChain,
|
||||
TimeSource: blockchain.NewMedianTime(),
|
||||
SigCache: sigCache,
|
||||
UtxoCache: blockchain.NewUtxoCache(&blockchain.UtxoCacheConfig{
|
||||
DB: db,
|
||||
DB: utxoDb,
|
||||
MaxSize: 100 * 1024 * 1024, // 100 MiB
|
||||
}),
|
||||
})
|
||||
|
||||
@ -111,7 +111,7 @@ func (b *BlockChain) TicketsWithAddress(address stdaddr.Address, isTreasuryEnabl
|
||||
|
||||
encodedAddr := address.String()
|
||||
var ticketsWithAddr []chainhash.Hash
|
||||
err := b.db.View(func(dbTx database.Tx) error {
|
||||
err := b.utxoDb.View(func(dbTx database.Tx) error {
|
||||
for _, hash := range tickets {
|
||||
outpoint := wire.OutPoint{Hash: hash, Index: 0, Tree: wire.TxTreeStake}
|
||||
utxo, err := b.utxoCache.FetchEntry(dbTx, outpoint)
|
||||
@ -223,7 +223,7 @@ func (b *BlockChain) TicketPoolValue() (dcrutil.Amount, error) {
|
||||
b.chainLock.RUnlock()
|
||||
|
||||
var amt int64
|
||||
err := b.db.View(func(dbTx database.Tx) error {
|
||||
err := b.utxoDb.View(func(dbTx database.Tx) error {
|
||||
for _, hash := range sn.LiveTickets() {
|
||||
outpoint := wire.OutPoint{Hash: hash, Index: 0, Tree: wire.TxTreeStake}
|
||||
utxo, err := b.utxoCache.FetchEntry(dbTx, outpoint)
|
||||
|
||||
@ -3183,6 +3183,169 @@ func upgradeToVersion10(ctx context.Context, db database.DB, dbInfo *databaseInf
|
||||
return nil
|
||||
}
|
||||
|
||||
// separateUtxoDatabase moves the UTXO set and state from the block database to
|
||||
// the UTXO database.
|
||||
func separateUtxoDatabase(ctx context.Context, b *BlockChain) error {
|
||||
// Hardcoded bucket and key names so updates do not affect old upgrades.
|
||||
v1UtxoSetStateKeyName := []byte("utxosetstate")
|
||||
v3UtxoSetBucketName := []byte("utxosetv3")
|
||||
|
||||
log.Info("Migrating UTXO database. This may take a while...")
|
||||
start := time.Now()
|
||||
|
||||
// doBatch contains the primary logic for migrating the UTXO database. This
|
||||
// is done because attempting to migrate in a single database transaction
|
||||
// could result in massive memory usage and could potentially crash on many
|
||||
// systems due to ulimits.
|
||||
//
|
||||
// It returns whether or not all entries have been updated.
|
||||
const maxEntries = 20000
|
||||
var resumeOffset uint32
|
||||
var totalMigrated uint64
|
||||
var err error
|
||||
doBatch := func(dbTx database.Tx, utxoDbTx database.Tx) (bool, error) {
|
||||
// Get the UTXO set bucket for both the old and the new database.
|
||||
v3UtxoSetBucketOldDb := dbTx.Metadata().Bucket(v3UtxoSetBucketName)
|
||||
if v3UtxoSetBucketOldDb == nil {
|
||||
// If the UTXO set doesn't exist in the old database, return immediately
|
||||
// as there is nothing to do.
|
||||
return true, nil
|
||||
}
|
||||
v3UtxoSetBucketNewDb := utxoDbTx.Metadata().Bucket(v3UtxoSetBucketName)
|
||||
if v3UtxoSetBucketNewDb == nil {
|
||||
return false, fmt.Errorf("bucket %s does not exist",
|
||||
v3UtxoSetBucketName)
|
||||
}
|
||||
|
||||
// Migrate UTXO set entries so long as the max number of entries for this
|
||||
// batch has not been exceeded.
|
||||
var logProgress bool
|
||||
var numMigrated, numIterated uint32
|
||||
cursor := v3UtxoSetBucketOldDb.Cursor()
|
||||
for ok := cursor.Last(); ok; ok = cursor.Prev() {
|
||||
// Reset err on each iteration.
|
||||
err = nil
|
||||
|
||||
if interruptRequested(ctx) {
|
||||
logProgress = true
|
||||
err = errInterruptRequested
|
||||
break
|
||||
}
|
||||
|
||||
if numMigrated >= maxEntries {
|
||||
logProgress = true
|
||||
err = errBatchFinished
|
||||
break
|
||||
}
|
||||
|
||||
// Skip entries that have already been migrated in previous batches.
|
||||
numIterated++
|
||||
if numIterated-1 < resumeOffset {
|
||||
continue
|
||||
}
|
||||
resumeOffset++
|
||||
|
||||
// Create the new entry in the V3 bucket.
|
||||
err = v3UtxoSetBucketNewDb.Put(cursor.Key(), cursor.Value())
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
numMigrated++
|
||||
}
|
||||
isFullyDone := err == nil
|
||||
if (isFullyDone || logProgress) && numMigrated > 0 {
|
||||
totalMigrated += uint64(numMigrated)
|
||||
log.Infof("Migrated %d entries (%d total)", numMigrated,
|
||||
totalMigrated)
|
||||
}
|
||||
return isFullyDone, err
|
||||
}
|
||||
|
||||
// Migrate all entries in batches for the reasons mentioned above.
|
||||
var isFullyDone bool
|
||||
for !isFullyDone {
|
||||
err := b.db.View(func(dbTx database.Tx) error {
|
||||
err := b.utxoDb.Update(func(utxoDbTx database.Tx) error {
|
||||
var err error
|
||||
isFullyDone, err = doBatch(dbTx, utxoDbTx)
|
||||
if errors.Is(err, errInterruptRequested) ||
|
||||
errors.Is(err, errBatchFinished) {
|
||||
|
||||
// No error here so the database transaction is not cancelled
|
||||
// and therefore outstanding work is written to disk. The outer
|
||||
// function will exit with an interrupted error below due to
|
||||
// another interrupted check.
|
||||
err = nil
|
||||
}
|
||||
return err
|
||||
})
|
||||
return err
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if interruptRequested(ctx) {
|
||||
return errInterruptRequested
|
||||
}
|
||||
}
|
||||
|
||||
elapsed := time.Since(start).Round(time.Millisecond)
|
||||
log.Infof("Done migrating UTXO database. Total entries: %d in %v",
|
||||
totalMigrated, elapsed)
|
||||
|
||||
if interruptRequested(ctx) {
|
||||
return errInterruptRequested
|
||||
}
|
||||
|
||||
// Drop UTXO set from the old database.
|
||||
log.Info("Removing old UTXO set...")
|
||||
start = time.Now()
|
||||
err = incrementalFlatDrop(ctx, b.db, v3UtxoSetBucketName, "old UTXO set")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
err = b.db.Update(func(dbTx database.Tx) error {
|
||||
return dbTx.Metadata().DeleteBucket(v3UtxoSetBucketName)
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
elapsed = time.Since(start).Round(time.Millisecond)
|
||||
log.Infof("Done removing old UTXO set in %v", elapsed)
|
||||
|
||||
// Get the UTXO set state from the old database.
|
||||
var serialized []byte
|
||||
err = b.db.View(func(dbTx database.Tx) error {
|
||||
serialized = dbTx.Metadata().Get(v1UtxoSetStateKeyName)
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if serialized != nil {
|
||||
// Set the UTXO set state in the new database.
|
||||
err = b.utxoDb.Update(func(utxoDbTx database.Tx) error {
|
||||
return utxoDbTx.Metadata().Put(v1UtxoSetStateKeyName, serialized)
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Delete the UTXO set state from the old database.
|
||||
err = b.db.Update(func(dbTx database.Tx) error {
|
||||
return dbTx.Metadata().Delete(v1UtxoSetStateKeyName)
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// checkDBTooOldToUpgrade returns an ErrDBTooOldToUpgrade error if the provided
|
||||
// database version can no longer be upgraded due to being too old.
|
||||
func checkDBTooOldToUpgrade(dbVersion uint32) error {
|
||||
@ -3306,7 +3469,8 @@ func upgradeDB(ctx context.Context, db database.DB, chainParams *chaincfg.Params
|
||||
// 3, so spend journal upgrades prior to version 3 are handled in the overall
|
||||
// block database upgrade path.
|
||||
//
|
||||
// NOTE: The passed database info will be updated with the latest versions.
|
||||
// NOTE: The database info housed by the passed blockchain instance will be
|
||||
// updated with the latest versions.
|
||||
func upgradeSpendJournal(ctx context.Context, b *BlockChain) error {
|
||||
// Update to a version 3 spend journal as needed.
|
||||
if b.dbInfo.stxoVer == 2 {
|
||||
@ -3324,7 +3488,8 @@ func upgradeSpendJournal(ctx context.Context, b *BlockChain) error {
|
||||
// so utxo set upgrades prior to version 3 are handled in the block database
|
||||
// upgrade path.
|
||||
//
|
||||
// NOTE: The passed database info will be updated with the latest versions.
|
||||
// NOTE: The database info housed by the passed blockchain instance will be
|
||||
// updated with the latest versions.
|
||||
func upgradeUtxoDb(ctx context.Context, b *BlockChain) error {
|
||||
// Update to a version 3 utxo set as needed.
|
||||
if b.utxoDbInfo.utxoVer == 2 {
|
||||
@ -3333,5 +3498,23 @@ func upgradeUtxoDb(ctx context.Context, b *BlockChain) error {
|
||||
}
|
||||
}
|
||||
|
||||
// Check if the block database contains the UTXO set or state. If it does,
|
||||
// move the UTXO set and state from the block database to the UTXO database.
|
||||
blockDbUtxoSetExists := false
|
||||
b.db.View(func(dbTx database.Tx) error {
|
||||
if dbTx.Metadata().Bucket([]byte("utxosetv3")) != nil ||
|
||||
dbTx.Metadata().Get([]byte("utxosetstate")) != nil {
|
||||
|
||||
blockDbUtxoSetExists = true
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if blockDbUtxoSetExists {
|
||||
err := separateUtxoDatabase(ctx, b)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@ -920,7 +920,7 @@ func (b *BlockChain) FetchUtxoEntry(outpoint wire.OutPoint) (*UtxoEntry, error)
|
||||
defer b.chainLock.RUnlock()
|
||||
|
||||
var entry *UtxoEntry
|
||||
err := b.db.View(func(dbTx database.Tx) error {
|
||||
err := b.utxoDb.View(func(dbTx database.Tx) error {
|
||||
var err error
|
||||
entry, err = b.utxoCache.FetchEntry(dbTx, outpoint)
|
||||
return err
|
||||
@ -949,7 +949,7 @@ type UtxoStats struct {
|
||||
// always be in sync with the best block.
|
||||
func (b *BlockChain) FetchUtxoStats() (*UtxoStats, error) {
|
||||
var stats *UtxoStats
|
||||
err := b.db.View(func(dbTx database.Tx) error {
|
||||
err := b.utxoDb.View(func(dbTx database.Tx) error {
|
||||
var err error
|
||||
stats, err = dbFetchUtxoStats(dbTx)
|
||||
return err
|
||||
|
||||
@ -1118,7 +1118,7 @@ func TestInitialize(t *testing.T) {
|
||||
// gets created and initialized at startup.
|
||||
resetTestUtxoCache := func() *testUtxoCache {
|
||||
testUtxoCache := newTestUtxoCache(&UtxoCacheConfig{
|
||||
DB: g.chain.db,
|
||||
DB: g.chain.utxoDb,
|
||||
MaxSize: 100 * 1024 * 1024, // 100 MiB
|
||||
})
|
||||
g.chain.utxoCache = testUtxoCache
|
||||
@ -1260,7 +1260,7 @@ func TestShutdownUtxoCache(t *testing.T) {
|
||||
// Replace the chain utxo cache with a test cache so that flushing can be
|
||||
// disabled.
|
||||
testUtxoCache := newTestUtxoCache(&UtxoCacheConfig{
|
||||
DB: g.chain.db,
|
||||
DB: g.chain.utxoDb,
|
||||
MaxSize: 100 * 1024 * 1024, // 100 MiB
|
||||
})
|
||||
g.chain.utxoCache = testUtxoCache
|
||||
|
||||
@ -298,7 +298,7 @@ func TestCheckBlockHeaderContext(t *testing.T) {
|
||||
ChainParams: params,
|
||||
TimeSource: NewMedianTime(),
|
||||
UtxoCache: NewUtxoCache(&UtxoCacheConfig{
|
||||
DB: db,
|
||||
DB: utxoDb,
|
||||
MaxSize: 100 * 1024 * 1024, // 100 MiB
|
||||
}),
|
||||
})
|
||||
|
||||
@ -3403,7 +3403,7 @@ func newServer(ctx context.Context, listenAddrs []string, db database.DB,
|
||||
|
||||
// Create a new block chain instance with the appropriate configuration.
|
||||
utxoCache := blockchain.NewUtxoCache(&blockchain.UtxoCacheConfig{
|
||||
DB: s.db,
|
||||
DB: utxoDb,
|
||||
MaxSize: uint64(cfg.UtxoCacheMaxSize) * 1024 * 1024,
|
||||
})
|
||||
s.chain, err = blockchain.New(ctx,
|
||||
|
||||
Loading…
Reference in New Issue
Block a user