From 92e6b9823119b5461efc304a81e79e186a928e6d Mon Sep 17 00:00:00 2001 From: Christophe Besson Date: Wed, 7 Oct 2026 21:31:10 +0200 Subject: perf(node): list a group's directories off the event loop Every full index walked all roots on the loop, and a node with several large roots stopped answering for seconds. Walk directories only, on the roots' disk thread; index_sync is spawned and still answers on failure. Co-Authored-By: Claude Opus 5.5 --- packages/meshbay-node/src/meshbay_node/daemon.py | 25 ++-- .../src/meshbay_node/transport/webrtc/dispatch.py | 2 +- .../src/meshbay_node/transport/webrtc/files.py | 20 +++- .../src/meshbay_node/transport/wire.py | 58 +++++++--- packages/meshbay-node/tests/golden/dispatch.json | 128 +++++++-------------- .../meshbay-node/tests/test_content_blocklist.py | 11 +- .../meshbay-node/tests/test_disk_io_off_loop.py | 43 +++++++ 7 files changed, 167 insertions(+), 120 deletions(-) (limited to 'packages/meshbay-node') diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index d6293ee..e1aac41 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -67,7 +67,11 @@ from meshbay_node.transport import ( WEBRTC_AVAILABLE, Denylist, ) -from meshbay_node.transport.wire import index_delta_message, index_sync_message +from meshbay_node.transport.wire import ( + index_delta_message, + index_sync_message, + list_dirs_off_loop, +) if QUIC_AVAILABLE: from meshbay_node.transport import QuicChunkServer @@ -1452,9 +1456,11 @@ class NodeDaemon(EnrichmentMixin): # "nobody is listening", not a case to send in clear for. if peers and idx.gek: hidden = self._blocklist if self._is_public(group_id) else () + dirs = (None if delta is not None + else await list_dirs_off_loop(indexer.roots)) msg = (index_delta_message(idx, delta, indexer.roots, hidden) if delta is not None - else index_sync_message(idx, indexer.roots, hidden)) + else index_sync_message(idx, indexer.roots, hidden, dirs)) pushed = 0 for session in peers: try: @@ -1488,16 +1494,18 @@ class NodeDaemon(EnrichmentMixin): return if self._blocklist.replace(hashes): log.info("Content blocklist: %d hashes", len(self._blocklist)) - self._push_public_indexes() + await self._push_public_indexes() - def _on_blocklist_update(self, add: list, remove: list) -> None: + def _on_blocklist_update(self, add: list, remove: list): + """Returns the push it started, if any, for whoever wants to wait on it.""" if not self._hosts_public_group(): - return + return None if self._blocklist.apply(add, remove): log.info("Content blocklist updated: +%d -%d", len(add), len(remove)) - self._push_public_indexes() + return spawn(self._push_public_indexes()) + return None - def _push_public_indexes(self) -> None: + async def _push_public_indexes(self) -> None: """Resend each public group's whole index, so what the list now hides leaves every connected member's view, and what it released comes back.""" if not self._webrtc: @@ -1506,7 +1514,8 @@ class NodeDaemon(EnrichmentMixin): idx = indexer.index if not self._is_public(idx.group_id) or not idx.gek: continue - msg = index_sync_message(idx, indexer.roots, self._blocklist) + dirs = await list_dirs_off_loop(indexer.roots) + msg = index_sync_message(idx, indexer.roots, self._blocklist, dirs) for session in list(self._webrtc._sessions.values()): if session._group_id == idx.group_id: try: diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/dispatch.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/dispatch.py index 4bf9dbb..a8b5767 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/dispatch.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/dispatch.py @@ -18,7 +18,7 @@ MAX_PRE_PROOF_FETCHES = 4 # records it for every type. SPAWNED, INLINE = True, False _HANDLERS = { - MNP.INDEX_SYNC: ("_do_index_sync", INLINE), + MNP.INDEX_SYNC: ("_do_index_sync", SPAWNED), # Spawned rather than answered inline: the reply waits for room # on the channel, and blocking the message loop for that would # stop everything else this peer is doing — including the diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py index f629563..0434554 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py @@ -19,7 +19,7 @@ from meshbay_node.transport.webrtc.disk import ( _rmdir_if_empty, ) from meshbay_node.transport.webrtc.limits import CHUNK_SIZE, LEASE_GRANTED, LEASE_QUEUED -from meshbay_node.transport.wire import index_sync_message +from meshbay_node.transport.wire import index_sync_message, list_dirs_off_loop log = logging.getLogger("meshbay_node.transport.webrtc_server") @@ -177,9 +177,21 @@ class FilesMixin: self._audit("dir_delete", rel) self._send({"type": MNP.DIR_DELETE_ACK, "v": MNP_VERSION, "dir": rel}) - def _do_index_sync(self) -> None: - ctx = self._group_ctx() - self._send(index_sync_message(ctx["index"], ctx.get("roots"), self._hidden_ids())) + async def _do_index_sync(self) -> None: + try: + dirs = await list_dirs_off_loop(self._group_ctx().get("roots")) + # Read again after the walk: the entries go out as they are when the + # reply is sent, so no delta pushed meanwhile is older than it. + ctx = self._group_ctx() + msg = index_sync_message(ctx["index"], ctx.get("roots"), self._hidden_ids(), + dirs) + except Exception as e: + # What the dispatcher answered when this ran inline. A spawned task + # that fails only logs, and the client would wait out its timeout. + log.error("Error building index_sync: %s", e, exc_info=True) + self._send({"type": "error", "detail": "Request failed"}) + return + self._send(msg) async def _try_serve_thumbnail( self, thumb_hash: str, chunk_index: int, gek: bytes | None, diff --git a/packages/meshbay-node/src/meshbay_node/transport/wire.py b/packages/meshbay-node/src/meshbay_node/transport/wire.py index 258edb5..77c5d62 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/wire.py +++ b/packages/meshbay-node/src/meshbay_node/transport/wire.py @@ -27,11 +27,14 @@ message we have not yet authenticated. from __future__ import annotations +import os +from pathlib import Path + from meshbay_common import MNP_VERSION from meshbay_common.groupbox import PURPOSE_INDEX, seal from meshbay_common.protocol import MNP, index_entry_wire -from meshbay_node.roots import RootSet +from meshbay_node.roots import RootSet, off_disk # A group with a deep tree can hold more directories than anyone will navigate in one # sitting, and the whole list rides on one message. @@ -47,25 +50,49 @@ def list_dirs(roots: RootSet | None) -> list[str]: is listed too — its content is frozen, not gone, and hiding it would look exactly like deletion. """ + return _walk_dirs(_dir_roots(roots)) if roots else [] + + +async def list_dirs_off_loop(roots: RootSet | None) -> list[str]: + """ + `list_dirs`, on the thread that serves `roots`. + + The walk is a syscall per directory over every root of the group, and it ran + on the event loop for each full index: on a library of several large roots + that was seconds in which the node answered nobody, not its members and not + the desktop client asking whether it was there. Found live: a tester adding a + sixth directory saw "Add a directory" vanish because the node missed a 3 s + check while it indexed the other five. + """ if not roots: return [] + return await off_disk(roots, _walk_dirs, _dir_roots(roots)) + + +def _dir_roots(roots: RootSet) -> list[tuple[str, Path, bool]]: + """What the walk needs from each root, read on the loop that owns them.""" + return [(root.name, root.path, root.available) for root in roots] + + +def _walk_dirs(dir_roots: list[tuple[str, Path, bool]]) -> list[str]: + """Blocking. Directories only: `os.walk` reads each one once and never + stats a file, which `rglob("*")` and an `is_dir()` per entry did.""" out: list[str] = [] - for root in roots: - out.append(root.name) - if not root.available: - continue - try: - for path in sorted(root.path.rglob("*")): - if path.is_dir() and not path.name.startswith("."): - rel = path.relative_to(root.path) - if not any(part.startswith(".") for part in rel.parts): - out.append(f"{root.name}/{rel.as_posix()}") - except OSError: + for name, path, available in dir_roots: + out.append(name) + if not available: continue + for top, subdirs, _files in os.walk(path): + # Pruned in place: nothing under a hidden directory is listed either. + subdirs[:] = [d for d in subdirs if not d.startswith(".")] + rel = Path(top).relative_to(path) + for d in subdirs: + out.append(f"{name}/{(rel / d).as_posix()}") return sorted(out)[:MAX_DIRS] -def index_sync_message(index, roots: RootSet | None, hidden=()) -> dict: +def index_sync_message(index, roots: RootSet | None, hidden=(), + dirs: list[str] | None = None) -> dict: """ The full `index_sync` message for one group. @@ -76,11 +103,14 @@ def index_sync_message(index, roots: RootSet | None, hidden=()) -> dict: `hidden` is the content blocklist in a public group (`meshbay_node.blocklist`): those entries are not sent at all. + + `dirs` is `list_dirs_off_loop(roots)`, awaited by a caller on the event loop; + left out, the walk runs here, where the caller is. """ payload = { "version": index.version, "entries": [index_entry_wire(e) for e in index.entries if e.id not in hidden], - "dirs": list_dirs(roots), + "dirs": list_dirs(roots) if dirs is None else dirs, "roots": roots.describe() if roots else [], } return { diff --git a/packages/meshbay-node/tests/golden/dispatch.json b/packages/meshbay-node/tests/golden/dispatch.json index 61bbc45..44d030a 100644 --- a/packages/meshbay-node/tests/golden/dispatch.json +++ b/packages/meshbay-node/tests/golden/dispatch.json @@ -9167,115 +9167,67 @@ }, "index_sync | member | bare": { "audit": [], - "log": [ - "ERROR Error handling %s on DataChannel: %s" - ], - "sent": [ - { - "detail": "Request failed", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] + "log": [], + "sent": [], + "spawned": [ + "_do_index_sync" + ] }, "index_sync | member | lists": { "audit": [], - "log": [ - "ERROR Error handling %s on DataChannel: %s" - ], - "sent": [ - { - "detail": "Request failed", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] + "log": [], + "sent": [], + "spawned": [ + "_do_index_sync" + ] }, "index_sync | member | numbers": { "audit": [], - "log": [ - "ERROR Error handling %s on DataChannel: %s" - ], - "sent": [ - { - "detail": "Request failed", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] + "log": [], + "sent": [], + "spawned": [ + "_do_index_sync" + ] }, "index_sync | member | strings": { "audit": [], - "log": [ - "ERROR Error handling %s on DataChannel: %s" - ], - "sent": [ - { - "detail": "Request failed", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] + "log": [], + "sent": [], + "spawned": [ + "_do_index_sync" + ] }, "index_sync | operator | bare": { "audit": [], - "log": [ - "ERROR Error handling %s on DataChannel: %s" - ], - "sent": [ - { - "detail": "Request failed", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] + "log": [], + "sent": [], + "spawned": [ + "_do_index_sync" + ] }, "index_sync | operator | lists": { "audit": [], - "log": [ - "ERROR Error handling %s on DataChannel: %s" - ], - "sent": [ - { - "detail": "Request failed", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] + "log": [], + "sent": [], + "spawned": [ + "_do_index_sync" + ] }, "index_sync | operator | numbers": { "audit": [], - "log": [ - "ERROR Error handling %s on DataChannel: %s" - ], - "sent": [ - { - "detail": "Request failed", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] + "log": [], + "sent": [], + "spawned": [ + "_do_index_sync" + ] }, "index_sync | operator | strings": { "audit": [], - "log": [ - "ERROR Error handling %s on DataChannel: %s" - ], - "sent": [ - { - "detail": "Request failed", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] + "log": [], + "sent": [], + "spawned": [ + "_do_index_sync" + ] }, "invite_cancel | challenged | bare": { "audit": [], diff --git a/packages/meshbay-node/tests/test_content_blocklist.py b/packages/meshbay-node/tests/test_content_blocklist.py index 0e07787..8e0b790 100644 --- a/packages/meshbay-node/tests/test_content_blocklist.py +++ b/packages/meshbay-node/tests/test_content_blocklist.py @@ -129,7 +129,7 @@ async def test_a_public_group_refuses_a_blocked_file_and_its_thumbnail(tmp_path) @pytest.mark.asyncio async def test_a_public_group_does_not_list_a_blocked_file(tmp_path): s, gek = _session("public", tmp_path) - s._do_index_sync() + await s._do_index_sync() payload = unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, s._channel.sent[-1]) assert [e["id"] for e in payload["entries"]] == [KEPT] @@ -137,7 +137,7 @@ async def test_a_public_group_does_not_list_a_blocked_file(tmp_path): @pytest.mark.asyncio async def test_a_private_group_is_untouched(tmp_path): s, gek = _session("private", tmp_path) - s._do_index_sync() + await s._do_index_sync() payload = unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, s._channel.sent[-1]) assert sorted(e["id"] for e in payload["entries"]) == [BLOCKED, KEPT] assert not s._refuse_blocked(BLOCKED) @@ -218,14 +218,15 @@ def _daemon(tmp_path, visibility: str): return daemon, session, gek -def test_a_pushed_block_resends_the_index_without_the_file(tmp_path): +@pytest.mark.asyncio +async def test_a_pushed_block_resends_the_index_without_the_file(tmp_path): daemon, session, gek = _daemon(tmp_path, "public") - daemon._on_blocklist_update([BLOCKED], []) + await daemon._on_blocklist_update([BLOCKED], []) msg = session._send.call_args[0][0] ids = [e["id"] for e in unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, msg)["entries"]] assert ids == [KEPT] - daemon._on_blocklist_update([], [BLOCKED]) + await daemon._on_blocklist_update([], [BLOCKED]) msg = session._send.call_args[0][0] ids = sorted(e["id"] for e in unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, msg)["entries"]) assert ids == [BLOCKED, KEPT] diff --git a/packages/meshbay-node/tests/test_disk_io_off_loop.py b/packages/meshbay-node/tests/test_disk_io_off_loop.py index 44864ba..c0ab052 100644 --- a/packages/meshbay-node/tests/test_disk_io_off_loop.py +++ b/packages/meshbay-node/tests/test_disk_io_off_loop.py @@ -34,6 +34,7 @@ from meshbay_node.indexer.group_index import GroupIndex from meshbay_node.indexer.indexer import DirectoryIndexer from meshbay_node.roots import Root, RootSet from meshbay_node.transfers import LeaselessReads +from meshbay_node.transport import wire from meshbay_node.transport.webrtc import files, media_tools, upload_handlers from meshbay_node.transport.webrtc_server import WebRTCPeerSession from node_source import webrtc_files @@ -412,3 +413,45 @@ async def test_a_slow_scratch_read_does_not_stop_the_loop(tmp_path, monkeypatch) assert blob == CONTENT assert ticker.ticks > FREE_TICKS / 2, ( f"the loop was blocked: {ticker.ticks} wake-ups during a {SLOW_S}s read") + + +# ── The directory listing of a full index ─────────────────────────────────── +# +# Every `index_sync` lists the group's directories by walking all of its roots. +# On the loop, a library of several large roots was seconds in which the node +# answered nobody: found live, a tester adding a sixth directory lost "Add a +# directory" because the node missed the desktop client's 3 s check. + +async def test_a_slow_directory_walk_does_not_stop_the_loop(tmp_path, monkeypatch): + session, _, _, _ = await _served(tmp_path) + monkeypatch.setattr(wire, "_walk_dirs", _slow(wire._walk_dirs)) + + with _Ticker() as ticker: + await session._do_index_sync() + + assert ticker.ticks > FREE_TICKS / 2, ( + f"the loop was blocked: {ticker.ticks} wake-ups during a {SLOW_S}s walk") + assert session.sent[-1]["type"] == MNP.INDEX_SYNC + + +async def test_an_index_sync_that_fails_still_answers(tmp_path): + """It is spawned now, and a spawned task that fails only logs: without its + own answer the client would wait out its timeout for nothing.""" + session, ctx, _, _ = await _served(tmp_path) + del ctx["index"] + + await session._do_index_sync() + + assert session.sent[-1] == {"type": "error", "detail": "Request failed"} + + +def test_the_walk_lists_directories_and_skips_hidden_ones(tmp_path): + root = tmp_path / "films" + (root / "a" / "b").mkdir(parents=True) + (root / ".hidden" / "inside").mkdir(parents=True) + (root / "a" / "file.bin").write_bytes(b"x") + gone = tmp_path / "gone" + + dirs = wire._walk_dirs([("films", root, True), ("gone", gone, False)]) + + assert dirs == ["films", "films/a", "films/a/b", "gone"] -- cgit v1.2.3