responsemanager.go 15.6 KB
Newer Older
1 2 3 4
package responsemanager

import (
	"context"
5
	"errors"
6
	"math"
7 8
	"time"

Hannah Howard's avatar
Hannah Howard committed
9 10 11 12
	logging "github.com/ipfs/go-log"
	"github.com/ipfs/go-peertaskqueue/peertask"
	ipld "github.com/ipld/go-ipld-prime"
	"github.com/libp2p/go-libp2p-core/peer"
13

14
	"github.com/ipfs/go-graphsync"
15
	"github.com/ipfs/go-graphsync/ipldutil"
16
	gsmsg "github.com/ipfs/go-graphsync/message"
17
	"github.com/ipfs/go-graphsync/notifications"
Hannah Howard's avatar
Hannah Howard committed
18
	"github.com/ipfs/go-graphsync/responsemanager/hooks"
19 20 21
	"github.com/ipfs/go-graphsync/responsemanager/peerresponsemanager"
)

22 23
var log = logging.Logger("graphsync")

24
const (
25
	thawSpeed = time.Millisecond * 100
26 27 28
)

type inProgressResponseStatus struct {
29 30 31 32 33 34 35 36 37
	ctx        context.Context
	cancelFn   func()
	request    gsmsg.GraphSyncRequest
	loader     ipld.Loader
	traverser  ipldutil.Traverser
	signals    signals
	updates    []gsmsg.GraphSyncRequest
	isPaused   bool
	subscriber notifications.MappableSubscriber
38 39 40 41
}

type responseKey struct {
	p         peer.ID
42
	requestID graphsync.RequestID
43 44
}

45
type signals struct {
46
	pauseSignal  chan struct{}
Hannah Howard's avatar
Hannah Howard committed
47
	updateSignal chan struct{}
48
	errSignal    chan error
49 50 51
}

type responseTaskData struct {
52 53 54 55 56 57 58
	empty      bool
	subscriber notifications.MappableSubscriber
	ctx        context.Context
	request    gsmsg.GraphSyncRequest
	loader     ipld.Loader
	traverser  ipldutil.Traverser
	signals    signals
59 60
}

61 62 63
// QueryQueue is an interface that can receive new selector query tasks
// and prioritize them as needed, and pop them off later
type QueryQueue interface {
64 65 66 67
	PushTasks(to peer.ID, tasks ...peertask.Task)
	PopTasks(targetMinWork int) (peer.ID, []*peertask.Task, int)
	Remove(topic peertask.Topic, p peer.ID)
	TasksDone(to peer.ID, tasks ...*peertask.Task)
68 69 70
	ThawRound()
}

71 72
// RequestHooks is an interface for processing request hooks
type RequestHooks interface {
Hannah Howard's avatar
Hannah Howard committed
73
	ProcessRequestHooks(p peer.ID, request graphsync.RequestData) hooks.RequestResult
74 75 76 77
}

// BlockHooks is an interface for processing block hooks
type BlockHooks interface {
Hannah Howard's avatar
Hannah Howard committed
78 79 80 81 82 83
	ProcessBlockHooks(p peer.ID, request graphsync.RequestData, blockData graphsync.BlockData) hooks.BlockResult
}

// UpdateHooks is an interface for processing update hooks
type UpdateHooks interface {
	ProcessUpdateHooks(p peer.ID, request graphsync.RequestData, update graphsync.RequestData) hooks.UpdateResult
84 85
}

86 87 88
// CompletedListeners is an interface for notifying listeners that responses are complete
type CompletedListeners interface {
	NotifyCompletedListeners(p peer.ID, request graphsync.RequestData, status graphsync.ResponseStatusCode)
89 90
}

91 92 93 94 95
// CancelledListeners is an interface for notifying listeners that requestor cancelled
type CancelledListeners interface {
	NotifyCancelledListeners(p peer.ID, request graphsync.RequestData)
}

96 97 98 99 100 101 102 103 104 105
// BlockSentListeners is an interface for notifying listeners that of a block send occuring over the wire
type BlockSentListeners interface {
	NotifyBlockSentListeners(p peer.ID, request graphsync.RequestData, block graphsync.BlockData)
}

// NetworkErrorListeners is an interface for notifying listeners that an error occurred sending a data on the wire
type NetworkErrorListeners interface {
	NotifyNetworkErrorListeners(p peer.ID, request graphsync.RequestData, err error)
}

