eth/fetcher: blob fetcher style edits

This commit is contained in:
Felix Lange 2026-06-12 19:57:04 +02:00
parent 4d107e77dc
commit 63fe3ff7ff

View file

@ -746,20 +746,20 @@ func (f *BlobFetcher) scheduleFetches(timer *mclock.Timer, timeout chan struct{}
if len(actives) == 0 { if len(actives) == 0 {
return return
} }
// For each active peer, try to schedule some payload fetches
idle := len(f.requests) == 0
wasIdle := len(f.requests) == 0
// For each active peer, try to schedule some payload fetches.
f.forEachPeer(actives, func(peer string) { f.forEachPeer(actives, func(peer string) {
if len(f.announces[peer]) == 0 || len(f.requests[peer]) != 0 { if len(f.announces[peer]) == 0 || len(f.requests[peer]) != 0 {
return // continue return // continue
} }
var ( var (
hashes = make([]common.Hash, 0, maxTxRetrievals) hashes []common.Hash
custodies = make([]*types.CustodyBitmap, 0, maxTxRetrievals) custodies []*types.CustodyBitmap
) )
f.forEachAnnounce(f.announces[peer], func(hash common.Hash, cells *types.CustodyBitmap) bool { f.forEachAnnounce(f.announces[peer], func(hash common.Hash, cells *types.CustodyBitmap) bool {
var unfetched *types.CustodyBitmap var unfetched *types.CustodyBitmap
if f.fetches[hash] == nil { if f.fetches[hash] == nil {
// tx is not being fetched // tx is not being fetched
unfetched = cells unfetched = cells
@ -795,6 +795,7 @@ func (f *BlobFetcher) scheduleFetches(timer *mclock.Timer, timeout chan struct{}
return len(hashes) < maxPayloadRetrievals return len(hashes) < maxPayloadRetrievals
}) })
// If any hashes were allocated, request them from the peer // If any hashes were allocated, request them from the peer
if len(hashes) > 0 { if len(hashes) > 0 {
// Group hashes by custody bitmap // Group hashes by custody bitmap
@ -817,7 +818,7 @@ func (f *BlobFetcher) scheduleFetches(timer *mclock.Timer, timeout chan struct{}
request = append(request, cr) request = append(request, cr)
} }
f.requests[peer] = request f.requests[peer] = request
go func(peer string, request []*cellRequest) { go func() {
for _, req := range request { for _, req := range request {
blobRequestOutMeter.Mark(int64(len(req.txs))) blobRequestOutMeter.Mark(int64(len(req.txs)))
if err := f.fn.FetchPayloads(peer, req.txs, req.cells); err != nil { if err := f.fn.FetchPayloads(peer, req.txs, req.cells); err != nil {
@ -826,11 +827,12 @@ func (f *BlobFetcher) scheduleFetches(timer *mclock.Timer, timeout chan struct{}
break break
} }
} }
}(peer, request) }()
} }
}) })
// If a new request was fired, schedule a timeout timer // If a new request was fired, schedule a timeout timer
if idle && len(f.requests) > 0 { if wasIdle && len(f.requests) > 0 {
f.rescheduleTimeout(timer, timeout) f.rescheduleTimeout(timer, timeout)
} }
} }