diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/transport')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py | 208 |
1 files changed, 166 insertions, 42 deletions
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 9ea70d8..8a5bbff 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -174,7 +174,8 @@ USER_BLOB_ACCOUNT_MAX = 8 * 1024 * 1024 _USER_BLOB_KIND_RE = re.compile( r"^(playlists|playlist:[A-Za-z0-9_-]{1,64})$") -# Chat link-preview results, kept in memory only (draft-v6 §2.7: the node +# Chat link-preview results, kept in memory only (docs/MESHBAY_DESIGN.md §6.5: +# the node # produces enrichment on demand and keeps nothing durable — the asking device # caches). Bounded and time-limited so a busy group cannot grow it without end # and a page that changed its card is picked up within the hour. @@ -212,6 +213,23 @@ _LINK_PREVIEW_RATE_WINDOW = 60.0 _LINK_PREVIEW_RATE_PER_CONN = 15 _LINK_PREVIEW_RATE_NODE = 60 +# A free-text TMDB search spends the *operator's* credential, which is rated by +# TMDB and shared by everyone in the group: one member typing in the search box +# can exhaust what every other member's automatic matching depends on, and the +# operator is the one who has to notice. §6.5's rule is a bound and a named +# adversary in the same commit; this one arrived without either. +# +# Per member rather than per connection, unlike link previews above: three tabs +# is one person, and a ceiling a tab can multiply is not a ceiling. Kept in the +# group context so it survives a reconnect, which is the other thing a per-session +# count cannot do. +# +# Generous next to what a person types — ten searches a minute is a search every +# six seconds, sustained — and small next to a loop. +_TMDB_SEARCH_WINDOW = 60.0 +_TMDB_SEARCH_PER_MEMBER = 10 +_TMDB_SEARCH_NODE = 30 + # Chat limits. A message is a member-supplied write onto the operator's disk # (`chat.db`, where retention is a manual CLI command — §6.6), relayed from there # to every other connected member and turned into a notification for every member @@ -327,6 +345,24 @@ SEEK_PROBE_MAX_BACKOFF_SECS = 60 # subtitle track, it is an ffmpeg that found something else to write, and it # would sit in the media cache for ever. SUBTITLE_MAX_BYTES = 8 * 1024 * 1024 + +# What a whole-file audio transcode may produce. The output is AAC at 192 kbit/s, +# so this is about forty-five minutes of source — past any track, any single +# piece, most sets. +# +# The bound is the media cache's, not memory's. `put_thumb` writes one SQLite row +# and the store is 512 MB with least-recently-used eviction, sized for what it +# holds: thumbnails, posters, subtitle tracks, short transcodes. A three-hour +# audiobook at this bitrate is ~260 MB — a single row that would evict most of +# the cache to make room for itself, and be evicted in turn by the next few +# thumbnails. It is not a size this store can hold usefully. +# +# It does not take away something that worked: `AUDIO_TRANSCODE_TIMEOUT_SECS` is +# 120, so a source long enough to reach this cap was already liable to be killed +# mid-transcode. What changes is that the refusal now says which limit was met. +# Serving audio of that length properly is streaming the transcode rather than +# buffering it, which is a different feature from this one. +AUDIO_TRANSCODE_MAX_BYTES = 64 * 1024 * 1024 # Bundle fetches are served in the pre-proof window (C4). Bounded and audited # until the native client removes remote keypair bundles entirely. MAX_PRE_PROOF_FETCHES = 4 @@ -1004,7 +1040,7 @@ class WebRTCPeerSession: self._ctx.get("daemon_state", {}).get("tmdb_token_customized", False)), "tmdb_language": str( self._ctx.get("daemon_state", {}).get("tmdb_language") or ""), - # Music app (docs/musicbay.md §6) — same shape as the TMDB + # Music app (docs/MESHBAY_DESIGN.md §9.8) — same shape as the TMDB # fields above. No language field: MusicBrainz search doesn't # take one the way TMDB does. "musicbrainz_enabled": bool(self._group_ctx().get("musicbrainz_enabled", True)), @@ -1180,7 +1216,7 @@ class WebRTCPeerSession: } # The recovery-wrapped copy (MNP 0.14) rides along when present, so a # client holding the recovery key can re-wrap it under a new - # passphrase — docs/auth-confirm.md §4.5. + # passphrase — docs/MESHBAY_DESIGN.md §3.6. if kp.get("bundle_enc_recovery"): resp["bundle_enc_recovery"] = kp["bundle_enc_recovery"] self._send(resp) @@ -1564,7 +1600,7 @@ class WebRTCPeerSession: # A person may hold several devices on one node. The authority admitting a # new one is a key the node already pinned — never the hub, which has stored # no user keys since 2026-08-14 and therefore cannot countersign anything. - # See docs/desktop-client-v1.md §4. + # See docs/MESHBAY_DESIGN.md §3.3. async def _do_device_request(self, msg: dict) -> None: """ @@ -1725,7 +1761,7 @@ class WebRTCPeerSession: # # This is what lets another member check for themselves that this device # belongs to an account whose earlier device they have already pinned, - # instead of taking the node's word (Tier 2, desktop-client-v1.md §4.8). + # instead of taking the node's word (Tier 2, docs/MESHBAY_DESIGN.md §3.3). await roster.pin_identity( user_id=self._user_id, username=self._username or "", pk_ed25519=pk_ed_b64, pk_x25519=pk_x_b64, via="device", @@ -1752,7 +1788,7 @@ class WebRTCPeerSession: What is checked, in order: the key is a live device *of this account* in the node's own roster (never a token claim — that is - `per-node-identity-v1.md`'s rule), the timestamp is fresh, and the + `docs/MESHBAY_DESIGN.md` §3.2's rule), the timestamp is fresh, and the signature verifies over a transcript naming this node, this group and this connection's nonce. A key that is merely well-formed proves nothing. @@ -1970,7 +2006,7 @@ class WebRTCPeerSession: Admission policy for a group, read from the node's own configuration. Never from the hub: a hub that could declare a group open would be handed - the key to it (§3.4 of docs/invite-pairing-v1.md). + the key to it (docs/MESHBAY_DESIGN.md §3.4). """ gctx = (self._ctx.get("groups") or {}).get(group_id) or {} return gctx.get("join_policy", "invite") @@ -2296,9 +2332,9 @@ class WebRTCPeerSession: # include "video" or "music" — both can make outbound third-party # network calls (TMDB, MusicBrainz) once enabled, so an operator opts a # group in explicitly rather than getting it for free - # (docs/mediacenter.md §5.6, docs/musicbay.md §4.4). - # `helloworld` is the reference implementation (docs/refactor-groups.md - # §4.1), hidden client-side behind `?dev=1`. It is here because the + # (docs/MESHBAY_DESIGN.md §9.7, §9.8). + # `helloworld` is the reference implementation (docs/MESHBAY_DESIGN.md + # §9.4), hidden client-side behind `?dev=1`. It is here because the # allow-list is server-side enforcement — a client that names an app this # node does not know is refused — and an app the node refused could not # demonstrate anything. This entry and the client's registry line are the @@ -2369,7 +2405,7 @@ class WebRTCPeerSession: per-group concern (see _do_tmdb_enabled for the per-group on/off switch). Signed like the rest: this changes outbound third-party network traffic the node did not have before the Videos app - (docs/mediacenter.md §5.5, §8) — an unsigned change would let any + (docs/MESHBAY_DESIGN.md §9.7, §6.5) — an unsigned change would let any member alter egress the operator never agreed to. """ token = msg.get("token") @@ -2752,7 +2788,7 @@ class WebRTCPeerSession: def _do_musicbrainz_enabled(self, msg: dict) -> None: """ Whether MusicBrainz lookups run for this group at all. Per-group - from the start (docs/musicbay.md §3.2/§6) — signed like + from the start (docs/MESHBAY_DESIGN.md §9.8) — signed like tmdb_enabled: it decides whether this group's members' Music tab ever makes outbound MusicBrainz traffic. """ @@ -3833,13 +3869,13 @@ class WebRTCPeerSession: self, thumb_hash: str, chunk_index: int, gek: bytes | None, ) -> dict | None: """ - docs/mediacenter.md §5.3: a thumbnail is served through the same + docs/MESHBAY_DESIGN.md §6.5: a thumbnail is served through the same chunked file_req path as a real file, resolved against the media cache instead of the index when the id doesn't match a file. Sliced by `chunk_index` like a real file's chunks, not just handed back whole: a thumbnail/poster/cover never approached CHUNK_SIZE so this used to be equivalent to "only chunk 0 exists", but an audio - transcode result (docs/musicbay.md, the WMA/Musepack exception) is + transcode result (docs/MESHBAY_DESIGN.md §9.8, the WMA/Musepack exception) is cached in the same media_cache blob store and can be several MB — genuinely multi-chunk, same as a file read straight off disk. """ @@ -3959,7 +3995,7 @@ class WebRTCPeerSession: async def _fetch_and_cache_poster(media_cache, tmdb_client, poster_path: str | None) -> str | None: """ Downloads a TMDB poster/backdrop once, caches it under its own - blake3 like a video thumbnail (docs/mediacenter.md §5.4), and + blake3 like a video thumbnail (docs/MESHBAY_DESIGN.md §9.7), and returns the hash a client then fetches via the normal file_req/ chunk path (§5.3) — no client ever contacts image.tmdb.org directly. @@ -4008,7 +4044,7 @@ class WebRTCPeerSession: async def _do_audio_transcode_request(self, msg: dict) -> None: """ - docs/musicbay.md's one exception to "no node-side transcode pool": + docs/MESHBAY_DESIGN.md §9.8's one exception to "no node-side transcode pool": WMA and Musepack tag/cover fine (enrich_audio.py) but decode in no mainstream browser's <audio> element at all. Transcoded to AAC/M4A once and cached under its own content hash — same "computed once, @@ -4023,6 +4059,25 @@ class WebRTCPeerSession: if not entry: self._send({"type": "error", "detail": "File not found"}) return + + # The gate `BROWSER_INCOMPATIBLE_AUDIO_EXTS` exists for, applied where it + # costs something. Nothing on the node read it: the player asks for these + # two extensions and no others, and `music-player.js` described itself as + # "kept in sync with the node's" constant — so the whole restriction lived + # in the caller, and a member's own message is not the caller. + # + # What that let through: this converts a *whole file* and holds a + # transcode slot shared with video streaming while it runs. Pointed at a + # two-hour film it spends minutes of the operator's CPU and a slot every + # other viewer is queued behind. `AUDIO_TRANSCODE_MAX_BYTES` catches the + # result, after the work; only this catches the work. + if Path(entry.name).suffix.lower() not in BROWSER_INCOMPATIBLE_AUDIO_EXTS: + self._send({ + "type": "error", + "detail": "This file does not need transcoding — play it directly.", + "code": "transcode_not_applicable", + }) + return file_path, refusal = await off_disk(ctx["roots"], _locate, ctx["roots"], entry) if refusal is not None: self._send({"type": "error", "detail": refusal}) @@ -4215,7 +4270,7 @@ class WebRTCPeerSession: async def _do_music_meta_request(self, msg: dict) -> None: """ - docs/musicbay.md §4.3: MusicBrainz metadata for one track, resolved + docs/MESHBAY_DESIGN.md §9.8: MusicBrainz metadata for one track, resolved from the group's index by its content id. Album-level (release), the direct analogue of Videos' show-level TMDB caching: one search per (artist, album) pair serves cover art and canonical naming to every @@ -4295,7 +4350,7 @@ class WebRTCPeerSession: async def _do_media_meta_request(self, msg: dict) -> None: """ - docs/mediacenter.md §5.4: TMDB metadata for one file, resolved from + docs/MESHBAY_DESIGN.md §9.7: TMDB metadata for one file, resolved from the group's index by its content id (root+relpath the client already knows from index_sync/index_delta identify the entry; its own `id` is what actually names one file — never a raw filesystem path off @@ -4320,7 +4375,7 @@ class WebRTCPeerSession: media_cache = self._ctx.get("media_cache") tmdb_client = self._ctx.get("tmdb_client") - # Per-group, not node-wide (docs/mediacenter.md §5.5, 2026-08-24): + # Per-group, not node-wide (docs/MESHBAY_DESIGN.md §9.7, 2026-08-24): # treated exactly like "no client configured" — same silent, no-error # degradation, since a member's Videos tab already has to handle "no # TMDB match" as the ordinary case. @@ -4429,7 +4484,7 @@ class WebRTCPeerSession: return media_cache = self._ctx.get("media_cache") tmdb_client = self._ctx.get("tmdb_client") - # Per-group, not node-wide (docs/mediacenter.md §5.5, 2026-08-24) — + # Per-group, not node-wide (docs/MESHBAY_DESIGN.md §9.7, 2026-08-24) — # same silent zero-confidence degradation as "no client configured". if (media_cache is None or tmdb_client is None or not self._group_ctx().get("tmdb_enabled", True)): @@ -4467,7 +4522,7 @@ class WebRTCPeerSession: async def _do_tmdb_search_request(self, msg: dict) -> None: """ Candidate TMDB matches for an operator correcting a wrong automatic - match (docs/mediacenter.md, §V-whatever this becomes) — a plain + match (docs/MESHBAY_DESIGN.md §9.7, §V-whatever this becomes) — a plain lookup, not a mutation, so unlike `tmdb_override` this needs no admin authority: any member can see what TMDB itself would offer, the same as the automatic search already silently does on their @@ -4480,7 +4535,7 @@ class WebRTCPeerSession: return media_cache = self._ctx.get("media_cache") tmdb_client = self._ctx.get("tmdb_client") - # Per-group, not node-wide (docs/mediacenter.md §5.5, 2026-08-24) — + # Per-group, not node-wide (docs/MESHBAY_DESIGN.md §9.7, 2026-08-24) — # same silent empty-results degradation as "no client configured": # a member with TMDB off for this group sees the same "type it in # yourself" affordance either way, never an error. @@ -4490,6 +4545,21 @@ class WebRTCPeerSession: "query": query, "media_type": media_type, "results": []}) return + # Refused out loud, not as an empty result: "no matches" is what the + # client draws for an empty list, and telling somebody their film is + # unknown when the node simply declined to ask is a worse answer than + # the truth. `video-app.js`'s `runSearch` puts `detail` on screen. + if not self._tmdb_search_rate_ok(): + log.info("tmdb_search_req: rate-limited (user=%s)", (self._user_id or "")[:8]) + self._send({ + "type": "error", + "detail": "Too many searches in the last minute. This spends the " + "operator's search quota, which everyone in the group " + "shares — try again shortly.", + "code": "tmdb_search_rate_limited", + }) + return + raw = (await tmdb_client.search_movie_results(query) if media_type == "movie" else await tmdb_client.search_tv_results(query)) results = [] @@ -4613,7 +4683,7 @@ class WebRTCPeerSession: def _do_tmdb_rematch(self, msg: dict) -> None: """ An operator dropping one file's cached TMDB match so it re-resolves - with the current matcher (§10.1/V13) — the one-click alternative to + with the current matcher (V13) — the one-click alternative to the full search-and-pick "Fix match" flow, and reachable without SSH (`meshbay-node video rematch` clears a whole group). Signed like `tmdb_override`: `media_cache` is shared node-wide. @@ -4655,7 +4725,8 @@ class WebRTCPeerSession: async def _tmdb_search(self, tmdb_client, entry, is_show: bool): """ - §3.3's retry ladder — same shape for movies and shows (§10.1/V8). + docs/MESHBAY_DESIGN.md §9.7's scored ladder — same shape for movies + and shows (V8). TMDB's own top result is still trusted per query (§3.3's last row — no local re-ranking of *its* list); what the ladder adds is that it *scores every candidate query* and keeps the best, instead of @@ -4701,7 +4772,7 @@ class WebRTCPeerSession: specific than a punctuation-normalised restatement of `primary` (an alternative_title, a sequel variant); when it does not and the primary hit is already decent, the remaining calls are skipped - (§10.1/V11 — they almost never win and cost a round trip each). + (V11 — they almost never win and cost a round trip each). """ def _year_of(res: dict) -> int | None: d = str(res.get("release_date") or res.get("first_air_date") or "") @@ -5121,6 +5192,33 @@ class WebRTCPeerSession: if isinstance(m.payload, bytes) else m.payload) return row + def _tmdb_search_rate_ok(self) -> bool: + """ + True when this search is within both the member's window and the node's; + records it when so, and trims both to the window on every call so neither + list can grow without bound. + + Both are checked because they answer different questions: the member's + keeps one person from spending everyone's quota, and the node's keeps a + group of them from doing it together. + """ + now = time.monotonic() + w = _TMDB_SEARCH_WINDOW + ctx = self._group_ctx() + by_member = ctx.setdefault("tmdb_search_hits", {}) + who = self._user_id or "" + mine = [t for t in by_member.get(who, []) if now - t < w] + node = [t for t in self._ctx.get("tmdb_search_hits_node", []) if now - t < w] + if len(mine) >= _TMDB_SEARCH_PER_MEMBER or len(node) >= _TMDB_SEARCH_NODE: + by_member[who] = mine + self._ctx["tmdb_search_hits_node"] = node + return False + mine.append(now) + node.append(now) + by_member[who] = mine + self._ctx["tmdb_search_hits_node"] = node + return True + def _link_preview_rate_ok(self) -> bool: """ True when this preview fetch is within both the per-connection and the @@ -5144,7 +5242,8 @@ class WebRTCPeerSession: async def _do_link_preview_request(self, msg: dict) -> None: """ - Unfurl a URL a member pasted into chat (draft-v6 §2.7 enrichment rule: + Unfurl a URL a member pasted into chat (docs/MESHBAY_DESIGN.md §6.5's + enrichment rule: the client asks, the node produces on demand, the asking device caches — nothing durable here). @@ -5915,13 +6014,14 @@ class WebRTCPeerSession: # anything failing loudly: a file uploaded from a phone could not be # deleted from the same person's laptop, and the only symptom was # "Signature verification failed" on their own file - # (docs/desktop-client-v1.md §4.8 A). + # (docs/MESHBAY_DESIGN.md §3.3). # # `uploader_pk` is kept, and stops being the authorization key: it is # now the audit record of *which device* did it. Authorization is by # account, through the roster — never through a token claim, which is - # the protection `per-node-identity-v1.md` added and which a lookup by - # `uploader_id` in the hub's world would give straight back. + # the protection per-node identity keys give (docs/MESHBAY_DESIGN.md + # §3.2) and which a lookup by `uploader_id` in the hub's world would + # give straight back. if not (await self._verify_admin_sig(transcript, sig) or await self._verify_uploader_sig(entry, transcript, sig)): self._send({"type": "error", "detail": "Signature verification failed"}) @@ -5935,7 +6035,7 @@ class WebRTCPeerSession: ) -> None: # Node operator only. A group admin who does not run the node has no # authority over who this node admits (deny by default). Delegation is - # designed but deferred — see §6.2 of docs/invite-pairing-v1.md. + # designed but deferred — see docs/MESHBAY_DESIGN.md §3.4. if not await self._verify_admin_sig(transcript, sig): self._send({"type": "error", "detail": "Signature verification failed"}) self._audit("admin_auth_failed", f"invite_create:{pending['subject'][:16]}") @@ -6240,7 +6340,7 @@ class WebRTCPeerSession: # node's answer to it, since ffmpeg re-encodes these in real time on # any machine that can run this daemon. Reported live against an # Xvid/MP3 .avi. `transcode_incompatible_video`'s own documentation - # (draft-v6 §2.11) already said "HEVC *and other browser- + # (docs/MESHBAY_DESIGN.md §6.8) already said "HEVC *and other browser- # incompatible video codecs*"; only HEVC was ever wired up. can_copy = (bool(codec_str) and raw_video_codec not in BROWSER_INCOMPATIBLE_VIDEO_CODECS) @@ -6678,6 +6778,30 @@ def _locate(roots: RootSet, entry) -> tuple[Path | None, str | None]: return path, None +def _read_scratch_capped(tmp_path: Path, cap: int, what: str) -> bytes: + """ + Stat ffmpeg's output, refuse it if it is too big, read it. Blocking. + + Run through `asyncio.to_thread` and not `off_disk`: this file is ffmpeg's + own, under `tempfile.mkstemp` on the system disk, so it is not a group root + and there is no spun-down platter to serialise against — it only has to be + off the event loop. A whole transcode read inline is tens of megabytes of + blocking read while nothing else in the node is served. + + The size is checked before the bytes are asked for, so an oversized result + costs a stat rather than the read *and* the memory. + """ + size = tmp_path.stat().st_size + if size > cap: + raise RuntimeError(f"{what} is {size} bytes, over the {cap} cap") + return tmp_path.read_bytes() + + +async def _discard_scratch(tmp_path: Path) -> None: + """Remove one of ffmpeg's temp files, off the loop like the read of it.""" + await asyncio.to_thread(tmp_path.unlink, True) + + 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: @@ -6734,9 +6858,11 @@ async def _transcode_audio_to_aac(file_path: Path) -> bytes: if proc.returncode != 0: raise RuntimeError( f"ffmpeg exited {proc.returncode}: {stderr.decode(errors='replace')[:300]}") - return tmp_path.read_bytes() + return await asyncio.to_thread( + _read_scratch_capped, tmp_path, AUDIO_TRANSCODE_MAX_BYTES, + "transcoded audio") finally: - tmp_path.unlink(missing_ok=True) + await _discard_scratch(tmp_path) async def _seek_lands_at(file_path: Path, t: float, map_args: list[str]) -> float | None: @@ -6790,7 +6916,7 @@ async def _seek_lands_at(file_path: Path, t: float, map_args: list[str]) -> floa log.warning("stream: seek probe failed at %.1fs: %r", t, e) return None finally: - tmp_path.unlink(missing_ok=True) + await _discard_scratch(tmp_path) text = stdout.decode(errors="replace").strip().rstrip(",") try: landed = float(text) @@ -6848,10 +6974,8 @@ async def _extract_subtitle_to_webvtt(file_path: Path, ordinal: int, if proc.returncode != 0: raise RuntimeError( f"ffmpeg exited {proc.returncode}: {stderr.decode(errors='replace')[:300]}") - size = tmp_path.stat().st_size - if size > SUBTITLE_MAX_BYTES: - raise RuntimeError(f"subtitle track is {size} bytes, over the {SUBTITLE_MAX_BYTES} cap") - blob = tmp_path.read_bytes() + blob = await asyncio.to_thread( + _read_scratch_capped, tmp_path, SUBTITLE_MAX_BYTES, "subtitle track") # A WebVTT file that is only its header has no cues in it. That is what # a bitmap track extracted by mistake produces, and what a text track # whose stream is empty produces; either way there is nothing to show, @@ -6861,7 +6985,7 @@ async def _extract_subtitle_to_webvtt(file_path: Path, ordinal: int, raise RuntimeError("extracted subtitle contains no cues") return blob finally: - tmp_path.unlink(missing_ok=True) + await _discard_scratch(tmp_path) class WebRTCTransport: @@ -6931,9 +7055,9 @@ class WebRTCTransport: `ctx["_transcode_sem"]`, and `hasattr(webrtc, "_stream_sem")` is always False. So the hot-swap was a no-op and **`max_concurrent_streams` has never taken effect from the Node page without a restart**, contrary to - draft-v6 §2.11. This is the one implementation, on the object that owns - the state, so the next two caps do not each grow their own copy of the - mistake. + docs/MESHBAY_DESIGN.md §6.8. This is the one implementation, on the + object that owns the state, so the next two caps do not each grow their + own copy of the mistake. What resizing means, stated because it is a decision and not a detail: **the new cap governs new streams; the ones already running are |