core/rawdb : freezer sync closed and timer loops

This commit is contained in:
ucwong 2020-04-27 05:36:58 +00:00
parent 1aa83290f5
commit a04ed31615
2 changed files with 135 additions and 112 deletions

View file

@ -29,17 +29,22 @@ import (
"github.com/ethereum/go-ethereum/ethdb/memorydb" "github.com/ethereum/go-ethereum/ethdb/memorydb"
"github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/log"
"github.com/olekukonko/tablewriter" "github.com/olekukonko/tablewriter"
"sync"
) )
// freezerdb is a database wrapper that enabled freezer data retrievals. // freezerdb is a database wrapper that enabled freezer data retrievals.
type freezerdb struct { type freezerdb struct {
ethdb.KeyValueStore ethdb.KeyValueStore
ethdb.AncientStore ethdb.AncientStore
quitChan chan struct{}
wg sync.WaitGroup
} }
// Close implements io.Closer, closing both the fast key-value store as well as // Close implements io.Closer, closing both the fast key-value store as well as
// the slow ancient tables. // the slow ancient tables.
func (frdb *freezerdb) Close() error { func (frdb *freezerdb) Close() error {
close(frdb.quitChan)
var errs []error var errs []error
if err := frdb.KeyValueStore.Close(); err != nil { if err := frdb.KeyValueStore.Close(); err != nil {
errs = append(errs, err) errs = append(errs, err)
@ -171,13 +176,22 @@ func NewDatabaseWithFreezer(db ethdb.KeyValueStore, freezer string, namespace st
// feezer. // feezer.
} }
} }
// Freezer is consistent with the key-value database, permit combining the two
go frdb.freeze(db)
return &freezerdb{ fdb := &freezerdb{
KeyValueStore: db, KeyValueStore: db,
AncientStore: frdb, //AncientStore: frdb,
}, nil quitChan: make(chan struct{}),
}
// Freezer is consistent with the key-value database, permit combining the two
fdb.wg.Add(1)
go func() {
defer fdb.wg.Done()
frdb.freeze(db, fdb.quitChan)
}()
fdb.AncientStore = frdb
return fdb, nil
} }
// NewMemoryDatabase creates an ephemeral in-memory key-value database without a // NewMemoryDatabase creates an ephemeral in-memory key-value database without a

View file

