2019-07-05 14:29:17 +00:00
|
|
|
package chain
|
|
|
|
|
|
|
|
import (
|
|
|
|
"bufio"
|
|
|
|
"context"
|
|
|
|
"fmt"
|
2019-07-08 15:14:36 +00:00
|
|
|
"math/rand"
|
|
|
|
"sync"
|
|
|
|
|
2019-07-28 05:35:32 +00:00
|
|
|
bserv "github.com/ipfs/go-blockservice"
|
2019-07-08 13:36:43 +00:00
|
|
|
"github.com/libp2p/go-libp2p-core/host"
|
2019-07-05 14:46:21 +00:00
|
|
|
"github.com/libp2p/go-libp2p-core/protocol"
|
2019-10-12 09:44:56 +00:00
|
|
|
"go.opencensus.io/trace"
|
2019-07-26 16:31:07 +00:00
|
|
|
"golang.org/x/xerrors"
|
2019-07-08 15:14:36 +00:00
|
|
|
|
2019-10-18 04:47:41 +00:00
|
|
|
"github.com/filecoin-project/lotus/chain/store"
|
|
|
|
"github.com/filecoin-project/lotus/chain/types"
|
|
|
|
"github.com/filecoin-project/lotus/lib/cborrpc"
|
|
|
|
"github.com/filecoin-project/lotus/node/modules/dtypes"
|
2019-07-05 14:29:17 +00:00
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
blocks "github.com/ipfs/go-block-format"
|
2019-07-05 14:29:17 +00:00
|
|
|
"github.com/ipfs/go-cid"
|
|
|
|
cbor "github.com/ipfs/go-ipld-cbor"
|
|
|
|
inet "github.com/libp2p/go-libp2p-core/network"
|
2019-07-08 12:51:45 +00:00
|
|
|
"github.com/libp2p/go-libp2p-core/peer"
|
2019-07-05 14:29:17 +00:00
|
|
|
)
|
|
|
|
|
2019-07-05 14:46:21 +00:00
|
|
|
type NewStreamFunc func(context.Context, peer.ID, ...protocol.ID) (inet.Stream, error)
|
|
|
|
|
2019-07-05 14:29:17 +00:00
|
|
|
const BlockSyncProtocolID = "/fil/sync/blk/0.0.1"
|
|
|
|
|
|
|
|
func init() {
|
|
|
|
cbor.RegisterCborType(BlockSyncRequest{})
|
|
|
|
cbor.RegisterCborType(BlockSyncResponse{})
|
|
|
|
cbor.RegisterCborType(BSTipSet{})
|
|
|
|
}
|
|
|
|
|
|
|
|
type BlockSyncService struct {
|
2019-07-26 04:54:22 +00:00
|
|
|
cs *store.ChainStore
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
type BlockSyncRequest struct {
|
|
|
|
Start []cid.Cid
|
|
|
|
RequestLength uint64
|
|
|
|
|
|
|
|
Options uint64
|
|
|
|
}
|
|
|
|
|
|
|
|
type BSOptions struct {
|
|
|
|
IncludeBlocks bool
|
|
|
|
IncludeMessages bool
|
|
|
|
}
|
|
|
|
|
|
|
|
func ParseBSOptions(optfield uint64) *BSOptions {
|
|
|
|
return &BSOptions{
|
|
|
|
IncludeBlocks: optfield&(BSOptBlocks) != 0,
|
|
|
|
IncludeMessages: optfield&(BSOptMessages) != 0,
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
const (
|
|
|
|
BSOptBlocks = 1 << 0
|
|
|
|
BSOptMessages = 1 << 1
|
|
|
|
)
|
|
|
|
|
|
|
|
type BlockSyncResponse struct {
|
|
|
|
Chain []*BSTipSet
|
|
|
|
|
2019-08-22 01:29:19 +00:00
|
|
|
Status uint64
|
2019-07-05 14:29:17 +00:00
|
|
|
Message string
|
|
|
|
}
|
|
|
|
|
|
|
|
type BSTipSet struct {
|
2019-07-25 22:15:03 +00:00
|
|
|
Blocks []*types.BlockHeader
|
2019-07-05 14:29:17 +00:00
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
BlsMessages []*types.Message
|
2019-08-22 01:29:19 +00:00
|
|
|
BlsMsgIncludes [][]uint64
|
2019-08-01 20:40:47 +00:00
|
|
|
|
|
|
|
SecpkMessages []*types.SignedMessage
|
2019-08-22 01:29:19 +00:00
|
|
|
SecpkMsgIncludes [][]uint64
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
2019-07-26 04:54:22 +00:00
|
|
|
func NewBlockSyncService(cs *store.ChainStore) *BlockSyncService {
|
2019-07-05 14:29:17 +00:00
|
|
|
return &BlockSyncService{
|
|
|
|
cs: cs,
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (bss *BlockSyncService) HandleStream(s inet.Stream) {
|
2019-10-12 09:44:56 +00:00
|
|
|
ctx, span := trace.StartSpan(context.Background(), "blocksync.HandleStream")
|
|
|
|
defer span.End()
|
|
|
|
|
2019-07-05 14:29:17 +00:00
|
|
|
defer s.Close()
|
|
|
|
|
|
|
|
var req BlockSyncRequest
|
2019-07-08 12:46:30 +00:00
|
|
|
if err := cborrpc.ReadCborRPC(bufio.NewReader(s), &req); err != nil {
|
2019-07-05 14:29:17 +00:00
|
|
|
log.Errorf("failed to read block sync request: %s", err)
|
|
|
|
return
|
|
|
|
}
|
2019-07-30 13:55:36 +00:00
|
|
|
log.Infof("block sync request for: %s %d", req.Start, req.RequestLength)
|
2019-07-05 14:29:17 +00:00
|
|
|
|
2019-10-12 09:44:56 +00:00
|
|
|
resp, err := bss.processRequest(ctx, &req)
|
2019-07-05 14:29:17 +00:00
|
|
|
if err != nil {
|
|
|
|
log.Error("failed to process block sync request: ", err)
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2019-07-08 12:46:30 +00:00
|
|
|
if err := cborrpc.WriteCborRPC(s, resp); err != nil {
|
2019-07-05 14:29:17 +00:00
|
|
|
log.Error("failed to write back response for handle stream: ", err)
|
|
|
|
return
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-10-12 09:44:56 +00:00
|
|
|
func (bss *BlockSyncService) processRequest(ctx context.Context, req *BlockSyncRequest) (*BlockSyncResponse, error) {
|
|
|
|
ctx, span := trace.StartSpan(ctx, "blocksync.ProcessRequest")
|
|
|
|
defer span.End()
|
|
|
|
|
2019-07-05 14:29:17 +00:00
|
|
|
opts := ParseBSOptions(req.Options)
|
2019-08-02 22:21:46 +00:00
|
|
|
if len(req.Start) == 0 {
|
|
|
|
return &BlockSyncResponse{
|
2019-08-02 22:32:02 +00:00
|
|
|
Status: 204,
|
|
|
|
Message: "no cids given in blocksync request",
|
2019-08-02 22:21:46 +00:00
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2019-10-13 00:41:24 +00:00
|
|
|
span.AddAttributes(
|
|
|
|
trace.BoolAttribute("blocks", opts.IncludeBlocks),
|
|
|
|
trace.BoolAttribute("messages", opts.IncludeMessages),
|
|
|
|
)
|
2019-10-12 09:44:56 +00:00
|
|
|
|
2019-07-05 14:29:17 +00:00
|
|
|
chain, err := bss.collectChainSegment(req.Start, req.RequestLength, opts)
|
|
|
|
if err != nil {
|
|
|
|
log.Error("encountered error while responding to block sync request: ", err)
|
|
|
|
return &BlockSyncResponse{
|
|
|
|
Status: 203,
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
return &BlockSyncResponse{
|
|
|
|
Chain: chain,
|
|
|
|
Status: 0,
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (bss *BlockSyncService) collectChainSegment(start []cid.Cid, length uint64, opts *BSOptions) ([]*BSTipSet, error) {
|
|
|
|
var bstips []*BSTipSet
|
|
|
|
cur := start
|
|
|
|
for {
|
|
|
|
var bst BSTipSet
|
|
|
|
ts, err := bss.cs.LoadTipSet(cur)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
if opts.IncludeMessages {
|
2019-08-01 20:40:47 +00:00
|
|
|
bmsgs, bmincl, smsgs, smincl, err := bss.gatherMessages(ts)
|
2019-07-05 14:29:17 +00:00
|
|
|
if err != nil {
|
2019-08-01 20:40:47 +00:00
|
|
|
return nil, xerrors.Errorf("gather messages failed: %w", err)
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
bst.BlsMessages = bmsgs
|
|
|
|
bst.BlsMsgIncludes = bmincl
|
|
|
|
bst.SecpkMessages = smsgs
|
|
|
|
bst.SecpkMsgIncludes = smincl
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
if opts.IncludeBlocks {
|
|
|
|
bst.Blocks = ts.Blocks()
|
|
|
|
}
|
|
|
|
|
|
|
|
bstips = append(bstips, &bst)
|
|
|
|
|
|
|
|
if uint64(len(bstips)) >= length || ts.Height() == 0 {
|
|
|
|
return bstips, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
cur = ts.Parents()
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-08-22 01:29:19 +00:00
|
|
|
func (bss *BlockSyncService) gatherMessages(ts *types.TipSet) ([]*types.Message, [][]uint64, []*types.SignedMessage, [][]uint64, error) {
|
|
|
|
blsmsgmap := make(map[cid.Cid]uint64)
|
|
|
|
secpkmsgmap := make(map[cid.Cid]uint64)
|
2019-08-01 20:40:47 +00:00
|
|
|
var secpkmsgs []*types.SignedMessage
|
|
|
|
var blsmsgs []*types.Message
|
2019-08-22 01:29:19 +00:00
|
|
|
var secpkincl, blsincl [][]uint64
|
2019-07-05 14:29:17 +00:00
|
|
|
|
|
|
|
for _, b := range ts.Blocks() {
|
2019-08-01 20:40:47 +00:00
|
|
|
bmsgs, smsgs, err := bss.cs.MessagesForBlock(b)
|
2019-07-05 14:29:17 +00:00
|
|
|
if err != nil {
|
2019-08-01 20:40:47 +00:00
|
|
|
return nil, nil, nil, nil, err
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
2019-08-22 01:29:19 +00:00
|
|
|
bmi := make([]uint64, 0, len(bmsgs))
|
2019-08-01 20:40:47 +00:00
|
|
|
for _, m := range bmsgs {
|
|
|
|
i, ok := blsmsgmap[m.Cid()]
|
2019-07-05 14:29:17 +00:00
|
|
|
if !ok {
|
2019-08-22 01:29:19 +00:00
|
|
|
i = uint64(len(blsmsgs))
|
2019-08-01 20:40:47 +00:00
|
|
|
blsmsgs = append(blsmsgs, m)
|
|
|
|
blsmsgmap[m.Cid()] = i
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
bmi = append(bmi, i)
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
2019-08-01 20:40:47 +00:00
|
|
|
blsincl = append(blsincl, bmi)
|
|
|
|
|
2019-08-22 01:29:19 +00:00
|
|
|
smi := make([]uint64, 0, len(smsgs))
|
2019-08-01 20:40:47 +00:00
|
|
|
for _, m := range smsgs {
|
|
|
|
i, ok := secpkmsgmap[m.Cid()]
|
|
|
|
if !ok {
|
2019-08-22 01:29:19 +00:00
|
|
|
i = uint64(len(secpkmsgs))
|
2019-08-01 20:40:47 +00:00
|
|
|
secpkmsgs = append(secpkmsgs, m)
|
|
|
|
secpkmsgmap[m.Cid()] = i
|
|
|
|
}
|
|
|
|
|
|
|
|
smi = append(smi, i)
|
|
|
|
}
|
|
|
|
secpkincl = append(secpkincl, smi)
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
return blsmsgs, blsincl, secpkmsgs, secpkincl, nil
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
type BlockSync struct {
|
2019-07-28 05:35:32 +00:00
|
|
|
bserv bserv.BlockService
|
2019-07-05 14:29:17 +00:00
|
|
|
newStream NewStreamFunc
|
|
|
|
|
|
|
|
syncPeersLk sync.Mutex
|
|
|
|
syncPeers map[peer.ID]struct{}
|
|
|
|
}
|
|
|
|
|
2019-08-01 14:14:16 +00:00
|
|
|
func NewBlockSyncClient(bserv dtypes.ChainBlockService, h host.Host) *BlockSync {
|
2019-07-05 14:29:17 +00:00
|
|
|
return &BlockSync{
|
2019-07-28 05:35:32 +00:00
|
|
|
bserv: bserv,
|
2019-07-08 13:36:43 +00:00
|
|
|
newStream: h.NewStream,
|
2019-07-05 14:29:17 +00:00
|
|
|
syncPeers: make(map[peer.ID]struct{}),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (bs *BlockSync) getPeers() []peer.ID {
|
|
|
|
bs.syncPeersLk.Lock()
|
|
|
|
defer bs.syncPeersLk.Unlock()
|
|
|
|
var out []peer.ID
|
|
|
|
for p := range bs.syncPeers {
|
|
|
|
out = append(out, p)
|
|
|
|
}
|
|
|
|
return out
|
|
|
|
}
|
|
|
|
|
2019-08-01 14:05:35 +00:00
|
|
|
func (bs *BlockSync) processStatus(req *BlockSyncRequest, res *BlockSyncResponse) error {
|
2019-07-26 16:47:18 +00:00
|
|
|
switch res.Status {
|
|
|
|
case 101: // Partial Response
|
|
|
|
panic("not handled")
|
|
|
|
case 201: // req.Start not found
|
2019-08-01 14:05:35 +00:00
|
|
|
return fmt.Errorf("not found")
|
2019-07-26 16:47:18 +00:00
|
|
|
case 202: // Go Away
|
|
|
|
panic("not handled")
|
|
|
|
case 203: // Internal Error
|
2019-08-01 14:05:35 +00:00
|
|
|
return fmt.Errorf("block sync peer errored: %s", res.Message)
|
2019-09-06 20:03:28 +00:00
|
|
|
case 204:
|
|
|
|
return fmt.Errorf("block sync request invalid: %s", res.Message)
|
2019-07-26 16:47:18 +00:00
|
|
|
default:
|
2019-09-06 20:03:28 +00:00
|
|
|
return fmt.Errorf("unrecognized response code: %d", res.Status)
|
2019-07-26 16:47:18 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-07-26 04:54:22 +00:00
|
|
|
func (bs *BlockSync) GetBlocks(ctx context.Context, tipset []cid.Cid, count int) ([]*types.TipSet, error) {
|
2019-10-12 09:44:56 +00:00
|
|
|
ctx, span := trace.StartSpan(ctx, "bsync.GetBlocks")
|
|
|
|
defer span.End()
|
|
|
|
if span.IsRecordingEvents() {
|
|
|
|
span.AddAttributes(
|
|
|
|
trace.StringAttribute("tipset", fmt.Sprint(tipset)),
|
|
|
|
trace.Int64Attribute("count", int64(count)),
|
|
|
|
)
|
|
|
|
}
|
|
|
|
|
2019-07-05 14:29:17 +00:00
|
|
|
peers := bs.getPeers()
|
|
|
|
perm := rand.Perm(len(peers))
|
|
|
|
// TODO: round robin through these peers on error
|
|
|
|
|
|
|
|
req := &BlockSyncRequest{
|
|
|
|
Start: tipset,
|
|
|
|
RequestLength: uint64(count),
|
|
|
|
Options: BSOptBlocks,
|
|
|
|
}
|
|
|
|
|
2019-09-06 20:03:28 +00:00
|
|
|
var oerr error
|
2019-07-26 16:31:07 +00:00
|
|
|
for _, p := range perm {
|
2019-08-01 14:05:35 +00:00
|
|
|
res, err := bs.sendRequestToPeer(ctx, peers[p], req)
|
2019-07-26 16:47:18 +00:00
|
|
|
if err != nil {
|
2019-09-06 20:03:28 +00:00
|
|
|
oerr = err
|
2019-07-26 16:47:18 +00:00
|
|
|
log.Warnf("BlockSync request failed for peer %s: %s", peers[p].String(), err)
|
|
|
|
continue
|
2019-07-26 16:31:07 +00:00
|
|
|
}
|
|
|
|
|
2019-08-01 14:05:35 +00:00
|
|
|
if res.Status == 0 {
|
|
|
|
return bs.processBlocksResponse(req, res)
|
2019-07-26 16:47:18 +00:00
|
|
|
}
|
2019-09-06 20:03:28 +00:00
|
|
|
oerr = bs.processStatus(req, res)
|
|
|
|
if oerr != nil {
|
2019-10-05 15:59:35 +00:00
|
|
|
log.Warnf("BlockSync peer %s response was an error: %s", peers[p].String(), oerr)
|
2019-08-01 17:19:18 +00:00
|
|
|
}
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
2019-09-06 20:03:28 +00:00
|
|
|
return nil, xerrors.Errorf("GetBlocks failed with all peers: %w", oerr)
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
2019-07-26 04:54:22 +00:00
|
|
|
func (bs *BlockSync) GetFullTipSet(ctx context.Context, p peer.ID, h []cid.Cid) (*store.FullTipSet, error) {
|
2019-07-05 14:29:17 +00:00
|
|
|
// TODO: round robin through these peers on error
|
|
|
|
|
|
|
|
req := &BlockSyncRequest{
|
|
|
|
Start: h,
|
|
|
|
RequestLength: 1,
|
|
|
|
Options: BSOptBlocks | BSOptMessages,
|
|
|
|
}
|
|
|
|
|
|
|
|
res, err := bs.sendRequestToPeer(ctx, p, req)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
switch res.Status {
|
|
|
|
case 0: // Success
|
|
|
|
if len(res.Chain) == 0 {
|
|
|
|
return nil, fmt.Errorf("got zero length chain response")
|
|
|
|
}
|
|
|
|
bts := res.Chain[0]
|
|
|
|
|
|
|
|
return bstsToFullTipSet(bts)
|
|
|
|
case 101: // Partial Response
|
|
|
|
panic("not handled")
|
|
|
|
case 201: // req.Start not found
|
|
|
|
return nil, fmt.Errorf("not found")
|
|
|
|
case 202: // Go Away
|
|
|
|
panic("not handled")
|
|
|
|
case 203: // Internal Error
|
2019-08-02 22:32:02 +00:00
|
|
|
return nil, fmt.Errorf("block sync peer errored: %q", res.Message)
|
|
|
|
case 204: // Invalid Request
|
|
|
|
return nil, fmt.Errorf("block sync request invalid: %q", res.Message)
|
2019-07-05 14:29:17 +00:00
|
|
|
default:
|
|
|
|
return nil, fmt.Errorf("unrecognized response code")
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-07-26 04:54:22 +00:00
|
|
|
func (bs *BlockSync) GetChainMessages(ctx context.Context, h *types.TipSet, count uint64) ([]*BSTipSet, error) {
|
2019-10-12 09:44:56 +00:00
|
|
|
ctx, span := trace.StartSpan(ctx, "GetChainMessages")
|
|
|
|
defer span.End()
|
|
|
|
|
2019-07-05 14:29:17 +00:00
|
|
|
peers := bs.getPeers()
|
|
|
|
perm := rand.Perm(len(peers))
|
|
|
|
// TODO: round robin through these peers on error
|
|
|
|
|
|
|
|
req := &BlockSyncRequest{
|
|
|
|
Start: h.Cids(),
|
|
|
|
RequestLength: count,
|
2019-08-02 00:13:57 +00:00
|
|
|
Options: BSOptMessages | BSOptBlocks,
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
2019-08-01 14:05:35 +00:00
|
|
|
var err error
|
|
|
|
for _, p := range perm {
|
|
|
|
res, err := bs.sendRequestToPeer(ctx, peers[p], req)
|
|
|
|
if err != nil {
|
|
|
|
log.Warnf("BlockSync request failed for peer %s: %s", peers[p].String(), err)
|
|
|
|
continue
|
|
|
|
}
|
2019-07-05 14:29:17 +00:00
|
|
|
|
2019-08-01 14:05:35 +00:00
|
|
|
if res.Status == 0 {
|
|
|
|
return res.Chain, nil
|
|
|
|
}
|
|
|
|
err = bs.processStatus(req, res)
|
2019-08-01 17:19:18 +00:00
|
|
|
if err != nil {
|
|
|
|
log.Warnf("BlockSync peer %s response was an error: %s", peers[p].String(), err)
|
|
|
|
}
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
2019-08-01 14:05:35 +00:00
|
|
|
|
|
|
|
// TODO: What if we have no peers (and err is nil)?
|
2019-10-03 18:20:29 +00:00
|
|
|
return nil, xerrors.Errorf("GetChainMessages failed with all peers(%d): %w", len(peers), err)
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
2019-07-26 04:54:22 +00:00
|
|
|
func bstsToFullTipSet(bts *BSTipSet) (*store.FullTipSet, error) {
|
|
|
|
fts := &store.FullTipSet{}
|
2019-07-05 14:29:17 +00:00
|
|
|
for i, b := range bts.Blocks {
|
2019-07-25 22:15:03 +00:00
|
|
|
fb := &types.FullBlock{
|
2019-07-05 14:29:17 +00:00
|
|
|
Header: b,
|
|
|
|
}
|
2019-08-01 20:40:47 +00:00
|
|
|
for _, mi := range bts.BlsMsgIncludes[i] {
|
|
|
|
fb.BlsMessages = append(fb.BlsMessages, bts.BlsMessages[mi])
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
2019-10-06 00:51:48 +00:00
|
|
|
for _, mi := range bts.SecpkMsgIncludes[i] {
|
|
|
|
fb.SecpkMessages = append(fb.SecpkMessages, bts.SecpkMessages[mi])
|
|
|
|
}
|
|
|
|
|
2019-07-05 14:29:17 +00:00
|
|
|
fts.Blocks = append(fts.Blocks, fb)
|
|
|
|
}
|
|
|
|
|
|
|
|
return fts, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (bs *BlockSync) sendRequestToPeer(ctx context.Context, p peer.ID, req *BlockSyncRequest) (*BlockSyncResponse, error) {
|
|
|
|
s, err := bs.newStream(inet.WithNoDial(ctx, "should already have connection"), p, BlockSyncProtocolID)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
2019-07-08 12:46:30 +00:00
|
|
|
if err := cborrpc.WriteCborRPC(s, req); err != nil {
|
2019-07-05 14:29:17 +00:00
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
var res BlockSyncResponse
|
2019-07-08 12:46:30 +00:00
|
|
|
if err := cborrpc.ReadCborRPC(bufio.NewReader(s), &res); err != nil {
|
2019-07-05 14:29:17 +00:00
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
return &res, nil
|
|
|
|
}
|
|
|
|
|
2019-07-26 04:54:22 +00:00
|
|
|
func (bs *BlockSync) processBlocksResponse(req *BlockSyncRequest, res *BlockSyncResponse) ([]*types.TipSet, error) {
|
|
|
|
cur, err := types.NewTipSet(res.Chain[0].Blocks)
|
2019-07-05 14:29:17 +00:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
2019-07-26 04:54:22 +00:00
|
|
|
out := []*types.TipSet{cur}
|
2019-07-05 14:29:17 +00:00
|
|
|
for bi := 1; bi < len(res.Chain); bi++ {
|
|
|
|
next := res.Chain[bi].Blocks
|
2019-07-26 04:54:22 +00:00
|
|
|
nts, err := types.NewTipSet(next)
|
2019-07-05 14:29:17 +00:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
2019-09-03 04:36:07 +00:00
|
|
|
if !types.CidArrsEqual(cur.Parents(), nts.Cids()) {
|
2019-07-05 14:29:17 +00:00
|
|
|
return nil, fmt.Errorf("parents of tipset[%d] were not tipset[%d]", bi-1, bi)
|
|
|
|
}
|
|
|
|
|
|
|
|
out = append(out, nts)
|
|
|
|
cur = nts
|
|
|
|
}
|
|
|
|
return out, nil
|
|
|
|
}
|
|
|
|
|
2019-07-25 22:15:03 +00:00
|
|
|
func (bs *BlockSync) GetBlock(ctx context.Context, c cid.Cid) (*types.BlockHeader, error) {
|
2019-07-28 05:35:32 +00:00
|
|
|
sb, err := bs.bserv.GetBlock(ctx, c)
|
2019-07-05 14:29:17 +00:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
2019-07-25 22:15:03 +00:00
|
|
|
return types.DecodeBlock(sb.RawData())
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func (bs *BlockSync) AddPeer(p peer.ID) {
|
|
|
|
bs.syncPeersLk.Lock()
|
|
|
|
defer bs.syncPeersLk.Unlock()
|
|
|
|
bs.syncPeers[p] = struct{}{}
|
|
|
|
}
|
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
func (bs *BlockSync) FetchMessagesByCids(ctx context.Context, cids []cid.Cid) ([]*types.Message, error) {
|
|
|
|
out := make([]*types.Message, len(cids))
|
|
|
|
|
|
|
|
err := bs.fetchCids(ctx, cids, func(i int, b blocks.Block) error {
|
|
|
|
msg, err := types.DecodeMessage(b.RawData())
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
if out[i] != nil {
|
|
|
|
return fmt.Errorf("received duplicate message")
|
|
|
|
}
|
|
|
|
|
|
|
|
out[i] = msg
|
|
|
|
return nil
|
|
|
|
})
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
return out, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (bs *BlockSync) FetchSignedMessagesByCids(ctx context.Context, cids []cid.Cid) ([]*types.SignedMessage, error) {
|
2019-07-25 22:15:03 +00:00
|
|
|
out := make([]*types.SignedMessage, len(cids))
|
2019-07-05 14:29:17 +00:00
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
err := bs.fetchCids(ctx, cids, func(i int, b blocks.Block) error {
|
|
|
|
smsg, err := types.DecodeSignedMessage(b.RawData())
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
if out[i] != nil {
|
|
|
|
return fmt.Errorf("received duplicate message")
|
|
|
|
}
|
|
|
|
|
|
|
|
out[i] = smsg
|
|
|
|
return nil
|
|
|
|
})
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
return out, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (bs *BlockSync) fetchCids(ctx context.Context, cids []cid.Cid, cb func(int, blocks.Block) error) error {
|
2019-07-28 05:35:32 +00:00
|
|
|
resp := bs.bserv.GetBlocks(context.TODO(), cids)
|
2019-07-05 14:29:17 +00:00
|
|
|
|
|
|
|
m := make(map[cid.Cid]int)
|
|
|
|
for i, c := range cids {
|
|
|
|
m[c] = i
|
|
|
|
}
|
|
|
|
|
|
|
|
for i := 0; i < len(cids); i++ {
|
|
|
|
select {
|
|
|
|
case v, ok := <-resp:
|
|
|
|
if !ok {
|
|
|
|
if i == len(cids)-1 {
|
|
|
|
break
|
|
|
|
}
|
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
return fmt.Errorf("failed to fetch all messages")
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
ix, ok := m[v.Cid()]
|
2019-07-05 14:29:17 +00:00
|
|
|
if !ok {
|
2019-08-01 20:40:47 +00:00
|
|
|
return fmt.Errorf("received message we didnt ask for")
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
if err := cb(ix, v); err != nil {
|
|
|
|
return err
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-08-01 20:40:47 +00:00
|
|
|
return nil
|
2019-07-05 14:29:17 +00:00
|
|
|
}
|