OnBlobsByRange is a callback method to invoke when a batch of Contract bytes codes are received from a remote peer.
(res *blobsByRangeResponse)
| 974 | // OnBlobsByRange is a callback method to invoke when a batch of Contract |
| 975 | // bytes codes are received from a remote peer. |
| 976 | func (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() |
no test coverage detected