@ -251,125 +251,134 @@ func (f *freezer) Sync() error {
// //
// This functionality is deliberately broken off from block importing to avoid // This functionality is deliberately broken off from block importing to avoid
// incurring additional data shuffling delays on block propagation. // incurring additional data shuffling delays on block propagation.
func (f *freezer) freeze(db ethdb.KeyValueStore) { func (f *freezer) freeze(db ethdb.KeyValueStore, quitChan chan struct{}) {
nfdb := &nofreezedb{KeyValueStore: db} nfdb := &nofreezedb{KeyValueStore: db}
timer := time.NewTimer(freezerRecheckInterval)
defer timer.Stop()
for { for {
// Retrieve the freezing threshold. select {
hash := ReadHeadBlockHash(nfdb) case <-quitChan:
if hash == (common.Hash{}) { log.Info("Database freezer quit")
log.Debug("Current full block hash unavailable") // new chain, empty database return
time.Sleep(freezerRecheckInterval) case <-timer.C:
continue // Retrieve the freezing threshold.
} hash := ReadHeadBlockHash(nfdb)
number := ReadHeaderNumber(nfdb, hash)
switch {
case number == nil:
log.Error("Current full block number unavailable", "hash", hash)
time.Sleep(freezerRecheckInterval)
continue
case *number < params.ImmutabilityThreshold:
log.Debug("Current full block not old enough", "number", *number, "hash", hash, "delay", params.ImmutabilityThreshold)
time.Sleep(freezerRecheckInterval)
continue
case *number-params.ImmutabilityThreshold <= f.frozen:
log.Debug("Ancient blocks frozen already", "number", *number, "hash", hash, "frozen", f.frozen)
time.Sleep(freezerRecheckInterval)
continue
}
head := ReadHeader(nfdb, hash, *number)
if head == nil {
log.Error("Current full block unavailable", "number", *number, "hash", hash)
time.Sleep(freezerRecheckInterval)
continue
}
// Seems we have data ready to be frozen, process in usable batches
limit := *number - params.ImmutabilityThreshold
if limit-f.frozen > freezerBatchLimit {
limit = f.frozen + freezerBatchLimit
}
var (
start = time.Now()
first = f.frozen
ancients = make([]common.Hash, 0, limit)
)
for f.frozen < limit {
// Retrieves all the components of the canonical block
hash := ReadCanonicalHash(nfdb, f.frozen)
if hash == (common.Hash{}) { if hash == (common.Hash{}) {
log.Error("Canonical hash missing, can't freeze", "number", f.frozen) log.Debug("Current full block hash unavailable") // new chain, empty database
break timer.Reset(freezerRecheckInterval)
continue
} }
header := ReadHeaderRLP(nfdb, hash, f.frozen) number := ReadHeaderNumber(nfdb, hash)
if len(header) == 0 { switch {
log.Error("Block header missing, can't freeze", "number", f.frozen, "hash", hash) case number == nil:
break log.Error("Current full block number unavailable", "hash", hash)
timer.Reset(freezerRecheckInterval)
continue
case *number < params.ImmutabilityThreshold:
log.Debug("Current full block not old enough", "number", *number, "hash", hash, "delay", params.ImmutabilityThreshold)
timer.Reset(freezerRecheckInterval)
continue
case *number-params.ImmutabilityThreshold <= f.frozen:
log.Debug("Ancient blocks frozen already", "number", *number, "hash", hash, "frozen", f.frozen)
timer.Reset(freezerRecheckInterval)
continue
} }
body := ReadBodyRLP(nfdb, hash, f.frozen) head := ReadHeader(nfdb, hash, *number)
if len(body) == 0 { if head == nil {
log.Error("Block body missing, can't freeze", "number", f.frozen, "hash", hash) log.Error("Current full block unavailable", "number", *number, "hash", hash)
break timer.Reset(freezerRecheckInterval)
continue
} }
receipts := ReadReceiptsRLP(nfdb, hash, f.frozen) // Seems we have data ready to be frozen, process in usable batches
if len(receipts) == 0 { limit := *number - params.ImmutabilityThreshold
log.Error("Block receipts missing, can't freeze", "number", f.frozen, "hash", hash) if limit-f.frozen > freezerBatchLimit {
break limit = f.frozen + freezerBatchLimit
} }
td := ReadTdRLP(nfdb, hash, f.frozen) var (
if len(td) == 0 { start = time.Now()
log.Error("Total difficulty missing, can't freeze", "number", f.frozen, "hash", hash) first = f.frozen
break ancients = make([]common.Hash, 0, limit)
)
for f.frozen < limit {
// Retrieves all the components of the canonical block
hash := ReadCanonicalHash(nfdb, f.frozen)
if hash == (common.Hash{}) {
log.Error("Canonical hash missing, can't freeze", "number", f.frozen)
break
}
header := ReadHeaderRLP(nfdb, hash, f.frozen)
if len(header) == 0 {
log.Error("Block header missing, can't freeze", "number", f.frozen, "hash", hash)
break
}
body := ReadBodyRLP(nfdb, hash, f.frozen)
if len(body) == 0 {
log.Error("Block body missing, can't freeze", "number", f.frozen, "hash", hash)
break
}
receipts := ReadReceiptsRLP(nfdb, hash, f.frozen)
if len(receipts) == 0 {
log.Error("Block receipts missing, can't freeze", "number", f.frozen, "hash", hash)
break
}
td := ReadTdRLP(nfdb, hash, f.frozen)
if len(td) == 0 {
log.Error("Total difficulty missing, can't freeze", "number", f.frozen, "hash", hash)
break
}
log.Trace("Deep froze ancient block", "number", f.frozen, "hash", hash)
// Inject all the components into the relevant data tables
if err := f.AppendAncient(f.frozen, hash[:], header, body, receipts, td); err != nil {
break
}
ancients = append(ancients, hash)
} }
log.Trace("Deep froze ancient block", "number", f.frozen, "hash", hash) // Batch of blocks have been frozen, flush them before wiping from leveldb
// Inject all the components into the relevant data tables if err := f.Sync(); err != nil {
if err := f.AppendAncient(f.frozen, hash[:], header, body, receipts, td); err != nil { log.Crit("Failed to flush frozen tables", "err", err)
break
} }
ancients = append(ancients, hash) // Wipe out all data from the active database
} batch := db.NewBatch()
// Batch of blocks have been frozen, flush them before wiping from leveldb for i := 0; i < len(ancients); i++ {
if err := f.Sync(); err != nil { // Always keep the genesis block in active database
log.Crit("Failed to flush frozen tables", "err", err) if first+uint64(i) != 0 {
} DeleteBlockWithoutNumber(batch, ancients[i], first+uint64(i))
// Wipe out all data from the active database DeleteCanonicalHash(batch, first+uint64(i))
batch := db.NewBatch()
for i := 0; i < len(ancients); i++ {
// Always keep the genesis block in active database
if first+uint64(i) != 0 {
DeleteBlockWithoutNumber(batch, ancients[i], first+uint64(i))
DeleteCanonicalHash(batch, first+uint64(i))
}
}
if err := batch.Write(); err != nil {
log.Crit("Failed to delete frozen canonical blocks", "err", err)
}
batch.Reset()
// Wipe out side chain also.
for number := first; number < f.frozen; number++ {
// Always keep the genesis block in active database
if number != 0 {
for _, hash := range ReadAllHashes(db, number) {
DeleteBlock(batch, hash, number)
} }
} }
} if err := batch.Write(); err != nil {
if err := batch.Write(); err != nil { log.Crit("Failed to delete frozen canonical blocks", "err", err)
log.Crit("Failed to delete frozen side blocks", "err", err) }
} batch.Reset()
// Log something friendly for the user // Wipe out side chain also.
context := []interface{}{ for number := first; number < f.frozen; number++ {
"blocks", f.frozen - first, "elapsed", common.PrettyDuration(time.Since(start)), "number", f.frozen - 1, // Always keep the genesis block in active database
} if number != 0 {
if n := len(ancients); n > 0 { for _, hash := range ReadAllHashes(db, number) {
context = append(context, []interface{}{"hash", ancients[n-1]}...) DeleteBlock(batch, hash, number)
} }
log.Info("Deep froze chain segment", context...) }
}
if err := batch.Write(); err != nil {
log.Crit("Failed to delete frozen side blocks", "err", err)
}
// Log something friendly for the user
context := []interface{}{
"blocks", f.frozen - first, "elapsed", common.PrettyDuration(time.Since(start)), "number", f.frozen - 1,
}
if n := len(ancients); n > 0 {
context = append(context, []interface{}{"hash", ancients[n-1]}...)
}
log.Info("Deep froze chain segment", context...)
// Avoid database thrashing with tiny writes // Avoid database thrashing with tiny writes
if f.frozen-first < freezerBatchLimit { if f.frozen-first < freezerBatchLimit {
time.Sleep(freezerRecheckInterval) timer.Reset(freezerRecheckInterval)
} else {
timer.Reset(0)
}
} }
} }
} }