summaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
Diffstat (limited to 'packages')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py9
-rw-r--r--packages/meshbay-node/src/meshbay_node/roots.py63
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py74
-rw-r--r--packages/meshbay-node/tests/test_added_root_does_not_wait_for_its_scan.py2
-rw-r--r--packages/meshbay-node/tests/test_disk_io_off_loop.py224
5 files changed, 338 insertions, 34 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py
index fa9ef81..69c7e66 100644
--- a/packages/meshbay-node/src/meshbay_node/daemon.py
+++ b/packages/meshbay-node/src/meshbay_node/daemon.py
@@ -1789,6 +1789,15 @@ class NodeDaemon:
for indexer in self._indexers:
await indexer.stop()
+ # The thread each group's roots are read from (roots.py `io_executor`).
+ # Created on the first read, so a group nobody downloaded from never
+ # started one and this is a no-op for it.
+ groups = self._webrtc._ctx.get("groups") if self._webrtc else None
+ for group in (groups or {}).values():
+ roots = group.get("roots")
+ if roots is not None:
+ roots.close_io()
+
if self._quic_server:
await self._quic_server.stop()
diff --git a/packages/meshbay-node/src/meshbay_node/roots.py b/packages/meshbay-node/src/meshbay_node/roots.py
index 7854e28..67778ca 100644
--- a/packages/meshbay-node/src/meshbay_node/roots.py
+++ b/packages/meshbay-node/src/meshbay_node/roots.py
@@ -28,8 +28,10 @@ are one directory.
from __future__ import annotations
+import asyncio
import logging
import re
+from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass, field
from pathlib import Path
@@ -171,6 +173,49 @@ class RootSet:
# like a deletion all over again on the pass after it.
auto_ejected: list[str] = field(default_factory=list)
+ # The thread every blocking filesystem call on these roots is made from.
+ #
+ # **Why it exists at all.** A root that has spun down, or that lives on a
+ # network mount, answers its first syscall in seconds rather than
+ # microseconds. Made from the event loop that is the whole node: no other
+ # group is served, no stream is fed, no chat message is delivered and the
+ # hub socket is not read, for as long as the platter takes to come back.
+ # A client's connection attempt times out and has to be made again, which
+ # is what this was found by. It is `AV9`'s lesson with the disk in the
+ # place of the mail server.
+ #
+ # **Why one worker and not a pool.** The same reason the indexer's executor
+ # has one: two interleaved reads on a spinning drive seek-thrash against
+ # each other rather than go faster, measured there on a USB disk. One
+ # worker also keeps every read of these roots in the order it was asked
+ # for, which costs nothing — the protocol addresses a chunk by index, so
+ # no caller depends on that order — and leaves no way for two threads to
+ # be inside the same file at once.
+ #
+ # **Why per root set and not one for the node.** A node serves several
+ # groups, and their roots are not all on the same volume. A single worker
+ # would put the sleeping USB drive of one group in front of the SSD of
+ # another, which is the symptom this removes, one level down.
+ #
+ # Created on first use, so a RootSet that never reads anything — most of
+ # them, in tests — never starts a thread. Not compared and not printed: it
+ # is machinery, not part of what a root set *is*, and `daemon.py` compares
+ # root sets to decide whether a reload changed anything.
+ _io: ThreadPoolExecutor | None = field(
+ default=None, init=False, repr=False, compare=False)
+
+ @property
+ def io_executor(self) -> ThreadPoolExecutor:
+ if self._io is None:
+ self._io = ThreadPoolExecutor(max_workers=1, thread_name_prefix="rootio")
+ return self._io
+
+ def close_io(self) -> None:
+ """Stop the disk thread. Safe to call twice, and on a set that never read."""
+ if self._io is not None:
+ self._io.shutdown(wait=False)
+ self._io = None
+
# ── Construction ─────────────────────────────────────────────────────────
@classmethod
@@ -401,6 +446,24 @@ def entry_abs_path(roots: RootSet, entry) -> Path | None:
return (parent / entry.name) if parent else None
+async def off_disk(roots: RootSet, fn, *args):
+ """
+ Run one blocking filesystem call on the thread that serves `roots`.
+
+ Every syscall against a group's content goes through here, `resolve()` and
+ `exists()` included: a `stat` is what *wakes* a sleeping disk, so a check
+ left on the event loop pays the spin-up in full and the read that follows
+ it finds the disk already awake. Offloading only the read would move the
+ stall, not remove it.
+
+ `fn` must not touch anything the loop also touches — it runs on another
+ thread. Reading and encrypting a chunk qualifies; updating a session's
+ state does not.
+ """
+ loop = asyncio.get_running_loop()
+ return await loop.run_in_executor(roots.io_executor, fn, *args)
+
+
def _refuse_nesting(new: Root, existing: list[Root]) -> None:
"""
No root may contain another, compared case-insensitively.
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
index d79674d..e748be9 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
@@ -139,7 +139,8 @@ from meshbay_node.media_probe import (
probe_video as _probe_video,
)
from meshbay_node.roots import (
- ROOT_NOT_SERVED, RootSet, entry_abs_path, SAFE_UPLOAD_NAME, safe_subdir, _free_name,
+ ROOT_NOT_SERVED, RootSet, entry_abs_path, off_disk, SAFE_UPLOAD_NAME, safe_subdir,
+ _free_name,
)
log = logging.getLogger(__name__)
@@ -3877,12 +3878,9 @@ class WebRTCPeerSession:
self._send({"type": "error", "detail": "File not found"})
return
- file_path = entry_abs_path(ctx["roots"], entry)
- if file_path is None:
- self._send({"type": "error", "detail": ROOT_NOT_SERVED})
- return
- if not file_path.exists():
- self._send({"type": "error", "detail": "File not on disk"})
+ file_path, refusal = await off_disk(ctx["roots"], _locate, ctx["roots"], entry)
+ if refusal is not None:
+ self._send({"type": "error", "detail": refusal})
return
# A real index entry, asked for without a lease: browsing, or a client
@@ -3913,7 +3911,8 @@ class WebRTCPeerSession:
file_id[:12], chunk_index,
getattr(self._channel, "bufferedAmount", "?"))
file_hash = bytes.fromhex(entry.id)
- chunk_data = _read_and_encrypt(
+ chunk_data = await off_disk(
+ ctx["roots"], _read_and_encrypt,
ctx["gek"], file_path, chunk_index, file_hash, entry.id)
# Backpressure. Without it the node hands the whole window to the
# channel at once and the reader sees the first chunk, then nothing for
@@ -4007,12 +4006,9 @@ class WebRTCPeerSession:
if not entry:
self._send({"type": "error", "detail": "File not found"})
return
- file_path = entry_abs_path(ctx["roots"], entry)
- if file_path is None:
- self._send({"type": "error", "detail": ROOT_NOT_SERVED})
- return
- if not file_path.exists():
- self._send({"type": "error", "detail": "File not on disk"})
+ file_path, refusal = await off_disk(ctx["roots"], _locate, ctx["roots"], entry)
+ if refusal is not None:
+ self._send({"type": "error", "detail": refusal})
return
media_cache = self._ctx.get("media_cache")
@@ -4074,12 +4070,9 @@ class WebRTCPeerSession:
if not entry:
self._send({"type": "error", "detail": "File not found"})
return
- file_path = entry_abs_path(ctx["roots"], entry)
- if file_path is None:
- self._send({"type": "error", "detail": ROOT_NOT_SERVED})
- return
- if not file_path.exists():
- self._send({"type": "error", "detail": "File not on disk"})
+ file_path, refusal = await off_disk(ctx["roots"], _locate, ctx["roots"], entry)
+ if refusal is not None:
+ self._send({"type": "error", "detail": refusal})
return
media_cache = self._ctx.get("media_cache")
@@ -5882,7 +5875,7 @@ class WebRTCPeerSession:
self._audit("admin_auth_failed", f"file_delete:{file_id[:16]}")
return
- self._exec_file_delete(ctx, file_id, entry)
+ await self._exec_file_delete(ctx, file_id, entry)
async def _admin_exec_invite_create(
self, pending: dict, transcript: bytes, sig: bytes,
@@ -5918,15 +5911,15 @@ class WebRTCPeerSession:
"username": result.get("username", ""),
})
- def _exec_file_delete(self, ctx: dict, file_id: str, entry) -> None:
- file_path = entry_abs_path(ctx["roots"], entry)
- if file_path is None:
+ async def _exec_file_delete(self, ctx: dict, file_id: str, entry) -> None:
+ file_path, refusal = await off_disk(ctx["roots"], _locate, ctx["roots"], entry)
+ if refusal == ROOT_NOT_SERVED:
# Frozen, not gone: removing the entry would lose a file that is
# still on a drive the node cannot read right now.
self._send({"type": "error", "detail": ROOT_NOT_SERVED})
return
- if file_path.exists():
- file_path.unlink()
+ if file_path is not None:
+ await off_disk(ctx["roots"], file_path.unlink)
log.info("File deleted: %s", entry.name)
self._audit("file_delete", entry.name)
@@ -6123,12 +6116,9 @@ class WebRTCPeerSession:
self._send({"type": "error", "detail": "File not found"})
return
- file_path = entry_abs_path(ctx["roots"], entry)
- if file_path is None:
- self._send({"type": "error", "detail": ROOT_NOT_SERVED})
- return
- if not file_path.exists():
- self._send({"type": "error", "detail": "File not on disk"})
+ file_path, refusal = await off_disk(ctx["roots"], _locate, ctx["roots"], entry)
+ if refusal is not None:
+ self._send({"type": "error", "detail": refusal})
return
gek = ctx.get("gek")
@@ -6581,6 +6571,24 @@ class WebRTCPeerSession:
await self._pc.close()
+def _locate(roots: RootSet, entry) -> tuple[Path | None, str | None]:
+ """
+ Where an entry is, and whether it is readable — or the refusal to send.
+
+ Both halves are syscalls: `resolve()` walks the path and `exists()` stats
+ it, and a stat is what *wakes* a sleeping disk. Leaving either on the event
+ loop and offloading only the read would move the stall rather than remove
+ it, and the read would then find the disk already awake. Blocking; called
+ through `off_disk`.
+ """
+ path = entry_abs_path(roots, entry)
+ if path is None:
+ return None, ROOT_NOT_SERVED
+ if not path.exists():
+ return None, "File not on disk"
+ return path, None
+
+
def _read_and_encrypt(
gek: bytes,
file_path: Path,
@@ -6588,7 +6596,7 @@ def _read_and_encrypt(
file_hash: bytes,
file_id: str = "",
) -> dict:
- """Read one chunk off disk and encrypt it. Blocking; the caller keeps it short."""
+ """Read one chunk off disk and encrypt it. Blocking; called through `off_disk`."""
with open(file_path, "rb") as f:
f.seek(chunk_index * CHUNK_SIZE)
plaintext = f.read(CHUNK_SIZE)
diff --git a/packages/meshbay-node/tests/test_added_root_does_not_wait_for_its_scan.py b/packages/meshbay-node/tests/test_added_root_does_not_wait_for_its_scan.py
index a80645a..4d4c4e8 100644
--- a/packages/meshbay-node/tests/test_added_root_does_not_wait_for_its_scan.py
+++ b/packages/meshbay-node/tests/test_added_root_does_not_wait_for_its_scan.py
@@ -235,6 +235,6 @@ async def test_a_file_request_is_refused_not_crashed(tmp_path):
async def test_a_delete_is_refused_and_the_entry_kept(tmp_path):
session, ctx, entry = await _served_without_two(tmp_path)
- session._exec_file_delete(ctx, entry.id, entry)
+ await session._exec_file_delete(ctx, entry.id, entry)
assert [m.get("detail") for m in session.sent] == [ROOT_NOT_SERVED]
assert ctx["index"].get_entry(entry.id) is not None
diff --git a/packages/meshbay-node/tests/test_disk_io_off_loop.py b/packages/meshbay-node/tests/test_disk_io_off_loop.py
new file mode 100644
index 0000000..206b77b
--- /dev/null
+++ b/packages/meshbay-node/tests/test_disk_io_off_loop.py
@@ -0,0 +1,224 @@
+"""
+A slow disk must cost the caller that touched it, and nobody else.
+
+A root that has spun down, or that lives on a network mount, answers its first
+syscall in seconds rather than microseconds. Made from the event loop, that
+stalls the entire node: no other group is served, no stream is fed, no chat
+message is delivered, and the hub socket is not read, for as long as the platter
+takes to come back. It was found from the other end — a client's connection
+attempt timing out on a node with one member, while the disk woke up.
+
+These tests do not read the source to check which thread a call is made from.
+They make the disk slow and **measure whether the loop kept running**: a ticker
+counts its own wake-ups beside the request, and a handler that blocks the loop
+takes every one of those wake-ups with it. Put the call back inline and the
+ticker count collapses to zero, which is the property being guarded.
+
+The third test is the one that is not about latency: one worker per root set
+means two reads of the same file can never be inside it at once, and that is
+what makes the single `f.seek()`/`f.read()` pair safe without a lock.
+"""
+
+import asyncio
+import re
+import threading
+import time
+from pathlib import Path
+
+from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
+from meshbay_common.crypto import generate_gek
+from meshbay_common.protocol import MNP
+from meshbay_common.webcrypto import chunk_key_aes, decrypt_chunk_aes
+from meshbay_node.indexer.indexer import DirectoryIndexer
+from meshbay_node.roots import RootSet
+from meshbay_node.transfers import LeaselessReads
+from meshbay_node.transport import webrtc_server
+from meshbay_node.transport.webrtc_server import WebRTCPeerSession
+
+GROUP = "g" * 32
+
+# Long enough that a blocked loop is unmistakable, short enough to keep the
+# suite quick. The ticker below wakes every 5 ms, so a loop that stays free
+# gets ~40 wake-ups inside one of these and a blocked one gets none.
+SLOW_S = 0.2
+TICK_S = 0.005
+CONTENT = b"a file worth waking a disk for" * 400
+
+
+class _Channel:
+ readyState = "open"
+ bufferedAmount = 0
+
+
+class _Session(WebRTCPeerSession):
+ def __init__(self, ctx):
+ self._ctx = ctx
+ self._registry_key = "s1"
+ self._user_id = "alice"
+ self._username = "alice"
+ self._group_id = GROUP
+ self._channel = _Channel()
+ self._leaseless = LeaselessReads()
+ self._unleased_noted = False
+ self.sent: list[dict] = []
+
+ def _send(self, msg):
+ self.sent.append(msg)
+
+ def _audit(self, event, detail=""):
+ pass
+
+ def _spawn(self, coro):
+ coro.close()
+ return None
+
+
+async def _served(tmp_path: Path):
+ """A session serving one real file out of one real root."""
+ root = tmp_path / "films"
+ root.mkdir()
+ (root / "clip.bin").write_bytes(CONTENT)
+ roots = RootSet.build([{"path": str(root), "name": "films"}])
+ gek = generate_gek()
+ idx = DirectoryIndexer(roots=roots, group_id=GROUP,
+ sk_node=Ed25519PrivateKey.generate(), gek=gek)
+ await idx.initial_scan()
+ entry = next(e for e in idx.index.entries if e.name == "clip.bin")
+ ctx = {"_peers": {}, "roots": roots, "index": idx.index,
+ "sk_node": idx.sk_node, "gek": gek}
+ return _Session(ctx), ctx, entry, gek
+
+
+class _Ticker:
+ """Counts how many times the event loop came back to it."""
+
+ def __init__(self):
+ self.ticks = 0
+ self._stop = False
+ self._task: asyncio.Task | None = None
+
+ async def _run(self):
+ while not self._stop:
+ await asyncio.sleep(TICK_S)
+ self.ticks += 1
+
+ def __enter__(self):
+ self._task = asyncio.get_running_loop().create_task(self._run())
+ return self
+
+ def __exit__(self, *exc):
+ self._stop = True
+ self._task.cancel()
+ return False
+
+
+def _slow(fn):
+ """`fn`, plus a sleep on whichever thread calls it."""
+
+ def wrapper(*args, **kwargs):
+ time.sleep(SLOW_S)
+ return fn(*args, **kwargs)
+
+ return wrapper
+
+
+async def test_a_slow_chunk_read_does_not_stop_the_loop(tmp_path, monkeypatch):
+ session, _, entry, gek = await _served(tmp_path)
+ monkeypatch.setattr(webrtc_server, "_read_and_encrypt",
+ _slow(webrtc_server._read_and_encrypt))
+
+ with _Ticker() as ticker:
+ await session._do_file_request(
+ {"type": MNP.FILE_REQUEST, "file_id": entry.id, "chunk_index": 0})
+
+ # The read really did take its time, and the loop really did keep running:
+ # both halves matter, because a wrapper that never ran would also leave the
+ # ticker free.
+ assert ticker.ticks > SLOW_S / TICK_S / 2, (
+ f"the loop was blocked: {ticker.ticks} wake-ups during a {SLOW_S}s read")
+
+ chunk = next(m for m in session.sent if m.get("type") == MNP.FILE_CHUNK)
+ key = chunk_key_aes(gek, bytes.fromhex(entry.id), 0)
+ assert decrypt_chunk_aes(key, chunk["nonce"], chunk["ct"]) == CONTENT
+
+
+async def test_a_slow_stat_does_not_stop_the_loop(tmp_path, monkeypatch):
+ """
+ The stat matters as much as the read: it is what *wakes* the disk.
+
+ Offloading only the read would leave the spin-up on the loop and the read
+ would then find the disk already awake — the stall moved, not removed.
+ """
+ session, _, entry, _ = await _served(tmp_path)
+ monkeypatch.setattr(webrtc_server, "_locate", _slow(webrtc_server._locate))
+
+ with _Ticker() as ticker:
+ await session._do_file_request(
+ {"type": MNP.FILE_REQUEST, "file_id": entry.id, "chunk_index": 0})
+
+ assert ticker.ticks > SLOW_S / TICK_S / 2, (
+ f"the loop was blocked: {ticker.ticks} wake-ups during a {SLOW_S}s stat")
+ assert any(m.get("type") == MNP.FILE_CHUNK for m in session.sent)
+
+
+async def test_one_root_set_reads_one_chunk_at_a_time(tmp_path):
+ """
+ Two reads of the same root are never in flight together.
+
+ Not a performance choice: interleaved reads of one spinning drive seek-thrash
+ (the indexer's executor carries the measurement), and a second thread inside
+ the same `open`/`seek`/`read` sequence is a correctness question this avoids
+ having to answer. A pool would reopen both.
+ """
+ session, _, entry, _ = await _served(tmp_path)
+ inside = 0
+ peak = 0
+ guard = threading.Lock()
+ real = webrtc_server._read_and_encrypt
+
+ def counting(*args, **kwargs):
+ nonlocal inside, peak
+ with guard:
+ inside += 1
+ peak = max(peak, inside)
+ try:
+ time.sleep(0.02)
+ return real(*args, **kwargs)
+ finally:
+ with guard:
+ inside -= 1
+
+ webrtc_server._read_and_encrypt = counting
+ try:
+ await asyncio.gather(*[
+ session._do_file_request(
+ {"type": MNP.FILE_REQUEST, "file_id": entry.id, "chunk_index": 0})
+ for _ in range(6)
+ ])
+ finally:
+ webrtc_server._read_and_encrypt = real
+
+ assert peak == 1, f"{peak} reads of one root set were inside the disk at once"
+
+
+def test_no_handler_resolves_a_path_on_the_loop():
+ """
+ `entry_abs_path` is `Path.resolve()`, which is syscalls — it belongs on the
+ disk thread with everything else.
+
+ This reads the source because there is nothing else to read: a handler added
+ later that resolves an entry inline would pass every test above, since those
+ only exercise the handlers that exist today. `_locate` is the one place
+ allowed to call it, and `off_disk` is how `_locate` is reached.
+ """
+ src = Path(webrtc_server.__file__).read_text()
+ # Every call site, not the first one: a guard that stops at the first
+ # occurrence stops guarding the moment a new call is inserted above it.
+ calls = [m.start() for m in re.finditer(r"\bentry_abs_path\(", src)]
+ body = re.search(r"\ndef _locate\(.*?\n(?=\n\ndef |\n\nclass )", src, re.S)
+ assert body, "_locate is gone or has been renamed — this guard needs rewriting"
+ allowed = range(body.start(), body.end())
+ stray = [c for c in calls if c not in allowed]
+ assert not stray, (
+ f"{len(stray)} call(s) to entry_abs_path outside _locate: "
+ f"resolve a path through `off_disk(roots, _locate, ...)` instead")