responsemanager.go 14.2 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
	"github.com/ipfs/go-graphsync/responsemanager/hooks"
10

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

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

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

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

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

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

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

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

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

79 80 81 82 83 84 85 86 87 88 89 90
// 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
91 92 93 94 95 96 97 98
	ctx                 context.Context
	cancelFn            context.CancelFunc
	loader              ipld.Loader
	peerManager         PeerManager
	queryQueue          QueryQueue
	requestHooks        RequestHooks
	blockHooks          BlockHooks
	updateHooks         UpdateHooks
99 100 101 102
	messages            chan responseManagerMessage
	workSignal          chan struct{}
	ticker              *time.Ticker
	inProgressResponses map[responseKey]*inProgressResponseStatus
103 104 105 106 107
}

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

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

// ProcessRequests processes incoming requests for the given peer
137
func (rm *ResponseManager) ProcessRequests(ctx context.Context, p peer.ID, requests []gsmsg.GraphSyncRequest) {
138 139 140
	select {
	case rm.messages <- &processRequestMessage{p, requests}:
	case <-rm.ctx.Done():
141
	case <-ctx.Done():
142 143 144
	}
}

145 146 147 148
type unpauseRequestMessage struct {
	p         peer.ID
	requestID graphsync.RequestID
	response  chan error
149 150
}

151 152 153 154 155 156 157 158 159 160 161 162 163
// 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
164 165 166
	}
}

167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188
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
}

189
type finishTaskRequest struct {
190
	key responseKey
191 192 193 194 195 196 197
	err error
}

type setResponseDataRequest struct {
	key       responseKey
	loader    ipld.Loader
	traverser ipldutil.Traverser
198 199
}

Hannah Howard's avatar
Hannah Howard committed
200 201 202 203 204
type responseUpdateRequest struct {
	key        responseKey
	updateChan chan []gsmsg.GraphSyncRequest
}

205
func (rm *ResponseManager) processQueriesWorker() {
206
	const targetWork = 1
207 208 209
	taskDataChan := make(chan *responseTaskData)
	var taskData *responseTaskData
	for {
210 211
		pid, tasks, _ := rm.queryQueue.PopTasks(targetWork)
		for len(tasks) == 0 {
212 213 214 215
			select {
			case <-rm.ctx.Done():
				return
			case <-rm.workSignal:
216
				pid, tasks, _ = rm.queryQueue.PopTasks(targetWork)
217 218
			case <-rm.ticker.C:
				rm.queryQueue.ThawRound()
219
				pid, tasks, _ = rm.queryQueue.PopTasks(targetWork)
220 221
			}
		}
222 223
		for _, task := range tasks {
			key := task.Topic.(responseKey)
224 225 226 227 228 229 230 231 232 233
			select {
			case rm.messages <- &responseDataRequest{key, taskDataChan}:
			case <-rm.ctx.Done():
				return
			}
			select {
			case taskData = <-taskDataChan:
			case <-rm.ctx.Done():
				return
			}
234
			err := rm.executeTask(key, taskData)
235
			select {
236
			case rm.messages <- &finishTaskRequest{key, err}:
237 238 239
			case <-rm.ctx.Done():
			}
		}
240
		rm.queryQueue.TasksDone(pid, tasks...)
241

242 243 244 245
	}

}

246 247 248 249 250 251 252 253 254 255 256 257 258 259
func (rm *ResponseManager) executeTask(key responseKey, taskData *responseTaskData) error {
	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 {
			return err
		}
		select {
		case <-rm.ctx.Done():
			return nil
		case rm.messages <- &setResponseDataRequest{key, loader, traverser}:
		}
260
	}
Hannah Howard's avatar
Hannah Howard committed
261
	return rm.executeQuery(key.p, taskData.request, loader, traverser, taskData.updateSignal)
262 263
}

264
func (rm *ResponseManager) prepareQuery(ctx context.Context,
265
	p peer.ID,
266 267
	request gsmsg.GraphSyncRequest) (ipld.Loader, ipldutil.Traverser, error) {
	result := rm.requestHooks.ProcessRequestHooks(p, request)
268
	peerResponseSender := rm.peerManager.SenderForPeer(p)
269 270
	for _, extension := range result.Extensions {
		peerResponseSender.SendExtensionData(request.ID(), extension)
271
	}
272
	if result.Err != nil || !result.IsValidated {
273
		peerResponseSender.FinishWithError(request.ID(), graphsync.RequestFailedUnknown)
274
		return nil, nil, errors.New("request not valid")
275
	}
276
	rootLink := cidlink.Link{Cid: request.Root()}
277 278 279 280 281 282 283 284 285 286 287 288
	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
}

