aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-10-07 21:31:10 +0200
committerChristophe Besson <cbesson@gmail.com>2026-10-07 21:31:10 +0200
commit92e6b9823119b5461efc304a81e79e186a928e6d (patch)
treea58a742140025b553b6872e3269ad54c3ac601b5 /packages/meshbay-node/src
parentd3bc50cca1c4dfe143c29a522b8cd57eb69ab92e (diff)
downloadmeshbay-92e6b9823119b5461efc304a81e79e186a928e6d.tar.gz
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 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-node/src')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py25
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/dispatch.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py20
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/wire.py58
4 files changed, 78 insertions, 27 deletions
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 {