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

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

9 10
	"github.com/ipfs/go-cid"
	"github.com/ipfs/go-graphsync/cidset"
Hannah Howard's avatar
Hannah Howard committed
11
	"github.com/ipfs/go-graphsync/responsemanager/hooks"
12

13
	"github.com/ipfs/go-graphsync"
14
	"github.com/ipfs/go-graphsync/ipldutil"
15 16
	gsmsg "github.com/ipfs/go-graphsync/message"
	"github.com/ipfs/go-graphsync/responsemanager/peerresponsemanager"
17 18
	"github.com/ipfs/go-graphsync/responsemanager/runtraversal"
	logging "github.com/ipfs/go-log"
19
	"github.com/ipfs/go-peertaskqueue/peertask"
20 21
	ipld "github.com/ipld/go-ipld-prime"
	cidlink "github.com/ipld/go-ipld-prime/linking/cid"
22
	"github.com/libp2p/go-libp2p-core/peer"
23 24
)

25 26
var log = logging.Logger("graphsync")

27 28 29 30 31 32
const (
	maxInProcessRequests = 6
	thawSpeed            = time.Millisecond * 100
)

type inProgressResponseStatus struct {
Hannah Howard's avatar
Hannah Howard committed
33 34 35 36 37 38 39 40
	ctx          context.Context
	cancelFn     func()
	request      gsmsg.GraphSyncRequest
	loader       ipld.Loader
	traverser    ipldutil.Traverser
	updateSignal chan struct{}
	updates      []gsmsg.GraphSyncRequest
	isPaused     bool
41 42 43 44
}

type responseKey struct {
	p         peer.ID
45
	requestID graphsync.RequestID
46 47 48
}

type responseTaskData struct {
Hannah Howard's avatar
Hannah Howard committed
49 50 51 52 53
	ctx          context.Context
	request      gsmsg.GraphSyncRequest
	loader       ipld.Loader
	traverser    ipldutil.Traverser
	updateSignal chan struct{}
54 55
}

56 57 58
// QueryQueue is an interface that can receive new selector query tasks
// and prioritize them as needed, and pop them off later
type QueryQueue interface {
59 60 61 62
	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)
63 64 65
	ThawRound()
}

66 67
// RequestHooks is an interface for processing request hooks
type RequestHooks interface {
Hannah Howard's avatar
Hannah Howard committed
68
	ProcessRequestHooks(p peer.ID, request graphsync.RequestData) hooks.RequestResult
69 70 71 72
}

// BlockHooks is an interface for processing block hooks
type BlockHooks interface {
Hannah Howard's avatar
Hannah Howard committed
73 74 75 76 77 78
	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
79 80
}

81 82 83 84 85
// CompletedListeners is an interface for notifying listeners that responses are complete
type CompletedListeners interface {
	NotifyCompletedListeners(p peer.ID, request graphsync.RequestData, status graphsync.ResponseStatusCode)
}

86 87 88 89 90 91 92 93 94 95 96 97
// 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 {
Hannah Howard's avatar
Hannah Howard committed
98 99 100 101 102 103 104 105
	ctx                 context.Context
	cancelFn            context.CancelFunc
	loader              ipld.Loader
	peerManager         PeerManager
	queryQueue          QueryQueue
	requestHooks        RequestHooks
	blockHooks          BlockHooks
	updateHooks         UpdateHooks
106
	completedListeners  CompletedListeners
107 108 109 110
	messages            chan responseManagerMessage
	workSignal          chan struct{}
	ticker              *time.Ticker
	inProgressResponses map[responseKey]*inProgressResponseStatus
111 112 113 114 115
}