Hannah Howard's avatar
Hannah Howard committed
289 290
func (rm *ResponseManager) executeQuery(
	p peer.ID,
291 292
	request gsmsg.GraphSyncRequest,
	loader ipld.Loader,
Hannah Howard's avatar
Hannah Howard committed
293 294 295
	traverser ipldutil.Traverser,
	updateSignal chan struct{}) error {
	updateChan := make(chan []gsmsg.GraphSyncRequest)
296 297
	peerResponseSender := rm.peerManager.SenderForPeer(p)
	err := runtraversal.RunTraversal(loader, traverser, func(link ipld.Link, data []byte) error {
Hannah Howard's avatar
Hannah Howard committed
298 299 300 301
		err := rm.checkForUpdates(p, request, updateSignal, updateChan, peerResponseSender)
		if err != nil {
			return err
		}
302 303 304 305 306 307 308 309 310 311 312 313
		blockData := peerResponseSender.SendResponse(request.ID(), link, data)
		if blockData.BlockSize() > 0 {
			result := rm.blockHooks.ProcessBlockHooks(p, request, blockData)
			for _, extension := range result.Extensions {
				peerResponseSender.SendExtensionData(request.ID(), extension)
			}
			if result.Err != nil {
				return result.Err
			}
		}
		return nil
	})
314
	if err != nil {
Hannah Howard's avatar
Hannah Howard committed
315
		if err != hooks.ErrPaused {
316
			peerResponseSender.FinishWithError(request.ID(), graphsync.RequestFailedUnknown)
Hannah Howard's avatar
Hannah Howard committed
317 318
		} else {
			peerResponseSender.PauseRequest(request.ID())
319 320
		}
		return err
321 322
	}
	peerResponseSender.FinishRequest(request.ID())
323
	return nil
324 325
}

Hannah Howard's avatar
Hannah Howard committed
326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355
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
}

356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390
// 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
391
		if request.IsCancel() {
392 393 394 395
			rm.queryQueue.Remove(key, key.p)
			response, ok := rm.inProgressResponses[key]
			if ok {
				response.cancelFn()
Hannah Howard's avatar
Hannah Howard committed
396 397 398 399 400 401 402 403 404 405 406 407 408 409 410
				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),
411
			}
Hannah Howard's avatar
Hannah Howard committed
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 446 447 448 449
		// 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)
	for _, extension := range result.Extensions {
		peerResponseSender.SendExtensionData(key.requestID, extension)
	}
	if result.Err != nil {
		peerResponseSender.FinishWithError(key.requestID, graphsync.RequestFailedUnknown)
		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())
450 451 452 453
		}
	}
}

Hannah Howard's avatar
Hannah Howard committed
454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471
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
}

472 473 474 475
func (rdr *responseDataRequest) handle(rm *ResponseManager) {
	response, ok := rm.inProgressResponses[rdr.key]
	var taskData *responseTaskData
	if ok {
Hannah Howard's avatar
Hannah Howard committed
476
		taskData = &responseTaskData{response.ctx, response.request, response.loader, response.traverser, response.updateSignal}
477 478 479 480 481 482 483 484 485
	} else {
		taskData = nil
	}
	select {
	case <-rm.ctx.Done():
	case rdr.taskDataChan <- taskData:
	}
}

486 487
func (ftr *finishTaskRequest) handle(rm *ResponseManager) {
	response, ok := rm.inProgressResponses[ftr.key]
488 489 490
	if !ok {
		return
	}
Hannah Howard's avatar
Hannah Howard committed
491
	if ftr.err == hooks.ErrPaused {
492 493 494 495 496 497 498
		response.isPaused = true
		return
	}
	if ftr.err != nil {
		log.Infof("response failed: %w", ftr.err)
	}
	delete(rm.inProgressResponses, ftr.key)
499 500 501
	response.cancelFn()
}

502 503 504 505 506 507 508 509 510
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
511 512 513 514 515 516 517 518 519
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
	}
520 521
	select {
	case <-rm.ctx.Done():
Hannah Howard's avatar
Hannah Howard committed
522
	case rur.updateChan <- updates:
523 524
	}
}
525

Hannah Howard's avatar
Hannah Howard committed
526
func (sm *synchronizeMessage) handle(rm *ResponseManager) {
527
	select {
Hannah Howard's avatar
Hannah Howard committed
528 529
	case <-rm.ctx.Done():
	case sm.sync <- struct{}{}:
530 531 532 533
	}
}

func (urm *unpauseRequestMessage) handle(rm *ResponseManager) {
Hannah Howard's avatar
Hannah Howard committed
534
	err := rm.unpauseRequest(urm.p, urm.requestID)
535 536 537 538 539
	select {
	case <-rm.ctx.Done():
	case urm.response <- err:
	}
}