(replyMsg, receivedMsg chan P2PMessage)
| 293 | } |
| 294 | |
| 295 | func (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) |
no test coverage detected