// New creates a new response manager from the given context, loader,
// bridge to IPLD interface, peerManager, and queryQueue.
func New(ctx context.Context,
116
	loader ipld.Loader,
117
	peerManager PeerManager,
118 119
	queryQueue QueryQueue,
	requestHooks RequestHooks,
Hannah Howard's avatar
Hannah Howard committed
120
	blockHooks BlockHooks,
121 122
	updateHooks UpdateHooks,
	completedListeners CompletedListeners) *ResponseManager {
123 124 125 126 127 128 129
	ctx, cancelFn := context.WithCancel(ctx)
	return &ResponseManager{
		ctx:                 ctx,
		cancelFn:            cancelFn,
		loader:              loader,
		peerManager:         peerManager,
		queryQueue:          queryQueue,
130 131
		requestHooks:        requestHooks,
		blockHooks:          blockHooks,
Hannah Howard's avatar
Hannah Howard committed
132
		updateHooks:         updateHooks,
133
		completedListeners:  completedListeners,
134 135 136
		messages:            make(chan responseManagerMessage, 16),
		workSignal:          make(chan struct{}, 1),
		ticker:              time.NewTicker(thawSpeed),
137
		inProgressResponses: make(map[responseKey]*inProgressResponseStatus),
138 139 140 141 142 143 144 145 146
	}
}

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

// ProcessRequests processes incoming requests for the given peer
147
func (rm *ResponseManager) ProcessRequests(ctx context.Context, p peer.ID, requests []gsmsg.GraphSyncRequest) {
148 149 150
	select {
	case rm.messages <- &processRequestMessage{p, requests}:
	case <-rm.ctx.Done():
151
	case <-ctx.Done():
152 153 154
	}
}

155 156 157 158
type unpauseRequestMessage struct {
	p         peer.ID
	requestID graphsync.RequestID
	response  chan error
159 160
}

161 162 163 164 165 166 167 168 169 170 171 172 173
// UnpauseResponse unpauses a response that was previously paused
func (rm *ResponseManager) UnpauseResponse(p peer.ID, requestID graphsync.RequestID) error {
	response := make(chan error, 1)
	select {
	case <-rm.ctx.Done():
		return errors.New("Context Cancelled")
	case rm.messages <- &unpauseRequestMessage{p, requestID, response}:
	}
	select {
	case <-rm.ctx.Done():
		return errors.New("Context Cancelled")
	case err := <-response:
		return err
174 175 176
	}
}

177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198
type synchronizeMessage struct {
	sync chan struct{}
}

// this is a test utility method to force all messages to get processed
func (rm *ResponseManager) synchronize() {
	sync := make(chan struct{})
	select {
	case rm.messages <- &synchronizeMessage{sync}:
	case <-rm.ctx.Done():
	}
	select {
	case <-sync:
	case <-rm.ctx.Done():
	}
}

type responseDataRequest struct {
	key          responseKey
	taskDataChan chan *responseTaskData
}

199
type finishTaskRequest struct {
200 201 202
	key    responseKey
	status graphsync.ResponseStatusCode
	err    error
203 204 205 206 207 208
}

type setResponseDataRequest struct {
	key       responseKey
	loader    ipld.Loader
	traverser ipldutil.Traverser
209 210
}

Hannah Howard's avatar
Hannah Howard committed
211 212 213 214 215
type responseUpdateRequest struct {
	key        responseKey
	updateChan chan []gsmsg.GraphSyncRequest
}

216
func (rm *ResponseManager) processQueriesWorker() {
217
	const targetWork = 1
218 219 220
	taskDataChan := make(chan *responseTaskData)
	var taskData *responseTaskData
	for {
221 222
		pid, tasks, _ := rm.queryQueue.PopTasks(targetWork)
		for len(tasks) == 0 {
223 224 225 226
			select {
			case <-rm.ctx.Done():
				return
			case <-rm.workSignal:
227
				pid, tasks, _ = rm.queryQueue.PopTasks(targetWork)
228 229
			case <-rm.ticker.C:
				rm.queryQueue.ThawRound()
230
				pid, tasks, _ = rm.queryQueue.PopTasks(targetWork)
231 232
			}
		}
233 234
		for _, task := range tasks {
			key := task.Topic.(responseKey)
235 236 237 238 239 240 241 242 243 244
			select {
			case rm.messages <- &responseDataRequest{key, taskDataChan}:
			case <-rm.ctx.Done():
				return
			}
			select {
			case taskData = <-taskDataChan:
			case <-rm.ctx.Done():
				return
			}
245 246 247 248
			if taskData == nil {
				log.Info("Empty task on peer request stack")
				continue
			}
249
			status, err := rm.executeTask(key, taskData)
250
			select {
251
			case rm.messages <- &finishTaskRequest{key, status, err}:
252 253 254
			case <-rm.ctx.Done():
			}
		}
255
		rm.queryQueue.TasksDone(pid, tasks...)
256

257 258 259 260
	}

}

