MCPcopy Create free account
hub / github.com/DOSNetwork/core / dispatch

Method dispatch

p2p/client.go:295–368  ·  view source on GitHub ↗
(replyMsg, receivedMsg chan P2PMessage)

Source from the content-addressed store, hash-verified

293}
294
295func (c *client) dispatch(replyMsg, receivedMsg chan P2PMessage) (out chan p2pRequest) {
296 out = make(chan p2pRequest)
297 requests := make(map[uint64]*p2pRequest)
298 var nonce uint64
299 idleTimer := time.NewTimer(idleTimeout)
300 var readTimer bool
301 go func() {
302 defer close(out)
303 for {
304 var (
305 ok bool
306 req p2pRequest
307 msg P2PMessage
308 )
309 select {
310 case <-c.ctx.Done():
311 for _, req := range requests {
312 req.replyResult(nil, errors.Errorf("client dispatch: %w", c.ctx.Err()))
313 }
314 if !idleTimer.Stop() && !readTimer {
315 <-idleTimer.C
316 }
317 return
318 case <-idleTimer.C:
319 readTimer = true
320 err := c.close()
321 if err != nil {
322 errors.Errorf("client close: %w", err)
323 }
324 case req, ok = <-c.peerSend:
325 if ok {
326 if !idleTimer.Stop() {
327 <-idleTimer.C
328 }
329 idleTimer.Reset(idleTimeout)
330 if req.rType != replyReq {
331 req.nonce = nonce
332 requests[nonce] = &req
333 nonce++
334 }
335 select {
336 case <-c.ctx.Done():
337 case out <- req:
338 }
339 }
340 case msg, ok = <-replyMsg:
341 if ok {
342 if !idleTimer.Stop() {
343 <-idleTimer.C
344 }
345 idleTimer.Reset(idleTimeout)
346 p2pRequest := requests[msg.RequestNonce]
347 if p2pRequest != nil {
348 delete(requests, msg.RequestNonce)
349 select {
350 case <-p2pRequest.ctx.Done():
351 default:
352 p2pRequest.replyResult(msg, nil)

Callers 1

runMethod · 0.95

Calls 4

closeMethod · 0.95
reportMsgMethod · 0.95
replyResultMethod · 0.80
ResetMethod · 0.45

Tested by

no test coverage detected