106 107 108 109 110 111 112 113 114 115 116 117
// PeerManager is an interface that returns sender interfaces for peer responses.
type PeerManager interface {
	SenderForPeer(p peer.ID) peerresponsemanager.PeerResponseSender
}

type responseManagerMessage interface {
	handle(rm *ResponseManager)
}

// ResponseManager handles incoming requests from the network, initiates selector
// traversals, and transmits responses
type ResponseManager struct {
118 119 120 121 122 123 124 125 126 127 128 129 130
	ctx                   context.Context
	cancelFn              context.CancelFunc
	peerManager           PeerManager
	queryQueue            QueryQueue
	updateHooks           UpdateHooks
	cancelledListeners    CancelledListeners
	completedListeners    CompletedListeners
	blockSentListeners    BlockSentListeners
	networkErrorListeners NetworkErrorListeners
	messages              chan responseManagerMessage
	workSignal            chan struct{}
	qe                    *queryExecutor
	inProgressResponses   map[responseKey]*inProgressResponseStatus
131
	maxInProcessRequests  uint64
132 133 134 135 136
}

// New creates a new response manager from the given context, loader,
// bridge to IPLD interface, peerManager, and queryQueue.
func New(ctx context.Context,
137
	loader ipld.Loader,
138
	peerManager PeerManager,
139 140
	queryQueue QueryQueue,
	requestHooks RequestHooks,
Hannah Howard's avatar
Hannah Howard committed
141
	blockHooks BlockHooks,
142
	updateHooks UpdateHooks,
143
	completedListeners CompletedListeners,
144
	cancelledListeners CancelledListeners,
145 146
	blockSentListeners BlockSentListeners,
	networkErrorListeners NetworkErrorListeners,
147
	maxInProcessRequests uint64,
148
) *ResponseManager {
149
	ctx, cancelFn := context.WithCancel(ctx)
150 151 152
	messages := make(chan responseManagerMessage, 16)
	workSignal := make(chan struct{}, 1)
	qe := &queryExecutor{
153 154 155 156 157 158 159 160 161 162 163
		requestHooks:       requestHooks,
		blockHooks:         blockHooks,
		updateHooks:        updateHooks,
		cancelledListeners: cancelledListeners,
		peerManager:        peerManager,
		loader:             loader,
		queryQueue:         queryQueue,
		messages:           messages,
		ctx:                ctx,
		workSignal:         workSignal,
		ticker:             time.NewTicker(thawSpeed),
164
	}
165
	return &ResponseManager{
166 167 168 169 170 171 172 173 174 175 176 177 178
		ctx:                   ctx,
		cancelFn:              cancelFn,
		peerManager:           peerManager,
		queryQueue:            queryQueue,
		updateHooks:           updateHooks,
		cancelledListeners:    cancelledListeners,
		completedListeners:    completedListeners,
		blockSentListeners:    blockSentListeners,
		networkErrorListeners: networkErrorListeners,
		messages:              messages,
		workSignal:            workSignal,
		qe:                    qe,
		inProgressResponses:   make(map[responseKey]*inProgressResponseStatus),
179
		maxInProcessRequests:  maxInProcessRequests,
180 181 182 183 184 185 186 187 188
	}
}

type processRequestMessage struct {
	p        peer.ID
	requests []gsmsg.GraphSyncRequest
}

// ProcessRequests processes incoming requests for the given peer
189
func (rm *ResponseManager) ProcessRequests(ctx context.Context, p peer.ID, requests []gsmsg.GraphSyncRequest) {
190 191 192
	select {
	case rm.messages <- &processRequestMessage{p, requests}:
	case <-rm.ctx.Done():
193
	case <-ctx.Done():
194 195 196
	}
}

197
type unpauseRequestMessage struct {
198 199 200 201 202 203 204 205 206 207 208 209 210
	p          peer.ID
	requestID  graphsync.RequestID
	response   chan error
	extensions []graphsync.ExtensionData
}

// UnpauseResponse unpauses a response that was previously paused
func (rm *ResponseManager) UnpauseResponse(p peer.ID, requestID graphsync.RequestID, extensions ...graphsync.ExtensionData) error {
	response := make(chan error, 1)
	return rm.sendSyncMessage(&unpauseRequestMessage{p, requestID, response, extensions}, response)
}

