diff --git a/backend/app/services/virtual_printer/mqtt_server.py b/backend/app/services/virtual_printer/mqtt_server.py index 4f93e569c..893ab833f 100644 --- a/backend/app/services/virtual_printer/mqtt_server.py +++ b/backend/app/services/virtual_printer/mqtt_server.py @@ -594,12 +594,45 @@ class SimpleMQTTServer: # Enable TCP keepalive so a hard network drop is detected # by the OS within a few minutes rather than waiting for # the next outbound write to ECONNRESET. + # + # Also tighten the Linux keepalive schedule. Defaults are + # tcp_keepalive_time=7200 s (2 h before first probe), + # tcp_keepalive_intvl=75, tcp_keepalive_probes=9 — so a + # macOS client that goes to sleep silently is only + # detected as dead ~2 h 11 min later, and until then the + # push loop keeps stalling on drain-timeouts to the + # zombie socket. #1872: a P1S sleep/wake left the pre- + # sleep session in _clients for 5+ min with no eviction + # signal. New settings (idle=60 s, interval=15 s, + # count=4) detect a dead peer in ~2 min. `getattr` guards + # keep this cross-platform — macOS has TCP_KEEPINTVL but + # not TCP_KEEPIDLE (uses TCP_KEEPALIVE); other platforms + # silently skip. sock = writer.get_extra_info("socket") if sock is not None: try: sock.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1) except OSError as e: logger.debug("%sFailed to set SO_KEEPALIVE on %s: %s", self._log_prefix, client_id, e) + for opt_name, opt_value in ( + ("TCP_KEEPIDLE", 60), + ("TCP_KEEPINTVL", 15), + ("TCP_KEEPCNT", 4), + ): + opt = getattr(socket, opt_name, None) + if opt is None: + continue + try: + sock.setsockopt(socket.IPPROTO_TCP, opt, opt_value) + except OSError as e: + logger.debug( + "%sFailed to set %s=%s on %s: %s", + self._log_prefix, + opt_name, + opt_value, + client_id, + e, + ) # Register client for periodic status pushes; start with # self.serial as the fallback until we learn the slicer's # preferred serial from the first SUBSCRIBE/PUBLISH. @@ -1145,11 +1178,30 @@ class SimpleMQTTServer: writer.write(packet) # Timeout the drain to prevent blocking the event loop if the - # MQTT client stops reading (e.g. slicer busy with FTP upload). + # MQTT client stops reading (e.g. slicer busy with FTP upload, + # macOS suspends the client mid-session — #1872). + # + # On timeout, close the writer and raise BrokenPipeError so the + # push-loop's ``except OSError`` at ``_periodic_status_push`` + # evicts the client from ``self._clients`` on this same tick. + # Before this, timeouts logged at DEBUG and returned silently, + # so the zombie writer sat in ``self._clients`` until SO_KEEPALIVE + # detected the dead peer (~2 h on Linux defaults). That kept the + # push loop spending 5 s per iteration on the stalled client and + # left the slicer's UI unaware the session was gone. try: await asyncio.wait_for(writer.drain(), timeout=5) - except TimeoutError: - logger.debug("MQTT drain timeout for %s — client may be busy", topic) + except TimeoutError as e: + logger.info( + "%sMQTT drain timeout for %s — closing stalled writer", + self._log_prefix, + topic, + ) + try: + writer.close() + except Exception: + pass # best-effort — writer may already be broken + raise BrokenPipeError(f"drain timeout on {topic}") from e async def _send_print_response( self, writer: asyncio.StreamWriter, sequence_id: str, filename: str, serial: str | None = None diff --git a/backend/tests/unit/test_vp_mqtt_server.py b/backend/tests/unit/test_vp_mqtt_server.py index 0eb97b161..849900bdb 100644 --- a/backend/tests/unit/test_vp_mqtt_server.py +++ b/backend/tests/unit/test_vp_mqtt_server.py @@ -621,3 +621,96 @@ class TestPendingRequestRouting: return None so the response broadcasts.""" assert server._lookup_pending_request_client(b"not valid json") is None assert server._lookup_pending_request_client(b'"a string, not a dict"') is None + + +class TestSendPublishDrainTimeoutEviction: + """#1872: a slicer client that stops draining (e.g. macOS sleeps the + machine mid-session) used to keep its writer in `self._clients` for + hours — drain timed out at DEBUG, returned silently, and the push loop + kept spending 5 s per iteration on the zombie until SO_KEEPALIVE + detected the dead peer. + + `_publish_to_report` now closes the writer and raises + `BrokenPipeError` on drain timeout so the push loop's existing + `except OSError` branch evicts the client on the same tick. + """ + + @pytest.mark.asyncio + async def test_drain_timeout_raises_broken_pipe_and_closes_writer(self): + """The core contract — drain > 5 s must raise BrokenPipeError AND + close the writer, not swallow the timeout.""" + server = _make_server() + + # Writer whose drain never completes — asyncio.wait_for should + # hit its 5 s ceiling. Use an unresolved future so the coroutine + # returned by drain() blocks indefinitely. + writer = MagicMock() + writer.write = MagicMock(return_value=None) + never = asyncio.Future() # deliberately never resolved + writer.drain = MagicMock(return_value=never) + writer.close = MagicMock() + + # Patch wait_for to raise TimeoutError immediately instead of + # actually waiting 5 s — we're testing our handler, not asyncio. + with pytest.MonkeyPatch.context() as mp: + + async def fake_wait_for(coro, timeout): + # Cancel the pending drain future so it doesn't leak. + if not never.done(): + never.cancel() + raise TimeoutError() + + mp.setattr(asyncio, "wait_for", fake_wait_for) + + with pytest.raises(BrokenPipeError, match="drain timeout"): + await server._publish_to_report(writer, {"x": 1}, serial="01P00A391800001") + + writer.close.assert_called_once() + + @pytest.mark.asyncio + async def test_drain_timeout_still_closes_writer_when_close_fails(self): + """Best-effort close: if the writer is already broken and + `.close()` raises, `_send_publish` must still raise + BrokenPipeError so the push loop evicts the client. Silent + swallowing here would put us right back to the #1872 zombie.""" + server = _make_server() + + writer = MagicMock() + writer.write = MagicMock(return_value=None) + never = asyncio.Future() + writer.drain = MagicMock(return_value=never) + writer.close = MagicMock(side_effect=OSError("already broken")) + + with pytest.MonkeyPatch.context() as mp: + + async def fake_wait_for(coro, timeout): + if not never.done(): + never.cancel() + raise TimeoutError() + + mp.setattr(asyncio, "wait_for", fake_wait_for) + + with pytest.raises(BrokenPipeError): + await server._publish_to_report(writer, {"x": 1}, serial="01P00A391800001") + + +class TestHandleClientTCPKeepaliveTuning: + """#1872: without tightening Linux TCP keepalive knobs, dead-peer + detection defaults to ~2 h. A macOS sleep leaves the pre-sleep socket + in `self._clients` until then. Tighten to detect within ~2 min. + """ + + def test_handle_client_source_names_the_tuning_constants(self): + """The tuning code needs the three TCP_KEEP* constants to be + referenced by name so a socket-module regression / a stripped-down + platform can be diagnosed from a support bundle. Inspecting the + source keeps this pinned without spinning up a real socket in + the unit test (that's covered separately by integration).""" + source = inspect.getsource(SimpleMQTTServer._handle_client) + for name in ("TCP_KEEPIDLE", "TCP_KEEPINTVL", "TCP_KEEPCNT"): + assert name in source, ( + f"_handle_client must reference {name} so the Linux " + "keepalive schedule is tightened (#1872). Without this " + "a macOS sleep leaves the pre-sleep socket in _clients " + "for ~2 h until the default SO_KEEPALIVE probes fire." + )