aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/tests/test_partial_uploads.py
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/tests/test_partial_uploads.py')
-rw-r--r--packages/meshbay-node/tests/test_partial_uploads.py489
1 files changed, 489 insertions, 0 deletions
diff --git a/packages/meshbay-node/tests/test_partial_uploads.py b/packages/meshbay-node/tests/test_partial_uploads.py
new file mode 100644
index 0000000..f5b6602
--- /dev/null
+++ b/packages/meshbay-node/tests/test_partial_uploads.py
@@ -0,0 +1,489 @@
+"""
+Two rules about an upload that stopped in the middle.
+
+**It belongs to the group, not to the connection.** Progress used to be kept on
+the session, so a dropped connection lost it and the client's next chunk was
+refused with `not_started` — an upload interrupted at 99% could only start again
+from zero, on a link flaky enough to have interrupted it once.
+
+**And what it leaves on disk has an owner or it has an end.** The state that was
+lost left a `.part` file nothing would ever finish, delete or look at again:
+invisible in the index, because `.part` is not an index entry, and a gigabyte of
+somebody else's disk for one abandoned film.
+
+The keying is a correctness property rather than a nicety: a shared directory
+means two members can be sending `IMG_1234.jpg` at the same moment, and neither
+may inherit — or overwrite the position of — the other's.
+"""
+
+import os
+import time
+import types
+from pathlib import Path
+
+from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
+from meshbay_common.crypto import generate_gek
+from meshbay_node.daemon import NodeDaemon
+from meshbay_node.indexer.group_index import GroupIndex
+from meshbay_node.roots import Root, RootSet
+from meshbay_node.transport.webrtc_server import WebRTCPeerSession
+
+from meshbay_common.protocol import (
+ UPLOAD_PROBE_INDEX, file_upload_ack_payload,
+)
+
+from conftest import one_root, sealed_upload
+
+from meshbay_node.uploads import (
+ ORPHAN_AFTER_SECS, PART_SUFFIX, PartialUploads, find_parts, orphaned_parts,
+)
+
+
+# ── the state ───────────────────────────────────────────────────────────────
+
+def test_an_upload_is_found_again_after_the_connection_went_away():
+ """The whole point: the store outlives the session, so the position is
+ still there when the client comes back."""
+ uploads = PartialUploads()
+ uploads.start("alice", "media", "film.mkv", "film.mkv")
+ uploads.advance("alice", "media", "film.mkv", chunk_index=0, nbytes=1024)
+ uploads.advance("alice", "media", "film.mkv", chunk_index=1, nbytes=1024)
+
+ state = uploads.get("alice", "media", "film.mkv")
+ assert state is not None
+ assert state.next_index == 2
+ assert state.bytes == 2048
+
+
+def test_two_members_uploading_the_same_name_do_not_share_a_position():
+ """A shared folder makes this ordinary, not adversarial: everyone's camera
+ produces the same filenames. Inheriting the other's position would append
+ one person's chunks to another person's file."""
+ uploads = PartialUploads()
+ uploads.start("alice", "photos", "IMG_1234.jpg", "IMG_1234.jpg")
+ uploads.start("bob", "photos", "IMG_1234.jpg", "IMG_1234 (2).jpg")
+ uploads.advance("alice", "photos", "IMG_1234.jpg", 0, 10)
+
+ assert uploads.get("alice", "photos", "IMG_1234.jpg").next_index == 1
+ assert uploads.get("bob", "photos", "IMG_1234.jpg").next_index == 0
+ assert uploads.get("bob", "photos", "IMG_1234.jpg").stored_name \
+ == "IMG_1234 (2).jpg"
+
+
+def test_the_same_name_in_two_directories_is_two_uploads():
+ uploads = PartialUploads()
+ uploads.start("alice", "media", "a.bin", "a.bin")
+ uploads.start("alice", "archive", "a.bin", "a.bin")
+ uploads.advance("alice", "media", "a.bin", 0, 5)
+ assert uploads.get("alice", "archive", "a.bin").next_index == 0
+
+
+def test_advancing_an_upload_nobody_started_says_so():
+ """The caller refuses the chunk on this; silently creating the state here
+ would let a client append to whatever `.part` is already on disk."""
+ assert PartialUploads().advance("alice", "media", "x", 0, 1) is None
+
+
+def test_starting_again_forgets_the_old_position():
+ """Chunk zero means "from the beginning" — the file is opened for writing,
+ not appending, so the position has to go with it."""
+ uploads = PartialUploads()
+ uploads.start("alice", "media", "a.bin", "a.bin")
+ uploads.advance("alice", "media", "a.bin", 0, 500)
+ uploads.start("alice", "media", "a.bin", "a.bin")
+ assert uploads.get("alice", "media", "a.bin").next_index == 0
+ assert uploads.get("alice", "media", "a.bin").bytes == 0
+
+
+# ── the reaper ──────────────────────────────────────────────────────────────
+
+def _old(seconds: float) -> float:
+ return 1_000_000.0 - seconds
+
+
+NOW = 1_000_000.0
+FILM = Path("/roots/media/film.mkv.part")
+
+
+def test_a_part_nobody_is_writing_and_nobody_has_touched_is_deleted():
+ """The leak this exists to close: an abandoned upload's file, kept for ever
+ and invisible because `.part` is not an index entry."""
+ doomed = orphaned_parts([(FILM, _old(ORPHAN_AFTER_SECS + 1))],
+ live=set(), now=NOW)
+ assert doomed == [FILM]
+
+
+def test_an_upload_in_progress_is_never_deleted():
+ """Even when its file is old: a large upload over a slow link is exactly the
+ one that has been on disk the longest, and it is the one that would hurt
+ most to lose."""
+ doomed = orphaned_parts([(FILM, _old(ORPHAN_AFTER_SECS * 3))],
+ live={FILM}, now=NOW)
+ assert doomed == []
+
+
+def test_a_recently_written_part_is_left_alone():
+ """No state and recent writes is a client that has just reconnected, or one
+ whose state this node has not seen yet. Waiting a day costs disk; being
+ wrong costs somebody their upload."""
+ doomed = orphaned_parts([(FILM, _old(60))], live=set(), now=NOW)
+ assert doomed == []
+
+
+def test_the_same_name_in_another_directory_does_not_protect_it():
+ """Matched on the whole path, so an upload to `media/` cannot keep an
+ orphan in `archive/` alive for ever. Comparing names would; comparing a
+ path rebuilt from a root and a relative directory would be a second
+ implementation that has to agree with the first for ever, and the state
+ records the path it is writing instead."""
+ other = Path("/roots/archive/film.mkv.part")
+ doomed = orphaned_parts([(other, _old(ORPHAN_AFTER_SECS + 1))],
+ live={FILM}, now=NOW)
+ assert doomed == [other]
+
+
+def test_a_finished_file_is_not_a_candidate():
+ """Only `.part` is ever deleted. A bug that let this touch a real file would
+ be the worst one in the project, so the check is here as well as at the call
+ site that only offers `.part` paths."""
+ doomed = orphaned_parts(
+ [(Path("/roots/media/film.mkv"), _old(ORPHAN_AFTER_SECS * 10))],
+ live=set(), now=NOW)
+ assert doomed == []
+
+
+def test_a_file_from_the_future_is_left_alone():
+ """A clock that went backwards is not evidence that a file is abandoned, and
+ deleting is not reversible."""
+ doomed = orphaned_parts([(Path("/roots/media/a.part"), NOW + 10_000)],
+ live=set(), now=NOW)
+ assert doomed == []
+
+
+def test_the_boundary_is_the_age_itself():
+ at = [(Path("/roots/media/a.part"), _old(ORPHAN_AFTER_SECS))]
+ just_under = [(Path("/roots/media/a.part"), _old(ORPHAN_AFTER_SECS - 1))]
+ assert orphaned_parts(at, set(), NOW) == [Path("/roots/media/a.part")]
+ assert orphaned_parts(just_under, set(), NOW) == []
+
+
+def test_an_upload_records_the_file_it_is_writing():
+ """What keeps the reaper honest. Without it the two sides would have to
+ agree on how a path is built from a root name and a relative directory —
+ two implementations of one rule, and the failure mode is deleting a live
+ upload."""
+ uploads = PartialUploads()
+ uploads.start("alice", "media", "film.mkv", "film.mkv", part_path=FILM)
+ assert uploads.live_paths() == {FILM}
+ uploads.drop("alice", "media", "film.mkv")
+ assert uploads.live_paths() == set()
+
+
+def test_the_suffix_is_named_once():
+ """Two spellings of `.part` would be a bug nobody could see: the writer
+ would produce one and the reaper would look for the other."""
+ assert PART_SUFFIX == ".part"
+
+
+# ── the walk, and the deletion ──────────────────────────────────────────────
+
+def _root(tmp_path, name, *, writable=True, available=True) -> Root:
+ path = tmp_path / name
+ path.mkdir(parents=True, exist_ok=True)
+ return Root(name=name, path=path, writable=writable, available=available)
+
+
+def _aged(path: Path, seconds: float, content: bytes = b"x") -> Path:
+ path.parent.mkdir(parents=True, exist_ok=True)
+ path.write_bytes(content)
+ when = time.time() - seconds
+ os.utime(path, (when, when))
+ return path
+
+
+def test_the_walk_finds_parts_in_subdirectories(tmp_path):
+ """Uploads go into the folder the sender was looking at, which is any
+ directory in the group — not a quarantine subfolder, since 2026-08-14."""
+ root = _root(tmp_path, "media")
+ _aged(root.path / "a.part", 10)
+ _aged(root.path / "series" / "b.part", 10)
+ _aged(root.path / "series" / "kept.mkv", 10)
+ found = {p.name for p, _ in find_parts([root])}
+ assert found == {"a.part", "b.part"}
+
+
+def test_a_read_only_root_is_not_walked(tmp_path):
+ """It cannot have received an upload, so anything `.part` in it belongs to
+ the operator and is none of this code's business."""
+ root = _root(tmp_path, "library", writable=False)
+ _aged(root.path / "theirs.part", ORPHAN_AFTER_SECS * 2)
+ assert find_parts([root]) == []
+
+
+def test_an_unavailable_root_is_not_walked(tmp_path):
+ """A drive that is not mounted. Walking it finds nothing, and "nothing
+ found" is the input from which a careless janitor concludes everything is
+ gone."""
+ root = _root(tmp_path, "external", available=False)
+ _aged(root.path / "x.part", ORPHAN_AFTER_SECS * 2)
+ assert find_parts([root]) == []
+
+
+def _daemon(groups: dict) -> NodeDaemon:
+ """A daemon with nothing but what `_reap_once` reads."""
+ daemon = NodeDaemon.__new__(NodeDaemon)
+ daemon._webrtc = types.SimpleNamespace(_ctx={"groups": groups})
+ return daemon
+
+
+def test_the_janitor_deletes_the_abandoned_and_keeps_the_rest(tmp_path):
+ """End to end on real files: the old orphan goes, the recent one and the
+ one somebody is still writing stay, and a finished file is never a
+ candidate."""
+ root = _root(tmp_path, "media")
+ old = _aged(root.path / "abandoned.mkv.part", ORPHAN_AFTER_SECS + 60)
+ recent = _aged(root.path / "fresh.mkv.part", 30)
+ live = _aged(root.path / "sending.mkv.part", ORPHAN_AFTER_SECS * 2)
+ finished = _aged(root.path / "done.mkv", ORPHAN_AFTER_SECS * 5)
+
+ uploads = PartialUploads()
+ uploads.start("alice", "media", "sending.mkv", "sending.mkv", part_path=live)
+
+ daemon = _daemon({"g1": {"roots": RootSet(roots=[root]),
+ "partial_uploads": uploads}})
+ assert daemon._reap_once() == 1
+ assert not old.exists()
+ assert recent.exists() and live.exists() and finished.exists()
+
+
+def test_a_group_that_has_never_uploaded_anything_is_handled(tmp_path):
+ """No `partial_uploads` in the context yet — it is created on first use, so
+ a node that has been up for five minutes has none."""
+ root = _root(tmp_path, "media")
+ old = _aged(root.path / "left.mkv.part", ORPHAN_AFTER_SECS + 1)
+ daemon = _daemon({"g1": {"roots": RootSet(roots=[root])}})
+ assert daemon._reap_once() == 1
+ assert not old.exists()
+
+
+def test_a_group_with_no_roots_is_skipped(tmp_path):
+ assert _daemon({"g1": {}})._reap_once() == 0
+
+
+# ── across two connections ──────────────────────────────────────────────────
+
+GROUP = "g" * 32
+
+
+def _peer(ctx: dict, user_id: str = "user-1") -> WebRTCPeerSession:
+ """One connection into a group whose context is shared, as it is on a node.
+
+ Two of these standing for the same member is the whole point: the second is
+ the reconnection, and it must find what the first was doing.
+ """
+ session = WebRTCPeerSession.__new__(WebRTCPeerSession)
+ session._ctx = ctx
+ session._group_id = GROUP
+ session._user_id = user_id
+ session._pk_user = ""
+ session.sent = []
+ session._send = session.sent.append
+ session._audit = lambda *a, **k: None
+ return session
+
+
+def _group_ctx(tmp_path) -> dict:
+ shared = tmp_path / "shared"
+ shared.mkdir(exist_ok=True)
+ return {"roots": one_root(shared),
+ "index": GroupIndex(group_id=GROUP,
+ sk_node=Ed25519PrivateKey.generate()),
+ "gek": generate_gek()}
+
+
+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):
+ """The defect this stage exists to fix.
+
+ The state used to live on the session, so the second connection saw no
+ upload at all and refused the chunk with `not_started`: an upload
+ interrupted at 99% could only be started again from zero, on a link flaky
+ enough to have interrupted it once.
+ """
+ ctx = _group_ctx(tmp_path)
+ first = _peer(ctx)
+ 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",
+ data=b"second-half",
+ chunk_index=1, total_chunks=2))
+ assert _errors(second) == [], _errors(second)
+
+ root = ctx["roots"].roots[0]
+ assert (root.path / "film.mkv").read_bytes() == b"first-halfsecond-half"
+
+
+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",
+ 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",
+ 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):
+ """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",
+ data=b"half", chunk_index=0,
+ total_chunks=2))
+ live = ctx["partial_uploads"].live_paths()
+ assert len(live) == 1
+ assert next(iter(live)).name == "film.mkv.part"
+ assert next(iter(live)).exists()
+
+
+# ── asking where to resume ──────────────────────────────────────────────────
+
+
+def _acks(session, ctx):
+ return [file_upload_ack_payload(ctx["gek"], GROUP, m)
+ for m in session.sent if m.get("type") == "file_upload_ack"]
+
+
+def _probe(session, filename: str) -> dict:
+ """The question, asked exactly as the client asks it: an ordinary sealed
+ upload chunk with no bytes and the probe index."""
+ return sealed_upload(session, filename=filename, data=b"",
+ chunk_index=UPLOAD_PROBE_INDEX, total_chunks=1)
+
+
+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"))
+ assert _errors(peer) == []
+ assert _acks(peer, ctx)[0]["resume_from"] == 0
+
+
+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",
+ data=b"xxxx", chunk_index=i,
+ total_chunks=5))
+ assert _errors(first) == []
+
+ reconnected = _peer(ctx)
+ 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):
+ """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"))
+ root = ctx["roots"].roots[0]
+ assert list(root.path.iterdir()) == []
+ assert len(ctx.get("partial_uploads") or []) == 0
+ # And it promises no destination it has not taken.
+ assert _acks(peer, ctx)[0]["stored_as"] == ""
+
+
+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",
+ data=b"xxxx", chunk_index=0,
+ total_chunks=5))
+ bob = _peer(ctx, "bob")
+ 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):
+ """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",
+ 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):
+ """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",
+ data=b"", chunk_index=UPLOAD_PROBE_INDEX,
+ total_chunks=1))
+ assert [m.get("code") for m in _errors(peer)] == ["invalid_filename"]
+ assert _acks(peer, ctx) == []
+
+
+# ── the slot an upload holds ────────────────────────────────────────────────
+
+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
+ read "waiting, 0 ahead" for a minute and a half while it was transferring,
+ with three reclaims logged against it.
+
+ The download twin of this was fixed a day earlier; the same omission was
+ still here, invisible until uploads took a real lease.
+ """
+ from meshbay_node.transfers import TransferSlots, UPLOAD
+
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ slots = TransferSlots()
+ peer._ctx = dict(ctx)
+ peer._ctx["_transfer_slots"] = slots
+ peer._registry_key = "session-1"
+ lease, err = slots.open(tr="up-1", kind=UPLOAD, session_key="session-1",
+ user_id="user-1", group_id=GROUP, bytes=10, chunks=2)
+ assert not err and lease.state == "granted"
+ assert lease.used is False
+
+ 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)
+
+ assert _errors(peer) == []
+ assert slots.leases["up-1"].used is True, (
+ "the node still believes nobody took this slot up, and will reclaim it")