type pauseRequestMessage struct {
211 212 213
	p         peer.ID
	requestID graphsync.RequestID
	response  chan error
214 215
}

216 217
// PauseResponse pauses an in progress response (may take 1 or more blocks to process)
func (rm *ResponseManager) PauseResponse(p peer.ID, requestID graphsync.RequestID) error {
218
	response := make(chan error, 1)
219 220 221
	return rm.sendSyncMessage(&pauseRequestMessage{p, requestID, response}, response)
}

222
type errorRequestMessage struct {
223 224
	p         peer.ID
	requestID graphsync.RequestID
225
	err       error
226 227 228 229 230 231
	response  chan error
}

// CancelResponse cancels an in progress response
func (rm *ResponseManager) CancelResponse(p peer.ID, requestID graphsync.RequestID) error {
	response := make(chan error, 1)
232
	return rm.sendSyncMessage(&errorRequestMessage{p, requestID, errCancelledByCommand, response}, response)
233 234 235
}

func (rm *ResponseManager) sendSyncMessage(message responseManagerMessage, response chan error) error {
236 237 238
	select {
	case <-rm.ctx.Done():
		return errors.New("Context Cancelled")
239
	case rm.messages <- message:
240 241 242 243 244 245
	}
	select {
	case <-rm.ctx.Done():
		return errors.New("Context Cancelled")
	case err := <-response:
		return err
246 247 248
	}
}

249
type synchronizeMessage struct {
250
	sync chan error
251 252 253 254
}

// this is a test utility method to force all messages to get processed
func (rm *ResponseManager) synchronize() {
255 256
	sync := make(chan error)
	_ = rm.sendSyncMessage(&synchronizeMessage{sync}, sync)
257 258 259 260
}

type responseDataRequest struct {
	key          responseKey
261
	taskDataChan chan responseTaskData
262 263
}

264
type finishTaskRequest struct {
265 266 267
	key    responseKey
	status graphsync.ResponseStatusCode
	err    error
268 269 270 271 272 273
}

type setResponseDataRequest struct {
	key       responseKey
	loader    ipld.Loader
	traverser ipldutil.Traverser
274 275
}

Hannah Howard's avatar
Hannah Howard committed
276 277 278 279 280
type responseUpdateRequest struct {
	key        responseKey
	updateChan chan []gsmsg.GraphSyncRequest
}

281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298
// Startup starts processing for the WantManager.
func (rm *ResponseManager) Startup() {
	go rm.run()
}

// Shutdown ends processing for the want manager.
func (rm *ResponseManager) Shutdown() {
	rm.cancelFn()
}

func (rm *ResponseManager) cleanupInProcessResponses() {
	for _, response := range rm.inProgressResponses {
		response.cancelFn()
	}
}

func (rm *ResponseManager) run() {
	defer rm.cleanupInProcessResponses()
299
	for i := uint64(0); i < rm.maxInProcessRequests; i++ {
300
		go rm.qe.processQueriesWorker()
301 302 303 304 305 306 307 308 309 310 311 312
	}

	for {
		select {
		case <-rm.ctx.Done():
			return
		case message := <-rm.messages:
			message.handle(rm)
		}
	}
}

Hannah Howard's avatar
Hannah Howard committed
313 314 315 316 317 318 319 320 321
func (rm *ResponseManager) processUpdate(key responseKey, update gsmsg.GraphSyncRequest) {
	response, ok := rm.inProgressResponses[key]
	if !ok {
		log.Warnf("received update for non existent request, peer %s, request ID %d", key.p.Pretty(), key.requestID)
		return
	}
	if !response.isPaused {
		response.updates = append(response.updates, update)
		select {
322
		case response.signals.updateSignal <- struct{}{}:
Hannah Howard's avatar
Hannah Howard committed
323 324 325 326 327 328
		default:
		}
		return
	}
	result := rm.updateHooks.ProcessUpdateHooks(key.p, response.request, update)
	peerResponseSender := rm.peerManager.SenderForPeer(key.p)
Hannah Howard's avatar
Hannah Howard committed
329 330 331 332 333 334
	err := peerResponseSender.Transaction(key.requestID, func(transaction peerresponsemanager.PeerResponseTransactionSender) error {
		for _, extension := range result.Extensions {
			transaction.SendExtensionData(extension)
		}
		if result.Err != nil {
			transaction.FinishWithError(graphsync.RequestFailedUnknown)
335
			transaction.AddNotifee(notifications.Notifee{Topic: graphsync.RequestFailedUnknown, Subscriber: response.subscriber})
Hannah Howard's avatar
Hannah Howard committed
336 337 338 339 340
		}
		return nil
	})
	if err != nil {
		log.Errorf("Error processing update: %s", err)
Hannah Howard's avatar
Hannah Howard committed
341 342 343 344 345 346 347 348 349 350
	}
	if result.Err != nil {
		delete(rm.inProgressResponses, key)
		response.cancelFn()
		return
	}
	if result.Unpause {
		err := rm.unpauseRequest(key.p, key.requestID)
		if err != nil {
			log.Warnf("error unpausing request: %s", err.Error())
351 352
		}
	}
Hannah Howard's avatar
Hannah Howard committed
353

354 355
}