261
func (rm *ResponseManager) executeTask(key responseKey, taskData *responseTaskData) (graphsync.ResponseStatusCode, error) {
262 263 264 265 266 267
	var err error
	loader := taskData.loader
	traverser := taskData.traverser
	if loader == nil || traverser == nil {
		loader, traverser, err = rm.prepareQuery(taskData.ctx, key.p, taskData.request)
		if err != nil {
268
			return graphsync.RequestFailedUnknown, err
269 270 271
		}
		select {
		case <-rm.ctx.Done():
272
			return graphsync.RequestFailedUnknown, errors.New("context cancelled")
273 274
		case rm.messages <- &setResponseDataRequest{key, loader, traverser}:
		}
275
	}
Hannah Howard's avatar
Hannah Howard committed
276
	return rm.executeQuery(key.p, taskData.request, loader, traverser, taskData.updateSignal)
277 278
}

279
func (rm *ResponseManager) prepareQuery(ctx context.Context,
280
	p peer.ID,
281 282
	request gsmsg.GraphSyncRequest) (ipld.Loader, ipldutil.Traverser, error) {
	result := rm.requestHooks.ProcessRequestHooks(p, request)
283
	peerResponseSender := rm.peerManager.SenderForPeer(p)
Hannah Howard's avatar
Hannah Howard committed
284 285 286 287 288 289 290 291 292 293 294 295 296
	var validationErr error
	err := peerResponseSender.Transaction(request.ID(), func(transaction peerresponsemanager.PeerResponseTransactionSender) error {
		for _, extension := range result.Extensions {
			transaction.SendExtensionData(extension)
		}
		if result.Err != nil || !result.IsValidated {
			transaction.FinishWithError(graphsync.RequestFailedUnknown)
			validationErr = errors.New("request not valid")
		}
		return nil
	})
	if err != nil {
		return nil, nil, err
297
	}
Hannah Howard's avatar
Hannah Howard committed
298 299
	if validationErr != nil {
		return nil, nil, validationErr
300
	}
301 302 303
	if err := rm.processDoNoSendCids(request, peerResponseSender); err != nil {
		return nil, nil, err
	}
304
	rootLink := cidlink.Link{Cid: request.Root()}
305 306 307 308 309 310 311 312 313 314 315 316
	traverser := ipldutil.TraversalBuilder{
		Root:     rootLink,
		Selector: request.Selector(),
		Chooser:  result.CustomChooser,
	}.Start(ctx)
	loader := result.CustomLoader
	if loader == nil {
		loader = rm.loader
	}
	return loader, traverser, nil
}

317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338
func (rm *ResponseManager) processDoNoSendCids(request gsmsg.GraphSyncRequest, peerResponseSender peerresponsemanager.PeerResponseSender) error {
	doNotSendCidsData, has := request.Extension(graphsync.ExtensionDoNotSendCIDs)
	if !has {
		return nil
	}
	cidSet, err := cidset.DecodeCidSet(doNotSendCidsData)
	if err != nil {
		peerResponseSender.FinishWithError(request.ID(), graphsync.RequestFailedUnknown)
		return err
	}
	links := make([]ipld.Link, 0, cidSet.Len())
	err = cidSet.ForEach(func(c cid.Cid) error {
		links = append(links, cidlink.Link{Cid: c})
		return nil
	})
	if err != nil {
		return err
	}
	peerResponseSender.IgnoreBlocks(request.ID(), links)
	return nil
}

