()
| 215 | } |
| 216 | |
| 217 | func (device *Device) RoutineReadFromTUN() { |
| 218 | defer func() { |
| 219 | device.Log.Verbosef("Routine: TUN reader - stopped") |
| 220 | device.state.stopping.Done() |
| 221 | device.queue.encryption.wg.Done() |
| 222 | }() |
| 223 | |
| 224 | device.Log.Verbosef("Routine: TUN reader - started") |
| 225 | |
| 226 | var ( |
| 227 | batchSize = device.BatchSize() |
| 228 | readErr error |
| 229 | rBufs = make([][]byte, batchSize) |
| 230 | bufs = make([]*[MaxMessageSize]byte, batchSize) |
| 231 | count = batchSize |
| 232 | sizes = make([]int, batchSize) |
| 233 | tcBufs = make([]*TCElement, 0, batchSize) |
| 234 | offset = MessageTransportHeaderSize |
| 235 | tcs = NewTCState() |
| 236 | ) |
| 237 | |
| 238 | for i := 0; i < batchSize; i++ { |
| 239 | bufs[i] = device.GetMessageBuffer() |
| 240 | rBufs[i] = bufs[i][:] |
| 241 | } |
| 242 | |
| 243 | for { |
| 244 | count, readErr = device.tun.device.Read(rBufs, sizes, offset) |
| 245 | |
| 246 | for i := 0; i < count; i++ { |
| 247 | if sizes[i] < 1 { |
| 248 | continue |
| 249 | } |
| 250 | tce := device.GetTCElement() |
| 251 | tce.Buffer = bufs[i] |
| 252 | tcBufs = append(tcBufs, tce) |
| 253 | |
| 254 | bufs[i] = device.GetMessageBuffer() |
| 255 | rBufs[i] = bufs[i][:] |
| 256 | } |
| 257 | |
| 258 | // pass to traffic control |
| 259 | device.TCBatch(tcBufs, tcs) |
| 260 | |
| 261 | tcBufs = tcBufs[:0] |
| 262 | |
| 263 | if readErr != nil { |
| 264 | if errors.Is(readErr, tun.ErrTooManySegments) { |
| 265 | // TODO: record stat for this |
| 266 | // This will happen if MSS is surprisingly small (< 576) |
| 267 | // coincident with reasonably high throughput. |
| 268 | device.Log.Verbosef("Dropped some packets from multi-segment read: %v", readErr) |
| 269 | continue |
| 270 | } |
| 271 | if !device.isClosed() { |
| 272 | if !errors.Is(readErr, os.ErrClosed) { |
| 273 | device.Log.Errorf("Failed to read packet from TUN device: %v", readErr) |
| 274 | } |
no test coverage detected