356
func (rm *ResponseManager) unpauseRequest(p peer.ID, requestID graphsync.RequestID, extensions ...graphsync.ExtensionData) error {
Hannah Howard's avatar
Hannah Howard committed
357 358 359 360 361 362 363 364 365
	key := responseKey{p, requestID}
	inProgressResponse, ok := rm.inProgressResponses[key]
	if !ok {
		return errors.New("could not find request")
	}
	if !inProgressResponse.isPaused {
		return errors.New("request is not paused")
	}
	inProgressResponse.isPaused = false
366 367 368 369 370 371 372 373 374
	if len(extensions) > 0 {
		peerResponseSender := rm.peerManager.SenderForPeer(key.p)
		_ = peerResponseSender.Transaction(requestID, func(transaction peerresponsemanager.PeerResponseTransactionSender) error {
			for _, extension := range extensions {
				transaction.SendExtensionData(extension)
			}
			return nil
		})
	}
Hannah Howard's avatar
Hannah Howard committed
375 376 377 378 379 380 381 382
	rm.queryQueue.PushTasks(p, peertask.Task{Topic: key, Priority: math.MaxInt32, Work: 1})
	select {
	case rm.workSignal <- struct{}{}:
	default:
	}
	return nil
}

383
func (rm *ResponseManager) abortRequest(p peer.ID, requestID graphsync.RequestID, err error) error {
384 385 386 387 388 389 390 391 392
	key := responseKey{p, requestID}
	rm.queryQueue.Remove(key, key.p)
	response, ok := rm.inProgressResponses[key]
	if !ok {
		return errors.New("could not find request")
	}

	if response.isPaused {
		peerResponseSender := rm.peerManager.SenderForPeer(key.p)
393 394
		if isContextErr(err) {

395 396
			rm.cancelledListeners.NotifyCancelledListeners(p, response.request)
			peerResponseSender.FinishWithCancel(requestID)
397 398
		} else if err != errNetworkError {
			peerResponseSender.FinishWithError(requestID, graphsync.RequestCancelled, notifications.Notifee{Topic: graphsync.RequestCancelled, Subscriber: response.subscriber})
399
		}
400 401 402 403 404
		delete(rm.inProgressResponses, key)
		response.cancelFn()
		return nil
	}
	select {
405
	case response.signals.errSignal <- err:
406 407 408 409 410
	default:
	}
	return nil
}

411 412 413 414
func (prm *processRequestMessage) handle(rm *ResponseManager) {
	for _, request := range prm.requests {
		key := responseKey{p: prm.p, requestID: request.ID()}
		if request.IsCancel() {
415
			_ = rm.abortRequest(prm.p, request.ID(), ipldutil.ContextCancelError{})
416 417 418 419 420 421 422
			continue
		}
		if request.IsUpdate() {
			rm.processUpdate(key, request)
			continue
		}
		ctx, cancelFn := context.WithCancel(rm.ctx)
423 424 425 426 427 428 429 430 431
		sub := notifications.NewMappableSubscriber(&subscriber{
			p:                     key.p,
			request:               request,
			ctx:                   rm.ctx,
			messages:              rm.messages,
			blockSentListeners:    rm.blockSentListeners,
			completedListeners:    rm.completedListeners,
			networkErrorListeners: rm.networkErrorListeners,
		}, notifications.IdentityTransform)
432 433
		rm.inProgressResponses[key] =
			&inProgressResponseStatus{
434 435 436 437
				ctx:        ctx,
				cancelFn:   cancelFn,
				subscriber: sub,
				request:    request,
438 439 440
				signals: signals{
					pauseSignal:  make(chan struct{}, 1),
					updateSignal: make(chan struct{}, 1),
441
					errSignal:    make(chan error, 1),
442
				},
443 444 445 446 447 448 449 450 451 452
			}
		// TODO: Use a better work estimation metric.
		rm.queryQueue.PushTasks(prm.p, peertask.Task{Topic: key, Priority: int(request.Priority()), Work: 1})
		select {
		case rm.workSignal <- struct{}{}:
		default:
		}
	}
}

