From 5fa158fab709d3d24a33318b3d910f75c051af2e Mon Sep 17 00:00:00 2001 From: Christophe Besson Date: Fri, 18 Sep 2026 15:59:24 +0200 Subject: fix(node): take the availability poll and every upload write off the loop MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The rest of AV9's disk half. Serving a file left the loop in the commit before this one; two paths were still on it. **The availability poll.** `RootSet.refresh_availability` stats every root, and eleven call sites reached it from `async def` — the reconcile loop among them, on a timer. On a sleeping disk that is a stall once per tick, and the stat is also what keeps the disk awake, so a node paid spin-up for a library nobody was reading. All eleven now go through `off_disk`, `Root.is_live` included. **The upload write.** `open`/`write`, and the resolve, the stat, the free-name search, the rename and the unlink around it. This one could not simply be awaited: the handler was synchronous, so nothing could come between the `chunk_index != state.next_index` check and the `advance` that answers it, and that is the whole of the chunk-ordering rule. Awaiting the write opens the gap — chunk 1 arriving while chunk 0 is in the disk thread reads a position that has not moved and is refused as out of order, so an upload would fail on a slow disk and nowhere else. Verified, not assumed: without the lock the new ordering test refuses three chunks of four. So the check, the write and the advance are one critical section again, under a lock held **per group**. Not per session: `partial_uploads` lives in the group context so a reconnecting client finds its upload where it left it, which means two sessions of one member share the position of one `.part` file. Arrival order is preserved by construction — the dispatcher creates one task per message as it arrives, tasks start in creation order, and the lock is the first thing each one waits on, so its waiters queue in arrival order too. `_do_file_upload` is a coroutine now, which is why forty-two test call sites gain an `await`. Their outcomes are unchanged, file by file, against the run before the change. `test_ops.py` asked which public coroutines `ops` exposes and got `off_disk`, imported rather than defined there. It now asks for the ones written in the module, which is what its own docstring means; all forty-three operations are still checked. Co-Authored-By: Claude Opus 5 --- packages/meshbay-node/src/meshbay_node/daemon.py | 8 +- .../src/meshbay_node/indexer/indexer.py | 15 ++-- packages/meshbay-node/src/meshbay_node/ops.py | 6 +- .../src/meshbay_node/transport/webrtc_server.py | 64 +++++++++++--- .../meshbay-node/tests/test_disk_io_off_loop.py | 98 +++++++++++++++++++++- .../meshbay-node/tests/test_lease_enforcement.py | 2 +- packages/meshbay-node/tests/test_ops.py | 9 +- .../meshbay-node/tests/test_partial_uploads.py | 48 +++++------ .../tests/test_root_writable_policy.py | 12 +-- .../tests/test_security_regressions.py | 44 +++++----- .../meshbay-node/tests/test_upload_attribution.py | 2 +- packages/meshbay-node/tests/test_upload_sealed.py | 42 +++++----- 12 files changed, 247 insertions(+), 103 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 69c7e66..26f73fc 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -45,7 +45,7 @@ from meshbay_node.audit import RETENTION_DAYS as AUDIT_RETENTION_DAYS, AuditStor from meshbay_node.bundle_store import BundleStore from meshbay_node.chat.store import ChatStore from meshbay_node.config import Config, DEFAULT_CONFIG_PATH, load_config, write_example_config -from meshbay_node.roots import RootSet, RootError, entry_abs_path +from meshbay_node.roots import RootSet, RootError, entry_abs_path, off_disk from meshbay_node.hub_client import HubClient, HubConfig from meshbay_node.indexer import DirectoryIndexer, IndexCache, GroupIndex from meshbay_node.indexer.enrich import Enricher @@ -331,7 +331,7 @@ class NodeDaemon: log.error("Group %r: %s — skipping", group_cfg.name, e) continue - roots.refresh_availability() + await off_disk(roots, roots.refresh_availability) if not any(r.available for r in roots): # Not skipped for being empty: a group whose only drive is # unplugged still exists, and its index is frozen rather @@ -826,7 +826,7 @@ class NodeDaemon: # reports "nothing changed" for exactly that edit. if _root_shape(ctx["roots"]) == _root_shape(roots): continue - roots.refresh_availability() + await off_disk(roots, roots.refresh_availability) indexer = next((i for i in self._indexers if i.group_id == group_cfg.id), None) if indexer is None: @@ -869,7 +869,7 @@ class NodeDaemon: except RootError as e: log.error("New group %r: %s — skipping", group_cfg.name, e) continue - roots.refresh_availability() + await off_disk(roots, roots.refresh_availability) gek = None if sk_x_raw and pk_x_raw: diff --git a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py index c1b151a..887d435 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py @@ -40,7 +40,7 @@ from meshbay_common.paths import fold, find_fold_collisions, long_path from meshbay_common.protocol import IndexEntry from meshbay_node.indexer.cache import IndexCache from meshbay_node.indexer.group_index import GroupIndex -from meshbay_node.roots import Root, RootSet +from meshbay_node.roots import Root, RootSet, off_disk log = logging.getLogger(__name__) @@ -387,7 +387,7 @@ class DirectoryIndexer: await self._initial_scan() async def _initial_scan(self) -> None: - self.roots.refresh_availability() + await off_disk(self.roots, self.roots.refresh_availability) total = 0 waiting = [r.name for r in self.roots if r.available] self._queue(waiting) @@ -695,7 +695,7 @@ class DirectoryIndexer: self._index.remove_entry(entry.id) self.roots = roots - roots.refresh_availability() + await off_disk(roots, roots.refresh_availability) added = [r for r in roots if r.folded not in old_names and r.available] self._index.roots = roots.describe() @@ -794,7 +794,7 @@ class DirectoryIndexer: return await self._reconcile() async def _reconcile(self) -> bool: - changed = self.roots.refresh_availability() + changed = await off_disk(self.roots, self.roots.refresh_availability) touched = False # Drained before the loop below, because persisting the flag is what @@ -1042,7 +1042,7 @@ class DirectoryIndexer: if root is None: return root.ejected = False - root.available = root.is_live() + root.available = await off_disk(self.roots, root.is_live) if not root.available: await self._finish_plug(None) return @@ -1071,7 +1071,8 @@ class DirectoryIndexer: # Ejected or removed again while it waited. `_rescan_root` # drops the entries before it walks, so going ahead would # empty a root that is not there to be read. - if self._holds(root) and not root.ejected and root.is_live(): + live = await off_disk(self.roots, root.is_live) + if self._holds(root) and not root.ejected and live: await self._rescan_root(root) finally: if waiting: @@ -1143,7 +1144,7 @@ class DirectoryIndexer: if root is None: return - if deleted and not root.is_live(): + if deleted and not await off_disk(self.roots, root.is_live): # The volume went away rather than the file. Freeze: mark the # root and touch nothing. Every other event for this root # will arrive here too and be dropped the same way, which is diff --git a/packages/meshbay-node/src/meshbay_node/ops.py b/packages/meshbay-node/src/meshbay_node/ops.py index 8cd893c..2221bea 100644 --- a/packages/meshbay-node/src/meshbay_node/ops.py +++ b/packages/meshbay-node/src/meshbay_node/ops.py @@ -39,7 +39,7 @@ from meshbay_common.crypto import ( ) from meshbay_node.config import DEFAULT_CONFIG_PATH from meshbay_common.join import ROLE_MEMBER, ROLE_OPERATOR -from meshbay_node.roots import RootError, RootSet +from meshbay_node.roots import RootError, RootSet, off_disk log = logging.getLogger(__name__) @@ -1130,7 +1130,7 @@ async def plug_root(state: dict, group_id: str, root_name: str) -> dict: if not root.ejected: return {"status": "already_plugged", "name": root_name, "group_id": group_id, "roots": roots.describe()} - if not root.is_live(): + if not await off_disk(roots, root.is_live): raise OpError( f"Directory not found: {root.path}. Is the device connected?", status=409) @@ -1145,7 +1145,7 @@ async def plug_root(state: dict, group_id: str, root_name: str) -> dict: if indexer: await indexer.plug_root(root_name) root.ejected = False - root.available = root.is_live() + root.available = await off_disk(roots, root.is_live) log.info("Root plugged: %s in group %s", root_name, group_id[:8]) return {"status": "plugged", "name": root_name, "group_id": group_id, 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 e748be9..92f2951 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -644,7 +644,7 @@ class WebRTCPeerSession: elif mtype == MNP.TRANSFER_CLOSE: self._do_transfer_close(msg) elif mtype == MNP.FILE_UPLOAD: - self._do_file_upload(msg) + self._spawn(self._do_file_upload(msg)) elif mtype == MNP.DIR_CREATE: self._do_dir_create(msg) elif mtype == MNP.DIR_DELETE: @@ -5205,7 +5205,43 @@ class WebRTCPeerSession: ctx["partial_uploads"] = store return store - def _do_file_upload(self, msg: dict) -> None: + @staticmethod + def _upload_lock(ctx: dict) -> asyncio.Lock: + """ + One lock per group, beside the state it protects. + + Not per session: `partial_uploads` lives in the group context so a + client that reconnects finds its upload where it left it, which means + two sessions of the same member share the position of one `.part` file. + A lock on the session would let them interleave — and the loop no longer + serializes them for free now that a chunk write is awaited. + """ + lock = ctx.get("upload_lock") + if lock is None: + lock = asyncio.Lock() + ctx["upload_lock"] = lock + return lock + + async def _do_file_upload(self, msg: dict) -> None: + """ + One chunk of an upload, in the order it arrived. + + The chunk ordering rule — `chunk_index != state.next_index` is refused — + used to hold for free: the handler was synchronous, so nothing could run + between the check and the `advance` that answers it. Awaiting the write + opens that gap, and two chunks of one upload racing through it is a + `.part` file with a hole in it or a chunk refused for arriving on time. + So the check, the write and the advance are one critical section again. + + The order is the arrival order: `_dispatch_message` runs per message as + it arrives and creates these tasks in that order, tasks start in + creation order, and this lock is the first thing each one waits on, so + its queue of waiters is in arrival order too. + """ + async with self._upload_lock(self._group_ctx()): + await self._upload_chunk(msg) + + async def _upload_chunk(self, msg: dict) -> None: """ One chunk of an upload, sealed under the group key (MNP 2.0). @@ -5386,8 +5422,8 @@ class WebRTCPeerSession: # client names *where among the group's own folders*, never a path on # the operator's filesystem. if target_rel: - target_dir = roots.resolve(target_rel) - if target_dir is None or not target_dir.is_dir(): + target_dir = await off_disk(roots, roots.resolve, target_rel) + if target_dir is None or not await off_disk(roots, target_dir.is_dir): _refuse("Not a directory in this group", "no_such_directory") return rel_dir = target_rel @@ -5396,7 +5432,7 @@ class WebRTCPeerSession: # one destination is. target_dir = upload_root.path rel_dir = upload_root.name - if not target_dir.is_dir(): + if not await off_disk(roots, target_dir.is_dir): _refuse("That directory is currently unavailable", "root_unavailable") return @@ -5419,7 +5455,8 @@ class WebRTCPeerSession: # A shared directory means two people can send the same name. Refusing the # second is safe but silly — everyone's camera produces IMG_1234.jpg — so # a free name is found instead. Never a replacement. - stored_name = state.stored_name if state else _free_name(target_dir, filename) + stored_name = (state.stored_name if state + else await off_disk(roots, _free_name, target_dir, filename)) tmp_path = target_dir / f"{stored_name}{uploads_mod.PART_SUFFIX}" final_path = target_dir / stored_name @@ -5449,7 +5486,7 @@ class WebRTCPeerSession: if chunk_index == 0: # Backstop: _free_name already guarantees this, and it stays because # it asserts the invariant where the write happens. - if final_path.exists(): + if await off_disk(roots, final_path.exists): _refuse("File already exists", "already_exists") return state = uploads.start(user_id, rel_dir, filename, stored_name, @@ -5466,12 +5503,11 @@ class WebRTCPeerSession: if state.bytes + len(chunk_bytes) > MAX_UPLOAD_BYTES: uploads.drop(user_id, rel_dir, filename) - tmp_path.unlink(missing_ok=True) + await off_disk(roots, tmp_path.unlink, True) _refuse("Upload exceeds size limit", "too_large") return - with open(tmp_path, "wb" if chunk_index == 0 else "ab") as f: - f.write(chunk_bytes) + await off_disk(roots, _append_chunk, tmp_path, chunk_bytes, chunk_index == 0) uploads.advance(user_id, rel_dir, filename, chunk_index, len(chunk_bytes)) self._send(file_upload_ack_wire( @@ -5487,7 +5523,7 @@ class WebRTCPeerSession: if chunk_index + 1 >= total_chunks: uploads.drop(user_id, rel_dir, filename) - tmp_path.rename(final_path) + await off_disk(roots, tmp_path.rename, final_path) log.info("Upload complete: %s (%d chunks, %d bytes)", stored_name, total_chunks, state.bytes) self._audit("file_upload", f"{rel_dir}/{stored_name}") @@ -6589,6 +6625,12 @@ def _locate(roots: RootSet, entry) -> tuple[Path | None, str | None]: return path, None +def _append_chunk(tmp_path: Path, chunk_bytes: bytes, first: bool) -> None: + """Add one chunk to a partial upload. Blocking; called through `off_disk`.""" + with open(tmp_path, "wb" if first else "ab") as f: + f.write(chunk_bytes) + + def _read_and_encrypt( gek: bytes, file_path: Path, 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 206b77b..2811836 100644 --- a/packages/meshbay-node/tests/test_disk_io_off_loop.py +++ b/packages/meshbay-node/tests/test_disk_io_off_loop.py @@ -29,12 +29,15 @@ 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 RootSet +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_server import WebRTCPeerSession +from conftest import one_root, sealed_upload + GROUP = "g" * 32 # Long enough that a blocked loop is unmistakable, short enough to keep the @@ -222,3 +225,96 @@ def test_no_handler_resolves_a_path_on_the_loop(): assert not stray, ( f"{len(stray)} call(s) to entry_abs_path outside _locate: " f"resolve a path through `off_disk(roots, _locate, ...)` instead") + + +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 _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) diff --git a/packages/meshbay-node/tests/test_lease_enforcement.py b/packages/meshbay-node/tests/test_lease_enforcement.py index c516b5a..2b1014b 100644 --- a/packages/meshbay-node/tests/test_lease_enforcement.py +++ b/packages/meshbay-node/tests/test_lease_enforcement.py @@ -236,7 +236,7 @@ async def test_a_queued_upload_is_refused_before_anything_is_written(group, msg = sealed_upload(second, filename="sent.bin", data=b"DATA", dir=Path(group["roots"].roots[0].path).name) msg["tr"] = "u2" - second._do_file_upload(msg) + await second._do_file_upload(msg) assert second.errors()[-1]["code"] == "lease_not_granted" assert not list(Path(group["roots"].roots[0].path).glob("sent*")), ( diff --git a/packages/meshbay-node/tests/test_ops.py b/packages/meshbay-node/tests/test_ops.py index c118b5a..083f31e 100644 --- a/packages/meshbay-node/tests/test_ops.py +++ b/packages/meshbay-node/tests/test_ops.py @@ -44,9 +44,14 @@ def test_operations_take_state_and_nothing_web_shaped(): second copy. Every public operation therefore takes `state` first and returns plain data. """ + # Defined here, not merely visible here: an async helper imported into the + # module (`off_disk`, say) is not an operation, and reporting it as one says + # nothing about a second implementation appearing. Every real operation is + # written in this module, so nothing is lost by asking. public = [(n, f) for n, f in vars(ops).items() - if inspect.iscoroutinefunction(f) and not n.startswith("_")] - assert public, "no operations found — did the module move?" + if inspect.iscoroutinefunction(f) and not n.startswith("_") + and getattr(f, "__module__", None) == ops.__name__] + assert len(public) > 30, "operations have gone missing — did the module move?" for name, fn in public: params = list(inspect.signature(fn).parameters) assert params and params[0] == "state", ( diff --git a/packages/meshbay-node/tests/test_partial_uploads.py b/packages/meshbay-node/tests/test_partial_uploads.py index f5b6602..0270a69 100644 --- a/packages/meshbay-node/tests/test_partial_uploads.py +++ b/packages/meshbay-node/tests/test_partial_uploads.py @@ -305,7 +305,7 @@ def _errors(session): return [m for m in session.sent if m.get("type") == "error"] -def test_an_upload_survives_the_connection_that_started_it(tmp_path): +async def test_an_upload_survives_the_connection_that_started_it(tmp_path): """The defect this stage exists to fix. The state used to live on the session, so the second connection saw no @@ -315,14 +315,14 @@ def test_an_upload_survives_the_connection_that_started_it(tmp_path): """ ctx = _group_ctx(tmp_path) first = _peer(ctx) - first._do_file_upload(sealed_upload(first, filename="film.mkv", + await first._do_file_upload(sealed_upload(first, filename="film.mkv", data=b"first-half", chunk_index=0, total_chunks=2)) assert _errors(first) == [] # The link drops; the client comes back on a new connection and carries on. second = _peer(ctx) - second._do_file_upload(sealed_upload(second, filename="film.mkv", + await second._do_file_upload(sealed_upload(second, filename="film.mkv", data=b"second-half", chunk_index=1, total_chunks=2)) assert _errors(second) == [], _errors(second) @@ -331,31 +331,31 @@ def test_an_upload_survives_the_connection_that_started_it(tmp_path): assert (root.path / "film.mkv").read_bytes() == b"first-halfsecond-half" -def test_another_member_cannot_continue_somebody_elses_upload(tmp_path): +async def test_another_member_cannot_continue_somebody_elses_upload(tmp_path): """The key includes the member for a reason. Without it, a second person sending the same name into the same folder would append their chunks to the first person's file — which a shared folder makes an ordinary accident, not only an attack.""" ctx = _group_ctx(tmp_path) alice = _peer(ctx, "alice") - alice._do_file_upload(sealed_upload(alice, filename="IMG_1234.jpg", + await alice._do_file_upload(sealed_upload(alice, filename="IMG_1234.jpg", data=b"hers", chunk_index=0, total_chunks=2)) assert _errors(alice) == [] bob = _peer(ctx, "bob") - bob._do_file_upload(sealed_upload(bob, filename="IMG_1234.jpg", + await bob._do_file_upload(sealed_upload(bob, filename="IMG_1234.jpg", data=b"his", chunk_index=1, total_chunks=2)) assert [m.get("code") for m in _errors(bob)] == ["not_started"] -def test_an_upload_in_flight_is_known_to_the_reaper(tmp_path): +async def test_an_upload_in_flight_is_known_to_the_reaper(tmp_path): """The two halves of this stage meeting: the state the node keeps is what stops the janitor deleting a file somebody is still sending.""" ctx = _group_ctx(tmp_path) peer = _peer(ctx) - peer._do_file_upload(sealed_upload(peer, filename="film.mkv", + await peer._do_file_upload(sealed_upload(peer, filename="film.mkv", data=b"half", chunk_index=0, total_chunks=2)) live = ctx["partial_uploads"].live_paths() @@ -379,38 +379,38 @@ def _probe(session, filename: str) -> dict: chunk_index=UPLOAD_PROBE_INDEX, total_chunks=1) -def test_a_probe_for_an_unknown_file_says_start_at_the_beginning(tmp_path): +async def test_a_probe_for_an_unknown_file_says_start_at_the_beginning(tmp_path): ctx = _group_ctx(tmp_path) peer = _peer(ctx) - peer._do_file_upload(_probe(peer, "film.mkv")) + await peer._do_file_upload(_probe(peer, "film.mkv")) assert _errors(peer) == [] assert _acks(peer, ctx)[0]["resume_from"] == 0 -def test_a_probe_reports_what_the_node_already_holds(tmp_path): +async def test_a_probe_reports_what_the_node_already_holds(tmp_path): """The point of the whole stage: the client learns it has 2 chunks there and sends the third, instead of sending a film again.""" ctx = _group_ctx(tmp_path) first = _peer(ctx) for i in range(2): - first._do_file_upload(sealed_upload(first, filename="film.mkv", + await first._do_file_upload(sealed_upload(first, filename="film.mkv", data=b"xxxx", chunk_index=i, total_chunks=5)) assert _errors(first) == [] reconnected = _peer(ctx) - reconnected._do_file_upload(_probe(reconnected, "film.mkv")) + await reconnected._do_file_upload(_probe(reconnected, "film.mkv")) ack = _acks(reconnected, ctx)[0] assert ack["resume_from"] == 2 assert ack["stored_as"] == "film.mkv" -def test_a_probe_writes_nothing_and_reserves_nothing(tmp_path): +async def test_a_probe_writes_nothing_and_reserves_nothing(tmp_path): """It has to be free of consequence: a client that asks and goes away must leave no file, no state and no name taken.""" ctx = _group_ctx(tmp_path) peer = _peer(ctx) - peer._do_file_upload(_probe(peer, "film.mkv")) + await peer._do_file_upload(_probe(peer, "film.mkv")) root = ctx["roots"].roots[0] assert list(root.path.iterdir()) == [] assert len(ctx.get("partial_uploads") or []) == 0 @@ -418,36 +418,36 @@ def test_a_probe_writes_nothing_and_reserves_nothing(tmp_path): assert _acks(peer, ctx)[0]["stored_as"] == "" -def test_a_probe_answers_only_about_the_member_who_asks(tmp_path): +async def test_a_probe_answers_only_about_the_member_who_asks(tmp_path): """Same keying as the upload itself. Otherwise one member could measure another's progress on a file they never sent — and worse, resume it.""" ctx = _group_ctx(tmp_path) alice = _peer(ctx, "alice") - alice._do_file_upload(sealed_upload(alice, filename="film.mkv", + await alice._do_file_upload(sealed_upload(alice, filename="film.mkv", data=b"xxxx", chunk_index=0, total_chunks=5)) bob = _peer(ctx, "bob") - bob._do_file_upload(_probe(bob, "film.mkv")) + await bob._do_file_upload(_probe(bob, "film.mkv")) assert _acks(bob, ctx)[0]["resume_from"] == 0 -def test_an_ordinary_ack_carries_no_resume_field(tmp_path): +async def test_an_ordinary_ack_carries_no_resume_field(tmp_path): """So a client can tell a probe's answer from a chunk's without looking at the index it echoed.""" ctx = _group_ctx(tmp_path) peer = _peer(ctx) - peer._do_file_upload(sealed_upload(peer, filename="a.bin", data=b"x", + await peer._do_file_upload(sealed_upload(peer, filename="a.bin", data=b"x", chunk_index=0, total_chunks=2)) assert "resume_from" not in _acks(peer, ctx)[0] -def test_a_probe_is_refused_where_an_upload_would_be(tmp_path): +async def test_a_probe_is_refused_where_an_upload_would_be(tmp_path): """Every check the write path makes has already run when the probe is answered, so it cannot be used to ask questions about somewhere the caller may not write.""" ctx = _group_ctx(tmp_path) peer = _peer(ctx) - peer._do_file_upload(sealed_upload(peer, filename="../escape", + await peer._do_file_upload(sealed_upload(peer, filename="../escape", data=b"", chunk_index=UPLOAD_PROBE_INDEX, total_chunks=1)) assert [m.get("code") for m in _errors(peer)] == ["invalid_filename"] @@ -456,7 +456,7 @@ def test_a_probe_is_refused_where_an_upload_would_be(tmp_path): # ── the slot an upload holds ──────────────────────────────────────────────── -def test_an_upload_chunk_says_its_slot_is_in_use(tmp_path): +async def test_an_upload_chunk_says_its_slot_is_in_use(tmp_path): """A grant nobody takes up is reclaimed after thirty seconds and abandoned on the third miss. Uploads are not gated by the lease, so the file arrived anyway — but the widget follows the lease, and a 3.5 GB upload therefore @@ -482,7 +482,7 @@ def test_an_upload_chunk_says_its_slot_is_in_use(tmp_path): msg = sealed_upload(peer, filename="film.mkv", data=b"xxxx", chunk_index=0, total_chunks=2) msg["tr"] = "up-1" - peer._do_file_upload(msg) + await peer._do_file_upload(msg) assert _errors(peer) == [] assert slots.leases["up-1"].used is True, ( diff --git a/packages/meshbay-node/tests/test_root_writable_policy.py b/packages/meshbay-node/tests/test_root_writable_policy.py index 730e636..23b55cb 100644 --- a/packages/meshbay-node/tests/test_root_writable_policy.py +++ b/packages/meshbay-node/tests/test_root_writable_policy.py @@ -64,8 +64,8 @@ def _session(tmp_path: Path, user_id: str, *, return session -def _upload(session, filename="clip.mp4", body=b"bytes"): - session._do_file_upload(sealed_upload( +async def _upload(session, filename="clip.mp4", body=b"bytes"): + await session._do_file_upload(sealed_upload( session, filename=filename, data=body, dir="shared")) @@ -79,7 +79,7 @@ def _uploads_dir(session) -> Path: async def test_a_member_cannot_upload_to_a_read_only_root(tmp_path): session = _session(tmp_path, "member-1", writable=False) - _upload(session) + await _upload(session) refusal = [m for m in session.sent if m.get("type") == "error"] assert refusal and refusal[0].get("code") == "root_read_only" @@ -88,7 +88,7 @@ async def test_a_member_cannot_upload_to_a_read_only_root(tmp_path): async def test_members_upload_normally_to_a_writable_root(tmp_path): session = _session(tmp_path, "member-1", writable=True) - _upload(session) + await _upload(session) assert not [m for m in session.sent if m.get("type") == "error"] assert (_uploads_dir(session) / "clip.mp4").read_bytes() == b"bytes" @@ -103,7 +103,7 @@ async def test_read_only_binds_the_operator_too(tmp_path): session = _session(tmp_path, "the-operator", writable=False, operator="the-operator") session._is_node_admin = lambda: True - _upload(session) + await _upload(session) refusal = [m for m in session.sent if m.get("type") == "error"] assert refusal and refusal[0].get("code") == "root_read_only" @@ -249,6 +249,6 @@ async def test_no_message_can_reopen_uploads_for_a_whole_group(tmp_path): "the node answered an instruction it does not implement") # And the door is still shut. - _upload(session) + await _upload(session) refusal = [m for m in session.sent if m.get("type") == "error"] assert refusal and refusal[0].get("code") == "root_read_only" diff --git a/packages/meshbay-node/tests/test_security_regressions.py b/packages/meshbay-node/tests/test_security_regressions.py index fa55ff6..e203811 100644 --- a/packages/meshbay-node/tests/test_security_regressions.py +++ b/packages/meshbay-node/tests/test_security_regressions.py @@ -174,7 +174,7 @@ def _session(tmp_path: Path, user_id: str) -> WebRTCPeerSession: return session -def test_upload_cannot_overwrite_another_members_file(tmp_path): +async def test_upload_cannot_overwrite_another_members_file(tmp_path): """ C5a: uploads used to land in the shared root under a client-chosen name and overwrite whatever was there. That let any member destroy the operator's @@ -192,7 +192,7 @@ def test_upload_cannot_overwrite_another_members_file(tmp_path): original.write_bytes(b"operator's original content") attacker = _session(tmp_path, "attacker-user") - attacker._do_file_upload(sealed_upload( + await attacker._do_file_upload(sealed_upload( attacker, filename="important.mp4", data=b"attacker content")) assert original.read_bytes() == b"operator's original content", ( @@ -200,18 +200,18 @@ def test_upload_cannot_overwrite_another_members_file(tmp_path): assert (uploads / "important (2).mp4").read_bytes() == b"attacker content" -def test_upload_second_attempt_cannot_replace_own_completed_file(tmp_path): +async def test_upload_second_attempt_cannot_replace_own_completed_file(tmp_path): """C5a: even the original uploader does not get to overwrite.""" session = _session(tmp_path, "user-1") - def _send_it(): + async def _send_it(): # Sealed afresh each time: a nonce is drawn per message, so re-sending # the same dict would be a replay rather than a second upload. - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="movie.mp4", data=b"first")) - _send_it() + await _send_it() session.sent.clear() - _send_it() + await _send_it() uploads = _uploads_dir(session) assert (uploads / "movie.mp4").read_bytes() == b"first", ( "the first upload was replaced") @@ -236,7 +236,7 @@ def test_dir_create_cannot_escape_the_shared_root(tmp_path, bad): assert set(tmp_path.rglob("*")) == before, f"created something via {bad!r}" -def test_the_client_names_a_folder_and_never_a_filesystem_path(tmp_path): +async def test_the_client_names_a_folder_and_never_a_filesystem_path(tmp_path): """ The destination is now the folder the sender is looking at, which means the client does choose it — and the whole of what keeps that safe is that the @@ -255,7 +255,7 @@ def test_the_client_names_a_folder_and_never_a_filesystem_path(tmp_path): for bad in ("../../etc", "/etc", "shared/../..", "shared/../../etc", "nope", "shared/missing"): session.sent.clear() - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="note.txt", data=b"x", dir=bad)) refusal = [m for m in session.sent if m.get("type") == "error"] assert refusal, f"{bad!r} was accepted" @@ -264,7 +264,7 @@ def test_the_client_names_a_folder_and_never_a_filesystem_path(tmp_path): assert set(tmp_path.rglob("*")) == before, "a refused upload still wrote" -def test_an_upload_lands_in_the_folder_it_names(tmp_path): +async def test_an_upload_lands_in_the_folder_it_names(tmp_path): """ And in that folder itself — the `uploads/` subdirectory the node used to create is gone. Somebody dropping a file into the folder they are looking @@ -274,7 +274,7 @@ def test_an_upload_lands_in_the_folder_it_names(tmp_path): root = session._ctx["roots"].roots[0] (root.path / "Albums").mkdir() - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="note.txt", data=b"x", dir=f"{root.name}/Albums")) assert (root.path / "Albums" / "note.txt").read_bytes() == b"x" @@ -283,7 +283,7 @@ def test_an_upload_lands_in_the_folder_it_names(tmp_path): assert not (root.path / "uploads").exists() -def test_an_upload_goes_to_the_root_it_names(tmp_path): +async def test_an_upload_goes_to_the_root_it_names(tmp_path): """ With two writable roots there is no defensible default, and the client is the only party that knows which directory the person is looking at. The @@ -301,14 +301,14 @@ def test_an_upload_goes_to_the_root_it_names(tmp_path): {"path": str(incoming), "writable": True}, ]) - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="note.txt", data=b"x", dir="Incoming")) assert (incoming / "note.txt").read_bytes() == b"x" assert not (media / "note.txt").exists(), "it went to the first root instead" -def test_a_read_only_root_refuses_an_upload(tmp_path): +async def test_a_read_only_root_refuses_an_upload(tmp_path): """ RO is the mechanism now, not a hidden button. It binds the operator too: "read-only for everyone" is what makes a published library one, and an @@ -321,7 +321,7 @@ def test_a_read_only_root_refuses_an_upload(tmp_path): session._ctx["roots"] = RootSet.build([{"path": str(published)}]) session._is_node_admin = lambda: True - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="note.txt", data=b"x", dir="Published")) refusal = [m for m in session.sent if m.get("type") == "error"] @@ -329,7 +329,7 @@ def test_a_read_only_root_refuses_an_upload(tmp_path): assert not (published / "note.txt").exists() -def test_a_fully_read_only_group_refuses_an_unaddressed_upload(tmp_path): +async def test_a_fully_read_only_group_refuses_an_unaddressed_upload(tmp_path): """ An MNP 1.0 client names no root, so the node falls back to the first writable one. There isn't one here, and the fallback must refuse rather @@ -340,7 +340,7 @@ def test_a_fully_read_only_group_refuses_an_unaddressed_upload(tmp_path): session = _session(tmp_path, "user-1") session._ctx["roots"] = RootSet.build([{"path": str(published)}]) - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="note.txt", data=b"x")) refusal = [m for m in session.sent if m.get("type") == "error"] @@ -348,7 +348,7 @@ def test_a_fully_read_only_group_refuses_an_unaddressed_upload(tmp_path): assert not (published / "note.txt").exists() -def test_an_ejected_root_refuses_an_upload(tmp_path): +async def test_an_ejected_root_refuses_an_upload(tmp_path): """ Writing to a drive somebody has their hand on is the thing eject exists to stop. `writable` is still true — that is configuration — so availability @@ -363,7 +363,7 @@ def test_an_ejected_root_refuses_an_upload(tmp_path): roots.roots[0].available = False session._ctx["roots"] = roots - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="note.txt", data=b"x", dir="USB")) refusal = [m for m in session.sent if m.get("type") == "error"] @@ -371,19 +371,19 @@ def test_an_ejected_root_refuses_an_upload(tmp_path): assert not (usb / "note.txt").exists() -def test_two_members_can_send_the_same_filename(tmp_path): +async def test_two_members_can_send_the_same_filename(tmp_path): """ One shared uploads/ means collisions are ordinary — every camera produces IMG_1234.jpg. The second gets a free name; neither replaces the other. """ first = _session(tmp_path, "user-1") - first._do_file_upload(sealed_upload( + await first._do_file_upload(sealed_upload( first, filename="IMG_1234.jpg", data=b"first")) second = _session(tmp_path, "user-2") # Same group, so the same key: `_session` builds one per call, and two # members of one group do not have two. second._ctx["gek"] = first._ctx["gek"] - second._do_file_upload(sealed_upload( + await second._do_file_upload(sealed_upload( second, filename="IMG_1234.jpg", data=b"second")) uploads = _uploads_dir(first) diff --git a/packages/meshbay-node/tests/test_upload_attribution.py b/packages/meshbay-node/tests/test_upload_attribution.py index dc8d989..63762e9 100644 --- a/packages/meshbay-node/tests/test_upload_attribution.py +++ b/packages/meshbay-node/tests/test_upload_attribution.py @@ -91,7 +91,7 @@ def _session(shared: Path, indexer: DirectoryIndexer, gek: bytes, async def _upload(session, shared: Path, name: str, data: bytes) -> None: root_name = shared.name - session._do_file_upload( + await session._do_file_upload( sealed_upload(session, filename=name, data=data, dir=root_name)) errors = [m for m in session.sent if m.get("type") == "error"] assert not errors, errors diff --git a/packages/meshbay-node/tests/test_upload_sealed.py b/packages/meshbay-node/tests/test_upload_sealed.py index 7c1be96..7f78ba7 100644 --- a/packages/meshbay-node/tests/test_upload_sealed.py +++ b/packages/meshbay-node/tests/test_upload_sealed.py @@ -80,14 +80,14 @@ def test_the_wire_message_carries_no_filename_and_no_plaintext(tmp_path): assert b"JPEGDATA" not in blob, "the file content is on the wire in clear" -def test_the_ack_carries_no_stored_name(tmp_path): +async def test_the_ack_carries_no_stored_name(tmp_path): """ `stored_as` is the name the node settled on — it finds a free one rather than replacing anything — and naming it in clear would hand back exactly what the request took the trouble to hide. """ session = _session(tmp_path) - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="holiday.jpg", data=b"x", dir=_root(session).name)) ack = [m for m in session.sent if m.get("type") == MNP.FILE_UPLOAD_ACK][-1] @@ -98,7 +98,7 @@ def test_the_ack_carries_no_stored_name(tmp_path): "dir": _root(session).name} -def test_the_ack_names_the_upload_so_one_refusal_fails_one_upload(tmp_path): +async def test_the_ack_names_the_upload_so_one_refusal_fails_one_upload(tmp_path): """ `filename` used to be the correlation key on both sides. It cannot be one any more, and `upload_id` replaces it — a client-chosen label, opaque to @@ -107,10 +107,10 @@ def test_the_ack_names_the_upload_so_one_refusal_fails_one_upload(tmp_path): for one file used to fail every upload in flight. """ session = _session(tmp_path) - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="a.txt", data=b"x", dir=_root(session).name, upload_id="upload-A")) - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="../evil", data=b"x", dir=_root(session).name, upload_id="upload-B")) @@ -121,14 +121,14 @@ def test_the_ack_names_the_upload_so_one_refusal_fails_one_upload(tmp_path): # ── What is refused ────────────────────────────────────────────────────────── -def test_a_plaintext_upload_is_refused(tmp_path): +async def test_a_plaintext_upload_is_refused(tmp_path): """ The MNP 1.x shape, which is what an un-updated client sends. Refused with a code and a message saying which side is old — never accepted "just this once", because a path that still takes plaintext is not a sealed path. """ session = _session(tmp_path) - session._do_file_upload({ + await session._do_file_upload({ "filename": "note.txt", "dir": _root(session).name, "chunk_index": 0, "total_chunks": 1, "data": b"x", }) @@ -137,7 +137,7 @@ def test_a_plaintext_upload_is_refused(tmp_path): assert not _wrote_anything(tmp_path) -def test_a_tampered_chunk_is_refused(tmp_path): +async def test_a_tampered_chunk_is_refused(tmp_path): """ AES-GCM's tag, asserted where it matters: a flipped bit in the ciphertext must stop the upload, not produce a corrupt file with a plausible name. @@ -146,13 +146,13 @@ def test_a_tampered_chunk_is_refused(tmp_path): msg = sealed_upload(session, filename="note.txt", data=b"x" * 64, dir=_root(session).name) msg["ct"] = bytes([msg["ct"][0] ^ 0x01]) + msg["ct"][1:] - session._do_file_upload(msg) + await session._do_file_upload(msg) assert _errors(session)[0]["code"] == "upload_not_sealed" assert not _wrote_anything(tmp_path) -def test_an_upload_sealed_for_another_group_is_refused(tmp_path): +async def test_an_upload_sealed_for_another_group_is_refused(tmp_path): """ The group is the AAD, so a node hosting two groups cannot have a chunk moved between them — and a member of one cannot write into the other by @@ -163,20 +163,20 @@ def test_an_upload_sealed_for_another_group_is_refused(tmp_path): session._ctx["gek"], "some-other-group", upload_id="u1", chunk_index=0, total_chunks=1, filename="note.txt", data=b"x", dir=_root(session).name) - session._do_file_upload(msg) + await session._do_file_upload(msg) assert _errors(session)[0]["code"] == "upload_not_sealed" assert not _wrote_anything(tmp_path) -def test_an_upload_under_another_key_is_refused(tmp_path): +async def test_an_upload_under_another_key_is_refused(tmp_path): """A peer past the handshake with the wrong GEK still writes nothing.""" session = _session(tmp_path) msg = file_upload_wire( generate_gek(), GROUP, upload_id="u1", chunk_index=0, total_chunks=1, filename="note.txt", data=b"x", dir=_root(session).name) - session._do_file_upload(msg) + await session._do_file_upload(msg) assert _errors(session)[0]["code"] == "upload_not_sealed" assert not _wrote_anything(tmp_path) @@ -199,7 +199,7 @@ def test_an_ack_replayed_as_a_request_does_not_open(tmp_path): file_upload_payload(session._ctx["gek"], GROUP, ack) -def test_a_group_with_no_key_refuses_rather_than_falling_back(tmp_path): +async def test_a_group_with_no_key_refuses_rather_than_falling_back(tmp_path): """ A node whose group has no GEK yet cannot open anything. It must say so, not read the message as though it were the old plaintext shape. @@ -208,7 +208,7 @@ def test_a_group_with_no_key_refuses_rather_than_falling_back(tmp_path): msg = sealed_upload(session, filename="note.txt", data=b"x", dir=_root(session).name) session._ctx["gek"] = b"" - session._do_file_upload(msg) + await session._do_file_upload(msg) assert _errors(session)[0]["code"] == "no_group_key" assert not _wrote_anything(tmp_path) @@ -216,7 +216,7 @@ def test_a_group_with_no_key_refuses_rather_than_falling_back(tmp_path): # ── What must still work ───────────────────────────────────────────────────── -def test_a_multi_chunk_upload_reassembles(tmp_path): +async def test_a_multi_chunk_upload_reassembles(tmp_path): """ Every chunk is sealed under its own nonce, and the node appends in order. Nothing about the seal may change what lands on disk. @@ -225,7 +225,7 @@ def test_a_multi_chunk_upload_reassembles(tmp_path): body = bytes(range(256)) * 40 parts = [body[i:i + 1024] for i in range(0, len(body), 1024)] for i, part in enumerate(parts): - session._do_file_upload(sealed_upload( + await session._do_file_upload(sealed_upload( session, filename="blob.bin", data=part, chunk_index=i, total_chunks=len(parts), dir=_root(session).name)) @@ -247,7 +247,7 @@ def test_two_identical_chunks_do_not_reuse_a_nonce(tmp_path): assert a["ct"] != b["ct"] -def test_a_sealed_payload_is_authenticated_not_validated(tmp_path): +async def test_a_sealed_payload_is_authenticated_not_validated(tmp_path): """ Opening a payload proves a member wrote it, not that they wrote something sensible. A member can seal anything, so the fields still need their types @@ -261,7 +261,7 @@ def test_a_sealed_payload_is_authenticated_not_validated(tmp_path): {"filename": "note.txt", "data": "not bytes"}, {"filename": "note.txt"}): session.sent.clear() - session._do_file_upload({ + await session._do_file_upload({ "type": MNP.FILE_UPLOAD, "v": "2.0", "upload_id": "u1", "chunk_index": 0, "total_chunks": 1, **seal(session._ctx["gek"], PURPOSE_UPLOAD, MNP.FILE_UPLOAD, @@ -273,13 +273,13 @@ def test_a_sealed_payload_is_authenticated_not_validated(tmp_path): assert not _wrote_anything(tmp_path) -def test_a_peer_controlled_chunk_index_cannot_crash_the_handler(tmp_path): +async def test_a_peer_controlled_chunk_index_cannot_crash_the_handler(tmp_path): """`chunk_index` is outside the seal by necessity, so it is unchecked input.""" session = _session(tmp_path) msg = sealed_upload(session, filename="note.txt", data=b"x", dir=_root(session).name) msg["chunk_index"] = "zero" - session._do_file_upload(msg) + await session._do_file_upload(msg) assert _errors(session)[0]["code"] == "bad_chunk_index" assert not _wrote_anything(tmp_path) -- cgit v1.2.3