""" 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 ast import asyncio import threading import time from pathlib import Path import pytest 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.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 webrtc_server from meshbay_node.transport.webrtc import media_tools from meshbay_node.transport.webrtc_server import WebRTCPeerSession from node_source import webrtc_files from conftest import one_root, sealed_upload 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" async def test_the_availability_poll_does_not_stop_the_loop(tmp_path, monkeypatch): """ The poll is the one that runs whether anybody asked for anything. `refresh_availability` stats every root, and reconcile calls it on a timer. On the loop, a node with a sleeping disk stalls once per tick for as long as the spin-up takes — and the stat is also what keeps waking the disk, so the node pays for a library nobody is reading. """ root = tmp_path / "films" root.mkdir() (root / "clip.bin").write_bytes(CONTENT) roots = RootSet.build([{"path": str(root), "name": "films"}]) monkeypatch.setattr(Root, "is_live", _slow(Root.is_live)) idx = DirectoryIndexer(roots=roots, group_id=GROUP, sk_node=Ed25519PrivateKey.generate(), gek=None) try: with _Ticker() as ticker: await idx.initial_scan() finally: await idx.stop() roots.close_io() assert ticker.ticks > SLOW_S / TICK_S / 2, ( f"the loop was blocked: {ticker.ticks} wake-ups during a {SLOW_S}s poll") def test_no_handler_touches_the_disk_on_the_loop(): """ The whole class, not the calls that were fixed. Every measured test above exercises a handler that exists today; a new one that stats a root inline would pass all of them. So this walks the syntax tree of every module of the WebRTC transport instead and fails on any filesystem call outside the few functions written to be run through `off_disk`. `entry_abs_path` and `safe_subdir` are in the list because both are `Path.resolve()` underneath, and a resolve is syscalls whatever it is called. """ blocking = {"is_dir", "exists", "mkdir", "unlink", "rename", "rmdir", "iterdir", "read_bytes", "write_bytes", "stat", "resolve", "entry_abs_path"} # Written to block, and reached only through `off_disk`. on_the_disk_thread = {"_locate", "_append_chunk", "_read_and_encrypt", "_mkdir_if_absent", "_is_empty_dir", "_rmdir_if_empty", "safe_subdir"} # ffmpeg's own output goes through `_read_scratch_capped` and # `_discard_scratch` on a worker thread — `asyncio.to_thread` and not # `off_disk`, because a temp file is not a group root and has no platter to # serialise against. Nothing is exempt here any more. allowed = on_the_disk_thread | {"_read_scratch_capped"} found = [] def visit(node, owner): for child in ast.iter_child_nodes(node): if isinstance(child, (ast.FunctionDef, ast.AsyncFunctionDef)): visit(child, child.name) continue if isinstance(child, ast.Call): fn = child.func name = (fn.attr if isinstance(fn, ast.Attribute) else getattr(fn, "id", "")) if name in blocking and owner not in allowed: found.append(f"{owner} calls {name}() at line {child.lineno}") visit(child, owner) for path in webrtc_files(): tree = ast.parse(path.read_text(encoding="utf-8")) before = len(found) for node in tree.body: if isinstance(node, ast.ClassDef): for member in node.body: if isinstance(member, (ast.FunctionDef, ast.AsyncFunctionDef)): visit(member, member.name) elif isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef)): visit(node, node.name) found[before:] = [f"{path.name}: {f}" for f in found[before:]] assert not found, ( "filesystem calls made from the event loop:\n " + "\n ".join(found) + "\nRun them through `off_disk(roots, ...)`, or put the call in a " "helper that is only reached that way.") def _upload_session(tmp_path): """One connection into a group with one writable root, as a node has.""" shared = tmp_path / "shared" shared.mkdir(exist_ok=True) ctx = {"roots": one_root(shared), "index": GroupIndex(group_id=GROUP, sk_node=Ed25519PrivateKey.generate()), "gek": generate_gek()} session = WebRTCPeerSession.__new__(WebRTCPeerSession) session._ctx = ctx session._group_id = GROUP session._user_id = "user-1" session._pk_user = "" session.sent = [] session._send = session.sent.append session._audit = lambda *a, **k: None return session, shared async def test_a_slow_upload_write_does_not_stop_the_loop(tmp_path, monkeypatch): session, shared = _upload_session(tmp_path) monkeypatch.setattr(webrtc_server, "_append_chunk", _slow(webrtc_server._append_chunk)) with _Ticker() as ticker: await session._do_file_upload(sealed_upload( session, filename="clip.bin", data=CONTENT)) assert ticker.ticks > SLOW_S / TICK_S / 2, ( f"the loop was blocked: {ticker.ticks} wake-ups during a {SLOW_S}s write") assert (shared / "clip.bin").read_bytes() == CONTENT async def test_chunks_of_one_upload_keep_their_order_under_a_slow_disk(tmp_path, monkeypatch): """ The rule that used to hold for free. `chunk_index != state.next_index` is refused, and nothing could come between that check and the `advance` answering it while the handler was synchronous. Awaiting the write opens the gap: chunk 1 arriving while chunk 0 is still in the disk thread reads a position that has not moved yet and is refused as out of order — an upload that fails on a slow disk and nowhere else. The lock closes it, and the disk made slow here is what makes the gap wide enough to fall into. """ session, shared = _upload_session(tmp_path) pieces = [b"first-", b"second-", b"third-", b"fourth"] monkeypatch.setattr(webrtc_server, "_append_chunk", _slow(webrtc_server._append_chunk)) # Fired together and in order, which is what the dispatcher does: it creates # one task per message as it arrives. await asyncio.gather(*[ session._do_file_upload(sealed_upload( session, filename="clip.bin", data=piece, chunk_index=i, total_chunks=len(pieces))) for i, piece in enumerate(pieces) ]) refusals = [m for m in session.sent if m.get("type") == "error"] assert not refusals, f"a chunk was refused: {refusals}" assert (shared / "clip.bin").read_bytes() == b"".join(pieces) def test_the_scratch_read_is_only_ever_reached_on_a_thread(): """ `_read_scratch_capped` blocks by design, so the guard above allows it — and that allowance is worth nothing if somebody calls it straight from a handler. Passed to `asyncio.to_thread` it appears in the syntax tree as a name; called inline it appears as a call, which is what this refuses. """ direct = [f"{path.name}:{n.lineno}" for path in webrtc_files() for n in ast.walk(ast.parse(path.read_text(encoding="utf-8"))) if isinstance(n, ast.Call) and isinstance(n.func, ast.Name) and n.func.id == "_read_scratch_capped"] assert not direct, ( f"_read_scratch_capped is called directly at line(s) {direct} — hand it " "to `asyncio.to_thread` instead, or the cap is paid for on the loop") async def test_ffmpeg_output_over_the_cap_is_refused_before_it_is_read(tmp_path): """ The stat comes first, so an oversized result costs a stat rather than the read and the memory. The number in the message is the one that was measured, not the cap, because an operator reading a log wants to know by how much. """ scratch = tmp_path / "out.m4a" scratch.write_bytes(b"x" * 5000) with pytest.raises(RuntimeError, match=r"5000 bytes, over the 1024 cap"): media_tools._read_scratch_capped(scratch, 1024, "transcoded audio") # And under the cap it simply reads. assert media_tools._read_scratch_capped(scratch, 8192, "x") == b"x" * 5000 async def test_a_slow_scratch_read_does_not_stop_the_loop(tmp_path, monkeypatch): """Measured like the others: the loop keeps its wake-ups during the read.""" scratch = tmp_path / "out.vtt" scratch.write_bytes(CONTENT) monkeypatch.setattr(media_tools, "_read_scratch_capped", _slow(media_tools._read_scratch_capped)) with _Ticker() as ticker: blob = await asyncio.to_thread( media_tools._read_scratch_capped, scratch, 1 << 20, "subtitle track") assert blob == CONTENT assert ticker.ticks > SLOW_S / TICK_S / 2, ( f"the loop was blocked: {ticker.ticks} wake-ups during a {SLOW_S}s read")