diff options
Diffstat (limited to 'packages/meshbay-node/src')
3 files changed, 289 insertions, 14 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index e1cab61..371c2f0 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -32,6 +32,7 @@ import logging import os import signal import sys +import time from pathlib import Path import uvicorn @@ -51,6 +52,7 @@ from meshbay_node.indexer.enrich_audio import AudioEnricher from meshbay_node.indexer.enrich_photo import PhotoEnricher from meshbay_node.media_cache import MediaCache from meshbay_node.tmdb import TmdbClient +from meshbay_node import uploads as uploads_mod from meshbay_node.musicbrainz import MusicBrainzClient from meshbay_node.keystore import create_keystore, load_keystore, load_or_create_keystore from meshbay_node.platform import chmod_private, config_dir, data_dir, state_dir @@ -249,6 +251,8 @@ class NodeDaemon: self._tasks.append(asyncio.create_task(ui_server.serve())) log.info("Control API on 127.0.0.1:%d", self._config.node.ui_port) + self._tasks.append(asyncio.create_task(self._reap_partial_uploads())) + # 3. Hub connection (Ed25519 auth — retries until node key is linked) hub_cfg = HubConfig( hub_url=self._config.hub.url, @@ -1057,6 +1061,60 @@ class NodeDaemon: group_id[:8], e) return 0 + async def _reap_partial_uploads(self, interval: float = 3600.0, + first_delay: float = 60.0) -> None: + """ + Delete `.part` files that no upload will ever finish. + + An upload interrupted for good leaves its partial file behind, and + nothing else ever looks at it: `.part` is not an index entry, so it is + invisible to every group member and to the operator's own file list. One + abandoned film is a gigabyte of their disk, kept for ever. + + Two conditions, both required, and `uploads.orphaned_parts` is where + they are stated and tested. What this adds is the walk and the deletion, + and one rule of its own: it runs a minute after start rather than at + once, so a client reconnecting to finish an upload that outlived a node + restart is not raced by the janitor that would have deleted it — the age + threshold makes that impossible in practice, and doing it anyway costs a + minute. + + `interval` and `first_delay` are parameters so a test can drive this + without waiting an hour. + """ + await asyncio.sleep(first_delay) + while True: + try: + self._reap_once() + except Exception as exc: # never let the janitor kill the node + log.warning("Reaping partial uploads failed: %s", exc) + await asyncio.sleep(interval) + + def _reap_once(self, now: float | None = None) -> int: + """One pass over every group. Returns how many files were deleted.""" + groups = (self._webrtc._ctx.get("groups") or {}) if self._webrtc else {} + when = time.time() if now is None else now + deleted = 0 + for gid, ctx in groups.items(): + roots = ctx.get("roots") + if roots is None: + continue + store = ctx.get("partial_uploads") + live = store.live_paths() if store is not None else set() + for path in uploads_mod.orphaned_parts( + uploads_mod.find_parts(roots.roots), live, when): + try: + size = path.stat().st_size + path.unlink() + except OSError as exc: + log.warning("Could not remove abandoned upload %s: %s", + path.name, exc) + continue + deleted += 1 + log.info("Removed abandoned upload %s (%d bytes, group %s)", + path.name, size, gid[:8]) + return deleted + async def _progress_pusher(self, indexer: DirectoryIndexer, interval: float = 2.0) -> None: """ 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 ce599cc..6bae27c 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -129,6 +129,7 @@ from meshbay_node.indexer import GroupIndex from meshbay_node.indexer.indexer import DirectoryIndexer from meshbay_node import linkpreview, ops, platform from meshbay_node import transfers as transfers_mod +from meshbay_node import uploads as uploads_mod from meshbay_node.transfers import TransferSlots # Re-imported under its original name: every call site and existing test in # this module still refers to it as `_probe_video`. The implementation lives @@ -426,7 +427,8 @@ class WebRTCPeerSession: self._join_attempts = 0 self._nonce_client: bytes = b"" self._admin_ops: dict[str, dict] = {} # op_id → pending admin operation - self._uploads: dict[str, dict] = {} # filename → {next_index, bytes} + # Uploads in progress live in the group context, not here: see + # `_partial_uploads` and `uploads.py`. # Diagnostics only (_WEBRTC_TRACE): when the last DataChannel message # arrived, so the heartbeat can report silence duration. self._last_msg_at: float = 0.0 @@ -4757,6 +4759,19 @@ class WebRTCPeerSession: if k not in ("type", "v")}) self._send(resp) + def _partial_uploads(self, ctx: dict) -> uploads_mod.PartialUploads: + """This group's uploads in progress, created on first use. + + In the group context rather than on the session, so a client that + reconnects finds its own upload where it left it — and so the reaper has + something to ask "is anyone still writing this?". + """ + store = ctx.get("partial_uploads") + if store is None: + store = uploads_mod.PartialUploads() + ctx["partial_uploads"] = store + return store + def _do_file_upload(self, msg: dict) -> None: """ One chunk of an upload, sealed under the group key (MNP 2.0). @@ -4912,13 +4927,26 @@ class WebRTCPeerSession: "root_unavailable") return - upload_key = f"{rel_dir}/{filename}" - state = self._uploads.get(upload_key) + # Held by the group, not by this connection. + # + # This used to be `self._uploads`, on the session. A dropped link threw + # the position away and the next chunk was refused with `not_started`: + # an upload interrupted at 99% could only be started again from zero, on + # a connection flaky enough to have interrupted it once. And the state + # it lost was the only thing that knew about the `.part` file left + # behind — see `uploads.orphaned_parts`, which is the other half of this. + # + # Keyed by member as well as by name, because a shared directory means + # two people can be sending IMG_1234.jpg at the same moment and neither + # may inherit the other's position. + uploads = self._partial_uploads(ctx) + user_id = self._user_id or "" + state = uploads.get(user_id, rel_dir, filename) # 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) - tmp_path = target_dir / f"{stored_name}.part" + stored_name = state.stored_name if state else _free_name(target_dir, filename) + tmp_path = target_dir / f"{stored_name}{uploads_mod.PART_SUFFIX}" final_path = target_dir / stored_name if chunk_index == 0: @@ -4927,28 +4955,27 @@ class WebRTCPeerSession: if final_path.exists(): _refuse("File already exists", "already_exists") return - state = {"next_index": 0, "bytes": 0, "stored_name": stored_name} - self._uploads[upload_key] = state + state = uploads.start(user_id, rel_dir, filename, stored_name, + part_path=tmp_path) elif state is None: _refuse("Upload not started", "not_started") return # Reject out-of-order or replayed chunks — otherwise chunk_index>0 appends # blindly to whatever .part file is already on disk. - if chunk_index != state["next_index"]: + if chunk_index != state.next_index: _refuse("Unexpected chunk index", "bad_chunk_index") return - if state["bytes"] + len(chunk_bytes) > MAX_UPLOAD_BYTES: - self._uploads.pop(upload_key, None) + if state.bytes + len(chunk_bytes) > MAX_UPLOAD_BYTES: + uploads.drop(user_id, rel_dir, filename) tmp_path.unlink(missing_ok=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) - state["next_index"] = chunk_index + 1 - state["bytes"] += len(chunk_bytes) + uploads.advance(user_id, rel_dir, filename, chunk_index, len(chunk_bytes)) self._send(file_upload_ack_wire( gek, self._group_id or "", @@ -4962,10 +4989,10 @@ class WebRTCPeerSession: )) if chunk_index + 1 >= total_chunks: - self._uploads.pop(upload_key, None) + uploads.drop(user_id, rel_dir, filename) tmp_path.rename(final_path) log.info("Upload complete: %s (%d chunks, %d bytes)", - stored_name, total_chunks, state["bytes"]) + stored_name, total_chunks, state.bytes) self._audit("file_upload", f"{rel_dir}/{stored_name}") self._register_uploader(ctx, rel_dir, stored_name) diff --git a/packages/meshbay-node/src/meshbay_node/uploads.py b/packages/meshbay-node/src/meshbay_node/uploads.py new file mode 100644 index 0000000..f8ae7f9 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/uploads.py @@ -0,0 +1,190 @@ +""" +Partial uploads: the state that must outlive a connection, and the files that +must not outlive their upload. + +Two defects live here, and they are the same defect seen from two sides. + +An upload's progress was kept on the **session** — `WebRTCSession._uploads`, +keyed by `rel_dir/filename`. A dropped connection therefore lost it, and the +client's next chunk was refused with `not_started`: an upload interrupted at +99% could only be started again from zero. The state belongs to the group, not +to the connection that happened to carry it, and it is keyed by member as well, +because a shared directory means two people can be sending `IMG_1234.jpg` at +the same time and neither may inherit the other's position. + +And what the lost state left behind was a `.part` file that nothing would ever +finish, delete or even look at again. One abandoned upload of a film is a +gigabyte of somebody else's disk, kept for ever, invisible in the index because +`.part` is not an index entry. That is the leak this module reaps. + +Pure logic, no asyncio and no transport — the same shape as `transfers.py`, and +for the same reason: the rules are worth testing without a WebRTC connection to +build first. +""" + +from __future__ import annotations + +import time +from dataclasses import dataclass, field +from pathlib import Path +from typing import Iterable + +# What an unfinished upload is called on disk while it is being written. The +# node has always used this; it is named here because the reaper below has to +# recognise one, and a second spelling of it would be a bug nobody could see. +PART_SUFFIX = ".part" + +# 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 +# deleting an upload somebody is still making, which is unrecoverable and looks +# to them like a transfer that failed for no reason. A day covers a laptop +# closed overnight, a phone that lost signal in a tunnel, and a client that +# reconnects on the next launch — all of which are resumable and none of which +# should be swept. +ORPHAN_AFTER_SECS = 24 * 3600 + + +@dataclass +class Partial: + """One upload in progress, as the node knows it between two chunks.""" + + stored_name: str + # The `.part` this upload is writing. Recorded rather than recomputed: the + # reaper compares paths, and a path rebuilt from a root name and a relative + # directory is a second implementation of something that must agree exactly + # with the first, for ever, or a live upload gets deleted. + part_path: Path | None = None + next_index: int = 0 + bytes: int = 0 + updated_at: float = field(default_factory=time.time) + + +class PartialUploads: + """ + Every upload this group has in flight, keyed by member. + + Held in the group context rather than on a session, so that a reconnecting + client finds its own upload exactly where it left it. The key is + `(user_id, rel_dir, filename)`: the directory and name alone would let one + member resume — or clobber the position of — another member's upload of the + same name, which a shared folder makes an ordinary occurrence rather than an + attack. + """ + + def __init__(self) -> None: + self._by_key: dict[tuple[str, str, str], Partial] = {} + + # ── the state itself ──────────────────────────────────────────────────── + + def start(self, user_id: str, rel_dir: str, filename: str, + stored_name: str, part_path: Path | None = None, + now: float | None = None) -> Partial: + """Begin (or begin again) an upload, discarding any earlier position.""" + state = Partial(stored_name=stored_name, part_path=part_path, + updated_at=time.time() if now is None else now) + self._by_key[(user_id, rel_dir, filename)] = state + return state + + def get(self, user_id: str, rel_dir: str, filename: str) -> Partial | None: + return self._by_key.get((user_id, rel_dir, filename)) + + def advance(self, user_id: str, rel_dir: str, filename: str, + chunk_index: int, nbytes: int, + now: float | None = None) -> Partial | None: + """Record one accepted chunk. Returns None if there is no such upload.""" + state = self._by_key.get((user_id, rel_dir, filename)) + if state is None: + return None + state.next_index = chunk_index + 1 + state.bytes += nbytes + # Touched on every chunk, because the reaper measures *silence*, not + # age: an upload that has been running for two days is not an orphan, + # and one that stopped two days ago is, whatever it started as. + state.updated_at = time.time() if now is None else now + return state + + 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 __len__(self) -> int: + return len(self._by_key) + + # ── what the reaper must not touch ────────────────────────────────────── + + def live_paths(self) -> set[Path]: + """The `.part` files that still have an upload behind them. + + Deliberately without the member's identity: a file on disk has no owner, + and the only question the reaper asks is whether anyone is writing it. + """ + return {state.part_path for state in self._by_key.values() + if state.part_path is not None} + + +def orphaned_parts(candidates: Iterable[tuple[Path, float]], + live: set[Path], + now: float, + older_than: float = ORPHAN_AFTER_SECS) -> list[Path]: + """ + Which `.part` files may be deleted. + + `candidates` is `(path, mtime)` for every `.part` found under the group's + writable roots. A file is an orphan when **both** are true: no upload in + `live` is writing it, and nothing has been written to it for `older_than` + seconds. + + Both conditions are load-bearing. The first alone would delete an upload + that is mid-flight but whose state is held elsewhere; the second alone would + keep a file for a day after the upload that owned it was abandoned, which is + correct but is also the entire reason this is bounded rather than immediate. + + A file with a future mtime — a clock that went backwards, a filesystem with + a different idea of now — is left alone rather than treated as infinitely + old, because deleting is not reversible and a wrong clock is not evidence. + """ + doomed: list[Path] = [] + for path, mtime in candidates: + if path.suffix != PART_SUFFIX: + continue + if path in live: + continue + age = now - mtime + if age < older_than: + continue + doomed.append(path) + return doomed + + +def find_parts(roots: Iterable) -> list[tuple[Path, float]]: + """ + Every `.part` under these roots, with its modification time. + + Only writable, available roots are walked: a read-only root cannot have + received an upload, and an unavailable one is a disk that is not mounted — + walking it would find nothing and reporting nothing found there is how a + reaper deletes an entire library the day a drive is unplugged. (It cannot + here, since it only ever deletes what it finds, but the shape of that + mistake is worth refusing at the source.) + + Errors are swallowed per entry rather than per walk: one unreadable + subdirectory must not stop the rest from being tidied. + """ + found: list[tuple[Path, float]] = [] + for root in roots: + if not getattr(root, "writable", False): + continue + if not getattr(root, "available", False): + continue + try: + candidates = root.path.rglob(f"*{PART_SUFFIX}") + except OSError: + continue + for path in candidates: + try: + if not path.is_file(): + continue + found.append((path, path.stat().st_mtime)) + except OSError: + continue + return found |