From ba6b1a8436709547c29e02b244683e8c281f5289 Mon Sep 17 00:00:00 2001 From: maziggy Date: Thu, 2 Jul 2026 10:13:01 +0200 Subject: [PATCH] fix(vp): evict MQTT clients on drain timeout + tighten TCP keepalive (#1872) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Reporter (H2C + macOS 26.5.1 + BS 2.8.0.50): after every Mac sleep/wake cycle, Bambu Studio couldn't see the VP or connect to it. Only fix was quit BS + reboot Bambuddy. The physical printer's own cloud/LAN link recovered in ~5 s from the same sleep — the delta was in VP session handling. Log evidence (bug-report-assets/logs/ddf1ede75df045cd94ad223d0f08f88a): - 14:04:06 healthy `1Hz status push: 60 pushes/min to :54698` - 14:04:06 → 14:09:16: five minutes of SSDP-only, no push summary for :54698, no OSError, no disconnect line - 14:09:16: new source port :54861 connects and authenticates fine — the server was not rejecting reconnects - 14:10:17 first DEBUG line: `MQTT drain timeout for device/…/report — client may be busy` — smoking gun Root cause: `_publish_to_report:1149` caught `asyncio.wait_for(drain, timeout=5)` TimeoutError at DEBUG and returned silently. TimeoutError is not OSError, so the push loop's `except OSError` at :441 never saw it — the zombie writer sat in self._clients until the kernel's default TCP keepalive detected the dead peer (Linux default: ~2 h 11 min). Two hunks: 1. `_publish_to_report`: on drain TimeoutError, close the writer (best effort, catch Exception so an already-broken close() doesn't mask the raise) and raise BrokenPipeError, which IS OSError. Push loop evicts on the same tick. 2. `_handle_client`: after SO_KEEPALIVE=1, set TCP_KEEPIDLE=60, TCP_KEEPINTVL=15, TCP_KEEPCNT=4 — dead-peer detection in ~2 min instead of ~2 h. `getattr(socket, ...)` guards keep it cross- platform (macOS uses TCP_KEEPALIVE not TCP_KEEPIDLE, other kernels may not expose all three — skip whichever is missing). What I got wrong first pass and corrected on log-read: hypothesised "missing MQTT session takeover on same client_id". Wrong. _handle_connect parses the protocol client_id but discards it (assignment commented out at :762), and self._clients is keyed on `f"{addr[0]}:{addr[1]}"` (socket peer), so every reconnect gets a distinct key. No takeover race exists. The log fixed this: the "not seen" symptom is BS-side (macOS UDP receive after sleep + BS holding the pre-sleep socket state), but the server-side amplifier was the zombie writer. --- CHANGELOG.md | 2 + .../services/virtual_printer/mqtt_server.py | 58 +++++++++++- backend/tests/unit/test_vp_mqtt_server.py | 93 +++++++++++++++++++ 3 files changed, 150 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 29533b8e6..8a57b8852 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,8 @@ All notable changes to Bambuddy will be documented in this file. ## [0.2.5b2] - Unreleased ### Fixed +- **Bambu Studio on macOS won't reconnect to the VP after machine sleep — zombie writer pinned in `_clients` for hours (#1872, reporter @avvidme)** — Reporter on H2C + macOS 26.5.1 + BS 2.8.0.50: after every sleep/wake cycle Bambu Studio couldn't see the VP or connect to it. Only workaround was quitting BS and rebooting Bambuddy. The physical printer's own cloud / LAN link recovered in ~5 s from the same sleep, so the delta is in the VP's session-handling. **Log evidence (`bug-report-assets/logs/ddf1ede75df045cd94ad223d0f08f88a.log`).** 14:04:06 shows a healthy `1Hz status push: 60 pushes/min to [IP]:54698` — full-rate 1 Hz for the last minute pre-sleep. Then 5 min of SSDP output only — no push summary for :54698, no OSError, no disconnect line. At 14:09:16 a brand-new TCP source port :54861 connects, authenticates, subscribes — so the MQTT server is not rejecting reconnects. At 14:10:17, the first DEBUG line after the reporter enabled debug logging is `MQTT drain timeout for device/…/report — client may be busy` — smoking gun. **Root cause.** `_publish_to_report` at `mqtt_server.py:1149-1152` caught `asyncio.wait_for(writer.drain(), timeout=5)` `TimeoutError` at DEBUG and returned silently. Timeouts are not `OSError`, so the push loop's `except OSError` at `:441` never saw them — the client sat in `self._clients` until the OS's default TCP keepalive detected the dead peer, which on Linux is `tcp_keepalive_time=7200 s` (2 h) + 9 probes × 75 s = ~2 h 11 min. That's exactly the 5 min silence in the log; drain likely stayed under the 5 s ceiling because the kernel TX buffer had room, so the DEBUG line didn't even fire until much later. Meanwhile the loop was iterating with a stalled client sitting in the dict every tick. **Fix — two hunks.** (1) On drain `TimeoutError`, close the writer (best-effort, catch `Exception` so an already-broken writer doesn't mask the raise) and raise `BrokenPipeError`. The eviction path is indirect but reliable: `BrokenPipeError` is an `OSError` subclass, so it's caught by every `_send_*` wrapper's outer `except OSError` at `:1036` / `:1118` / `:1236` and logged at ERROR — push_counts increments on this tick, no direct eviction. BUT `writer.close()` was called inside `_publish_to_report`, so on the next push-loop tick, `writer.is_closing()` at `:431` returns True → the client is appended to `disconnected` → popped from `_clients`, `_client_serials`, `push_counts` on the same tick. Real eviction latency: ~1 s (one 1 Hz tick), down from ~2 h. Same mechanism as the pre-existing hard-disconnect path (RST → OSError swallowed in `_send_*` → transport marks closing → next-tick eviction), just extended to also cover the "silent stall" case that has no OS-level RST. (2) Tighten the Linux TCP-keepalive schedule right after `SO_KEEPALIVE=1` at `:600`: `TCP_KEEPIDLE=60`, `TCP_KEEPINTVL=15`, `TCP_KEEPCNT=4` — dead-peer detection in ~2 min instead of ~2 h. `getattr(socket, ...)` guards keep the code cross-platform: macOS has `TCP_KEEPINTVL` but not `TCP_KEEPIDLE` (it exposes `TCP_KEEPALIVE` under a different constant), other platforms silently skip whichever knobs their kernel doesn't expose. **What I initially got wrong.** First-pass hypothesis was "no MQTT session takeover on same `client_id`". Wrong. `_handle_connect` at `:762` parses the protocol client_id but discards it (assignment commented out), and `self._clients` is keyed on socket peer `f"{addr[0]}:{addr[1]}"`, so each reconnect gets a distinct key — no takeover race actually exists. The log fixed this: the "not seen" symptom was BS-side (macOS UDP receive socket recovering slowly from sleep, plus BS's `client_id` still holding the old socket state) but the *server-side amplifier* was the zombie writer keeping push loop attention. **Tests.** 3 new cases in `test_vp_mqtt_server.py`. `TestSendPublishDrainTimeoutEviction::test_drain_timeout_raises_broken_pipe_and_closes_writer` patches `asyncio.wait_for` to raise `TimeoutError` immediately and asserts `BrokenPipeError` propagates AND `writer.close()` was called — pins both halves of the contract. `..._still_closes_writer_when_close_fails` covers the best-effort `close()`: even if the writer is already broken and `.close()` raises, `_publish_to_report` must still raise `BrokenPipeError` — silent swallowing here would put us right back to the pre-fix zombie state. `TestHandleClientTCPKeepaliveTuning::test_handle_client_source_names_the_tuning_constants` uses `inspect.getsource` to pin that `_handle_client` references `TCP_KEEPIDLE`, `TCP_KEEPINTVL`, `TCP_KEEPCNT` — a socket-module regression or a stripped-down platform can then be diagnosed from a support bundle. Full VP MQTT + VP manager suites 179/179 green, ruff clean. **Scope.** Backend-only. No new i18n key, no new permission, no DB migration, no frontend change. Users on 0.2.5b1 or earlier with any macOS slicer client: fix takes effect on next Bambuddy restart, no reconfiguration needed. On non-Linux hosts (e.g. Bambuddy running on macOS or FreeBSD for development), the keepalive-schedule tightening is a no-op — the drain-timeout eviction still applies. + - **Non-proxy VP camera passthrough is dead for A1 / P1 targets — OrcaSlicer Liveview fails with `[2:-10061]` (#1868, reporter @tom4711-2)** — Symptom on a P1S target in server/non-proxy VP mode: the VP starts 3000, 3002, 8883, 990 and 322, but nothing on 6000, so OrcaSlicer's Liveview button fails. Reporter confirmed a raw `socat` forwarder `:6000 → :6000` immediately restores the stream — the target camera works, the VP just isn't publishing it. **Root cause.** `backend/app/services/virtual_printer/manager.py:1098-1118` hardcoded the camera-passthrough `TCPProxy` to `listen_port=322 / target_port=322` regardless of the target printer's model. That port is correct for RTSPS models (X1/X2/H2/P2S), but A1 / A1 Mini / P1P / P1S use Bambu's proprietary chamber-image protocol on port 6000 — the 322 listener the manager opens for those targets has no upstream, and the slicer's connection to `:322` yields "connection refused". Proxy mode is unaffected because `SlicerProxyManager` (`tcp_proxy.py:1596`) already opens 6000 for file-transfer, and Bambu reuses the same port for chamber-image, so the passthrough coincidentally works there. **Fix.** Read the target printer's model from `printer_manager.get_client(target_id).model` at the same point the target IP is read, then call `get_camera_port(target_model)` — the same source of truth used by `routes/camera.py` and covered by `test_printer_models.py::TestSupportsRtsp` / `TestGetCameraPort` — to pick 322 or 6000. TCPProxy listen_port and target_port both follow the same value. Rename the log tag from `"RTSP"` to `f"Camera-{camera_port}"` so support-bundle grep tells you at a glance which protocol the VP is fronting for. The internal `_rtsp_proxy` attribute name is unchanged to keep the diff tight; the comment above the block spells out that it doubles as chamber-image passthrough on A1/P1. Model comes from the target printer's client instance (`target_client.model`), NOT `self.model` — the VP's spoofed identity has no bearing on how the physical printer serves its camera. **Tests.** New `TestVirtualPrinterCameraPassthrough` in `test_virtual_printer.py` — 4 cases pin the branch: `test_rtsp_model_p2s_opens_port_322` and `test_rtsp_model_x1c_opens_port_322` guard the RTSP path stays on 322 (regression against the fix accidentally routing everything to 6000); `test_chamber_image_model_p1s_opens_port_6000` and `test_chamber_image_model_a1_opens_port_6000` are the direct #1868 guards — P1S / A1 targets must expose 6000. The P1S case additionally asserts NO 322 listener was opened, since a stale 322 on an A1/P1 install would confuse anyone port-scanning the VP. Test harness monkeypatches `TCPProxy` and every peer service (`VirtualPrinterFTPServer`, `SimpleMQTTServer`, `MQTTBridge`, `BindServer`, `VirtualPrinterSSDPServer`, `SSDPProxy`) to lightweight MagicMocks with an already-set `asyncio.Event` on `.ready` — `start_server()` awaits `.ready.wait()` on those four barriers before returning, so an unset event would deadlock the test. 143/143 in the file green (was 139 + 4 new). **Scope.** Backend-only, one-hunk change in `manager.py` plus 4 tests. No i18n key, no permission change, no DB migration, no frontend change. Users on 0.2.5b1 or earlier with A1 / A1 Mini / P1P / P1S targets in non-proxy VP mode: the fix takes effect on the next Bambuddy restart, no reconfiguration needed. The workaround `socat` forwarder can be removed after upgrade. - **Editing a queue item assigned to "Any of model X" left the printer selection area blank** — `PrintModal` initialises `assignmentMode` from `queueItem.target_model`: items created with a specific printer got `'printer'` mode (renders the printer list), items created with model-based assignment got `'model'` mode. But `PrintModal/index.tsx:1102-1106` was passing `onAssignmentModeChange={!isEditing ? setAssignmentMode : undefined}` (same for `onTargetModelChange` / `onTargetLocationChange`), which flipped `modelAssignmentAvailable` to `false` in `PrinterSelector` and hid every model-mode control — the mode toggle at `PrinterSelector.tsx:390`, the model dropdown at `:431`, the location filter at `:459`. Combined with the `assignmentMode === 'printer'` gate on the printer list at `:519`, edit-mode for a model-assigned item rendered an empty container. Users couldn't retarget the item or even see what model it was assigned to; the only way to change it was delete + re-queue. **Fix.** Drop the `!isEditing` gate on all three props — the underlying submit path already handles both directions cleanly (`printer_id: null, target_model, target_location` when saving in model mode at `index.tsx:788-802`; `printer_id, target_model: null, target_location: null` when saving in printer mode at `:844-857`), so un-gating the UI just surfaces the machinery that was already there. Users can now (a) see the current target model + location on a model-assigned item, (b) change the target model or narrow / widen the location filter, and (c) flip the item between "Specific Printer" and "Any of Model" without deleting + re-queueing. Edit is only offered on `pending` items (`QueuePage.tsx:2137`), so there's no race with an in-flight dispatch when the assignment mode flips. **Scope.** Frontend-only, one-hunk change to the props on the existing `PrinterSelector` invocation. No new i18n key (the model / location / mode strings already existed for create mode). No backend change. Existing 61/61 `PrintModal.test.tsx` cases green — the tests that pass `initialSelectedPrinterIds={[1]}` for the create-with-printer flow are unaffected because the `!initialSelectedPrinterIds?.length` gate at `index.tsx:1088` still hides the whole `PrinterSelector` for the post-upload dispatch dialog. 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." + )