diff --git a/deapi/client.py b/deapi/client.py index ff30b1fc..93c1ea9c 100644 --- a/deapi/client.py +++ b/deapi/client.py @@ -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: @@ -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: @@ -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): diff --git a/deapi/simulated_server/initialize_server.py b/deapi/simulated_server/initialize_server.py index ce533217..251fbaa4 100644 --- a/deapi/simulated_server/initialize_server.py +++ b/deapi/simulated_server/initialize_server.py @@ -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 diff --git a/deapi/tests/conftest.py b/deapi/tests/conftest.py index cd0dfbc7..f22b0d93 100644 --- a/deapi/tests/conftest.py +++ b/deapi/tests/conftest.py @@ -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): diff --git a/deapi/tests/test_client.py b/deapi/tests/test_client.py index bad9fb80..175698c7 100644 --- a/deapi/tests/test_client.py +++ b/deapi/tests/test_client.py @@ -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) diff --git a/deapi/tests/test_fake_server/test_partial_sends.py b/deapi/tests/test_fake_server/test_partial_sends.py new file mode 100644 index 00000000..cb57773f --- /dev/null +++ b/deapi/tests/test_fake_server/test_partial_sends.py @@ -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() diff --git a/upcoming_changes/60.bugfix.rst b/upcoming_changes/60.bugfix.rst new file mode 100644 index 00000000..7f0846b1 --- /dev/null +++ b/upcoming_changes/60.bugfix.rst @@ -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. diff --git a/upcoming_changes/60.maintenance.rst b/upcoming_changes/60.maintenance.rst new file mode 100644 index 00000000..c7f391a9 --- /dev/null +++ b/upcoming_changes/60.maintenance.rst @@ -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.