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
34 changes: 11 additions & 23 deletions deapi/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -1527,7 +1527,7 @@ def set_xy_array(self, positions, width=None, height=None):
command = self._addSingleCommand(self.SET_SCAN_XY_ARRAY, None, vals_to_send)
try:
packet = struct.pack("I", command.ByteSize()) + command.SerializeToString()
self.socket.send(packet)
self.socket.sendall(packet)
ret = self.__ReceiveResponseForCommand(command) != False
log.info(f"response {ret}")
except socket.error:
Expand Down Expand Up @@ -1879,7 +1879,7 @@ def set_virtual_mask(self, id, w, h, mask):
command = self._addSingleCommand(self.SET_VIRTUAL_MASK, None, [id, w, h])
ret = True
packet = struct.pack("I", command.ByteSize()) + command.SerializeToString()
self.socket.send(packet)
self.socket.sendall(packet)

if ret:
if mask.dtype != np.uint8:
Expand Down Expand Up @@ -3092,29 +3092,17 @@ def _recvFromSocket(self, sock, bytes):
return buffer

def __sendToSocket(self, sock, buffer, bytes):
"""Send all of *buffer*, or raise.

``socket.send`` may accept only part of a buffer (routinely on macOS's loopback once
the socket has a timeout), so the return value must be honoured: dropping the rest
left the server waiting for bytes that never came, and every later command on the
connection hung. ``sendall`` loops until every byte is sent and raises
``socket.timeout`` if the server stops reading for *timeout* seconds.
"""
timeout = self.exposureTime * 10 + 30
startTime = self.GetTime()
self.socket.settimeout(timeout)

retval = True
chunkSize = 4096
for i in range(0, len(buffer), chunkSize):
try:
sock.send(buffer[i : min(len(buffer), i + chunkSize)])
except socket.timeout:

log.debug(f" __sendToSocket : timeout in trying to send {bytes} bytes")
if self.GetTime() - startTime > timeout:
log.error(" __recvFromSocket: max timeout %d seconds", timeout)
retval = False
break
else:
pass # continue further
except socket.error as e:
log.error(f"Error during send: {e}")
# Handle the error as needed, e.g., close the connection
retval = False
break
sock.sendall(buffer)
return buffer

def __saveText(self, image, fileName, textSize):
Expand Down
3 changes: 3 additions & 0 deletions deapi/simulated_server/initialize_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,9 @@ def main(port=13240):
except Exception:
traceback.print_exc(file=sys.stderr)
connected = False
# Close it now, not when the next accept() replaces it: a connection left
# open here would take the client's next commands and never answer them.
conn.close()


# Using the special variable
Expand Down
18 changes: 16 additions & 2 deletions deapi/tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,25 @@


def close_port(port):
"""Stop a simulated server left listening on *port* by an earlier run.

Only deapi's own simulated server is stopped: on a development machine the port
may belong to a real DE-Server, which a test run must never kill.
"""
for conn in psutil.net_connections(kind="inet"):
if conn.laddr.port == port:
print(f"Closing port {port} by terminating PID {conn.pid}")
if conn.laddr.port != port or conn.status != psutil.CONN_LISTEN or not conn.pid:
continue
try:
process = psutil.Process(conn.pid)
if "initialize_server" not in " ".join(process.cmdline()):
pytest.exit(
f"port {port} is in use by {process.name()} (PID {conn.pid}), "
"not a simulated server; stop it or run with --server"
)
print(f"Closing port {port} by terminating PID {conn.pid}")
process.terminate()
except (psutil.NoSuchProcess, psutil.AccessDenied):
pass


def wait_for_idle(client, timeout: float = 30, interval: float = 0.1):
Expand Down
3 changes: 3 additions & 0 deletions deapi/tests/test_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,9 @@ def test_start_acquisition(self, client):

def test_start_acquisition_scan_disabled(self, client):
client.scan(enable="Off")
# Without a scan the acquisition is a single frame: at the suite's 1000 fps it
# is over in 1 ms, often before `acquiring` is read (macOS CI). 0.5 s is not.
client["Frames Per Second"] = 2
client.start_acquisition(1)
assert client.acquiring
wait_for_idle(client)
Expand Down
92 changes: 92 additions & 0 deletions deapi/tests/test_fake_server/test_partial_sends.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
"""The client must deliver every byte even when ``socket.send`` sends only part of a buffer.

``send`` is allowed to accept fewer bytes than it was given, and does so whenever the
socket has a timeout and the kernel's send buffer fills up -- routinely on macOS's
loopback. Dropping the rest of a buffer leaves the server waiting for bytes that never
come, and every later command on that connection hangs.
"""

import pathlib
import socket
import sys
import threading

import numpy as np
import pytest
from xprocess import ProcessStarter

from deapi import Client


class _TrickleSocket:
"""A socket whose ``send`` accepts at most ``limit`` bytes per call."""

def __init__(self, sock, limit=1000):
self._sock = sock
self._limit = limit

def send(self, data, *args):
return self._sock.send(bytes(data[: self._limit]), *args)

def sendall(self, data, *args):
view = memoryview(data)
while len(view):
view = view[self.send(view) :]

def __getattr__(self, name):
return getattr(self._sock, name)


@pytest.fixture
def trickle_client(xprocess):
# its own simulated server: a hung connection must not leak into other tests
port = int(np.random.randint(10000, 12000))
script = pathlib.Path(__file__).parents[2] / "simulated_server/initialize_server.py"

class Starter(ProcessStarter):
timeout = 60
pattern = "started"
args = [sys.executable, "-u", script, port]

xprocess.ensure(f"trickle-server-{port}", Starter)
c = Client()
c.usingMmf = False
c.connect(port=port)
c.socket = _TrickleSocket(c.socket)
yield c
try:
c.disconnect()
finally:
xprocess.getinfo(f"trickle-server-{port}").terminate()


@pytest.mark.timeout(30)
def test_virtual_mask_survives_partial_sends(trickle_client):
trickle_client.virtual_masks[2][:] = 2
np.testing.assert_allclose(trickle_client.virtual_masks[2][:], 2)
# the connection is still in step: the next command gets its own answer
assert trickle_client["Scan - Enable"] in ("On", "Off")


@pytest.mark.timeout(30)
def test_send_helper_delivers_every_byte():
"""The helper behind virtual masks and XY scan arrays, without a server."""
a, b = socket.socketpair()
try:
c = Client()
c.socket = _TrickleSocket(a, limit=777)
payload = np.arange(300_000, dtype=np.uint32).tobytes()
received = bytearray()

def read():
while len(received) < len(payload):
received.extend(b.recv(65536))

reader = threading.Thread(target=read)
reader.start()
c._Client__sendToSocket(c.socket, payload, len(payload))
reader.join(10)
assert bytes(received) == payload
finally:
a.close()
b.close()
1 change: 1 addition & 0 deletions upcoming_changes/60.bugfix.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Fixed virtual masks and XY scan arrays sometimes arriving incomplete, which hung the connection (seen on macOS): the client now sends every byte of them.
1 change: 1 addition & 0 deletions upcoming_changes/60.maintenance.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
The simulated server closes a connection after an error, so the client gets an error instead of waiting; the test suite only stops deapi's own simulated server on port 13240, never a real DE-Server.
Loading