peerresponsemanager.go 4.33 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161
package peerresponsemanager

import (
	"context"
	"sync"

	"github.com/ipld/go-ipld-prime/linking/cid"

	"github.com/ipfs/go-graphsync/ipldbridge"

	logging "github.com/ipfs/go-log"
	"github.com/ipld/go-ipld-prime"

	"github.com/ipfs/go-block-format"
	gsmsg "github.com/ipfs/go-graphsync/message"
	"github.com/ipfs/go-graphsync/responsemanager/linktracker"
	"github.com/ipfs/go-graphsync/responsemanager/responsebuilder"
	peer "github.com/libp2p/go-libp2p-peer"
)

var log = logging.Logger("graphsync")

// PeerHandler is an interface that can send a response for a given peer across
// the network.
type PeerHandler interface {
	SendResponse(peer.ID, []gsmsg.GraphSyncResponse, []blocks.Block) <-chan struct{}
}

// PeerResponseManager handles batching, deduping, and sending responses for
// a given peer across multiple requests.
type PeerResponseManager struct {
	p            peer.ID
	ctx          context.Context
	cancel       context.CancelFunc
	peerHandler  PeerHandler
	ipldBridge   ipldbridge.IPLDBridge
	outgoingWork chan struct{}

	linkTrackerLk     sync.RWMutex
	linkTracker       *linktracker.LinkTracker
	responseBuilderLk sync.RWMutex
	responseBuilder   *responsebuilder.ResponseBuilder
}

// New generates a new PeerResponse manager for the given context, peer ID,
// using the given peer handler and bridge to IPLD.
func New(ctx context.Context, p peer.ID, peerHandler PeerHandler, ipldBridge ipldbridge.IPLDBridge) *PeerResponseManager {
	ctx, cancel := context.WithCancel(ctx)
	return &PeerResponseManager{
		p:            p,
		ctx:          ctx,
		cancel:       cancel,
		peerHandler:  peerHandler,
		ipldBridge:   ipldBridge,
		outgoingWork: make(chan struct{}, 1),
		linkTracker:  linktracker.New(),
	}
}

// Startup initiates message sending for a peer
func (prm *PeerResponseManager) Startup() {
	go prm.run()
}

// Shutdown stops sending messages for a peer
func (prm *PeerResponseManager) Shutdown() {
	prm.cancel()
}

// SendResponse sends a given link for a given
// requestID across the wire, as well as its corresponding
// block if the block is present and has not already been sent
func (prm *PeerResponseManager) SendResponse(
	requestID gsmsg.GraphSyncRequestID,
	link ipld.Link,
	data []byte,
) {
	hasBlock := data != nil
	prm.linkTrackerLk.Lock()
	sendBlock := hasBlock && prm.linkTracker.ShouldSendBlockFor(link)
	prm.linkTracker.RecordLinkTraversal(requestID, link, hasBlock)
	prm.linkTrackerLk.Unlock()

	if prm.buildResponse(func(responseBuilder *responsebuilder.ResponseBuilder) {
		if sendBlock {
			cidLink := link.(cidlink.Link)
			block, err := blocks.NewBlockWithCid(data, cidLink.Cid)
			if err != nil {
				log.Errorf("Data did not match cid when sending link for %s", cidLink.String())
			}
			responseBuilder.AddBlock(block)
		}
		responseBuilder.AddLink(requestID, link, hasBlock)
	}) {
		prm.signalWork()
	}
}

// FinishRequest marks the given requestID as having sent all responses
func (prm *PeerResponseManager) FinishRequest(requestID gsmsg.GraphSyncRequestID) {
	prm.linkTrackerLk.Lock()
	isComplete := prm.linkTracker.FinishRequest(requestID)
	prm.linkTrackerLk.Unlock()

	if prm.buildResponse(func(responseBuilder *responsebuilder.ResponseBuilder) {
		responseBuilder.AddCompletedRequest(requestID, isComplete)
	}) {
		prm.signalWork()
	}
}

func (prm *PeerResponseManager) buildResponse(buildResponseFn func(*responsebuilder.ResponseBuilder)) bool {
	prm.responseBuilderLk.Lock()
	defer prm.responseBuilderLk.Unlock()
	if prm.responseBuilder == nil {
		prm.responseBuilder = responsebuilder.New()
	}
	buildResponseFn(prm.responseBuilder)
	return !prm.responseBuilder.Empty()
}

func (prm *PeerResponseManager) signalWork() {
	select {
	case prm.outgoingWork <- struct{}{}:
	default:
	}
}

func (prm *PeerResponseManager) run() {
	for {
		select {
		case <-prm.ctx.Done():
			return
		case <-prm.outgoingWork:
			prm.sendResponseMessage()
		}
	}
}

func (prm *PeerResponseManager) sendResponseMessage() {
	prm.responseBuilderLk.Lock()
	builder := prm.responseBuilder
	prm.responseBuilder = nil
	prm.responseBuilderLk.Unlock()

	if builder == nil || builder.Empty() {
		return
	}
	responses, blks, err := builder.Build(prm.ipldBridge)
	if err != nil {
		log.Errorf("Unable to assemble GraphSync response: %s", err.Error())
	}

	done := prm.peerHandler.SendResponse(prm.p, responses, blks)

	// wait for message to be processed
	select {
	case <-done:
	case <-prm.ctx.Done():
	}
}