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

Method OnBlobsByList

ethstorage/p2p/protocol/syncclient.go:1060–1120  ·  view source on GitHub ↗

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

(res *blobsByListResponse)

Source from the content-addressed store, hash-verified

1058// OnBlobsByList is a callback method to invoke when a batch of Contract
1059// bytes codes are received from a remote peer.
1060func (s *SyncClient) OnBlobsByList(res *blobsByListResponse) {
1061 var (
1062 size common.StorageSize
1063 req = res.req
1064 start = time.Now()
1065 )
1066 for _, blob := range res.Blobs {
1067 if blob != nil {
1068 size += common.StorageSize(len(blob.EncodedBlob))
1069 }
1070 }
1071 s.lg.Debug("OnBlobsByList: static", "reqId", req.id, "blobCount", len(res.Blobs), "bytes", size)
1072
1073 startIdx, endIdx := s.storageManager.KvEntries()*req.shardId, s.storageManager.KvEntries()*(req.shardId+1)-1
1074 blobsInRange := make([]*BlobPayload, 0)
1075 for _, blob := range res.Blobs {
1076 if startIdx <= blob.BlobIndex && endIdx >= blob.BlobIndex {
1077 blobsInRange = append(blobsInRange, blob)
1078 }
1079 }
1080 if len(res.Blobs) > len(blobsInRange) {
1081 s.lg.Trace("Drop unexpected kvs", "count", len(res.Blobs)-len(blobsInRange))
1082 }
1083
1084 // Response is valid, but check if peer is signalling that it does not have
1085 // the requested Data. For kv range queries that means the peer is not
1086 // yet synced.
1087 if len(blobsInRange) == 0 {
1088 s.lg.Info("Peer rejected get blobs by list request")
1089 s.lock.Lock()
1090 if _, ok := s.peers[req.peer]; ok {
1091 req.healTask.task.statelessPeers[req.peer] = struct{}{}
1092 }
1093 s.lock.Unlock()
1094 s.metrics.ClientOnBlobsByList(req.peer.String(), uint64(len(req.indexes)), uint64(len(res.Blobs)),
1095 0, time.Since(start))
1096 return
1097 }
1098
1099 synced, syncedBytes, inserted, err := s.onResult(blobsInRange)
1100 if err != nil {
1101 s.lg.Error("OnBlobsByList fail", "err", err.Error())
1102 return
1103 }
1104
1105 s.metrics.ClientOnBlobsByList(req.peer.String(), uint64(len(req.indexes)), uint64(len(res.Blobs)),
1106 synced, time.Since(start))
1107 s.lg.Debug("Persisted set of kvs", "count", synced, "bytes", syncedBytes)
1108
1109 s.lock.Lock()
1110 state := req.healTask.task.state
1111 state.BlobsSynced += uint64(len(inserted))
1112 // set peer to stateless peer if fail too much
1113 if len(inserted) == 0 {
1114 if _, ok := s.peers[req.peer]; ok {
1115 req.healTask.task.statelessPeers[req.peer] = struct{}{}
1116 }
1117 }

Callers 1

assignBlobHealTasksMethod · 0.95

Calls 6

onResultMethod · 0.95
removeMethod · 0.80
KvEntriesMethod · 0.65
ClientOnBlobsByListMethod · 0.65
StringMethod · 0.45
ErrorMethod · 0.45

Tested by

no test coverage detected