aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py8
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/indexer.py15
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops.py6
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py64
-rw-r--r--packages/meshbay-node/tests/test_disk_io_off_loop.py98
-rw-r--r--packages/meshbay-node/tests/test_lease_enforcement.py2
-rw-r--r--packages/meshbay-node/tests/test_ops.py9
-rw-r--r--packages/meshbay-node/tests/test_partial_uploads.py48
-rw-r--r--packages/meshbay-node/tests/test_root_writable_policy.py12
-rw-r--r--packages/meshbay-node/tests/test_security_regressions.py44
-rw-r--r--packages/meshbay-node/tests/test_upload_attribution.py2
-rw-r--r--packages/meshbay-node/tests/test_upload_sealed.py42
12 files changed, 247 insertions, 103 deletions
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)