eth/protocols/snap: address comments from peter

This commit is contained in:
Gary Rong 2024-04-10 09:59:22 +08:00
parent a6629dfa96
commit e66eb890af

View file

@ -327,9 +327,14 @@ type accountTask struct {
done bool // Flag whether the task can be removed done bool // Flag whether the task can be removed
} }
// effectiveSubTasks returns a set of storage tasks covered by the responded // activeSubTasks returns the set of storage tasks covered by the current account
// account range. Nil is returned if the account range is empty. // range. Normally this would be the entire subTask set, but on a sync interrupt
func (task *accountTask) effectiveSubTasks() map[common.Hash][]*storageTask { // and later resume it can happen that a shorter account range is retrieved. This
// method ensures that we only start up the subtasks covered by the latest account
// response.
//
// Nil is returned if the account range is empty.
func (task *accountTask) activeSubTasks() map[common.Hash][]*storageTask {
if len(task.res.hashes) == 0 { if len(task.res.hashes) == 0 {
return nil return nil
} }
@ -338,6 +343,7 @@ func (task *accountTask) effectiveSubTasks() map[common.Hash][]*storageTask {
last = task.res.hashes[len(task.res.hashes)-1] last = task.res.hashes[len(task.res.hashes)-1]
) )
for hash, subTasks := range task.SubTasks { for hash, subTasks := range task.SubTasks {
subTasks := subTasks // closure
if hash.Cmp(last) <= 0 { if hash.Cmp(last) <= 0 {
tasks[hash] = subTasks tasks[hash] = subTasks
} }
@ -777,6 +783,8 @@ func (s *Syncer) loadSyncStatus() {
for _, hash := range task.StorageCompleted { for _, hash := range task.StorageCompleted {
task.stateCompleted[hash] = struct{}{} task.stateCompleted[hash] = struct{}{}
} }
task.StorageCompleted = nil
// Allocate batch for account trie generation // Allocate batch for account trie generation
task.genBatch = ethdb.HookedBatch{ task.genBatch = ethdb.HookedBatch{
Batch: s.db.NewBatch(), Batch: s.db.NewBatch(),
@ -1257,7 +1265,7 @@ func (s *Syncer) assignStorageTasks(success chan *storageResponse, fail chan *st
continue continue
} }
// Skip tasks that are already retrieving (or done with) all small states // Skip tasks that are already retrieving (or done with) all small states
storageTasks := task.effectiveSubTasks() storageTasks := task.activeSubTasks()
if len(storageTasks) == 0 && len(task.stateTasks) == 0 { if len(storageTasks) == 0 && len(task.stateTasks) == 0 {
continue continue
} }
@ -1946,7 +1954,7 @@ func (s *Syncer) processAccountResponse(res *accountResponse) {
// The leftover storage tasks are not expected, unless system is // The leftover storage tasks are not expected, unless system is
// very wrong. // very wrong.
if _, ok := res.task.SubTasks[res.hashes[i]]; ok { if _, ok := res.task.SubTasks[res.hashes[i]]; ok {
log.Crit("Unexpected leftover storage tasks", "owner", res.hashes[i]) panic(fmt.Errorf("unexpected leftover storage tasks, owner: %x", res.hashes[i]))
} }
// Mark the healing tag if storage root node is inconsistent, or // Mark the healing tag if storage root node is inconsistent, or
// it's non-existent due to storage chunking. // it's non-existent due to storage chunking.
@ -2501,7 +2509,7 @@ func (s *Syncer) forwardAccountTask(task *accountTask) {
// Error out if there is any leftover completion flag. // Error out if there is any leftover completion flag.
if task.done && len(task.stateCompleted) != 0 { if task.done && len(task.stateCompleted) != 0 {
log.Crit("Storage completion flags should be emptied", "number", len(task.stateCompleted)) panic(fmt.Errorf("storage completion flags should be emptied, %d left", len(task.stateCompleted)))
} }
// Stack trie could have generated trie nodes, push them to disk (we need to // Stack trie could have generated trie nodes, push them to disk (we need to
// flush after finalizing task.done. It's fine even if we crash and lose this // flush after finalizing task.done. It's fine even if we crash and lose this