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/roots.py48
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py30
-rw-r--r--packages/meshbay-node/src/meshbay_node/uploads.py20
-rw-r--r--packages/meshbay-node/tests/test_partial_uploads.py65
4 files changed, 151 insertions, 12 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/roots.py b/packages/meshbay-node/src/meshbay_node/roots.py
index 89b0441..a4f9cc9 100644
--- a/packages/meshbay-node/src/meshbay_node/roots.py
+++ b/packages/meshbay-node/src/meshbay_node/roots.py
@@ -30,6 +30,7 @@ from __future__ import annotations
import asyncio
import logging
+import os
import re
from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass, field
@@ -52,25 +53,62 @@ SAFE_UPLOAD_NAME = re.compile(
re.UNICODE)
-def _free_name(directory: Path, filename: str) -> str:
+def _free_name(directory: Path, filename: str,
+ taken: frozenset[str] | set[str] = frozenset()) -> str:
"""
`filename`, or the first "name (n).ext" that is not taken.
- Never returns the name of a file that exists, so an upload cannot replace
- one — the property the per-user quarantine used to provide (C5a).
+ Never returns the name of a file that exists, nor one in `taken` — names
+ uploads in flight will publish under — so an upload cannot replace a file
+ or another upload (C5a).
"""
- if not (directory / filename).exists():
+ def free(name: str) -> bool:
+ return name not in taken and not (directory / name).exists()
+
+ if free(filename):
return filename
stem, dot, ext = filename.rpartition(".")
if not dot:
stem, ext = filename, ""
for n in range(2, 1000):
candidate = f"{stem} ({n}){dot}{ext}"
- if not (directory / candidate).exists():
+ if free(candidate):
return candidate
raise FileExistsError(filename)
+def publish_upload(part: Path, directory: Path, stored_name: str, filename: str,
+ taken: frozenset[str] | set[str] = frozenset()) -> str:
+ """
+ Move a finished `.part` to its name without ever replacing a file. The name
+ it was published under, which may not be `stored_name`.
+
+ A rename replaces whatever is at the target, and the target can appear
+ while the upload runs — the operator copying a file in, another group's
+ upload into a shared folder. A hard link refuses an existing target, so it
+ is the publication; where the filesystem has none (FAT, exFAT, some network
+ shares), the existence check and the rename are as close as it gets. A
+ taken name moves on to the next free one rather than failing the upload.
+ """
+ name = stored_name
+ for _ in range(8):
+ target = directory / name
+ try:
+ os.link(part, target)
+ except FileExistsError:
+ name = _free_name(directory, filename, taken)
+ continue
+ except OSError:
+ if target.exists():
+ name = _free_name(directory, filename, taken)
+ continue
+ part.rename(target)
+ return name
+ part.unlink()
+ return name
+ raise FileExistsError(stored_name)
+
+
def safe_subdir(roots: RootSet, rel: str) -> Path | None:
"""
Resolve a client-supplied directory inside one of the group's roots, or refuse.
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
index 1c6d1ce..f1ed139 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
@@ -8,7 +8,7 @@ from pathlib import Path
from meshbay_common.protocol import UPLOAD_PROBE_INDEX, file_upload_ack_wire, file_upload_payload
from meshbay_node import uploads as uploads_mod
-from meshbay_node.roots import SAFE_UPLOAD_NAME, RootSet, _free_name, off_disk
+from meshbay_node.roots import SAFE_UPLOAD_NAME, RootSet, _free_name, off_disk, publish_upload
from meshbay_node.transport.webrtc.disk import _append_chunk
from meshbay_node.transport.webrtc.limits import LEASE_NONE, LEASE_QUEUED
@@ -306,9 +306,13 @@ class UploadMixin:
# 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.
+ # Names other uploads into this directory will publish under are taken
+ # too: none of them is on disk yet.
+ reserved = uploads.reserved_names(rel_dir)
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}"
+ else await off_disk(roots, _free_name, target_dir, filename, reserved))
+ tmp_path = (state.part_path if state and state.part_path
+ else target_dir / uploads_mod.part_name(stored_name))
final_path = target_dir / stored_name
if chunk_index == UPLOAD_PROBE_INDEX:
@@ -361,6 +365,22 @@ class UploadMixin:
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))
+ last = chunk_index + 1 >= total_chunks
+ if last:
+ # Published before the last ack, so the ack names the file as it is
+ # on disk: publication never replaces a file, and may have had to
+ # take another free name for this one.
+ uploads.drop(user_id, rel_dir, filename)
+ try:
+ stored_name = await off_disk(roots, publish_upload, tmp_path, target_dir,
+ stored_name, filename,
+ uploads.reserved_names(rel_dir))
+ except OSError as e:
+ log.warning("Upload %s could not be published: %s", stored_name, e)
+ _refuse("The file could not be stored", "store_failed")
+ return
+ final_path = target_dir / stored_name
+
self._send(file_upload_ack_wire(
gek, self._group_id or "",
upload_id=upload_id,
@@ -372,9 +392,7 @@ class UploadMixin:
dir=rel_dir,
))
- if chunk_index + 1 >= total_chunks:
- uploads.drop(user_id, rel_dir, filename)
- await off_disk(roots, tmp_path.rename, final_path)
+ if last:
log.info("Upload complete: %s (%d chunks, %d bytes)",
stored_name, total_chunks, state.bytes)
self._audit("file_upload", f"{rel_dir}/{stored_name}")
diff --git a/packages/meshbay-node/src/meshbay_node/uploads.py b/packages/meshbay-node/src/meshbay_node/uploads.py
index dd31e1e..92b854e 100644
--- a/packages/meshbay-node/src/meshbay_node/uploads.py
+++ b/packages/meshbay-node/src/meshbay_node/uploads.py
@@ -24,6 +24,7 @@ build first.
from __future__ import annotations
+import secrets
import time
from collections.abc import Iterable
from dataclasses import dataclass, field
@@ -34,6 +35,16 @@ from pathlib import Path
# recognise one, and a second spelling of it would be a bug nobody could see.
PART_SUFFIX = ".part"
+
+def part_name(stored_name: str) -> str:
+ """The `.part` one upload writes: its final name, a tag of its own, `.part`.
+
+ Its own, because the final name alone is shared: two uploads that settled
+ on one name — two groups hosting one folder, each with its own lock —
+ would write one file, the second truncating the first.
+ """
+ return f"{stored_name}.{secrets.token_hex(4)}{PART_SUFFIX}"
+
# How long a `.part` with no upload behind it is kept before it is deleted.
#
# Generous on purpose. The cost of waiting is disk; the cost of being wrong is
@@ -107,6 +118,15 @@ class PartialUploads:
def drop(self, user_id: str, rel_dir: str, filename: str) -> Partial | None:
return self._by_key.pop((user_id, rel_dir, filename), None)
+ def reserved_names(self, rel_dir: str) -> set[str]:
+ """The final names uploads in flight into `rel_dir` will take.
+
+ None of them exists on disk yet, so a name check that looked only at
+ the directory would hand the same name to a second upload.
+ """
+ return {state.stored_name for (_u, d, _f), state in self._by_key.items()
+ if d == rel_dir}
+
def __len__(self) -> int:
return len(self._by_key)
diff --git a/packages/meshbay-node/tests/test_partial_uploads.py b/packages/meshbay-node/tests/test_partial_uploads.py
index 80ab454..138ee3b 100644
--- a/packages/meshbay-node/tests/test_partial_uploads.py
+++ b/packages/meshbay-node/tests/test_partial_uploads.py
@@ -17,6 +17,7 @@ may inherit — or overwrite the position of — the other's.
"""
import os
+import re
import time
import types
from pathlib import Path
@@ -362,7 +363,7 @@ async def test_an_upload_in_flight_is_known_to_the_reaper(tmp_path):
total_chunks=2))
live = ctx["partial_uploads"].live_paths()
assert len(live) == 1
- assert next(iter(live)).name == "film.mkv.part"
+ assert re.fullmatch(r"film\.mkv\.[0-9a-f]{8}\.part", next(iter(live)).name)
assert next(iter(live)).exists()
@@ -489,3 +490,65 @@ async def test_an_upload_chunk_says_its_slot_is_in_use(tmp_path):
assert _errors(peer) == []
assert slots.leases["up-1"].used is True, (
"the node still believes nobody took this slot up, and will reclaim it")
+
+
+
+# ── never replacing a file ──────────────────────────────────────────────────
+
+async def test_two_members_sending_one_name_at_once_get_two_files(tmp_path):
+ """
+ Both chose a free name at chunk 0, and the free name was the same: the
+ final file did not exist yet, only the first one's `.part`. They then wrote
+ one `.part`, the second truncating the first, and published it twice.
+ """
+ ctx = _group_ctx(tmp_path)
+ alice, bob = _peer(ctx, "alice"), _peer(ctx, "bob")
+ for who, data in ((alice, b"hers-1"), (bob, b"his-1")):
+ await who._do_file_upload(sealed_upload(who, filename="IMG_1234.jpg", data=data,
+ chunk_index=0, total_chunks=2))
+ for who, data in ((alice, b"hers-2"), (bob, b"his-2")):
+ await who._do_file_upload(sealed_upload(who, filename="IMG_1234.jpg", data=data,
+ chunk_index=1, total_chunks=2))
+ assert _errors(alice) == [] and _errors(bob) == []
+
+ root = ctx["roots"].roots[0].path
+ assert (root / "IMG_1234.jpg").read_bytes() == b"hers-1hers-2"
+ assert (root / "IMG_1234 (2).jpg").read_bytes() == b"his-1his-2"
+ assert _acks(bob, ctx)[-1]["stored_as"] == "IMG_1234 (2).jpg"
+
+
+async def test_a_file_that_appears_during_an_upload_is_not_replaced(tmp_path):
+ """The operator copies a file in under the same name while a member's
+ upload is running. The rename at the end used to replace it."""
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ await peer._do_file_upload(sealed_upload(peer, filename="film.mkv", data=b"up-1",
+ chunk_index=0, total_chunks=2))
+ root = ctx["roots"].roots[0].path
+ (root / "film.mkv").write_bytes(b"the operator's")
+ await peer._do_file_upload(sealed_upload(peer, filename="film.mkv", data=b"up-2",
+ chunk_index=1, total_chunks=2))
+
+ assert _errors(peer) == []
+ assert (root / "film.mkv").read_bytes() == b"the operator's"
+ assert (root / "film (2).mkv").read_bytes() == b"up-1up-2"
+ assert _acks(peer, ctx)[-1]["stored_as"] == "film (2).mkv"
+ assert not list(root.glob("*.part")), "the part was left behind"
+
+
+def test_without_hard_links_a_file_is_still_not_replaced(tmp_path, monkeypatch):
+ """FAT, exFAT and some network shares have no hard links."""
+ from meshbay_node import roots as roots_mod
+
+ def no_links(*a, **k):
+ raise PermissionError("operation not permitted")
+ monkeypatch.setattr(roots_mod.os, "link", no_links)
+ (tmp_path / "a.txt").write_bytes(b"there first")
+ part = tmp_path / "a.txt.0123abcd.part"
+ part.write_bytes(b"uploaded")
+
+ name = roots_mod.publish_upload(part, tmp_path, "a.txt", "a.txt")
+ assert name == "a (2).txt"
+ assert (tmp_path / "a.txt").read_bytes() == b"there first"
+ assert (tmp_path / "a (2).txt").read_bytes() == b"uploaded"
+ assert not part.exists()