responsemanager.go 14 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"
Hannah Howard's avatar
Hannah Howard committed
17
	"github.com/ipfs/go-graphsync/responsemanager/hooks"
18 19 20
	"github.com/ipfs/go-graphsync/responsemanager/peerresponsemanager"
)

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

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

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

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

44
type signals struct {
45
	pauseSignal  chan struct{}
Hannah Howard's avatar
Hannah Howard committed
46
	updateSignal chan struct{}
47 48 49 50 51 52 53 54 55 56
	stopSignal   chan bool
}

type responseTaskData struct {
	empty     bool
	ctx       context.Context
	request   gsmsg.GraphSyncRequest
	loader    ipld.Loader
	traverser ipldutil.Traverser
	signals   signals
57 58
}

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

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

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

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

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

94 95 96 97 98 99 100 101 102 103 104 105
// 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
106 107 108 109 110
	ctx                 context.Context
	cancelFn            context.CancelFunc
	peerManager         PeerManager
	queryQueue          QueryQueue
	updateHooks         UpdateHooks
111
	cancelledListeners  CancelledListeners
112
	completedListeners  CompletedListeners
113 114
	messages            chan responseManagerMessage
	workSignal          chan struct{}
115
	qe                  *queryExecutor
116
	inProgressResponses map[responseKey]*inProgressResponseStatus
117 118 119 120 121
}

// New creates a new response manager from the given context, loader,
// bridge to IPLD interface, peerManager, and queryQueue.
func New(ctx context.Context,
122
	loader ipld.Loader,
123
	peerManager PeerManager,
124 125
	queryQueue QueryQueue,
	requestHooks RequestHooks,
Hannah Howard's avatar
Hannah Howard committed
126
	blockHooks BlockHooks,
127
	updateHooks UpdateHooks,
128
	completedListeners CompletedListeners,
129 130
	cancelledListeners CancelledListeners,
) *ResponseManager {
131
	ctx, cancelFn := context.WithCancel(ctx)
132 133 134
	messages := make(chan responseManagerMessage, 16)
	workSignal := make(chan struct{}, 1)
	qe := &queryExecutor{
135 136 137
		requestHooks:       requestHooks,
		blockHooks:         blockHooks,
		updateHooks:        updateHooks,
138
		completedListeners: completedListeners,
139 140 141 142 143 144 145 146
		cancelledListeners: cancelledListeners,
		peerManager:        peerManager,
		loader:             loader,
		queryQueue:         queryQueue,
		messages:           messages,
		ctx:                ctx,
		workSignal:         workSignal,
		ticker:             time.NewTicker(thawSpeed),
147
	}
148 149 150 151 152
	return &ResponseManager{
		ctx:                 ctx,
		cancelFn:            cancelFn,
		peerManager:         peerManager,
		queryQueue:          queryQueue,
Hannah Howard's avatar
Hannah Howard committed
153
		updateHooks:         updateHooks,
154
		completedListeners:  completedListeners,
155
		cancelledListeners:  cancelledListeners,
156 157 158
		messages:            messages,
		workSignal:          workSignal,
		qe:                  qe,
159
		inProgressResponses: make(map[responseKey]*inProgressResponseStatus),
160 161 162 163 164 165 166 167 168
	}
}

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

// ProcessRequests processes incoming requests for the given peer
169
func (rm *ResponseManager) ProcessRequests(ctx context.Context, p peer.ID, requests []gsmsg.GraphSyncRequest) {
170 171 172
	select {
	case rm.messages <- &processRequestMessage{p, requests}:
	case <-rm.ctx.Done():
173
	case <-ctx.Done():
174 175 176
	}
}

177
type unpauseRequestMessage struct {
178 179 180 181 182 183 184 185 186 187 188 189 190
	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 {
191 192 193
	p         peer.ID
	requestID graphsync.RequestID
	response  chan error
194 195
}

196 197
// 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 {
198
	response := make(chan error, 1)
199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214
	return rm.sendSyncMessage(&pauseRequestMessage{p, requestID, response}, response)
}

type cancelRequestMessage struct {
	p         peer.ID
	requestID graphsync.RequestID
	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)
	return rm.sendSyncMessage(&cancelRequestMessage{p, requestID, response}, response)
}

func (rm *ResponseManager) sendSyncMessage(message responseManagerMessage, response chan error) error {
215 216 217
	select {
	case <-rm.ctx.Done():
		return errors.New("Context Cancelled")
218
	case rm.messages <- message:
219 220 221 222 223 224
	}
	select {
	case <-rm.ctx.Done():
		return errors.New("Context Cancelled")
	case err := <-response:
		return err
225 226 227
	}
}

228
type synchronizeMessage struct {
229
	sync chan error
230 231 232 233
}

// this is a test utility method to force all messages to get processed
func (rm *ResponseManager) synchronize() {
234 235
	sync := make(chan error)
	_ = rm.sendSyncMessage(&synchronizeMessage{sync}, sync)
236 237 238 239
}

type responseDataRequest struct {
	key          responseKey
240
	taskDataChan chan responseTaskData
241 242
}

243
type finishTaskRequest struct {
244 245 246
	key    responseKey
	status graphsync.ResponseStatusCode
	err    error
247 248 249 250 251 252
}

type setResponseDataRequest struct {
	key       responseKey
	loader    ipld.Loader
	traverser ipldutil.Traverser
253 254
}

Hannah Howard's avatar
Hannah Howard committed
255 256 257 258 259
type responseUpdateRequest struct {
	key        responseKey
	updateChan chan []gsmsg.GraphSyncRequest
}

260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278
// 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++ {
279
		go rm.qe.processQueriesWorker()
