From 3e8cd94f2ae1995402e91f7191730d6f2c35e582 Mon Sep 17 00:00:00 2001 From: Carter Francis Date: Thu, 24 Sep 2026 16:14:50 -0500 Subject: [PATCH 1/2] Send every byte of virtual masks and XY arrays; don't hang on a broken connection socket.send may accept only part of a buffer, and does on macOS's loopback once the socket has a timeout. __sendToSocket ignored the count (and skipped a chunk on timeout), so the server waited for mask bytes that never came. The fake server then went back to accept() with the connection still open, so every later command hung: one lost chunk became 60 s per remaining test until the job hit its 20-minute limit. - __sendToSocket uses sendall; set_xy_array and set_virtual_mask send their command packets with sendall too - the simulated server closes a connection after an error - close_port stops only deapi's simulated server; a real DE-Server on port 13240 ends the run with a message instead of being terminated - tests: a socket that sends at most 1000 bytes per call reproduces the hang on any OS --- deapi/client.py | 34 +++---- deapi/simulated_server/initialize_server.py | 3 + deapi/tests/conftest.py | 18 +++- .../test_fake_server/test_partial_sends.py | 92 +++++++++++++++++++ upcoming_changes/60.bugfix.rst | 1 + upcoming_changes/60.maintenance.rst | 1 + 6 files changed, 124 insertions(+), 25 deletions(-) create mode 100644 deapi/tests/test_fake_server/test_partial_sends.py create mode 100644 upcoming_changes/60.bugfix.rst create mode 100644 upcoming_changes/60.maintenance.rst 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_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. From d3cb5e244a37a334ed5863404e86dc17a4e6e7b2 Mon Sep 17 00:00:00 2001 From: Carter Francis Date: Thu, 24 Sep 2026 16:18:07 -0500 Subject: [PATCH 2/2] tests: give the scan-disabled acquisition time to be seen Without a scan it is one frame, 1 ms at 1000 fps, and could finish before the test read acquiring (the other macOS py3.12 failure). --- deapi/tests/test_client.py | 3 +++ 1 file changed, 3 insertions(+) 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)