453 454
func (rdr *responseDataRequest) handle(rm *ResponseManager) {
	response, ok := rm.inProgressResponses[rdr.key]
455
	var taskData responseTaskData
456
	if ok {
457
		taskData = responseTaskData{false, response.subscriber, response.ctx, response.request, response.loader, response.traverser, response.signals}
458
	} else {
459
		taskData = responseTaskData{empty: true}
460 461 462 463 464 465 466
	}
	select {
	case <-rm.ctx.Done():
	case rdr.taskDataChan <- taskData:
	}
}

467 468
func (ftr *finishTaskRequest) handle(rm *ResponseManager) {
	response, ok := rm.inProgressResponses[ftr.key]
469 470 471
	if !ok {
		return
	}
472
	if _, ok := ftr.err.(hooks.ErrPaused); ok {
473 474 475 476 477 478 479
		response.isPaused = true
		return
	}
	if ftr.err != nil {
		log.Infof("response failed: %w", ftr.err)
	}
	delete(rm.inProgressResponses, ftr.key)
480 481 482
	response.cancelFn()
}

483 484 485 486 487 488 489 490 491
func (srdr *setResponseDataRequest) handle(rm *ResponseManager) {
	response, ok := rm.inProgressResponses[srdr.key]
	if !ok {
		return
	}
	response.loader = srdr.loader
	response.traverser = srdr.traverser
}

Hannah Howard's avatar
Hannah Howard committed
492 493 494 495 496 497 498 499 500
func (rur *responseUpdateRequest) handle(rm *ResponseManager) {
	response, ok := rm.inProgressResponses[rur.key]
	var updates []gsmsg.GraphSyncRequest
	if ok {
		updates = response.updates
		response.updates = nil
	} else {
		updates = nil
	}
501 502
	select {
	case <-rm.ctx.Done():
Hannah Howard's avatar
Hannah Howard committed
503
	case rur.updateChan <- updates:
504 505
	}
}
506

Hannah Howard's avatar
Hannah Howard committed
507
func (sm *synchronizeMessage) handle(rm *ResponseManager) {
508
	select {
Hannah Howard's avatar
Hannah Howard committed
509
	case <-rm.ctx.Done():
510
	case sm.sync <- nil:
511 512 513 514
	}
}

func (urm *unpauseRequestMessage) handle(rm *ResponseManager) {
515
	err := rm.unpauseRequest(urm.p, urm.requestID, urm.extensions...)
516 517 518 519 520
	select {
	case <-rm.ctx.Done():
	case urm.response <- err:
	}
}
521 522 523 524 525 526 527 528 529 530 531

func (prm *pauseRequestMessage) pauseRequest(rm *ResponseManager) error {
	key := responseKey{prm.p, prm.requestID}
	inProgressResponse, ok := rm.inProgressResponses[key]
	if !ok {
		return errors.New("could not find request")
	}
	if inProgressResponse.isPaused {
		return errors.New("request is already paused")
	}
	select {
532
	case inProgressResponse.signals.pauseSignal <- struct{}{}:
533 534 535 536 537 538 539 540 541 542 543 544 545
	default:
	}
	return nil
}

func (prm *pauseRequestMessage) handle(rm *ResponseManager) {
	err := prm.pauseRequest(rm)
	select {
	case <-rm.ctx.Done():
	case prm.response <- err:
	}
}

546 547
func (crm *errorRequestMessage) handle(rm *ResponseManager) {
	err := rm.abortRequest(crm.p, crm.requestID, crm.err)
548 549
	select {
	case <-rm.ctx.Done():
550
	case crm.response <- err:
551 552
	}
}