280 281 282 283 284 285 286 287 288 289 290 291
	}

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

Hannah Howard's avatar
Hannah Howard committed
292 293 294 295 296 297 298 299 300
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 {
301
		case response.signals.updateSignal <- struct{}{}:
Hannah Howard's avatar
Hannah Howard committed
302 303 304 305 306 307
		default:
		}
		return
	}
	result := rm.updateHooks.ProcessUpdateHooks(key.p, response.request, update)
	peerResponseSender := rm.peerManager.SenderForPeer(key.p)
Hannah Howard's avatar
Hannah Howard committed
308 309 310 311 312 313 314 315 316 317 318
	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
319 320 321 322 323 324 325 326 327 328
	}
	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())
329 330
		}
	}
Hannah Howard's avatar
Hannah Howard committed
331

332 333
}

334
func (rm *ResponseManager) unpauseRequest(p peer.ID, requestID graphsync.RequestID, extensions ...graphsync.ExtensionData) error {
Hannah Howard's avatar
Hannah Howard committed
335 336 337 338 339 340 341 342 343
	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
344 345 346 347 348 349 350 351 352
	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
353 354 355 356 357 358 359 360
	rm.queryQueue.PushTasks(p, peertask.Task{Topic: key, Priority: math.MaxInt32, Work: 1})
	select {
	case rm.workSignal <- struct{}{}:
	default:
	}
	return nil
}

361 362 363 364 365 366 367 368 369 370
func (rm *ResponseManager) cancelRequest(p peer.ID, requestID graphsync.RequestID, selfCancel bool) error {
	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)
371 372 373 374 375 376 377
		if selfCancel {
			rm.completedListeners.NotifyCompletedListeners(p, response.request, graphsync.RequestCancelled)
			peerResponseSender.FinishWithError(requestID, graphsync.RequestCancelled)
		} else {
			rm.cancelledListeners.NotifyCancelledListeners(p, response.request)
			peerResponseSender.FinishWithCancel(requestID)
		}
378 379 380 381 382 383 384 385 386 387 388
		delete(rm.inProgressResponses, key)
		response.cancelFn()
		return nil
	}
	select {
	case response.signals.stopSignal <- selfCancel:
	default:
	}
	return nil
}

389 390 391 392
func (prm *processRequestMessage) handle(rm *ResponseManager) {
	for _, request := range prm.requests {
		key := responseKey{p: prm.p, requestID: request.ID()}
		if request.IsCancel() {
393
			_ = rm.cancelRequest(prm.p, request.ID(), false)
394 395 396 397 398 399 400 401 402
			continue
		}
		if request.IsUpdate() {
			rm.processUpdate(key, request)
			continue
		}
		ctx, cancelFn := context.WithCancel(rm.ctx)
		rm.inProgressResponses[key] =
			&inProgressResponseStatus{
403 404 405 406 407 408 409 410
				ctx:      ctx,
				cancelFn: cancelFn,
				request:  request,
				signals: signals{
					pauseSignal:  make(chan struct{}, 1),
					updateSignal: make(chan struct{}, 1),
					stopSignal:   make(chan bool, 1),
				},
411 412 413 414 415 416 417 418 419 420
			}
		// 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:
		}
	}
}

421 422
func (rdr *responseDataRequest) handle(rm *ResponseManager) {
	response, ok := rm.inProgressResponses[rdr.key]
423
	var taskData responseTaskData
424
	if ok {
425
		taskData = responseTaskData{false, response.ctx, response.request, response.loader, response.traverser, response.signals}
426
	} else {
427
		taskData = responseTaskData{empty: true}
428 429 430 431 432 433 434
	}
	select {
	case <-rm.ctx.Done():
	case rdr.taskDataChan <- taskData:
	}
}

435 436
func (ftr *finishTaskRequest) handle(rm *ResponseManager) {
	response, ok := rm.inProgressResponses[ftr.key]
437 438 439
	if !ok {
		return
	}
440
	if _, ok := ftr.err.(hooks.ErrPaused); ok {
441 442 443 444 445 446 447
		response.isPaused = true
		return
	}
	if ftr.err != nil {
		log.Infof("response failed: %w", ftr.err)
	}
	delete(rm.inProgressResponses, ftr.key)
448 449 450
	response.cancelFn()
}

451 452 453 454 455 456 457 458 459
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
460 461 462 463 464 465 466 467 468
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
	}
469 470
	select {
	case <-rm.ctx.Done():
Hannah Howard's avatar
Hannah Howard committed
471
	case rur.updateChan <- updates:
472 473
	}
}
474

Hannah Howard's avatar
Hannah Howard committed
475
func (sm *synchronizeMessage) handle(rm *ResponseManager) {
476
	select {
Hannah Howard's avatar
Hannah Howard committed
477
	case <-rm.ctx.Done():
478
	case sm.sync <- nil:
479 480 481 482
	}
}

func (urm *unpauseRequestMessage) handle(rm *ResponseManager) {
483
	err := rm.unpauseRequest(urm.p, urm.requestID, urm.extensions...)
484 485 486 487 488
	select {
	case <-rm.ctx.Done():
	case urm.response <- err:
	}
}
489 490 491 492 493 494 495 496 497 498 499

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 {
500
	case inProgressResponse.signals.pauseSignal <- struct{}{}:
501 502 503 504 505 506 507 508 509 510 511 512 513 514
	default:
	}
	return nil
}

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

func (crm *cancelRequestMessage) handle(rm *ResponseManager) {
515
	err := rm.cancelRequest(crm.p, crm.requestID, true)
516 517
	select {
	case <-rm.ctx.Done():
518
	case crm.response <- err:
519 520
	}
}