aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/transport/webrtc
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/transport/webrtc')
-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
2 files changed, 17 insertions, 5 deletions
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,