MCPcopy Create free account
hub / github.com/ethstorage/es-node / OnBlobsByRange

Method OnBlobsByRange

ethstorage/p2p/protocol/syncclient.go:976–1056  ·  view source on GitHub ↗

OnBlobsByRange is a callback method to invoke when a batch of Contract bytes codes are received from a remote peer.

(res *blobsByRangeResponse)

Source from the content-addressed store, hash-verified

974// OnBlobsByRange is a callback method to invoke when a batch of Contract
975// bytes codes are received from a remote peer.
976func (s *SyncClient) OnBlobsByRange(res *blobsByRangeResponse) {
977 var (
978 size common.StorageSize
979 req = res.req
980 start = time.Now()
981 reqCount = req.limit - req.origin + 1
982 )
983
984 if reqCount > maxKvCountPerReq {
985 reqCount = maxKvCountPerReq
986 }
987 for _, blob := range res.Blobs {
988 if blob != nil {
989 size += common.StorageSize(len(blob.EncodedBlob))
990 }
991 }
992 s.lg.Debug("OnBlobsByRange: static", "reqId", req.id, "blobCount", len(res.Blobs), "bytes", size)
993
994 blobsInRange := make([]*BlobPayload, 0)
995 for _, blob := range res.Blobs {
996 if req.origin <= blob.BlobIndex && req.limit >= blob.BlobIndex {
997 blobsInRange = append(blobsInRange, blob)
998 }
999 }
1000 if len(res.Blobs) > len(blobsInRange) {
1001 s.lg.Trace("Drop unexpected kvs", "count", len(res.Blobs)-len(blobsInRange))
1002 }
1003
1004 // Response is valid, but check if peer is signalling that it does not have
1005 // the requested Data. For blob range queries that means the peer is not
1006 // yet synced.
1007 if len(blobsInRange) == 0 {
1008 s.lg.Info("Peer rejected get blob by range request")
1009 s.lock.Lock()
1010 if _, ok := s.peers[req.peer]; ok {
1011 req.subTask.task.statelessPeers[req.peer] = struct{}{}
1012 }
1013 s.lock.Unlock()
1014 s.metrics.ClientOnBlobsByRange(req.peer.String(), reqCount, uint64(len(res.Blobs)), 0, time.Since(start))
1015 return
1016 }
1017
1018 synced, syncedBytes, inserted, err := s.onResult(blobsInRange)
1019 if err != nil {
1020 s.lg.Error("OnBlobsByRange fail", "err", err.Error())
1021 return
1022 }
1023
1024 s.metrics.ClientOnBlobsByRange(req.peer.String(), reqCount, uint64(len(res.Blobs)), synced, time.Since(start))
1025 s.lg.Debug("Persisted set of kvs", "count", synced, "bytes", syncedBytes)
1026
1027 // set peer to stateless peer if fail too much
1028 if len(inserted) == 0 {
1029 s.lock.Lock()
1030 if _, ok := s.peers[req.peer]; ok {
1031 req.subTask.task.statelessPeers[req.peer] = struct{}{}
1032 }
1033 s.lock.Unlock()

Callers 1

assignBlobRangeTasksMethod · 0.95

Calls 5

onResultMethod · 0.95
insertMethod · 0.80
ClientOnBlobsByRangeMethod · 0.65
StringMethod · 0.45
ErrorMethod · 0.45

Tested by

no test coverage detected