MCPcopy Create free account
hub / github.com/encodeous/nylon / RoutineSequentialReceiver

Method RoutineSequentialReceiver

polyamide/device/receive.go:428–510  ·  view source on GitHub ↗
(maxBatchSize int)

Source from the content-addressed store, hash-verified

426}
427
428func (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 }

Callers 1

StartMethod · 0.95

Calls 15

ReceivedWithKeypairMethod · 0.95
SetEndpointFromPacketMethod · 0.95
SendStagedPacketsMethod · 0.95
keepKeyFreshReceivingMethod · 0.95
timersDataReceivedMethod · 0.95
NewTCStateFunction · 0.85
ValidateCounterMethod · 0.80
GetTCElementMethod · 0.80
TCBatchMethod · 0.80

Tested by

no test coverage detected