Hannah Howard's avatar
Hannah Howard committed
339 340
func (rm *ResponseManager) executeQuery(
	p peer.ID,
341 342
	request gsmsg.GraphSyncRequest,
	loader ipld.Loader,
Hannah Howard's avatar
Hannah Howard committed
343
	traverser ipldutil.Traverser,
344
	updateSignal chan struct{}) (graphsync.ResponseStatusCode, error) {
Hannah Howard's avatar
Hannah Howard committed
345
	updateChan := make(chan []gsmsg.GraphSyncRequest)
346 347
	peerResponseSender := rm.peerManager.SenderForPeer(p)
	err := runtraversal.RunTraversal(loader, traverser, func(link ipld.Link, data []byte) error {
Hannah Howard's avatar
Hannah Howard committed
348 349 350 351
		err := rm.checkForUpdates(p, request, updateSignal, updateChan, peerResponseSender)
		if err != nil {
			return err
		}
Hannah Howard's avatar
Hannah Howard committed
352 353 354 355 356 357 358 359 360 361 362
		var result hooks.BlockResult
		err = peerResponseSender.Transaction(request.ID(), func(transaction peerresponsemanager.PeerResponseTransactionSender) error {
			blockData := transaction.SendResponse(link, data)
			if blockData.BlockSize() > 0 {
				result = rm.blockHooks.ProcessBlockHooks(p, request, blockData)
				for _, extension := range result.Extensions {
					transaction.SendExtensionData(extension)
				}
				if result.Err == hooks.ErrPaused {
					transaction.PauseRequest()
				}
363
			}
Hannah Howard's avatar
Hannah Howard committed
364 365 366 367
			return nil
		})
		if err != nil {
			return err
368
		}
Hannah Howard's avatar
Hannah Howard committed
369
		return result.Err
370
	})
371
	if err != nil {
Hannah Howard's avatar
Hannah Howard committed
372
		if err != hooks.ErrPaused {
373
			peerResponseSender.FinishWithError(request.ID(), graphsync.RequestFailedUnknown)
374
			return graphsync.RequestFailedUnknown, err
375
		}
376
		return graphsync.RequestPaused, err
377
	}
378
	return peerResponseSender.FinishRequest(request.ID()), nil
379 380
}

Hannah Howard's avatar
Hannah Howard committed
381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410
func (rm *ResponseManager) checkForUpdates(
	p peer.ID,
	request gsmsg.GraphSyncRequest,
	updateSignal chan struct{},
	updateChan chan []gsmsg.GraphSyncRequest,
	peerResponseSender peerresponsemanager.PeerResponseSender) error {
	select {
	case <-updateSignal:
		select {
		case rm.messages <- &responseUpdateRequest{responseKey{p, request.ID()}, updateChan}:
		case <-rm.ctx.Done():
		}
		select {
		case updates := <-updateChan:
			for _, update := range updates {
				result := rm.updateHooks.ProcessUpdateHooks(p, request, update)
				for _, extension := range result.Extensions {
					peerResponseSender.SendExtensionData(request.ID(), extension)
				}
				if result.Err != nil {
					return result.Err
				}
			}
		case <-rm.ctx.Done():
		}
	default:
	}
	return nil
}

411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445
// 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()
	for i := 0; i < maxInProcessRequests; i++ {
		go rm.processQueriesWorker()
	}

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

