From 6897b645b66892c71bf2805051ab32842a16ad7b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ga=C3=ABtan=20Trellu?= Date: Wed, 16 Sep 2026 07:47:29 -0400 Subject: [PATCH 1/2] Keep the lifecycle callbacks, and the query string out of the exception log Three from CodeRabbit on #21 and one on #15. At shutdown every executor was drained with `cancel_futures=True`, including the connect and disconnect lifecycle pair. A disconnect callback already queued by `on_close()` was therefore dropped, leaving a client marked connected on the runtime bus after the server that admitted it was gone. Admission and inbound work is still discardable -- nobody is left to receive what it would produce -- but those two now drain. `_positive_int` caught TypeError and ValueError. `int(float("inf"))` raises OverflowError, so a non-finite worker count from configuration aborted startup rather than falling back to the documented default. `_request_summary` only covers the ordinary request line. Tornado's exception logger prints `self.request`, and `HTTPServerRequest.__repr__` includes the URI, so a credential passed as `?authorization=` survived into the uncaught exception log that the normal path already redacted. `log_exception` is overridden to use the same summary; HTTPError keeps Tornado's own handling. The architecture doc now also describes the embedded-handler fallback and the admission/disconnect ordering, both of which transport integrators depend on. Co-Authored-By: Claude Opus 5 (1M context) --- docs/architecture.md | 11 +++++ hivemind_websocket_protocol/__init__.py | 37 +++++++++++++--- tests/test_protocol_unit.py | 59 +++++++++++++++++++++++++ 3 files changed, 101 insertions(+), 6 deletions(-) diff --git a/docs/architecture.md b/docs/architecture.md index 3fa8d9c..31230d0 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -59,6 +59,17 @@ to an ordered disconnect executor. Runtime-bus disconnect notifications can perform synchronous I/O, so they must not run on Tornado's event-loop thread or delay unrelated handshakes. +A disconnect that arrives while admission is still running is deferred until +the connect lifecycle finishes, so a client can never be published after its +own cleanup has run. At shutdown, the connect and disconnect executors are +drained rather than cancelled: dropping either would leave a client marked +connected on the runtime bus after the server that admitted it is gone. + +When no executor is installed -- an embedded or test harness that mounts the +handler without `HiveMindWebsocketProtocol.run()` -- the callback runs +synchronously on the caller's thread instead. Transport integrators relying on +that compatibility path keep the historical behaviour. + ## Authorization Clients connect with a URL query parameter: diff --git a/hivemind_websocket_protocol/__init__.py b/hivemind_websocket_protocol/__init__.py index d78ad51..6ef3dcf 100644 --- a/hivemind_websocket_protocol/__init__.py +++ b/hivemind_websocket_protocol/__init__.py @@ -24,6 +24,8 @@ from ovos_utils.log import LOG from ovos_utils.xdg_utils import xdg_data_home from poorman_handshake import HandShake, PasswordHandShake, check_password_strength +import traceback + from tornado import ioloop from tornado import web from tornado.iostream import StreamClosedError @@ -313,7 +315,10 @@ def _positive_int(value: Any, default: int, name: str) -> int: return default try: parsed = int(value) - except (TypeError, ValueError): + except (TypeError, ValueError, OverflowError): + # OverflowError with the others: int(float("inf")) raises it, and a + # non-finite worker count reaching this from configuration aborted + # startup instead of falling back to the documented default. LOG.warning(f"Ignoring invalid {name}: {value!r}") return default if parsed < 1: @@ -677,14 +682,17 @@ def start_listener() -> None: HiveMindTornadoWebSocket.slow_admission_log_ms = ( DEFAULT_SLOW_ADMISSION_LOG_MS ) + # Admission and inbound work is discardable at teardown: whatever + # it would have produced, nobody is left to receive. auth_executor.shutdown(wait=True, cancel_futures=True) handshake_executor.shutdown(wait=True, cancel_futures=True) inbound_executor.shutdown(wait=True, cancel_futures=True) - connect_lifecycle_executor.shutdown( - wait=True, - cancel_futures=True, - ) - disconnect_executor.shutdown(wait=True, cancel_futures=True) + # The lifecycle pair is not. A connect callback already queued owns + # the presence its disconnect has to clear, and dropping either + # leaves a client marked connected on the runtime bus after the + # server that admitted it is gone. Let them finish. + connect_lifecycle_executor.shutdown(wait=True) + disconnect_executor.shutdown(wait=True) if startup_error is not None: raise startup_error @@ -982,6 +990,23 @@ def _cancel_inbound_processing(self) -> None: def _peer_label(self, peer: str) -> str: return f"{peer} ({self.source_ip})" if self.source_ip else peer + def log_exception(self, typ, value, tb) -> None: + """Keep a query string out of the uncaught-exception log too. + + ``_request_summary`` only covers the ordinary request line. Tornado's + exception logger prints ``self.request`` directly, and + ``HTTPServerRequest.__repr__`` includes the URI -- so a credential + passed as a query parameter survived the redaction that the normal path + already applied. + """ + if isinstance(value, web.HTTPError): + return super().log_exception(typ, value, tb) + LOG.error( + "Uncaught exception %s\n%s", + self._request_summary(), + "".join(traceback.format_exception(typ, value, tb)), + ) + def _request_summary(self) -> str: """Keep query-string credentials out of Tornado request logs.""" return ( diff --git a/tests/test_protocol_unit.py b/tests/test_protocol_unit.py index dcdeca6..7e2111f 100644 --- a/tests/test_protocol_unit.py +++ b/tests/test_protocol_unit.py @@ -1765,3 +1765,62 @@ def _run(): assert not t.is_alive() assert (cert_dir / "gen-me.crt").exists() assert (cert_dir / "gen-me.key").exists() + + +def test_a_non_finite_worker_count_falls_back_instead_of_aborting_startup(): + """`int(float("inf"))` raises OverflowError, not ValueError. + + A non-finite worker count reaching this from configuration used to abort + startup rather than fall back to the documented default -- the one outcome + the fallback exists to prevent. + """ + from hivemind_websocket_protocol import _positive_int + + for value in (float("inf"), float("-inf"), float("nan")): + assert _positive_int(value, 7, "workers") == 7, value + # The ordinary invalid inputs keep behaving as they did. + assert _positive_int("nonsense", 7, "workers") == 7 + assert _positive_int(0, 7, "workers") == 7 + assert _positive_int(None, 7, "workers") == 7 + assert _positive_int(3, 7, "workers") == 3 + + +def test_an_uncaught_exception_log_carries_no_query_string(): + """Tornado prints `self.request`, whose repr includes the URI. + + `_request_summary` only covers the ordinary request line, so a credential + passed as a query parameter survived into the exception log that the normal + path already redacted. + """ + import logging + from unittest.mock import MagicMock + + from hivemind_websocket_protocol import HiveMindTornadoWebSocket + + handler = HiveMindTornadoWebSocket.__new__(HiveMindTornadoWebSocket) + handler.request = MagicMock() + handler.request.remote_ip = "203.0.113.7" + handler.request.method = "GET" + handler.request.uri = "/?authorization=c2VjcmV0OnRva2Vu" + handler.request.__repr__ = lambda _self: ( + "HTTPServerRequest(uri='/?authorization=c2VjcmV0OnRva2Vu')" + ) + + records = [] + import hivemind_websocket_protocol as module + + class _Recorder: + def error(self, message, *args): + records.append(message % args if args else message) + + original = module.LOG + module.LOG = _Recorder() + try: + handler.log_exception(ValueError, ValueError("boom"), None) + finally: + module.LOG = original + + assert records, "the exception was not logged at all" + joined = "\n".join(records) + assert "c2VjcmV0OnRva2Vu" not in joined, joined + assert "authorization=" not in joined, joined From 2c124766570622505f73895c907bfdfb59f18a38 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ga=C3=ABtan=20Trellu?= Date: Wed, 16 Sep 2026 09:25:03 -0400 Subject: [PATCH 2/2] Deliver a deferred disconnect without the IOLoop loop.start() returns before the executors drain, so a disconnect deferred behind admission and hopped through loop.add_callback sat in a queue until some later loop.start() that never comes -- leaving the client marked connected on the runtime bus after the server that admitted it was gone. It is submitted straight from the lifecycle future completion now; _submit_disconnect_callback only touches the executor, which is safe from that thread. And the two executor references are cleared after they drain rather than before, so a lifecycle callback finishing mid-drain still has somewhere to submit. Found by CodeRabbit on #33. Co-Authored-By: Claude Opus 5 (1M context) --- hivemind_websocket_protocol/__init__.py | 20 +++++++---- tests/test_protocol_unit.py | 44 +++++++++++++++++++++++++ 2 files changed, 58 insertions(+), 6 deletions(-) diff --git a/hivemind_websocket_protocol/__init__.py b/hivemind_websocket_protocol/__init__.py index 6ef3dcf..f14651d 100644 --- a/hivemind_websocket_protocol/__init__.py +++ b/hivemind_websocket_protocol/__init__.py @@ -676,8 +676,6 @@ def start_listener() -> None: HiveMindTornadoWebSocket.inbound_admission_capacity = None HiveMindTornadoWebSocket.inbound_client_queue_size = None HiveMindTornadoWebSocket.inbound_pending = 0 - HiveMindTornadoWebSocket.disconnect_executor = None - HiveMindTornadoWebSocket.connect_lifecycle_executor = None HiveMindTornadoWebSocket.password_min_bits = None HiveMindTornadoWebSocket.slow_admission_log_ms = ( DEFAULT_SLOW_ADMISSION_LOG_MS @@ -691,8 +689,14 @@ def start_listener() -> None: # the presence its disconnect has to clear, and dropping either # leaves a client marked connected on the runtime bus after the # server that admitted it is gone. Let them finish. + # Cleared only once they have drained: a lifecycle callback + # finishing during the drain submits its disconnect through these, + # and nulling them first sent it down the embedded-handler path + # instead -- or, with the loop already stopped, nowhere at all. connect_lifecycle_executor.shutdown(wait=True) disconnect_executor.shutdown(wait=True) + HiveMindTornadoWebSocket.disconnect_executor = None + HiveMindTornadoWebSocket.connect_lifecycle_executor = None if startup_error is not None: raise startup_error @@ -1390,11 +1394,15 @@ def on_close(self): LOG.debug(f"disconnecting client: {self._peer_label(client.peer)}") lifecycle = self._connect_lifecycle_future if lifecycle is not None and not lifecycle.done(): + # Submitted straight from the lifecycle future's completion, not + # hopped through the IOLoop: at shutdown `loop.start()` returns + # before the executors drain, and a callback added after that + # simply sits in the queue until some later `loop.start()` that + # never comes -- leaving the client marked connected on the runtime + # bus. `_submit_disconnect_callback` only touches the executor, + # which is safe from the lifecycle thread. lifecycle.add_done_callback( - lambda _future: self.loop.add_callback( - self._submit_disconnect_callback, - client, - ) + lambda _future: self._submit_disconnect_callback(client) ) return self._submit_disconnect_callback(client) diff --git a/tests/test_protocol_unit.py b/tests/test_protocol_unit.py index 7e2111f..9146d63 100644 --- a/tests/test_protocol_unit.py +++ b/tests/test_protocol_unit.py @@ -1824,3 +1824,47 @@ def error(self, message, *args): joined = "\n".join(records) assert "c2VjcmV0OnRva2Vu" not in joined, joined assert "authorization=" not in joined, joined + + +def test_a_disconnect_deferred_behind_admission_does_not_need_the_ioloop(): + """`loop.start()` returns before the executors drain. + + A deferred disconnect hopped through `loop.add_callback` would then sit in + a queue until some later `loop.start()` that never comes, leaving the + client marked connected on the runtime bus after the server is gone. + """ + from concurrent.futures import Future + from unittest.mock import MagicMock + + from hivemind_websocket_protocol import HiveMindTornadoWebSocket + + handler = HiveMindTornadoWebSocket.__new__(HiveMindTornadoWebSocket) + handler.request = MagicMock() + handler.request.remote_ip = "203.0.113.7" + handler.hm_protocol = MagicMock() + handler._disconnect_submitted = False + handler._auth_lookup_future = None + handler._auth_task = None + handler._cancel_inbound_processing = lambda: None + handler._release_auth_admission = lambda: None + handler._peer_label = lambda peer: str(peer) + handler.client = MagicMock() + handler.client.peer = "peer" + # No executor: the embedded/test path calls through synchronously. + handler.disconnect_executor = None + # An IOLoop that has already stopped drops anything added to it. + handler.loop = MagicMock() + handler.loop.add_callback.side_effect = AssertionError( + "the disconnect must not depend on the IOLoop still running" + ) + + lifecycle: Future = Future() + lifecycle.set_running_or_notify_cancel() + handler._connect_lifecycle_future = lifecycle + + handler.on_close() + handler.hm_protocol.handle_client_disconnected.assert_not_called() + + # Admission finishes after the socket closed -- the disconnect follows it. + lifecycle.set_result(None) + handler.hm_protocol.handle_client_disconnected.assert_called_once_with(handler.client)