OnBlobsByList is a callback method to invoke when a batch of Contract bytes codes are received from a remote peer.
(res *blobsByListResponse)
| 1058 | // OnBlobsByList is a callback method to invoke when a batch of Contract |
| 1059 | // bytes codes are received from a remote peer. |
| 1060 | func (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 | } |
no test coverage detected