Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 6 additions & 5 deletions pydivert/ebpf.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import socket
import threading
import time
from collections import deque
from typing import Any, cast

from .base import BaseDivert
Expand Down Expand Up @@ -71,7 +72,7 @@ def __init__(
if libbpf is None:
raise ImportError("libbpf missing on system.")
self._obj = self._ringbuf = self._raw_sock = self._raw_sock6 = None
self._queue: list[Packet] = []
self._queue: deque[Packet] = deque()
self._hooks: list[tuple[BpfTcHook, BpfTcOpts]] = []
self._interfaces = kwargs.get("interfaces", None)
self._tc_priority = 0
Expand Down Expand Up @@ -433,7 +434,7 @@ def _recv_impl(self, bufsize: int = DEFAULT_PACKET_BUFFER_SIZE, timeout: float |
if not self._queue:
raise OSError(socket.EBADF, "Handle closed while receiving")

return self._queue.pop(0)
return self._queue.popleft()

def _recv_batch_impl(self, count: int, bufsize: int, timeout: float | None) -> list[Packet]:
if Flag.SEND_ONLY in self.flags: # pragma: no cover
Expand All @@ -444,7 +445,7 @@ def _recv_batch_impl(self, count: int, bufsize: int, timeout: float | None) -> l
p = self._recv_impl(bufsize, timeout) # pragma: no cover
packets.append(p) # pragma: no cover
while len(packets) < count and self._queue: # pragma: no cover
packets.append(self._queue.pop(0)) # pragma: no cover
packets.append(self._queue.popleft()) # pragma: no cover
except TimeoutError: # pragma: no cover
if not packets: # pragma: no cover
raise # pragma: no cover
Expand Down Expand Up @@ -533,7 +534,7 @@ async def _recv_batch_async_impl(self, count: int, bufsize: int, timeout: float
p = await self._recv_async_impl(bufsize, timeout)
packets.append(p)
while len(packets) < count and self._queue:
packets.append(self._queue.pop(0))
packets.append(self._queue.popleft())
except (TimeoutError, Exception): # pragma: no cover
if not packets: # pragma: no cover
raise # pragma: no cover
Expand All @@ -545,7 +546,7 @@ def _recv_batch_impl(self, count: int, bufsize: int, timeout: float | None) -> l
p = self._recv_impl(bufsize, timeout)
packets.append(p)
while len(packets) < count and self._queue:
packets.append(self._queue.pop(0)) # pragma: no cover
packets.append(self._queue.popleft()) # pragma: no cover
except (TimeoutError, Exception):
if not packets:
raise
Expand Down
4 changes: 3 additions & 1 deletion pydivert/tests/test_coverage_complete.py
Original file line number Diff line number Diff line change
Expand Up @@ -1167,13 +1167,15 @@ def send_side_effect(*args):


def test_ebpf_recv_batch_linux():
from collections import deque

from pydivert.ebpf import EBPFDivert

with patch("pydivert.ebpf.libbpf", MagicMock()):
d = EBPFDivert()
d._is_open = True
p = MagicMock(spec=pydivert.Packet)
d._queue = [p]
d._queue = deque([p])
res = d.recv_batch(count=1)
assert res == [p]
assert len(d._queue) == 0
Expand Down