* Receives incoming datagrams for the device * * Every time the bind is updated a new routine is started for * IPv4 and IPv6 (separately) */
(maxBatchSize int, recv conn.ReceiveFunc)
| 69 | * IPv4 and IPv6 (separately) |
| 70 | */ |
| 71 | func (device *Device) RoutineReceiveIncoming(maxBatchSize int, recv conn.ReceiveFunc) { |
| 72 | recvName := recv.PrettyName() |
| 73 | defer func() { |
| 74 | device.Log.Verbosef("Routine: receive incoming %s - stopped", recvName) |
| 75 | device.queue.decryption.wg.Done() |
| 76 | device.queue.handshake.wg.Done() |
| 77 | device.net.stopping.Done() |
| 78 | }() |
| 79 | |
| 80 | device.Log.Verbosef("Routine: receive incoming %s - started", recvName) |
| 81 | |
| 82 | // receive datagrams until conn is closed |
| 83 | |
| 84 | var ( |
| 85 | bufsArrs = make([]*[MaxMessageSize]byte, maxBatchSize) |
| 86 | bufs = make([][]byte, maxBatchSize) |
| 87 | err error |
| 88 | sizes = make([]int, maxBatchSize) |
| 89 | count int |
| 90 | endpoints = make([]conn.Endpoint, maxBatchSize) |
| 91 | deathSpiral int |
| 92 | elemsByPeer = make(map[*Peer]*QueueInboundElementsContainer, maxBatchSize) |
| 93 | ) |
| 94 | |
| 95 | for i := range bufsArrs { |
| 96 | bufsArrs[i] = device.GetMessageBuffer() |
| 97 | bufs[i] = bufsArrs[i][:] |
| 98 | } |
| 99 | |
| 100 | defer func() { |
| 101 | for i := 0; i < maxBatchSize; i++ { |
| 102 | if bufsArrs[i] != nil { |
| 103 | device.PutMessageBuffer(bufsArrs[i]) |
| 104 | } |
| 105 | } |
| 106 | }() |
| 107 | |
| 108 | for { |
| 109 | count, err = recv(bufs, sizes, endpoints) |
| 110 | if err != nil { |
| 111 | if errors.Is(err, net.ErrClosed) { |
| 112 | return |
| 113 | } |
| 114 | device.Log.Verbosef("Failed to receive %s packet: %v", recvName, err) |
| 115 | if neterr, ok := err.(net.Error); ok && !neterr.Temporary() { |
| 116 | return |
| 117 | } |
| 118 | if deathSpiral < 10 { |
| 119 | deathSpiral++ |
| 120 | time.Sleep(time.Second / 3) |
| 121 | continue |
| 122 | } |
| 123 | return |
| 124 | } |
| 125 | deathSpiral = 0 |
| 126 | |
| 127 | // handle each packet in the batch |
| 128 | for i, size := range sizes[:count] { |
no test coverage detected