(maxBatchSize int)
| 426 | } |
| 427 | |
| 428 | func (peer *Peer) RoutineSequentialReceiver(maxBatchSize int) { |
| 429 | device := peer.device |
| 430 | defer func() { |
| 431 | device.Log.Verbosef("%v - Routine: sequential receiver - stopped", peer) |
| 432 | peer.stopping.Done() |
| 433 | }() |
| 434 | device.Log.Verbosef("%v - Routine: sequential receiver - started", peer) |
| 435 | |
| 436 | var ( |
| 437 | tcBufs = make([]*TCElement, 0, maxBatchSize) |
| 438 | tcs = NewTCState() |
| 439 | ) |
| 440 | |
| 441 | for elemsContainer := range peer.queue.inbound.c { |
| 442 | if elemsContainer == nil { |
| 443 | return |
| 444 | } |
| 445 | elemsContainer.Lock() |
| 446 | validTailPacket := -1 |
| 447 | dataPacketReceived := false |
| 448 | rxBytesLen := uint64(0) |
| 449 | rxPkts := 0 |
| 450 | |
| 451 | for i, elem := range elemsContainer.elems { |
| 452 | if elem.packet == nil { |
| 453 | // decryption failed |
| 454 | continue |
| 455 | } |
| 456 | |
| 457 | if !elem.keypair.replayFilter.ValidateCounter(elem.counter, RejectAfterMessages) { |
| 458 | continue |
| 459 | } |
| 460 | |
| 461 | validTailPacket = i |
| 462 | if peer.ReceivedWithKeypair(elem.keypair) { |
| 463 | peer.SetEndpointFromPacket(elem.endpoint) |
| 464 | peer.timersHandshakeComplete() |
| 465 | peer.SendStagedPackets() |
| 466 | } |
| 467 | rxBytesLen += uint64(len(elem.packet) + MinMessageSize) |
| 468 | rxPkts++ |
| 469 | |
| 470 | if len(elem.packet) == 0 { |
| 471 | device.Log.Verbosef("%v - Receiving keepalive packet", peer) |
| 472 | continue |
| 473 | } |
| 474 | dataPacketReceived = true |
| 475 | |
| 476 | tce := device.GetTCElement() |
| 477 | tce.Packet = elem.packet |
| 478 | tce.Buffer = elem.buffer |
| 479 | elem.buffer = nil |
| 480 | elem.packet = nil |
| 481 | tce.FromEp = elem.endpoint |
| 482 | tce.FromPeer = peer |
| 483 | |
| 484 | tcBufs = append(tcBufs, tce) |
| 485 | } |
no test coverage detected