eth/protocols/snap: attempt at refactoring for clarity

This commit is contained in:
Martin Holst Swende 2023-10-06 09:51:03 +02:00
parent 052355f5e2
commit c302a35772
No known key found for this signature in database
GPG key ID: 683B438C05A5DDF0

View file

@ -226,6 +226,37 @@ type storageResponse struct {
cont bool // Whether the last storage range has a continuation cont bool // Whether the last storage range has a continuation
} }
// lookupAccountRootFromRequest returns the account corresponding to the hash h, or nil if not present.
// It also marks whether healing is needed or not.
func (res *storageResponse) lookupAccountRootFromRequest(h common.Hash) (chunked bool, storageRoot common.Hash) {
for j, hash := range res.mainTask.res.hashes {
if h != hash {
continue
}
// If it is the last account, AND the 'cont' is true, the this is chunked
chunked = (j == len(res.hashes)-1) && res.cont
// State was delivered, if complete mark as not needed any more, otherwise
// mark the account as needing healing
// If the packet contains multiple contract storage slots, all
// but the last are surely complete. The last contract may be
// chunked, so check it's continuation flag.
if res.subTask == nil && res.mainTask.needState[j] && !chunked {
res.mainTask.needState[j] = false
res.mainTask.pend--
}
// If the last contract was chunked, mark it as needing healing
// to avoid writing it out to disk prematurely.
if res.subTask == nil && !res.mainTask.needHeal[j] && chunked {
res.mainTask.needHeal[j] = true
}
return chunked, res.mainTask.res.accounts[j].Root
}
return false, common.Hash{}
}
// trienodeHealRequest tracks a pending state trie request to ensure responses // trienodeHealRequest tracks a pending state trie request to ensure responses
// are to actual requests and to validate any security constraints. // are to actual requests and to validate any security constraints.
// //
@ -1947,149 +1978,72 @@ func (s *Syncer) processStorageResponse(res *storageResponse) {
res.mainTask.stateTasks[account] = res.roots[i] res.mainTask.stateTasks[account] = res.roots[i]
continue continue
} }
// State was delivered, if complete mark as not needed any more, otherwise // Look up the storage root, and also process healing status
// mark the account as needing healing chunked, storageRoot := res.lookupAccountRootFromRequest(account)
for j, hash := range res.mainTask.res.hashes { if storageRoot == (common.Hash{}) {
if account != hash { // Programming error - this should never happen
continue log.Warn("Delivered account not found in request", "account", account)
} continue
acc := res.mainTask.res.accounts[j] }
// If the last contract was chunked, we need to switch to large
// If the packet contains multiple contract storage slots, all // contract handling mode if we're not already doing that
// but the last are surely complete. The last contract may be if chunked {
// chunked, so check it's continuation flag. s.processChunkMode(res, account, storageRoot, res.hashes[i])
if res.subTask == nil && res.mainTask.needState[j] && (i < len(res.hashes)-1 || !res.cont) { }
res.mainTask.needState[j] = false // If we're in large contract delivery mode, forward the subtask
res.mainTask.pend-- if res.subTask != nil {
} // Ensure the response doesn't overflow into the subsequent task
// If the last contract was chunked, mark it as needing healing last := res.subTask.Last.Big()
// to avoid writing it out to disk prematurely. // Find the first overflowing key. While at it, mark res as complete
if res.subTask == nil && !res.mainTask.needHeal[j] && i == len(res.hashes)-1 && res.cont { // if we find the range to include or pass the 'last'
res.mainTask.needHeal[j] = true index := sort.Search(len(res.hashes[i]), func(k int) bool {
} cmp := res.hashes[i][k].Big().Cmp(last)
// If the last contract was chunked, we need to switch to large if cmp >= 0 {
// contract handling mode res.cont = false
if res.subTask == nil && i == len(res.hashes)-1 && res.cont {
// If we haven't yet started a large-contract retrieval, create
// the subtasks for it within the main account task
if tasks, ok := res.mainTask.SubTasks[account]; !ok {
var (
keys = res.hashes[i]
chunks = uint64(storageConcurrency)
lastKey common.Hash
)
if len(keys) > 0 {
lastKey = keys[len(keys)-1]
}
// If the number of slots remaining is low, decrease the
// number of chunks. Somewhere on the order of 10-15K slots
// fit into a packet of 500KB. A key/slot pair is maximum 64
// bytes, so pessimistically maxRequestSize/64 = 8K.
//
// Chunk so that at least 2 packets are needed to fill a task.
if estimate, err := estimateRemainingSlots(len(keys), lastKey); err == nil {
if n := estimate / (2 * (maxRequestSize / 64)); n+1 < chunks {
chunks = n + 1
}
log.Debug("Chunked large contract", "initiators", len(keys), "tail", lastKey, "remaining", estimate, "chunks", chunks)
} else {
log.Debug("Chunked large contract", "initiators", len(keys), "tail", lastKey, "chunks", chunks)
}
r := newHashRange(lastKey, chunks)
// Our first task is the one that was just filled by this response.
batch := ethdb.HookedBatch{
Batch: s.db.NewBatch(),
OnPut: func(key []byte, value []byte) {
s.storageBytes += common.StorageSize(len(key) + len(value))
},
}
tasks = append(tasks, &storageTask{
Next: common.Hash{},
Last: r.End(),
root: acc.Root,
genBatch: batch,
genTrie: trie.NewStackTrieWithOwner(func(owner common.Hash, path []byte, hash common.Hash, val []byte) {
rawdb.WriteTrieNode(batch, owner, path, hash, val, s.scheme)
}, account),
})
for r.Next() {
batch := ethdb.HookedBatch{
Batch: s.db.NewBatch(),
OnPut: func(key []byte, value []byte) {
s.storageBytes += common.StorageSize(len(key) + len(value))
},
}
tasks = append(tasks, &storageTask{
Next: r.Start(),
Last: r.End(),
root: acc.Root,
genBatch: batch,
genTrie: trie.NewStackTrieWithOwner(func(owner common.Hash, path []byte, hash common.Hash, val []byte) {
rawdb.WriteTrieNode(batch, owner, path, hash, val, s.scheme)
}, account),
})
}
for _, task := range tasks {
log.Debug("Created storage sync task", "account", account, "root", acc.Root, "from", task.Next, "last", task.Last)
}
res.mainTask.SubTasks[account] = tasks
// Since we've just created the sub-tasks, this response
// is surely for the first one (zero origin)
res.subTask = tasks[0]
} }
return cmp > 0
})
if index >= 0 {
// cut off excess
res.hashes[i] = res.hashes[i][:index]
res.slots[i] = res.slots[i][:index]
} }
// If we're in large contract delivery mode, forward the subtask // Forward the relevant storage chunk (even if created just now)
if res.subTask != nil { if res.cont {
// Ensure the response doesn't overflow into the subsequent task res.subTask.Next = incHash(res.hashes[i][len(res.hashes[i])-1])
last := res.subTask.Last.Big() } else {
// Find the first overflowing key. While at it, mark res as complete res.subTask.done = true
// if we find the range to include or pass the 'last'
index := sort.Search(len(res.hashes[i]), func(k int) bool {
cmp := res.hashes[i][k].Big().Cmp(last)
if cmp >= 0 {
res.cont = false
}
return cmp > 0
})
if index >= 0 {
// cut off excess
res.hashes[i] = res.hashes[i][:index]
res.slots[i] = res.slots[i][:index]
}
// Forward the relevant storage chunk (even if created just now)
if res.cont {
res.subTask.Next = incHash(res.hashes[i][len(res.hashes[i])-1])
} else {
res.subTask.done = true
}
} }
} }
// Iterate over all the complete contracts, reconstruct the trie nodes and // Iterate over all the complete contracts, reconstruct the trie nodes and
// push them to disk. If the contract is chunked, the trie nodes will be // push them to disk. If the contract is chunked, the trie nodes will be
// reconstructed later. // reconstructed later.
slots += len(res.hashes[i]) slots += len(res.hashes[i])
var (
keys = res.hashes[i]
values = res.slots[i]
)
if i < len(res.hashes)-1 || res.subTask == nil { if i < len(res.hashes)-1 || res.subTask == nil {
tr := trie.NewStackTrieWithOwner(func(owner common.Hash, path []byte, hash common.Hash, val []byte) { tr := trie.NewStackTrieWithOwner(func(owner common.Hash, path []byte, hash common.Hash, val []byte) {
rawdb.WriteTrieNode(batch, owner, path, hash, val, s.scheme) rawdb.WriteTrieNode(batch, owner, path, hash, val, s.scheme)
}, account) }, account)
for j := 0; j < len(res.hashes[i]); j++ { for j := 0; j < len(keys); j++ {
tr.Update(res.hashes[i][j][:], res.slots[i][j]) tr.Update(keys[j][:], values[j])
} }
tr.Commit() tr.Commit()
} }
// Persist the received storage segments. These flat state maybe // Persist the received storage segments. These flat state maybe
// outdated during the sync, but it can be fixed later during the // outdated during the sync, but it can be fixed later during the
// snapshot generation. // snapshot generation.
for j := 0; j < len(res.hashes[i]); j++ { for j := 0; j < len(keys); j++ {
rawdb.WriteStorageSnapshot(batch, account, res.hashes[i][j], res.slots[i][j]) rawdb.WriteStorageSnapshot(batch, account, keys[j], values[j])
// If we're storing large contracts, generate the trie nodes // If we're storing large contracts, generate the trie nodes
// on the fly to not trash the gluing points // on the fly to not trash the gluing points
if i == len(res.hashes)-1 && res.subTask != nil { if i == len(res.hashes)-1 && res.subTask != nil {
res.subTask.genTrie.Update(res.hashes[i][j][:], res.slots[i][j]) res.subTask.genTrie.Update(keys[j][:], values[j])
} }
} }
} }
@ -2132,6 +2086,83 @@ func (s *Syncer) processStorageResponse(res *storageResponse) {
// task assigners to pick up and fill. // task assigners to pick up and fill.
} }
func (s *Syncer) processChunkMode(res *storageResponse, account common.Hash, storageRoot common.Hash, keys []common.Hash) {
if res.subTask != nil {
return // Already started
}
// If we haven't yet started a large-contract retrieval, create
// the subtasks for it within the main account task
tasks, ok := res.mainTask.SubTasks[account]
if ok {
return // Already started
}
var (
chunks = uint64(storageConcurrency)
lastKey common.Hash
)
if len(keys) > 0 {
lastKey = keys[len(keys)-1]
}
// If the number of slots remaining is low, decrease the
// number of chunks. Somewhere on the order of 10-15K slots
// fit into a packet of 500KB. A key/slot pair is maximum 64
// bytes, so pessimistically maxRequestSize/64 = 8K.
//
// Chunk so that at least 2 packets are needed to fill a task.
if estimate, err := estimateRemainingSlots(len(keys), lastKey); err == nil {
if n := estimate / (2 * (maxRequestSize / 64)); n+1 < chunks {
chunks = n + 1
}
log.Debug("Chunked large contract", "initiators", len(keys), "tail", lastKey, "remaining", estimate, "chunks", chunks)
} else {
log.Debug("Chunked large contract", "initiators", len(keys), "tail", lastKey, "chunks", chunks)
}
r := newHashRange(lastKey, chunks)
// Our first task is the one that was just filled by this response.
batch := ethdb.HookedBatch{
Batch: s.db.NewBatch(),
OnPut: func(key []byte, value []byte) {
s.storageBytes += common.StorageSize(len(key) + len(value))
},
}
tasks = append(tasks, &storageTask{
Next: common.Hash{},
Last: r.End(),
root: storageRoot,
genBatch: batch,
genTrie: trie.NewStackTrieWithOwner(func(owner common.Hash, path []byte, hash common.Hash, val []byte) {
rawdb.WriteTrieNode(batch, owner, path, hash, val, s.scheme)
}, account),
})
for r.Next() {
batch := ethdb.HookedBatch{
Batch: s.db.NewBatch(),
OnPut: func(key []byte, value []byte) {
s.storageBytes += common.StorageSize(len(key) + len(value))
},
}
tasks = append(tasks, &storageTask{
Next: r.Start(),
Last: r.End(),
root: storageRoot,
genBatch: batch,
genTrie: trie.NewStackTrieWithOwner(func(owner common.Hash, path []byte, hash common.Hash, val []byte) {
rawdb.WriteTrieNode(batch, owner, path, hash, val, s.scheme)
}, account),
})
}
for _, task := range tasks {
log.Debug("Created storage sync task", "account", account, "root", storageRoot, "from", task.Next, "last", task.Last)
}
res.mainTask.SubTasks[account] = tasks
// Since we've just created the sub-tasks, this response
// is surely for the first one (zero origin)
res.subTask = tasks[0]
}
// processTrienodeHealResponse integrates an already validated trienode response // processTrienodeHealResponse integrates an already validated trienode response
// into the healer tasks. // into the healer tasks.
func (s *Syncer) processTrienodeHealResponse(res *trienodeHealResponse) { func (s *Syncer) processTrienodeHealResponse(res *trienodeHealResponse) {