func (prm *processRequestMessage) handle(rm *ResponseManager) {
	for _, request := range prm.requests {
		key := responseKey{p: prm.p, requestID: request.ID()}
Hannah Howard's avatar
Hannah Howard committed
446
		if request.IsCancel() {
447 448 449 450
			rm.queryQueue.Remove(key, key.p)
			response, ok := rm.inProgressResponses[key]
			if ok {
				response.cancelFn()
Hannah Howard's avatar
Hannah Howard committed
451 452 453 454 455 456 457 458 459 460 461 462 463 464 465
				delete(rm.inProgressResponses, key)
			}
			continue
		}
		if request.IsUpdate() {
			rm.processUpdate(key, request)
			continue
		}
		ctx, cancelFn := context.WithCancel(rm.ctx)
		rm.inProgressResponses[key] =
			&inProgressResponseStatus{
				ctx:          ctx,
				cancelFn:     cancelFn,
				request:      request,
				updateSignal: make(chan struct{}, 1),
466
			}
Hannah Howard's avatar
Hannah Howard committed
467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491
		// 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:
		}
	}
}

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 {
		case response.updateSignal <- struct{}{}:
		default:
		}
		return
	}
	result := rm.updateHooks.ProcessUpdateHooks(key.p, response.request, update)
	peerResponseSender := rm.peerManager.SenderForPeer(key.p)
Hannah Howard's avatar
Hannah Howard committed
492 493 494 495 496 497 498 499 500 501 502
	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)
		}
		return nil
	})
	if err != nil {
		log.Errorf("Error processing update: %s", err)
Hannah Howard's avatar
Hannah Howard committed
503 504 505 506 507 508 509 510 511 512
	}
	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())
513 514
		}
	}
Hannah Howard's avatar
Hannah Howard committed
515

516 517
}

Hannah Howard's avatar
Hannah Howard committed
518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535
func (rm *ResponseManager) unpauseRequest(p peer.ID, requestID graphsync.RequestID) error {
	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
	rm.queryQueue.PushTasks(p, peertask.Task{Topic: key, Priority: math.MaxInt32, Work: 1})
	select {
	case rm.workSignal <- struct{}{}:
	default:
	}
	return nil
}

536 537 538 539
func (rdr *responseDataRequest) handle(rm *ResponseManager) {
	response, ok := rm.inProgressResponses[rdr.key]
	var taskData *responseTaskData
	if ok {
Hannah Howard's avatar
Hannah Howard committed
540
		taskData = &responseTaskData{response.ctx, response.request, response.loader, response.traverser, response.updateSignal}
541 542 543 544 545 546 547 548 549
	} else {
		taskData = nil
	}
	select {
	case <-rm.ctx.Done():
	case rdr.taskDataChan <- taskData:
	}
}

550 551
func (ftr *finishTaskRequest) handle(rm *ResponseManager) {
	response, ok := rm.inProgressResponses[ftr.key]
552 553 554
	if !ok {
		return
	}
Hannah Howard's avatar
Hannah Howard committed
555
	if ftr.err == hooks.ErrPaused {
556 557 558
		response.isPaused = true
		return
	}
559
	rm.completedListeners.NotifyCompletedListeners(ftr.key.p, response.request, ftr.status)
560 561 562 563
	if ftr.err != nil {
		log.Infof("response failed: %w", ftr.err)
	}
	delete(rm.inProgressResponses, ftr.key)
564 565 566
	response.cancelFn()
}

567 568 569 570 571 572 573 574 575
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
576 577 578 579 580 581 582 583 584
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
	}
585 586
	select {
	case <-rm.ctx.Done():
Hannah Howard's avatar
Hannah Howard committed
587
	case rur.updateChan <- updates:
588 589
	}
}
590

Hannah Howard's avatar
Hannah Howard committed
591
func (sm *synchronizeMessage) handle(rm *ResponseManager) {
592
	select {
Hannah Howard's avatar
Hannah Howard committed
593 594
	case <-rm.ctx.Done():
	case sm.sync <- struct{}{}:
595 596 597 598
	}
}

func (urm *unpauseRequestMessage) handle(rm *ResponseManager) {
Hannah Howard's avatar
Hannah Howard committed
599
	err := rm.unpauseRequest(urm.p, urm.requestID)
600 601 602 603 604
	select {
	case <-rm.ctx.Done():
	case urm.response <- err:
	}
}