mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-18 18:02:24 +00:00
live: skip if freeze failed
Signed-off-by: jsvisa <delweng@gmail.com>
This commit is contained in:
parent
36c98a31e2
commit
c8dfa31422
1 changed files with 60 additions and 40 deletions
|
|
@ -15,19 +15,27 @@ const (
|
||||||
kvdbTailKey = "FilterFreezerTail"
|
kvdbTailKey = "FilterFreezerTail"
|
||||||
)
|
)
|
||||||
|
|
||||||
func (f *live) freeze(maxKeepBlocks uint64) {
|
func (l *live) freeze(maxKeepBlocks uint64) {
|
||||||
var lastFinalized uint64
|
var lastFinalized uint64
|
||||||
|
var freezeErr error
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-f.stopCh:
|
case <-l.stopCh:
|
||||||
return
|
return
|
||||||
case finalizedBlock := <-f.freezeCh:
|
case finalizedBlock := <-l.freezeCh:
|
||||||
|
// Skip if the finalized block is not increasing
|
||||||
if finalizedBlock <= lastFinalized {
|
if finalizedBlock <= lastFinalized {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
lastFinalized = finalizedBlock
|
lastFinalized = finalizedBlock
|
||||||
|
|
||||||
tail := f.getFreezerTail()
|
// Check if error occurred in previous iteration
|
||||||
|
if freezeErr != nil {
|
||||||
|
log.Error("Error occurred in previous freezing, checking the log for more detail", "error", freezeErr)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
tail := l.getFreezerTail()
|
||||||
|
|
||||||
// Freeze at most freezeThreshold blocks
|
// Freeze at most freezeThreshold blocks
|
||||||
freezeUpTo := min(finalizedBlock, tail+freezeThreshold)
|
freezeUpTo := min(finalizedBlock, tail+freezeThreshold)
|
||||||
|
|
@ -36,75 +44,87 @@ func (f *live) freeze(maxKeepBlocks uint64) {
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Info("Move traces from kvdb to frdb", "from", tail, "to", freezeUpTo-1)
|
log.Info("Move traces from kvdb to frdb", "from", tail, "to", freezeUpTo-1)
|
||||||
|
|
||||||
for blknum := tail; blknum < freezeUpTo; blknum++ {
|
for blknum := tail; blknum < freezeUpTo; blknum++ {
|
||||||
if err := f.moveBlockToFreezer(blknum); err != nil {
|
freezeErr = l.moveBlockToFreezer(blknum)
|
||||||
log.Error("Failed to move block to freezer", "block", blknum, "error", err)
|
if freezeErr != nil {
|
||||||
|
log.Error("Failed to move block to freezer", "block", blknum, "error", freezeErr)
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
if freezeErr != nil {
|
||||||
// Update the tail of the freezer
|
continue
|
||||||
tail = freezeUpTo
|
|
||||||
if err := f.updateFreezerTail(tail); err != nil {
|
|
||||||
log.Error("Failed to update freezer tail", "error", err)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
frozen, _ := f.frdb.Ancients()
|
// Update the tail of the freezer
|
||||||
offset := f.offset.Load()
|
if err := l.updateFreezerTail(freezeUpTo); err != nil {
|
||||||
|
log.Warn("Failed to update freezer tail", "old", tail, "new", freezeUpTo, "error", err)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
// No pruning
|
// No need to prune
|
||||||
if maxKeepBlocks == 0 || frozen <= maxKeepBlocks {
|
if maxKeepBlocks == 0 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
frozen, err := l.frdb.Ancients()
|
||||||
|
if err != nil {
|
||||||
|
log.Error("Failed to get number of ancient items", "error", err)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
// Not enough blocks to prune
|
||||||
|
if frozen <= maxKeepBlocks {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
// Prune old blocks if necessary
|
// Prune old blocks if necessary
|
||||||
itemsToPrune := min(freezeThreshold, frozen-maxKeepBlocks)
|
itemsToPrune := min(freezeThreshold, frozen-maxKeepBlocks)
|
||||||
head := offset + itemsToPrune - 1
|
from := l.offset.Load()
|
||||||
log.Info("Prune old blocks", "pruned", itemsToPrune, "from", offset, "to", head)
|
head := from + itemsToPrune
|
||||||
if err := f.pruneBlocksFromFreezer(frozen-itemsToPrune, head); err != nil {
|
log.Info("Prune old blocks", "pruned", itemsToPrune, "from", from, "to", head)
|
||||||
|
if err := l.pruneBlocksFromFreezer(frozen-itemsToPrune, head); err != nil {
|
||||||
log.Error("Failed to prune blocks from freezer", "error", err)
|
log.Error("Failed to prune blocks from freezer", "error", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *live) getFreezerTail() (tail uint64) {
|
func (l *live) getFreezerTail() (tail uint64) {
|
||||||
tailBytes, _ := f.kvdb.Get([]byte(kvdbTailKey))
|
tailBytes, _ := l.kvdb.Get([]byte(kvdbTailKey))
|
||||||
|
|
||||||
if len(tailBytes) > 0 {
|
if len(tailBytes) > 0 {
|
||||||
tail = binary.BigEndian.Uint64(tailBytes)
|
tail = binary.BigEndian.Uint64(tailBytes)
|
||||||
} else {
|
} else {
|
||||||
// If tail is 0 (not found in kvdb), use the offset
|
// If tail is 0 (not found in kvdb), use the offset
|
||||||
tail = f.offset.Load()
|
tail = l.offset.Load()
|
||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *live) updateFreezerTail(tail uint64) error {
|
func (l *live) updateFreezerTail(tail uint64) error {
|
||||||
tailBytes := make([]byte, 8)
|
tailBytes := make([]byte, 8)
|
||||||
binary.BigEndian.PutUint64(tailBytes, tail)
|
binary.BigEndian.PutUint64(tailBytes, tail)
|
||||||
return f.kvdb.Put([]byte(kvdbTailKey), tailBytes)
|
return l.kvdb.Put([]byte(kvdbTailKey), tailBytes)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *live) moveBlockToFreezer(blknum uint64) error {
|
func (l *live) moveBlockToFreezer(blknum uint64) error {
|
||||||
header, err := f.backend.HeaderByNumber(context.Background(), rpc.BlockNumber(blknum))
|
header, err := l.backend.HeaderByNumber(context.Background(), rpc.BlockNumber(blknum))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
offset := f.offset.Load()
|
offset := l.offset.Load()
|
||||||
|
|
||||||
size, err := f.frdb.ModifyAncients(func(op ethdb.AncientWriteOp) error {
|
size, err := l.frdb.ModifyAncients(func(op ethdb.AncientWriteOp) error {
|
||||||
for name := range f.tracer.Tracers() {
|
for name := range l.tracer.Tracers() {
|
||||||
kvKey := toKVKey(name, blknum, header.Hash())
|
kvKey := toKVKey(name, blknum, header.Hash())
|
||||||
data, err := f.kvdb.Get(kvKey)
|
data, err := l.kvdb.Get(kvKey)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
table := toTraceTable(name)
|
if err := op.AppendRaw(toTraceTable(name), blknum-offset, data); err != nil {
|
||||||
err = op.AppendRaw(table, blknum-offset, data)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -116,17 +136,17 @@ func (f *live) moveBlockToFreezer(blknum uint64) error {
|
||||||
log.Info("Move from kvdb to frdb", "blknum", blknum, "size", size)
|
log.Info("Move from kvdb to frdb", "blknum", blknum, "size", size)
|
||||||
|
|
||||||
// Delete all entries for this prefix from kvdb, ignore error
|
// Delete all entries for this prefix from kvdb, ignore error
|
||||||
if err := f.deleteKVDBEntriesWithPrefix(blknum); err != nil {
|
if err := l.deleteKVDBEntriesWithPrefix(blknum); err != nil {
|
||||||
log.Error("Failed to delete entries from kvdb", "error", err)
|
log.Error("Failed to delete entries from kvdb", "error", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *live) deleteKVDBEntriesWithPrefix(blknum uint64) error {
|
func (l *live) deleteKVDBEntriesWithPrefix(blknum uint64) error {
|
||||||
prefix := encodeNumber(blknum)
|
prefix := encodeNumber(blknum)
|
||||||
batch := f.kvdb.NewBatch()
|
batch := l.kvdb.NewBatch()
|
||||||
it := f.kvdb.NewIterator(prefix, nil)
|
it := l.kvdb.NewIterator(prefix, nil)
|
||||||
defer it.Release()
|
defer it.Release()
|
||||||
|
|
||||||
for it.Next() {
|
for it.Next() {
|
||||||
|
|
@ -157,11 +177,11 @@ func (f *live) deleteKVDBEntriesWithPrefix(blknum uint64) error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *live) pruneBlocksFromFreezer(items, head uint64) error {
|
func (l *live) pruneBlocksFromFreezer(items, newHead uint64) error {
|
||||||
if _, err := f.frdb.TruncateHead(items); err != nil {
|
if _, err := l.frdb.TruncateHead(items); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
// Head should be in sync with the on-mem offset
|
// Set the offset of the new head
|
||||||
f.offset.Store(head)
|
l.offset.Store(newHead)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue