diff options
Diffstat (limited to 'packages/meshbay-node')
32 files changed, 805 insertions, 1033 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/audit.py b/packages/meshbay-node/src/meshbay_node/audit.py index b382107..81ee600 100644 --- a/packages/meshbay-node/src/meshbay_node/audit.py +++ b/packages/meshbay-node/src/meshbay_node/audit.py @@ -33,20 +33,6 @@ CREATE INDEX IF NOT EXISTS idx_audit_user ON audit_log(user_id); CREATE INDEX IF NOT EXISTS idx_audit_event ON audit_log(event); """ -EVENTS = { - "connect", - "disconnect", - "handshake", - "file_download", - "file_upload", - "file_delete", - "stream_video", - "chat_message", - "chat_history", - "index_sync", - "auth_failed", -} - RETENTION_DAYS = 365 diff --git a/packages/meshbay-node/src/meshbay_node/blocklist.py b/packages/meshbay-node/src/meshbay_node/blocklist.py new file mode 100644 index 0000000..29d8666 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/blocklist.py @@ -0,0 +1,83 @@ +""" +The hub's content blocklist, as this node applies it (docs/MESHBAY_DESIGN.md §7.5). + +Hashes of public content that the hub's moderation has blocked. A node hosting a +public group stops serving those files there: they leave the index members are +sent, and a request for one is refused. Nothing is deleted from the operator's +disk — the node stops serving, and what the operator keeps is theirs to decide. + +Private groups are untouched. The hub never learns what a private group holds, +so there is nothing it could have blocked in one, and nothing here reaches them. + +Kept on disk as well as in memory, so a node that restarts while the hub is +unreachable does not serve again, meanwhile, what it had already stopped serving. +The hub's list is the authority: a full sync replaces this one. + +It is an exact match on the content id. A file changed by one byte is another +id, and files over the partial-hash threshold are identified by a sample of their +bytes (§6.3) — a moderation tool, not a guarantee. +""" + +import json +import logging +import os +import re +from pathlib import Path + +log = logging.getLogger(__name__) + +_HASH = re.compile(r"^[0-9a-f]{64}$") + + +class ContentBlocklist: + def __init__(self, path: Path | None = None): + self._path = path + self._hashes: set[str] = set() + if path is not None: + try: + data = json.loads(path.read_text(encoding="utf-8")) + self._hashes = {h for h in data.get("hashes", []) if _HASH.match(h)} + except FileNotFoundError: + pass + except (OSError, ValueError) as e: + # A damaged file is not a reason to refuse to start; the next + # sync with the hub rewrites it. + log.warning("Content blocklist %s unreadable: %s", path, e) + + def __contains__(self, content_hash: object) -> bool: + return content_hash in self._hashes + + def __len__(self) -> int: + return len(self._hashes) + + def __iter__(self): + return iter(self._hashes) + + def replace(self, hashes) -> bool: + """The hub's full list. True when it differs from what was applied.""" + new = {h for h in hashes if isinstance(h, str) and _HASH.match(h)} + if new == self._hashes: + return False + self._hashes = new + self._save() + return True + + def apply(self, add=(), remove=()) -> bool: + """One pushed change. True when it changed anything.""" + before = set(self._hashes) + self._hashes |= {h for h in add if isinstance(h, str) and _HASH.match(h)} + self._hashes -= {h for h in remove if isinstance(h, str)} + if self._hashes == before: + return False + self._save() + return True + + def _save(self) -> None: + if self._path is None: + return + tmp = self._path.with_suffix(".tmp") + try: + tmp.write_text(json.dumps({"hashes": sorted(self._hashes)}), encoding="utf-8") + os.replace(tmp, self._path) + except OSError as e: + log.warning("Content blocklist %s not saved: %s", self._path, e) diff --git a/packages/meshbay-node/src/meshbay_node/config.py b/packages/meshbay-node/src/meshbay_node/config.py index 563e395..2cd8a8f 100644 --- a/packages/meshbay-node/src/meshbay_node/config.py +++ b/packages/meshbay-node/src/meshbay_node/config.py @@ -505,10 +505,3 @@ def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config: "MESHBAY_MAX_CONCURRENT_STREAMS") return cfg - - -def write_example_config(path: Path = DEFAULT_CONFIG_PATH) -> None: - """Write an example config file if none exists.""" - if not path.exists(): - path.parent.mkdir(parents=True, exist_ok=True) - path.write_text(EXAMPLE_CONFIG, encoding="utf-8", newline="\n") diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index 220a908..45a6b9a 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -44,6 +44,7 @@ from meshbay_common.protocol import MNP from meshbay_node import uploads as uploads_mod from meshbay_node.audit import RETENTION_DAYS as AUDIT_RETENTION_DAYS from meshbay_node.audit import AuditStore +from meshbay_node.blocklist import ContentBlocklist from meshbay_node.bundle_store import BundleStore from meshbay_node.chat.store import ChatStore from meshbay_node.cli.dispatch import run, start @@ -114,6 +115,10 @@ class NodeDaemon(EnrichmentMixin): # Persisted so a restart does not silently un-revoke everyone (H4) self._denylist = ( Denylist(path=config.data_dir / "denylist.json") if Denylist else None) + # The hub's content blocklist, applied in public groups only (§7.5). + # Persisted for the same reason: a restart while the hub is unreachable + # must not serve again what had stopped being served. + self._blocklist = ContentBlocklist(config.data_dir / "blocklist.json") self._chat_stores: dict[str, ChatStore] = {} # One instance, shared by every group's DirectoryIndexer — see # indexer/cache.py's docstring for why this stopped being per-group. @@ -525,6 +530,7 @@ class NodeDaemon(EnrichmentMixin): self._webrtc._ctx["pk_x25519_b64"] = keys.pk_x25519_b64 self._webrtc._ctx["roster"] = self._roster + self._webrtc._ctx["blocklist"] = self._blocklist # The MNP adapter calls the same operations as the loopback API # (meshbay_node.ops), and those take the daemon's state. Handing # the transport a second set of lookups is how two paths to one @@ -611,11 +617,16 @@ class NodeDaemon(EnrichmentMixin): except Exception as e: log.warning("Invalid revocation token: %s", e) + async def on_connected(): + await self._sync_blocklist(hub) + ws_task = asyncio.create_task(hub.maintain_ws( on_incoming=on_incoming, on_revocation=on_revocation, on_webrtc_offer=on_webrtc_offer, group_ids=lambda: list((self._state.get("groups_ctx") or {}).keys()), + on_connected=on_connected, + on_blocklist=self._on_blocklist_update, )) self._tasks.append(ws_task) log.info("Hub WS task started") @@ -667,12 +678,6 @@ class NodeDaemon(EnrichmentMixin): await indexer.initial_scan() log.info("Background scan complete for %s: %d files", name, indexer.index.count) - # Swarm registration for public groups (after files are known). - if gctx.get("visibility") == "public": - endpoint = f"webrtc:{self._config.node.quic_port}" - hashes = [e.id for e in gctx["index"].entries] - if hashes: - await self._register_swarm(hashes, endpoint) # initial_scan() itself never calls on_change (it predates # the concept — every existing caller only cared about the # scan finishing, not about notifying anyone) — but Videos @@ -1454,9 +1459,10 @@ class NodeDaemon(EnrichmentMixin): # node refuses every handshake while the GEK is None (NS8) — so this is # "nobody is listening", not a case to send in clear for. if peers and idx.gek: - msg = (index_delta_message(idx, delta, indexer.roots) + hidden = self._blocklist if self._is_public(group_id) else () + msg = (index_delta_message(idx, delta, indexer.roots, hidden) if delta is not None - else index_sync_message(idx, indexer.roots)) + else index_sync_message(idx, indexer.roots, hidden)) pushed = 0 for session in peers: try: @@ -1468,20 +1474,53 @@ class NodeDaemon(EnrichmentMixin): log.info("Index %s pushed to %d WebRTC peers", "delta" if delta is not None else "sync", pushed) - # 11.9 — Register file hashes with hub swarm table (public groups only, H7) - group_cfg = next( - (g for g in self._config.groups if g.id == group_id), None) - if (self._hub and self._state.get("endpoint_hint") - and group_cfg and group_cfg.visibility == "public"): - # Only the newly added hashes once there is a delta to know them - # from — registering the whole library again on every change is - # the same O(changes x library size) cost the delta above exists - # to avoid. - hashes = ([e.id for e in delta.additions] if delta is not None - else [e.id for e in idx.entries]) - if hashes: - endpoint = f"webrtc:{self._config.node.quic_port}" - spawn(self._register_swarm(hashes, endpoint)) + # ── Content blocklist (public groups, §7.5) ────────────────────────────── + + def _is_public(self, group_id: str) -> bool: + return any(g.id == group_id and g.visibility == "public" + for g in self._config.groups) + + def _hosts_public_group(self) -> bool: + return any(g.visibility == "public" for g in self._config.groups) + + async def _sync_blocklist(self, hub) -> None: + """The hub's whole list, on every connection. Only a node hosting a + public group asks: nothing else here could have been blocked.""" + if not self._hosts_public_group(): + return + try: + hashes = await hub.fetch_blocklist() + except Exception as e: + # Keep applying the list on disk; the next connection tries again. + log.warning("Content blocklist sync failed: %s", e) + return + if self._blocklist.replace(hashes): + log.info("Content blocklist: %d hashes", len(self._blocklist)) + self._push_public_indexes() + + def _on_blocklist_update(self, add: list, remove: list) -> None: + if not self._hosts_public_group(): + return + if self._blocklist.apply(add, remove): + log.info("Content blocklist updated: +%d -%d", len(add), len(remove)) + self._push_public_indexes() + + def _push_public_indexes(self) -> None: + """Resend each public group's whole index, so what the list now hides + leaves every connected member's view, and what it released comes back.""" + if not self._webrtc: + return + for indexer in self._indexers: + idx = indexer.index + if not self._is_public(idx.group_id) or not idx.gek: + continue + msg = index_sync_message(idx, indexer.roots, self._blocklist) + for session in list(self._webrtc._sessions.values()): + if session._group_id == idx.group_id: + try: + session._send(msg) + except Exception: + pass def _drop_group_sessions(self, group_id: str) -> None: """Close live sessions for a revoked group (H4).""" @@ -1492,13 +1531,6 @@ class NodeDaemon(EnrichmentMixin): spawn(session.close()) log.info("Dropped session for revoked group %s", group_id[:8]) - async def _register_swarm(self, hashes: list[str], endpoint: str) -> None: - try: - n = await self._hub.register_swarm(hashes, endpoint) - log.info("Swarm: registered %d/%d hashes", n, len(hashes)) - except Exception as e: - log.warning("Swarm registration failed: %s", e) - async def _shutdown(self) -> None: log.info("Shutting down...") self._state["status"] = "stopping" diff --git a/packages/meshbay-node/src/meshbay_node/hub_client.py b/packages/meshbay-node/src/meshbay_node/hub_client.py index 346a4cd..2ce8993 100644 --- a/packages/meshbay-node/src/meshbay_node/hub_client.py +++ b/packages/meshbay-node/src/meshbay_node/hub_client.py @@ -6,7 +6,6 @@ Handles all communication from the node to a Mesh Hub: - JWT offline verification and auto-refresh - Node announcement (endpoint_hint) - User public key lookup (for GEK wrapping) - - Swarm hash registration The node authenticates via Ed25519 challenge-response (/v1/nodes/auth). No auth_key or password is ever stored on or transmitted from the node. @@ -335,6 +334,8 @@ class HubClient: on_revocation: Any = None, on_webrtc_offer: Any = None, group_ids: list[str] | None = None, # static list or callable returning one + on_connected: Any = None, + on_blocklist: Any = None, ) -> None: """ Maintain a persistent WebSocket connection to the hub. @@ -394,6 +395,12 @@ class HubClient: self._ws = ws log.info("Hub WS connected") + if on_connected: + # Off the read loop, for the reason offers are: a slow + # fetch must not stop this socket being read. + task = asyncio.create_task(on_connected()) + pending.add(task) + task.add_done_callback(pending.discard) async for raw in ws: msg = json.loads(raw) @@ -406,6 +413,9 @@ class HubClient: elif mtype == "revocation" and on_revocation: on_revocation(msg.get("token", "")) + elif mtype == "blocklist_update" and on_blocklist: + on_blocklist(msg.get("add") or [], msg.get("remove") or []) + elif mtype == "webrtc_offer" and on_webrtc_offer: # Answered off the read loop on purpose. Awaiting the # handler here meant one slow negotiation stopped the @@ -461,26 +471,25 @@ class HubClient: log.warning("Could not deliver WebRTC answer to %s: %s", str(msg.get("peer_id"))[:8], e) - # ── Swarm registration ───────────────────────────────────────────────── + # ── Content blocklist (public groups) ───────────────────────────────── - async def register_swarm(self, content_hashes: list[str], endpoint: str) -> int: - """Register file hashes in the hub swarm table. Returns count registered.""" + async def fetch_blocklist(self, max_pages: int = 100) -> set[str]: + """The hub's content blocklist, every page of it (docs/MESHBAY_DESIGN.md §7.5).""" if self._session is None: raise RuntimeError("Not logged in") await self.ensure_fresh_token() - - registered = 0 - for h in content_hashes: - try: - r = await self._http.post("/v1/swarm/register", json={ - "content_hash": h, - "endpoint": endpoint, - }, headers=self._session.auth_headers) - if r.status_code in (201, 200): - registered += 1 - except Exception: - pass - return registered + hashes: set[str] = set() + after = "" + for _ in range(max_pages): + r = await self._http.get("/v1/blocklist", params={"after": after}, + headers=self._session.auth_headers) + r.raise_for_status() + page = r.json() + hashes.update(page.get("hashes") or []) + after = page.get("next") or "" + if not after: + return hashes + raise RuntimeError("Content blocklist longer than expected; not applied") # ── Convenience: full startup sequence ─────────────────────────────────── diff --git a/packages/meshbay-node/src/meshbay_node/indexer/group_index.py b/packages/meshbay-node/src/meshbay_node/indexer/group_index.py index 9e4e780..85bc13e 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/group_index.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/group_index.py @@ -1,67 +1,41 @@ """ -Mesh Group Index — encrypted file listing for a group. +Mesh Group Index — the file listing for one group, and the deltas between versions. -Wire format (private group): - msgpack({entries, version, group_id}) → zstd compress → GEK ChaCha20 encrypt → sign - -Wire format (public group): - msgpack({entries, version, group_id}) → sign (no encryption) +What travels on the wire is built from it by `transport/wire.py` and sealed under +the group key (`meshbay_common.groupbox`); this module holds no encoding of its own. Delta format: {base_version, version, additions: [...], deletions: [id, ...]} """ -import base64 import logging -from dataclasses import asdict, dataclass, field +from dataclasses import dataclass, field -import blake3 -import msgpack -import zstandard as zstd from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey -from meshbay_common.crypto import ( - pk_to_b64, - sign_chunk, - verify_chunk_signature, -) from meshbay_common.protocol import IndexDelta, IndexEntry -from meshbay_common.webcrypto import ( - chunk_key_aes as derive_chunk_key, -) -from meshbay_common.webcrypto import ( - decrypt_chunk_aes as decrypt_chunk, -) -from meshbay_common.webcrypto import ( - encrypt_chunk_aes as encrypt_chunk, -) log = logging.getLogger(__name__) -ZSTD_LEVEL = 3 # fast compression -INDEX_CHUNK = 0 # the index itself is treated as chunk 0 of a virtual "index file" - @dataclass class GroupIndex: """ - Encrypted, signed Mesh Group Index for one group. + Mesh Group Index for one group. Usage: idx = GroupIndex(group_id="...", sk_node=sk, gek=gek_bytes) idx.add_entry(entry) - wire_bytes = idx.serialize() # for sending to members - recovered = GroupIndex.deserialize(wire_bytes, sk_node=sk, gek=gek_bytes) """ group_id: str sk_node: Ed25519PrivateKey gek: bytes | None = None # None → public group (no encryption) version: int = 1 - # The group's roots and whether each is readable right now. Travels inside - # the encrypted payload because it names the operator's directories, and a - # member needs it to tell "temporarily unavailable" from "deleted" — a - # distinction the entries alone cannot carry, since an unavailable root's - # files are still listed. Absent in an index written before roots existed. + # The group's roots and whether each is readable right now, as the indexer + # last described them. A member needs this to tell "temporarily unavailable" + # from "deleted" — a distinction the entries alone cannot carry, since an + # unavailable root's files are still listed. What members receive is built + # from the RootSet by `transport/wire.py`, not from this copy. roots: list = field(default_factory=list) _entries: dict = field(default_factory=dict, repr=False) # id → IndexEntry @@ -101,111 +75,6 @@ class GroupIndex: idx._entries = dict(entries_by_id) return idx - # ── Serialisation ───────────────────────────────────────────────────────── - - def serialize(self) -> bytes: - """ - Produce a signed index envelope: - msgpack → zstd → [GEK encrypt if private] → sign → length-prefixed envelope - - **This is not an MNP message.** It was the payload of `index_sync` on the QUIC - transport, while WebRTC sent plain entries under the same type — one message - type, two encodings (2026-09-03). Both transports now build `index_sync` from - `transport/wire.py`. This stays as a correct at-rest/interchange format, and - as the only thing that signs and encrypts a whole index; read it as that, not - as a wire contract. - """ - payload = msgpack.packb({ - "group_id": self.group_id, - "version": self.version, - "roots": list(self.roots), - "entries": [asdict(e) for e in self.entries], - }, use_bin_type=True) - - compressed = zstd.compress(payload, level=ZSTD_LEVEL) - - if self.gek is not None: - # Private group: encrypt with GEK-derived key - idx_hash = blake3.blake3(compressed).digest() - ckey = derive_chunk_key(self.gek, idx_hash, INDEX_CHUNK) - nonce, ct = encrypt_chunk(ckey, compressed) - sig = sign_chunk(self.sk_node, INDEX_CHUNK, nonce, blake3.blake3(ct).digest()) - envelope = msgpack.packb({ - "type": "index", - "encrypted": True, - "version": self.version, - "group_id": self.group_id, - "idx_hash_b64": base64.b64encode(idx_hash).decode(), - "nonce_b64": base64.b64encode(nonce).decode(), - "ct_b64": base64.b64encode(ct).decode(), - "sig_b64": base64.b64encode(sig).decode(), - "pk_node_b64": pk_to_b64(self.sk_node.public_key()), - }, use_bin_type=True) - else: - # Public group: just sign the compressed payload - payload_hash = blake3.blake3(compressed).digest() - sig = sign_chunk(self.sk_node, INDEX_CHUNK, - bytes(12), # zero nonce for plaintext - payload_hash) - envelope = msgpack.packb({ - "type": "index", - "encrypted": False, - "version": self.version, - "group_id": self.group_id, - "data_b64": base64.b64encode(compressed).decode(), - "hash_b64": base64.b64encode(payload_hash).decode(), - "sig_b64": base64.b64encode(sig).decode(), - "pk_node_b64": pk_to_b64(self.sk_node.public_key()), - }, use_bin_type=True) - - return envelope - - @classmethod - def deserialize( - cls, - data: bytes, - sk_node: Ed25519PrivateKey, - gek: bytes | None = None, - ) -> "GroupIndex": - """Deserialize, verify signature, and decrypt (if private).""" - from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey - - envelope = msgpack.unpackb(data, raw=False) - - pk_node_raw = base64.b64decode(envelope["pk_node_b64"]) - pk_node = Ed25519PublicKey.from_public_bytes(pk_node_raw) - sig = base64.b64decode(envelope["sig_b64"]) - - if envelope["encrypted"]: - if gek is None: - raise ValueError("GEK required to decrypt private group index") - ct = base64.b64decode(envelope["ct_b64"]) - nonce = base64.b64decode(envelope["nonce_b64"]) - ct_hash = blake3.blake3(ct).digest() - verify_chunk_signature(pk_node, INDEX_CHUNK, nonce, ct_hash, sig) - - idx_hash = base64.b64decode(envelope["idx_hash_b64"]) - ckey = derive_chunk_key(gek, idx_hash, INDEX_CHUNK) - compressed = decrypt_chunk(ckey, nonce, ct) - else: - compressed = base64.b64decode(envelope["data_b64"]) - payload_hash = base64.b64decode(envelope["hash_b64"]) - verify_chunk_signature(pk_node, INDEX_CHUNK, bytes(12), payload_hash, sig) - - payload = msgpack.unpackb(zstd.decompress(compressed), raw=False) - idx = cls( - group_id=payload["group_id"], - sk_node=sk_node, - gek=gek, - version=payload["version"], - # Absent from an index written before roots existed; an empty list - # reads as "nothing known about availability", not "no roots". - roots=payload.get("roots") or [], - ) - for e in payload["entries"]: - idx.add_entry(IndexEntry(**e)) - return idx - # ── Delta ───────────────────────────────────────────────────────────────── def diff(self, previous: "GroupIndex") -> IndexDelta: diff --git a/packages/meshbay-node/src/meshbay_node/modules/__init__.py b/packages/meshbay-node/src/meshbay_node/modules/__init__.py deleted file mode 100644 index e69de29..0000000 --- a/packages/meshbay-node/src/meshbay_node/modules/__init__.py +++ /dev/null diff --git a/packages/meshbay-node/src/meshbay_node/musicbrainz.py b/packages/meshbay-node/src/meshbay_node/musicbrainz.py index ce2d233..e3681ea 100644 --- a/packages/meshbay-node/src/meshbay_node/musicbrainz.py +++ b/packages/meshbay-node/src/meshbay_node/musicbrainz.py @@ -183,10 +183,6 @@ class MusicBrainzClient: results = (data or {}).get("releases", []) return _best_match_release(artist, album, results) - async def release_details(self, mbid: str) -> dict | None: - """Full release details, including recordings (tracklist).""" - return await self._get(_BASE_URL + f"release/{mbid}", {"inc": "recordings+artist-credits"}) - async def fetch_cover_art(self, mbid: str) -> bytes | None: """ The release's front cover, or None if Cover Art Archive has nothing diff --git a/packages/meshbay-node/src/meshbay_node/platform.py b/packages/meshbay-node/src/meshbay_node/platform.py index 09397a2..2255036 100644 --- a/packages/meshbay-node/src/meshbay_node/platform.py +++ b/packages/meshbay-node/src/meshbay_node/platform.py @@ -226,6 +226,14 @@ def _startup_vbs() -> Path: def _pid_alive(pid: int) -> bool: """Whether a process with this pid exists, in any session (tasklist lists session 0 too, where a service-mode daemon runs; opening it would not).""" + if sys.platform != "win32": + # A zombie counts as gone: it has exited, and only its parent has not + # reaped it yet -- os.kill(pid, 0) would still find it. + try: + stat = Path(f"/proc/{pid}/stat").read_text(encoding="ascii") + except OSError: + return False + return stat.rsplit(")", 1)[1].split()[0] != "Z" r = subprocess.run(["tasklist", "/FI", f"PID eq {pid}", "/NH", "/FO", "CSV"], capture_output=True, text=True) return f'"{pid}"' in r.stdout diff --git a/packages/meshbay-node/src/meshbay_node/replication.py b/packages/meshbay-node/src/meshbay_node/replication.py deleted file mode 100644 index 297ed0b..0000000 --- a/packages/meshbay-node/src/meshbay_node/replication.py +++ /dev/null @@ -1,140 +0,0 @@ -""" -MeshBay Node — content replication (node-to-node, admin-authorized). - -A replication node downloads files from a source node and stores them -locally, then registers itself as an additional swarm source in the hub. -This provides redundancy and improves availability for public content. - -Only public content is replicated (no GEK needed). -Private content replication requires the GEK and is admin-controlled. - -Usage: - replicator = ContentReplicator( - hub_url=..., access_token=..., source_endpoint=..., - local_dir=Path("/data/replicated"), node_pk_b64=..., - ) - await replicator.replicate_file(file_id, file_name, file_size) -""" - -import logging -from pathlib import Path - -import blake3 -import httpx - -log = logging.getLogger(__name__) - -CHUNK_SIZE = 1024 * 1024 # 1 MB - - -class ContentReplicator: - """ - Downloads public files from a source node and registers as swarm source. - """ - - def __init__( - self, - hub_url: str, - access_token: str, - source_endpoint: str, # "http://ip:port" of source node HTTP API - local_dir: Path, - node_pk_b64: str, - ): - self._hub_url = hub_url.rstrip("/") - self._access_token = access_token - self._source_endpoint = source_endpoint.rstrip("/") - self._local_dir = local_dir - self._node_pk_b64 = node_pk_b64 - self._local_dir.mkdir(parents=True, exist_ok=True) - - @property - def _auth_headers(self) -> dict: - return {"Authorization": f"Bearer {self._access_token}"} - - async def fetch_index(self) -> list[dict]: - """Fetch the public Mesh Group Index from the source node.""" - async with httpx.AsyncClient(timeout=30) as c: - r = await c.get(f"{self._source_endpoint}/index") - r.raise_for_status() - return r.json()["entries"] - - async def replicate_file( - self, - file_id: str, - file_name: str, - file_size: int, - progress_cb=None, - ) -> Path: - """ - Download a public file from the source node, verify integrity, - save locally, and register as swarm source in the hub. - Returns the local file path. - """ - local_path = self._local_dir / file_name - if local_path.exists(): - # Verify hash - existing_hash = blake3.blake3(local_path.read_bytes()).hexdigest() - if existing_hash == file_id: - log.info("Already have %s, skipping", file_name) - await self._register_swarm(file_id, local_path) - return local_path - - log.info("Replicating %s (%d bytes) from %s", file_name, file_size, self._source_endpoint) - - # Stream download chunk by chunk - n_chunks = max(1, (file_size + CHUNK_SIZE - 1) // CHUNK_SIZE) - with open(local_path, "wb") as f: - async with httpx.AsyncClient(timeout=60) as c: - for chunk_idx in range(n_chunks): - # Download full file (simpler for public content) - if chunk_idx == 0: - r = await c.get( - f"{self._source_endpoint}/file/{file_id}", - headers=self._auth_headers, - ) - r.raise_for_status() - f.write(r.content) - if progress_cb: - progress_cb(len(r.content), file_size) - break # full file downloaded in one request - - # Verify hash - actual_hash = blake3.blake3(local_path.read_bytes()).hexdigest() - if actual_hash != file_id: - local_path.unlink(missing_ok=True) - raise ValueError(f"Hash mismatch: expected {file_id[:16]}, got {actual_hash[:16]}") - - log.info("Replicated %s (hash OK)", file_name) - await self._register_swarm(file_id, local_path) - return local_path - - async def _register_swarm(self, content_hash: str, local_path: Path) -> None: - """Register this node as a swarm source for the content hash in the hub.""" - try: - async with httpx.AsyncClient(timeout=10) as c: - r = await c.post( - f"{self._hub_url}/v1/swarm/register", - json={"content_hash": content_hash, "endpoint": self._source_endpoint}, - headers=self._auth_headers, - ) - r.raise_for_status() - log.debug("Registered as swarm source for %s", content_hash[:16]) - except Exception as e: - log.warning("Failed to register swarm source: %s", e) - - async def replicate_all(self, progress_cb=None) -> list[Path]: - """Replicate all public files from the source node.""" - entries = await self.fetch_index() - results = [] - for entry in entries: - try: - path = await self.replicate_file( - file_id=entry["id"], - file_name=entry["name"], - file_size=entry["size"], - progress_cb=progress_cb, - ) - results.append(path) - except Exception as e: - log.error("Failed to replicate %s: %s", entry["name"], e) - return results diff --git a/packages/meshbay-node/src/meshbay_node/transport/quic_client.py b/packages/meshbay-node/src/meshbay_node/transport/quic_client.py index 8f3fb74..d548103 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/quic_client.py +++ b/packages/meshbay-node/src/meshbay_node/transport/quic_client.py @@ -212,16 +212,19 @@ class QuicChunkClient: "session — refusing to handshake without channel binding") binding = quic_binding(self._peer_cert_der) - # MNP 3.4: a signed challenge proves the node key before we send anything - # else. A wrong signature is refused; an absent one is an older node. - if reply.get("sig"): - try: - Ed25519PublicKey.from_public_bytes( - base64.b64decode(reply.get("node_pk", "")) - ).verify(base64.b64decode(reply["sig"]), challenge_transcript( - self._group_id, nonce_c, nonce_s, binding)) - except Exception as exc: - raise ConnectionError(f"Node challenge signature invalid: {exc}") from exc + # A signed challenge proves the node key before we send anything else. + # Every node we can reach signs (the floor is 4.0, signing is 3.4), and + # one without a certificate to bind could not complete the proof anyway, + # so a missing signature is refused like a wrong one. + if not reply.get("sig"): + raise ConnectionError("Node challenge is not signed") + try: + Ed25519PublicKey.from_public_bytes( + base64.b64decode(reply.get("node_pk", "")) + ).verify(base64.b64decode(reply["sig"]), challenge_transcript( + self._group_id, nonce_c, nonce_s, binding)) + except Exception as exc: + raise ConnectionError(f"Node challenge signature invalid: {exc}") from exc self._proto._send(self._ctrl_stream, { "type": MNP.HANDSHAKE_RESPONSE, diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py index f42d8e8..daaee62 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py @@ -5,6 +5,7 @@ import base64 import os import time +import msgpack from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey from meshbay_common import MNP_VERSION from meshbay_common.adminop import ( @@ -42,6 +43,16 @@ from meshbay_common.adminop import ( from meshbay_common.crypto import pk_to_b64 from meshbay_common.protocol import MNP +# What one connection may have waiting for a signature. Anyone authenticated can +# ask for a challenge — the signature is what is checked, and it comes later — so +# without a bound a member who never answers makes the node keep every request, +# payload and all, for the life of the connection (§13.5b). A person signs one +# operation at a time; a handful covers a settings page saving several at once. +MAX_PENDING_ADMIN_OPS = 8 +# The subject and payload of one pending operation, packed. A path is at most a +# few KiB, and the largest field a legitimate request carries is a directory list. +MAX_ADMIN_OP_BYTES = 64 * 1024 + # Which executor runs each signed operation once its signature has been # checked. Every one runs as a task of the session. _ADMIN_EXECUTORS = { @@ -98,8 +109,22 @@ class AdminMixin: (e.g. root management from a NodePage connection). """ gid = group_id if group_id is not None else (self._group_id or "") + now = time.time() + for op_id, pending in list(self._admin_ops.items()): + if now - pending["ts"] > ADMIN_CHALLENGE_TTL: + del self._admin_ops[op_id] + if len(self._admin_ops) >= MAX_PENDING_ADMIN_OPS: + self._send({"type": "error", "detail": "Too many operations waiting for a " + "signature", "code": "too_many_pending"}) + self._audit("admin_pending_flood", op) + return + if len(msgpack.packb([subject, payload or {}], use_bin_type=True)) \ + > MAX_ADMIN_OP_BYTES: + self._send({"type": "error", "detail": "Request too large", + "code": "too_large"}) + return nonce = os.urandom(32) - ts = int(time.time()) + ts = int(now) op_id = base64.b64encode(os.urandom(16)).decode() self._admin_ops[op_id] = { "op": op, "subject": subject, "nonce": nonce, "ts": ts, diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py index b02e59a..e678a13 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py @@ -8,7 +8,12 @@ import time from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey from meshbay_common import MNP_VERSION -from meshbay_common.adminop import OP_INVITE_CANCEL, OP_INVITE_CREATE, OP_INVITE_LINK_CREATE +from meshbay_common.adminop import ( + OP_INVITE_CANCEL, + OP_INVITE_CREATE, + OP_INVITE_LINK_CREATE, + invite_create_subject, +) from meshbay_common.crypto import wrap_gek_aes from meshbay_common.device import ( DEVICE_TTL, @@ -71,11 +76,13 @@ class AdmissionMixin: }) return - self._issue_admin_challenge(OP_INVITE_CREATE, invitee_id, { - "group_id": group_id, - "user_id": invitee_id, - "username": str(msg.get("username", ""))[:64], - }) + username = str(msg.get("username", ""))[:64] + self._issue_admin_challenge( + OP_INVITE_CREATE, invite_create_subject(invitee_id, username), { + "group_id": group_id, + "user_id": invitee_id, + "username": username, + }) def _do_invite_link_create(self, msg: dict) -> None: """ diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py index f2b6fa5..857db33 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py @@ -64,6 +64,8 @@ class MusicMixin: """ ctx = self._group_ctx() file_id = msg.get("file_id", "") + if self._refuse_blocked(file_id): + return entry = ctx["index"].get_entry(file_id) if not entry: self._send({"type": "error", "detail": "File not found"}) diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py index a425c65..4337e24 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py @@ -228,6 +228,8 @@ class StreamingMixin: async def _stream_video_inner(self, msg: dict) -> None: ctx = self._group_ctx() file_id = msg.get("file_id", "") + if self._refuse_blocked(file_id): + return entry = ctx["index"].get_entry(file_id) if not entry: self._send({"type": "error", "detail": "File not found"}) diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py index bd2ab52..70781eb 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py @@ -39,6 +39,8 @@ class SubtitlesMixin: """ ctx = self._group_ctx() file_id = msg.get("file_id", "") + if self._refuse_blocked(file_id): + return entry = ctx["index"].get_entry(file_id) if not entry: self._send({"type": "error", "detail": "File not found"}) diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/video_meta.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/video_meta.py index 834bac4..a9b35dc 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/video_meta.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/video_meta.py @@ -11,6 +11,7 @@ from meshbay_common.adminop import ( OP_TMDB_ENABLED, OP_TMDB_OVERRIDE, OP_TMDB_REMATCH, + tmdb_config_subject, ) from meshbay_common.protocol import MNP @@ -60,12 +61,12 @@ class VideoMetaMixin: if not self._has_admin_authority(): self._send({"type": "error", "detail": "No authorized key for this"}) return - # The subject is the signed, audited, human-shown string — it must - # never contain the token itself (it would end up in the audit log - # in plaintext). The actual token travels only in `payload`, which - # is node-side context, never re-sent or re-verified from the wire. - # The language is not a secret, so it travels in the subject itself. - subject = f"custom_token={'yes' if token else 'no'},language={language or 'default'}" + # The subject is the signed, audited string, so it must never contain + # the token itself (it would end up in the audit log in plaintext); it + # carries the token's SHA-256 instead, which binds the signature to this + # token without writing it down. `None` (unchanged) and `""` (clear) + # stay distinct, for the token and the language alike. + subject = tmdb_config_subject(token, language) self._issue_admin_challenge( OP_TMDB_CONFIG, subject, payload={"token": token, "language": language}, diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py index c3fb623..882ddbc 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py @@ -246,6 +246,25 @@ class SessionCore: return self._ctx["groups"].get(self._group_id) or {} return self._ctx + def _hidden_ids(self): + """The content blocklist, in a public group; nothing anywhere else. + + A private group's content never reaches the hub, so nothing in it can have + been blocked there (docs/MESHBAY_DESIGN.md §7.5, `meshbay_node.blocklist`). + """ + if self._group_ctx().get("visibility") != "public": + return () + return self._ctx.get("blocklist") or () + + def _refuse_blocked(self, file_id) -> bool: + """Refuse a file the blocklist names, in a public group. True if refused.""" + if not isinstance(file_id, str) or file_id not in self._hidden_ids(): + return False + self._send({"type": "error", "detail": "This file is not available here.", + "code": "content_blocked", "file_id": file_id}) + self._audit("content_blocked", file_id[:16]) + return True + def _register_peer(self) -> None: """Add this connection to its group's peer set. diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py index e76e23e..f629563 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py @@ -179,7 +179,7 @@ class FilesMixin: def _do_index_sync(self) -> None: ctx = self._group_ctx() - self._send(index_sync_message(ctx["index"], ctx.get("roots"))) + self._send(index_sync_message(ctx["index"], ctx.get("roots"), self._hidden_ids())) async def _try_serve_thumbnail( self, thumb_hash: str, chunk_index: int, gek: bytes | None, @@ -236,8 +236,18 @@ class FilesMixin: self._note_unleased(tr) file_id = msg["file_id"] chunk_index = msg["chunk_index"] + if self._refuse_blocked(file_id): + return entry = ctx["index"].get_entry(file_id) if not entry: + # A blocked file's own thumbnail is a preview of it. + hidden = self._hidden_ids() + if hidden and file_id in { + e.thumb_hash for e in map(ctx["index"].get_entry, hidden) + if e is not None and e.thumb_hash}: + self._send({"type": "error", "detail": "This file is not available here.", + "code": "content_blocked", "file_id": file_id}) + return thumb = await self._try_serve_thumbnail(file_id, chunk_index, ctx.get("gek")) if thumb is not None: log.debug("file_req file_id=%s chunk=%s: served as thumbnail", diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/handshake.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/handshake.py index eceb2b4..f701072 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/handshake.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/handshake.py @@ -24,7 +24,6 @@ from meshbay_common.handshake import ( ) from meshbay_common.protocol import MNP -from meshbay_node import transfers as transfers_mod from meshbay_node.indexer.indexer import DirectoryIndexer from meshbay_node.transport.webrtc.channel import _extract_dtls_fingerprint, _get_remote_ip from meshbay_node.transport.webrtc.limits import MAX_MSG @@ -266,19 +265,6 @@ class HandshakeMixin: # No `chat_encrypted` beside it: there is no switch. A peer that # reached this point speaks MNP 2.0, and 2.0 has no plaintext chat. "chat_epoch": int(self._group_ctx().get("chat_epoch", 0) or 0), - # This member's own transfer caps in this group, so the interface - # can say "2 of 2 of your slots are busy" rather than draw a bare - # spinner. Absent reads as "no limit known" and the hint is simply - # not drawn — never as "unlimited", which would have the interface - # contradicting the node. - "transfer_limits": { - "download": self._slots().member_cap( - transfers_mod.DOWNLOAD, - (self._group_id or "", self._user_id or "")), - "upload": self._slots().member_cap( - transfers_mod.UPLOAD, - (self._group_id or "", self._user_id or "")), - }, # So a client that connects mid-scan shows the indexing state # immediately, instead of waiting for the next periodic # INDEX_PROGRESS push. Never a path or filename — see diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py index 9d58d8c..f7bbbfa 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py @@ -14,6 +14,8 @@ from meshbay_common.adminop import ( OP_ROOT_UPDATE, OP_SET_SCAN_SETTINGS, OP_TRANSFER_LIMITS, + group_attach_subject, + root_add_subject, ) from meshbay_common.protocol import MNP @@ -246,10 +248,11 @@ class NodeOpsMixin: # is ignored rather than obeyed: on load it forces every other root # read-only, which is the model the RO/RW one replaced. A second # writable directory is `root_add` with `writable`. + writable = bool(msg.get("writable", True)) + # The directory being exposed is signed, not only the group's name. self._issue_admin_challenge( - OP_GROUP_ATTACH, name, - payload={"name": name, "shared_dir": shared_dir, - "writable": bool(msg.get("writable", True))}, + OP_GROUP_ATTACH, group_attach_subject(name, shared_dir, writable), + payload={"name": name, "shared_dir": shared_dir, "writable": writable}, group_id="") async def _admin_exec_group_attach( @@ -342,16 +345,20 @@ class NodeOpsMixin: if not self._has_admin_authority(): self._send({"type": "error", "detail": "No authorized key for this"}) return + payload = { + "group_id": target_group, "path": path, + "name": str(msg.get("name", ""))[:128], + "kind": str(msg.get("kind", "generic"))[:16], + "writable": bool(msg.get("writable", msg.get("upload", False))), + "removable": bool(msg.get("removable", False)), + } + # Everything the executor acts on is signed — `writable` decides whether + # every member may write there. The group is in the transcript itself. self._issue_admin_challenge( - OP_ROOT_ADD, path, - payload={ - "group_id": target_group, "path": path, - "name": str(msg.get("name", ""))[:128], - "kind": str(msg.get("kind", "generic"))[:16], - "writable": bool(msg.get("writable", msg.get("upload", False))), - "removable": bool(msg.get("removable", False)), - }, - group_id=target_group) + OP_ROOT_ADD, + root_add_subject(path, payload["name"], payload["kind"], + payload["writable"], payload["removable"]), + payload=payload, group_id=target_group) async def _admin_exec_root_add( self, pending: dict, transcript: bytes, sig: bytes, diff --git a/packages/meshbay-node/src/meshbay_node/transport/wire.py b/packages/meshbay-node/src/meshbay_node/transport/wire.py index 6986b01..258edb5 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/wire.py +++ b/packages/meshbay-node/src/meshbay_node/transport/wire.py @@ -12,11 +12,10 @@ envelope produced by `GroupIndex.serialize()`. Same message type, two encodings, consumer each and nothing asserting they matched. Same failure mode as the two `file_chunk` encoders, and the same fix: one builder, used by both. -`GroupIndex.serialize()`/`deserialize()` are unchanged and still tested — they remain -a correct signed index envelope — but they no longer describe any MNP message. Read -them as an at-rest/interchange format, not as a wire contract. It is also not a -candidate for reuse below: it compresses with zstd, which no browser can decompress -(`DecompressionStream` offers gzip and deflate only). +That envelope, and `GroupIndex.serialize()`/`deserialize()` which produced it, are +gone: once both transports built `index_sync` here, nothing stored or exchanged it. +It was never a candidate for reuse below either — it compressed with zstd, which no +browser can decompress (`DecompressionStream` offers gzip and deflate only). Since MNP 1.0 both messages carry their payload **sealed under a GEK-derived subkey** (`meshbay_common.groupbox`). Only the routing fields — `type`, `v`, `group_id` — stay @@ -66,7 +65,7 @@ def list_dirs(roots: RootSet | None) -> list[str]: return sorted(out)[:MAX_DIRS] -def index_sync_message(index, roots: RootSet | None) -> dict: +def index_sync_message(index, roots: RootSet | None, hidden=()) -> dict: """ The full `index_sync` message for one group. @@ -74,10 +73,13 @@ def index_sync_message(index, roots: RootSet | None) -> dict: without them a folder someone just created, or one they emptied, does not exist as far as a client is concerned, and a member cannot tell "the drive is unplugged" from "it is all still there". + + `hidden` is the content blocklist in a public group (`meshbay_node.blocklist`): + those entries are not sent at all. """ payload = { "version": index.version, - "entries": [index_entry_wire(e) for e in index.entries], + "entries": [index_entry_wire(e) for e in index.entries if e.id not in hidden], "dirs": list_dirs(roots), "roots": roots.describe() if roots else [], } @@ -89,7 +91,7 @@ def index_sync_message(index, roots: RootSet | None) -> dict: } -def index_delta_message(index, delta, roots=None) -> dict: +def index_delta_message(index, delta, roots=None, hidden=()) -> dict: """ One `index_delta` — what changed since the last thing this node broadcast. @@ -104,13 +106,16 @@ def index_delta_message(index, delta, roots=None) -> dict: page: the delta that told them something had changed was the one message that could not say what. It is a handful of dicts, bounded by the number of directories a group has, and it is sealed with the rest. + + `hidden`, as for `index_sync_message`: a blocked entry is never added or + updated; its deletion still goes out. """ payload = { "base_version": delta.base_version, "version": delta.version, - "additions": [index_entry_wire(e) for e in delta.additions], + "additions": [index_entry_wire(e) for e in delta.additions if e.id not in hidden], "deletions": list(delta.deletions), - "updates": [index_entry_wire(e) for e in delta.updates], + "updates": [index_entry_wire(e) for e in delta.updates if e.id not in hidden], } if roots is not None: payload["roots"] = roots.describe() diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py index 50beb26..bbc4649 100644 --- a/packages/meshbay-node/src/meshbay_node/ui/app.py +++ b/packages/meshbay-node/src/meshbay_node/ui/app.py @@ -20,7 +20,6 @@ from fastapi.responses import JSONResponse from meshbay_common.background import spawn from meshbay_node import __version__, ops -from meshbay_node.indexer.indexer import DirectoryIndexer log = logging.getLogger(__name__) @@ -493,50 +492,6 @@ def create_ui_app(state: dict) -> FastAPI: raise HTTPException(400, "apps must be a non-empty list") return await _op(lambda: ops.set_enabled_apps(state, group_id, apps)) - # ── App directories (operator only, localhost) ──────────────────────── - # - # The loopback twin of the `app_directories` MNP op. One endpoint for every - # application, keyed by the app's own name, so adding one needs no route - # here — the same reason the op is generic. `ALLOWED_APPS` is checked on - # the MNP path; here the caller is already on localhost holding the run - # token, and `ops` refuses a directory outside the group's roots either - # way, so an unknown key writes one unread settings row and nothing else. - - @app.put("/api/groups/{group_id}/app-directories/{app_key}") - async def set_app_directories(group_id: str, app_key: str, payload: dict): - dirs = payload.get("directories") - if not isinstance(dirs, list): - raise HTTPException(400, "directories must be a list") - return await _op(lambda: ops.set_app_directories( - state, group_id, app_key, [str(d) for d in dirs])) - - @app.put("/api/groups/{group_id}/chat-directory") - async def set_chat_directory(group_id: str, payload: dict): - return await _op(lambda: ops.set_chat_directory( - state, group_id, str(payload.get("path") or ""))) - - @app.put("/api/groups/{group_id}/chat-link-preview") - async def set_chat_link_preview(group_id: str, payload: dict): - return await _op(lambda: ops.set_chat_link_preview( - state, group_id, bool(payload.get("enabled", True)))) - - @app.put("/api/groups/{group_id}/search-listed") - async def set_search_listed(group_id: str, payload: dict): - return await _op(lambda: ops.set_search_listed( - state, group_id, bool(payload.get("listed", True)))) - - # ── Scan settings (operator only, localhost) ────────────────────────── - - @app.put("/api/groups/{group_id}/scan-settings") - async def set_scan_settings(group_id: str, payload: dict): - return await _op(lambda: ops.set_scan_settings( - state, group_id, - float(payload.get("reconcile_interval_secs", - DirectoryIndexer.DEFAULT_RECONCILE_SECS)), - float(payload.get("debounce_secs", - DirectoryIndexer.DEFAULT_DEBOUNCE_SECS)), - )) - # ── Reload config ──────────────────────────────────────────────────── @app.post("/api/reload") diff --git a/packages/meshbay-node/tests/golden/dispatch.json b/packages/meshbay-node/tests/golden/dispatch.json index 5d15057..a7bb543 100644 --- a/packages/meshbay-node/tests/golden/dispatch.json +++ b/packages/meshbay-node/tests/golden/dispatch.json @@ -1874,166 +1874,6 @@ "sent": [], "spawned": [] }, - "chat_attach | challenged | bare": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | challenged | lists": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | challenged | numbers": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | challenged | strings": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | fresh | bare": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | fresh | lists": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | fresh | numbers": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | fresh | strings": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "chat_attach | member | bare": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | member | lists": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | member | numbers": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | member | strings": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | operator | bare": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | operator | lists": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | operator | numbers": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "chat_attach | operator | strings": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, "chat_directory | challenged | bare": { "audit": [], "log": [], @@ -2228,7 +2068,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2551,7 +2391,7 @@ "subject": "gggggggggggggggggggggggggggggggg", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2570,7 +2410,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2589,7 +2429,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2608,7 +2448,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2878,7 +2718,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2892,7 +2732,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2906,7 +2746,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2920,7 +2760,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2934,7 +2774,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2948,7 +2788,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2962,7 +2802,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -2976,7 +2816,7 @@ "messages": [], "req_id": 4242, "type": "chat_hist_resp", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -7229,166 +7069,6 @@ "sent": [], "spawned": [] }, - "ephemeral_stream | challenged | bare": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | challenged | lists": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | challenged | numbers": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | challenged | strings": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | fresh | bare": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | fresh | lists": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | fresh | numbers": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | fresh | strings": { - "audit": [], - "log": [], - "sent": [ - { - "detail": "Handshake required", - "req_id": 4242, - "type": "error" - } - ], - "spawned": [] - }, - "ephemeral_stream | member | bare": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | member | lists": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | member | numbers": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | member | strings": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | operator | bare": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | operator | lists": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | operator | numbers": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, - "ephemeral_stream | operator | strings": { - "audit": [], - "log": [ - "WARNING Unknown MNP message type on DataChannel: %s" - ], - "sent": [], - "spawned": [] - }, "file_chunk | challenged | bare": { "audit": [], "log": [], @@ -8855,7 +8535,7 @@ "subject": "gggggggggggggggggggggggggggggggg", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -8874,7 +8554,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -8893,7 +8573,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -8912,7 +8592,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9244,10 +8924,10 @@ "op": "group_attach", "op_id": "<volatile>", "req_id": 4242, - "subject": "['x']", + "subject": "{\"name\":\"['x']\",\"shared_dir\":\"['x']\",\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9263,10 +8943,10 @@ "op": "group_attach", "op_id": "<volatile>", "req_id": 4242, - "subject": "7", + "subject": "{\"name\":\"7\",\"shared_dir\":\"7\",\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9282,10 +8962,10 @@ "op": "group_attach", "op_id": "<volatile>", "req_id": 4242, - "subject": "x", + "subject": "{\"name\":\"x\",\"shared_dir\":\"x\",\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9620,7 +9300,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9639,7 +9319,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -9658,7 +9338,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -11961,7 +11641,7 @@ "subject": "link:gggggggggggggggggggggggggggggggg", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14060,7 +13740,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14079,7 +13759,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14098,7 +13778,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14433,7 +14113,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14452,7 +14132,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -14471,7 +14151,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16532,7 +16212,7 @@ "req_id": 4242, "token": null, "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16547,7 +16227,7 @@ "x" ], "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16560,7 +16240,7 @@ "req_id": 4242, "token": 7, "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16573,7 +16253,7 @@ "req_id": 4242, "token": "x", "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16586,7 +16266,7 @@ "req_id": 4242, "token": null, "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16601,7 +16281,7 @@ "x" ], "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16614,7 +16294,7 @@ "req_id": 4242, "token": 7, "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16627,7 +16307,7 @@ "req_id": 4242, "token": "x", "type": "pong", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16959,10 +16639,10 @@ "op": "root_add", "op_id": "<volatile>", "req_id": 4242, - "subject": "['x']", + "subject": "{\"kind\":\"['x']\",\"name\":\"['x']\",\"path\":\"['x']\",\"removable\":true,\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16978,10 +16658,10 @@ "op": "root_add", "op_id": "<volatile>", "req_id": 4242, - "subject": "7", + "subject": "{\"kind\":\"7\",\"name\":\"7\",\"path\":\"7\",\"removable\":true,\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -16997,10 +16677,10 @@ "op": "root_add", "op_id": "<volatile>", "req_id": 4242, - "subject": "x", + "subject": "{\"kind\":\"x\",\"name\":\"x\",\"path\":\"x\",\"removable\":true,\"writable\":true}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17335,7 +17015,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17354,7 +17034,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17373,7 +17053,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17708,7 +17388,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17727,7 +17407,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -17746,7 +17426,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18081,7 +17761,7 @@ "subject": "['x']", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18100,7 +17780,7 @@ "subject": "7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18119,7 +17799,7 @@ "subject": "x", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18454,7 +18134,7 @@ "subject": "['x']:rw=on,rem=on", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18473,7 +18153,7 @@ "subject": "7:rw=on,rem=on", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -18492,7 +18172,7 @@ "subject": "x:rw=on,rem=on", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -21436,10 +21116,10 @@ "op": "tmdb_config", "op_id": "<volatile>", "req_id": 4242, - "subject": "custom_token=no,language=default", + "subject": "{\"language\":null,\"token\":null}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -21479,10 +21159,10 @@ "op": "tmdb_config", "op_id": "<volatile>", "req_id": 4242, - "subject": "custom_token=yes,language=x", + "subject": "{\"language\":\"x\",\"token\":\"sha256:2d711642b726b04401627ca9fbac32f5c8530fb1903cc4db02258717921a4881\"}", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] @@ -23369,7 +23049,7 @@ "subject": "d=7,u=7", "ts": "<volatile>", "type": "admin_challenge", - "v": "4.0" + "v": "5.0" } ], "spawned": [] diff --git a/packages/meshbay-node/tests/test_admin_challenge_bounds.py b/packages/meshbay-node/tests/test_admin_challenge_bounds.py new file mode 100644 index 0000000..fd3b2f1 --- /dev/null +++ b/packages/meshbay-node/tests/test_admin_challenge_bounds.py @@ -0,0 +1,135 @@ +""" +What a connection may leave waiting for a signature (docs/MESHBAY_DESIGN.md §13.5b). + +Anyone authenticated can ask for an admin challenge — the signature is checked +later — so a member who never answers must not make the node keep every request. +Measured before the bound: 200 `root_add` of 1 MiB each from a plain member held +200 pending operations and ~400 MiB for the life of the connection. +""" + +import struct +import time + +import msgpack +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey +from meshbay_node.transport.webrtc.admin import MAX_ADMIN_OP_BYTES, MAX_PENDING_ADMIN_OPS +from meshbay_node.transport.webrtc_server import WebRTCPeerSession + +GROUP = "g" * 32 + + +class _Channel: + readyState = "open" + + def __init__(self): + self.sent = [] + + def send(self, data: bytes) -> None: + (n,) = struct.unpack(">I", data[:4]) + self.sent.append(msgpack.unpackb(data[4:4 + n], raw=False)) + + +class _PC: + connectionState = "connected" + iceConnectionState = "connected" + remoteDescription = None + localDescription = None + sctp = None + + +def _member_session(): + """An authenticated member — not the operator — on a node that has one.""" + ctx = {"sk_node": Ed25519PrivateKey.from_private_bytes(b"\x01" * 32), + "groups": {GROUP: {}}, "has_admin_authority": True} + s = WebRTCPeerSession(_PC(), ctx, peer_id="peer") + s._channel = _Channel() + s._audit = lambda *a, **k: None + s._user_id, s._group_id = "member-1", GROUP + return s + + +def _root_add(s, path: str) -> dict: + s._dispatch_message({"type": "root_add", "group_id": GROUP, "path": path}) + return s._channel.sent[-1] + + +def test_a_member_cannot_pile_up_challenges(): + s = _member_session() + for i in range(MAX_PENDING_ADMIN_OPS): + assert _root_add(s, f"/srv/{i}")["type"] == "admin_challenge" + refused = _root_add(s, "/srv/one-too-many") + assert refused["type"] == "error" and refused["code"] == "too_many_pending" + assert len(s._admin_ops) == MAX_PENDING_ADMIN_OPS + + +def test_an_oversized_request_is_not_kept(): + s = _member_session() + refused = _root_add(s, "x" * (MAX_ADMIN_OP_BYTES + 1)) + assert refused["type"] == "error" and refused["code"] == "too_large" + assert s._admin_ops == {} + + +def test_an_expired_challenge_frees_its_place(): + s = _member_session() + for i in range(MAX_PENDING_ADMIN_OPS): + _root_add(s, f"/srv/{i}") + for pending in s._admin_ops.values(): + pending["ts"] -= 10_000 + assert _root_add(s, "/srv/after-expiry")["type"] == "admin_challenge" + assert len(s._admin_ops) == 1 + + +def test_answering_a_challenge_frees_its_place(): + s = _member_session() + for i in range(MAX_PENDING_ADMIN_OPS): + _root_add(s, f"/srv/{i}") + op_id = next(iter(s._admin_ops)) + s._dispatch_message({"type": "admin_response", "op_id": op_id, "signature": "!!"}) + assert len(s._admin_ops) == MAX_PENDING_ADMIN_OPS - 1 + assert _root_add(s, "/srv/next")["type"] == "admin_challenge" + assert all(time.time() - p["ts"] < 5 for p in s._admin_ops.values()) + + +# ── What a challenge covers (docs/MESHBAY_DESIGN.md §5.4) ──────────────────── +# +# The signature covers the subject and nothing else of a request, so every value +# the executor acts on has to be in it. + +def test_root_add_signs_whether_members_may_write(): + from meshbay_common.adminop import root_add_subject + s = _member_session() + s._dispatch_message({"type": "root_add", "group_id": GROUP, "path": "/srv/drop", + "name": "Drop", "writable": True, "removable": False}) + challenge = s._channel.sent[-1] + assert challenge["subject"] == root_add_subject("/srv/drop", "Drop", "generic", + True, False) + assert challenge["subject"] != root_add_subject("/srv/drop", "Drop", "generic", + False, False) + + +def test_group_attach_signs_the_directory_it_exposes(): + from meshbay_common.adminop import group_attach_subject + s = _member_session() + s._dispatch_message({"type": "group_attach", "name": "photos", + "shared_dir": "/home/me/Photos"}) + assert s._channel.sent[-1]["subject"] == group_attach_subject( + "photos", "/home/me/Photos", True) + + +def test_invite_create_signs_the_name_it_records(): + from meshbay_common.adminop import invite_create_subject + s = _member_session() + s._ctx["roster"] = object() # only its presence is checked before the challenge + s._dispatch_message({"type": "invite_create", "group_id": GROUP, + "user_id": "u-1", "username": "alice"}) + assert s._channel.sent[-1]["subject"] == invite_create_subject("u-1", "alice") + + +def test_tmdb_config_signs_the_token_without_writing_it(): + from meshbay_common.adminop import tmdb_config_subject + s = _member_session() + s._dispatch_message({"type": "tmdb_config", "token": "secret-token", + "language": "fr-FR"}) + subject = s._channel.sent[-1]["subject"] + assert subject == tmdb_config_subject("secret-token", "fr-FR") + assert "secret-token" not in subject diff --git a/packages/meshbay-node/tests/test_content_blocklist.py b/packages/meshbay-node/tests/test_content_blocklist.py new file mode 100644 index 0000000..0e07787 --- /dev/null +++ b/packages/meshbay-node/tests/test_content_blocklist.py @@ -0,0 +1,238 @@ +""" +The hub's content blocklist, applied by a node in its public groups +(docs/MESHBAY_DESIGN.md §7.5, `meshbay_node.blocklist`). + +A blocked file leaves the index members are sent and is refused if asked for — +in a public group, and nowhere else: a private group's content never reaches the +hub, so nothing there can have been blocked. +""" + +import struct + +import msgpack +import pytest +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey +from meshbay_common.crypto import generate_gek +from meshbay_common.groupbox import PURPOSE_INDEX, unseal +from meshbay_common.protocol import IndexEntry +from meshbay_node.blocklist import ContentBlocklist +from meshbay_node.hub_client import HubClient +from meshbay_node.indexer import GroupIndex +from meshbay_node.transport.webrtc_server import WebRTCPeerSession +from meshbay_node.transport.wire import index_delta_message, index_sync_message + +GROUP = "g" * 32 +BLOCKED = "b" * 64 +KEPT = "c" * 64 +THUMB = "d" * 64 + + +# ── The list ────────────────────────────────────────────────────────────────── + +def test_the_list_survives_a_restart(tmp_path): + path = tmp_path / "blocklist.json" + assert ContentBlocklist(path).replace([BLOCKED]) + assert BLOCKED in ContentBlocklist(path) + + +def test_a_full_sync_replaces_and_says_whether_anything_changed(tmp_path): + bl = ContentBlocklist(tmp_path / "blocklist.json") + assert bl.replace([BLOCKED, KEPT]) + assert not bl.replace([KEPT, BLOCKED]) + assert bl.replace([KEPT]) + assert BLOCKED not in bl and KEPT in bl + + +def test_a_pushed_change_applies_and_ignores_what_is_not_a_hash(tmp_path): + bl = ContentBlocklist(tmp_path / "blocklist.json") + assert bl.apply(add=[BLOCKED, "../etc", 7, "B" * 64]) + assert list(bl) == [BLOCKED] + assert not bl.apply(add=[BLOCKED]) + assert bl.apply(remove=[BLOCKED]) + assert len(bl) == 0 + + +def test_a_damaged_file_does_not_stop_the_node(tmp_path): + path = tmp_path / "blocklist.json" + path.write_text("{not json", encoding="utf-8") + assert len(ContentBlocklist(path)) == 0 + + +# ── What members are sent ───────────────────────────────────────────────────── + +def _index(gek) -> GroupIndex: + idx = GroupIndex(group_id=GROUP, sk_node=Ed25519PrivateKey.generate(), gek=gek) + idx.add_entry(IndexEntry(id=BLOCKED, name="a.jpg", path="r", size=1, type="image", + added_at=0, thumb_hash=THUMB)) + idx.add_entry(IndexEntry(id=KEPT, name="b.jpg", path="r", size=1, type="image", + added_at=0)) + return idx + + +def test_a_hidden_entry_is_not_in_the_index_or_its_deltas(): + gek = generate_gek() + idx = _index(gek) + msg = index_sync_message(idx, None, {BLOCKED}) + ids = [e["id"] for e in unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, msg)["entries"]] + assert ids == [KEPT] + + delta = idx.diff(GroupIndex(group_id=GROUP, sk_node=idx.sk_node, gek=gek, version=0)) + payload = unseal(gek, PURPOSE_INDEX, "index_delta", GROUP, + index_delta_message(idx, delta, None, {BLOCKED})) + assert [e["id"] for e in payload["additions"]] == [KEPT] + + +# ── What a session serves ──────────────────────────────────────────────────── + +class _Channel: + readyState = "open" + + def __init__(self): + self.sent = [] + + def send(self, data: bytes) -> None: + (n,) = struct.unpack(">I", data[:4]) + self.sent.append(msgpack.unpackb(data[4:4 + n], raw=False)) + + +class _PC: + connectionState = "connected" + iceConnectionState = "connected" + remoteDescription = None + localDescription = None + sctp = None + + +def _session(visibility: str, tmp_path): + gek = generate_gek() + bl = ContentBlocklist(tmp_path / "blocklist.json") + bl.replace([BLOCKED]) + ctx = {"sk_node": Ed25519PrivateKey.from_private_bytes(b"\x01" * 32), + "groups": {GROUP: {"visibility": visibility, "index": _index(gek), + "gek": gek, "roots": None}}, + "blocklist": bl} + s = WebRTCPeerSession(_PC(), ctx, peer_id="peer") + s._channel = _Channel() + s._audit = lambda *a, **k: None + s._user_id, s._group_id = "member-1", GROUP + return s, gek + + +@pytest.mark.asyncio +async def test_a_public_group_refuses_a_blocked_file_and_its_thumbnail(tmp_path): + s, _ = _session("public", tmp_path) + for file_id in (BLOCKED, THUMB): + await s._do_file_request({"file_id": file_id, "chunk_index": 0}) + assert s._channel.sent[-1]["code"] == "content_blocked" + + +@pytest.mark.asyncio +async def test_a_public_group_does_not_list_a_blocked_file(tmp_path): + s, gek = _session("public", tmp_path) + s._do_index_sync() + payload = unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, s._channel.sent[-1]) + assert [e["id"] for e in payload["entries"]] == [KEPT] + + +@pytest.mark.asyncio +async def test_a_private_group_is_untouched(tmp_path): + s, gek = _session("private", tmp_path) + s._do_index_sync() + payload = unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, s._channel.sent[-1]) + assert sorted(e["id"] for e in payload["entries"]) == [BLOCKED, KEPT] + assert not s._refuse_blocked(BLOCKED) + + +@pytest.mark.asyncio +async def test_streaming_subtitles_and_transcoding_are_refused_too(tmp_path): + s, _ = _session("public", tmp_path) + await s._stream_video_inner({"file_id": BLOCKED}) + assert s._channel.sent[-1]["code"] == "content_blocked" + await s._do_subtitle_request({"file_id": BLOCKED, "track": 0}) + assert s._channel.sent[-1]["code"] == "content_blocked" + await s._do_audio_transcode_request({"file_id": BLOCKED}) + assert s._channel.sent[-1]["code"] == "content_blocked" + + +# ── How the node gets the list ─────────────────────────────────────────────── + +@pytest.mark.asyncio +async def test_the_whole_list_is_fetched_page_by_page(): + pages = {"": {"hashes": [BLOCKED], "next": BLOCKED}, + BLOCKED: {"hashes": [KEPT], "next": None}} + + class _Resp: + def __init__(self, body): + self._body = body + + def raise_for_status(self): + pass + + def json(self): + return self._body + + class _Http: + async def get(self, path, params, headers): + assert path == "/v1/blocklist" + return _Resp(pages[params["after"]]) + + class _Session: + auth_headers = {} + + hub = HubClient.__new__(HubClient) + hub._session, hub._http = _Session(), _Http() + + async def fresh(): + return None + hub.ensure_fresh_token = fresh + assert await hub.fetch_blocklist() == {BLOCKED, KEPT} + + +# ── A change reaches members already connected ─────────────────────────────── + +def _daemon(tmp_path, visibility: str): + from unittest.mock import MagicMock + + from meshbay_node.config import Config, GroupConfig, HubConfig, KeystoreConfig, NodeConfig + from meshbay_node.daemon import NodeDaemon + + shared = tmp_path / "shared" + shared.mkdir() + config = Config( + hub=HubConfig(url="http://localhost:9999", username="t"), + node=NodeConfig(), + groups=[GroupConfig(id=GROUP, name="g", shared_dir=str(shared), + visibility=visibility)], + keystore=KeystoreConfig(path=tmp_path / "keystore.enc"), + data_dir=tmp_path / "data", + ) + daemon = NodeDaemon(config) + gek = generate_gek() + indexer = MagicMock() + indexer.index, indexer.roots = _index(gek), None + daemon._indexers = [indexer] + session = MagicMock() + session._group_id = GROUP + daemon._webrtc = MagicMock() + daemon._webrtc._sessions = {"p": session} + return daemon, session, gek + + +def test_a_pushed_block_resends_the_index_without_the_file(tmp_path): + daemon, session, gek = _daemon(tmp_path, "public") + daemon._on_blocklist_update([BLOCKED], []) + msg = session._send.call_args[0][0] + ids = [e["id"] for e in unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, msg)["entries"]] + assert ids == [KEPT] + + daemon._on_blocklist_update([], [BLOCKED]) + msg = session._send.call_args[0][0] + ids = sorted(e["id"] for e in unseal(gek, PURPOSE_INDEX, "index_sync", GROUP, msg)["entries"]) + assert ids == [BLOCKED, KEPT] + + +def test_a_node_with_only_private_groups_ignores_the_list(tmp_path): + daemon, session, _ = _daemon(tmp_path, "private") + daemon._on_blocklist_update([BLOCKED], []) + session._send.assert_not_called() + assert BLOCKED not in daemon._blocklist diff --git a/packages/meshbay-node/tests/test_daemon.py b/packages/meshbay-node/tests/test_daemon.py index 8c4da2d..acaafac 100644 --- a/packages/meshbay-node/tests/test_daemon.py +++ b/packages/meshbay-node/tests/test_daemon.py @@ -2,7 +2,7 @@ Integration test: Node daemon wires all components correctly. Phase 11 — verifies that NodeDaemon creates chat stores, WebRTC transport, -index push on change, swarm registration, and shuts down cleanly. +index push on change, and shuts down cleanly. Hub interaction is mocked. """ @@ -241,7 +241,6 @@ async def test_daemon_index_change_pushes_to_peers(tmp_path, shared_dir, gek, hu daemon = NodeDaemon(config) daemon._broadcast_coalesce_secs = 0.01 # real value would make this test wait 0.5s daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=2) daemon._state["endpoint_hint"] = "node123" sk_node = Ed25519PrivateKey.generate() @@ -268,47 +267,10 @@ async def test_daemon_index_change_pushes_to_peers(tmp_path, shared_dir, gek, hu payload = unseal(gek, PURPOSE_INDEX, "index_sync", "a" * 32, msg) assert len(payload["entries"]) == indexer.index.count - # Finding H7: this group is private, so its content hashes must NOT be - # registered with the hub. The test previously asserted the opposite — - # publishing a fingerprint of every private file was treated as expected - # behaviour. Index push to members is unaffected (asserted above). + # Finding H7: a change to the index tells the hub nothing — no content hash + # of any group reaches it. Index push to members is unaffected (above). await asyncio.sleep(0.1) - daemon._hub.register_swarm.assert_not_called() - -@pytest.mark.asyncio -async def test_daemon_index_change_registers_swarm_for_public_group( - tmp_path, shared_dir, gek, hub_pk_pem): - """Public groups still register content hashes with the hub swarm (H7).""" - config = Config( - hub=HubConfig(url="http://localhost:9999", username="testuser"), - node=NodeConfig(quic_port=_free_port(), ui_port=_free_port()), - groups=[GroupConfig( - id="a" * 32, - name="public-group", - shared_dir=str(shared_dir), - visibility="public", - quic_port=29010, - )], - keystore=KeystoreConfig(path=tmp_path / "keystore.enc"), - data_dir=tmp_path / "data", - ) - daemon = NodeDaemon(config) - daemon._broadcast_coalesce_secs = 0.01 - daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=2) - daemon._state["endpoint_hint"] = "node123" - - indexer = DirectoryIndexer( - roots=one_root(shared_dir), group_id="a" * 32, - sk_node=Ed25519PrivateKey.generate(), gek=gek) - await indexer.initial_scan() - - await daemon._on_index_change(indexer) - - await asyncio.sleep(0.1) - daemon._hub.register_swarm.assert_called_once() - assert len(daemon._hub.register_swarm.call_args[0][0]) == indexer.index.count - + assert daemon._hub.mock_calls == [] @pytest.mark.asyncio async def test_daemon_index_change_skips_other_group_peers( @@ -325,7 +287,6 @@ async def test_daemon_index_change_skips_other_group_peers( daemon = NodeDaemon(config) daemon._broadcast_coalesce_secs = 0.01 daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=0) daemon._state["endpoint_hint"] = "node123" sk_node = Ed25519PrivateKey.generate() @@ -370,7 +331,6 @@ def _new_daemon_for_group(tmp_path, shared_dir, gek, group_id="a" * 32, daemon = NodeDaemon(config) daemon._broadcast_coalesce_secs = 0.01 daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=0) daemon._state["endpoint_hint"] = "node123" return daemon @@ -532,29 +492,3 @@ async def test_a_burst_of_changes_produces_one_broadcast(tmp_path, shared_dir, g await asyncio.sleep(0.05) session._send.assert_called_once() - - -@pytest.mark.asyncio -async def test_swarm_registration_only_sends_new_hashes_after_the_first( - tmp_path, shared_dir, gek): - daemon = _new_daemon_for_group(tmp_path, shared_dir, gek, visibility="public") - indexer = DirectoryIndexer( - roots=one_root(shared_dir), group_id="a" * 32, - sk_node=Ed25519PrivateKey.generate(), gek=gek) - await indexer.initial_scan() - total_files = indexer.index.count - - await daemon._on_index_change(indexer) - await asyncio.sleep(0.05) - assert len(daemon._hub.register_swarm.call_args_list[0].args[0]) == total_files - - from meshbay_common.protocol import IndexEntry - indexer.index.add_entry(IndexEntry(id="new-file-id", name="new.mp4", - path="shared", size=10, type="video", - added_at=0)) - await daemon._on_index_change(indexer) - await asyncio.sleep(0.05) - - assert daemon._hub.register_swarm.call_count == 2 - assert daemon._hub.register_swarm.call_args_list[1].args[0] == ["new-file-id"], \ - "only the newly added hash must be (re-)registered, not the whole library" diff --git a/packages/meshbay-node/tests/test_indexer.py b/packages/meshbay-node/tests/test_indexer.py index e97e2ef..c60600f 100644 --- a/packages/meshbay-node/tests/test_indexer.py +++ b/packages/meshbay-node/tests/test_indexer.py @@ -42,58 +42,6 @@ def shared_dir(tmp_path): # ── GroupIndex tests ────────────────────────────────────────────────────────── -def test_group_index_serialize_deserialize_private(sk_node, gek, shared_dir): - idx = GroupIndex(group_id="grp-001", sk_node=sk_node, gek=gek) - from meshbay_common.protocol import IndexEntry - idx.add_entry(IndexEntry( - id="abc123", name="video.mkv", path="", size=1024, - type="video", added_at=int(time.time()), duration=120)) - - wire = idx.serialize() - recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) - - assert recovered.group_id == "grp-001" - assert recovered.count == 1 - assert recovered.entries[0].name == "video.mkv" - assert recovered.entries[0].type == "video" - - -def test_group_index_serialize_deserialize_public(sk_node): - idx = GroupIndex(group_id="pub-001", sk_node=sk_node, gek=None) - from meshbay_common.protocol import IndexEntry - idx.add_entry(IndexEntry( - id="xyz789", name="readme.txt", path="", size=42, - type="document", added_at=int(time.time()))) - - wire = idx.serialize() - recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=None) - assert recovered.count == 1 - assert recovered.entries[0].id == "xyz789" - - -def test_group_index_wrong_gek_rejected(sk_node, gek): - idx = GroupIndex(group_id="grp-002", sk_node=sk_node, gek=gek) - from meshbay_common.protocol import IndexEntry - idx.add_entry(IndexEntry(id="a", name="f.mp3", path="", size=1, - type="audio", added_at=0)) - wire = idx.serialize() - - wrong_gek = generate_gek() - with pytest.raises(Exception): # InvalidTag from AEAD - GroupIndex.deserialize(wire, sk_node=sk_node, gek=wrong_gek) - - -def test_group_index_tampered_rejected(sk_node, gek): - idx = GroupIndex(group_id="grp-003", sk_node=sk_node, gek=gek) - from meshbay_common.protocol import IndexEntry - idx.add_entry(IndexEntry(id="b", name="f.mp4", path="", size=1, - type="video", added_at=0)) - wire = bytearray(idx.serialize()) - wire[-5] ^= 0xFF # flip bytes at the end - with pytest.raises(Exception): - GroupIndex.deserialize(bytes(wire), sk_node=sk_node, gek=gek) - - def test_group_index_diff(sk_node, gek): from meshbay_common.protocol import IndexEntry v1 = GroupIndex(group_id="g", sk_node=sk_node, gek=gek, version=1) @@ -229,16 +177,6 @@ async def test_on_change_callback(shared_dir, sk_node, gek): assert len(changes) >= 1, "on_change should have been called" -@pytest.mark.asyncio -async def test_index_roundtrip_after_scan(shared_dir, sk_node, gek): - indexer = DirectoryIndexer(roots=one_root(shared_dir), group_id="g", sk_node=sk_node, gek=gek) - await indexer.initial_scan() - - wire = indexer.index.serialize() - recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) - assert recovered.count == indexer.index.count - - # ── Cache-aware scanning ─────────────────────────────────────────────────────── @pytest.fixture @@ -796,35 +734,13 @@ def test_partial_hash_is_deterministic(tmp_path): assert e1.id == e2.id -def test_group_index_roundtrip_preserves_hash_version(sk_node, gek): - from meshbay_common.protocol import IndexEntry - idx = GroupIndex(group_id="hv-test", sk_node=sk_node, gek=gek) - idx.add_entry(IndexEntry( - id="aaa", name="small.mp4", path="root", size=1024, - type="video", added_at=100, hash_version=1)) - idx.add_entry(IndexEntry( - id="bbb", name="big.mkv", path="root", size=50_000_000, - type="video", added_at=200, hash_version=2)) - - wire = idx.serialize() - recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) - - by_id = {e.id: e for e in recovered.entries} - assert by_id["aaa"].hash_version == 1 - assert by_id["bbb"].hash_version == 2 - - -def test_deserialize_without_hash_version_defaults_to_1(sk_node, gek): - """Entries serialized by old code (no hash_version field) must deserialize - as hash_version=1.""" +def test_an_entry_without_hash_version_defaults_to_1(): + """An entry built without the field (written before it existed) reads as a + full-read hash.""" from meshbay_common.protocol import IndexEntry - idx = GroupIndex(group_id="compat", sk_node=sk_node, gek=gek) - idx.add_entry(IndexEntry( - id="old", name="f.mp4", path="root", size=1024, - type="video", added_at=100)) - wire = idx.serialize() - recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) - assert recovered.entries[0].hash_version == 1 + e = IndexEntry(id="old", name="f.mp4", path="root", size=1024, + type="video", added_at=100) + assert e.hash_version == 1 def test_index_entry_wire_includes_hash_version(): diff --git a/packages/meshbay-node/tests/test_security_regressions.py b/packages/meshbay-node/tests/test_security_regressions.py index c0c2c2c..43e88c2 100644 --- a/packages/meshbay-node/tests/test_security_regressions.py +++ b/packages/meshbay-node/tests/test_security_regressions.py @@ -602,19 +602,18 @@ def test_denylist_persists_and_honours_groups(tmp_path): assert not reloaded.is_denied("someone", "other", "g-allowed") -def test_swarm_registration_skips_private_groups(): +def test_the_node_registers_no_content_hash_with_the_hub(): """ - H7: the daemon registered content hashes for every group, private included, - handing the hub a fingerprint of every private file. The bug was masked by a - mis-mounted route, so fixing the route without this filter would have turned a - dormant leak into a live one. + H7: the daemon once registered content hashes for every group, private + included, handing the hub a fingerprint of every private file. The swarm that + received them is gone, so no group's hashes are sent to the hub at all. """ - source = daemon_source() - assert 'visibility' in source and '_register_swarm' in source - # Both registration sites must gate on public visibility. - for marker in ['gctx.get("visibility") == "public"', - 'group_cfg.visibility == "public"']: - assert marker in source, f"swarm registration not gated: {marker}" + from pathlib import Path + + import meshbay_node.hub_client as hub_client + for source in (daemon_source(), Path(hub_client.__file__).read_text(encoding="utf-8")): + assert "/v1/swarm" not in source + assert "register_swarm" not in source def test_keystore_argon2_is_production_strength(): diff --git a/packages/meshbay-node/tests/test_tmdb_config_policy.py b/packages/meshbay-node/tests/test_tmdb_config_policy.py index 6a51eb0..c671ac8 100644 --- a/packages/meshbay-node/tests/test_tmdb_config_policy.py +++ b/packages/meshbay-node/tests/test_tmdb_config_policy.py @@ -22,7 +22,7 @@ from pathlib import Path import pytest from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey -from meshbay_common.adminop import OP_TMDB_CONFIG +from meshbay_common.adminop import OP_TMDB_CONFIG, tmdb_config_subject from meshbay_node.indexer.group_index import GroupIndex from meshbay_node.roster import Roster from meshbay_node.transport.webrtc_server import WebRTCPeerSession @@ -134,7 +134,9 @@ async def test_subject_reflects_whether_a_token_was_supplied(tmp_path): session._do_tmdb_config({"token": "x"}) _, subject, _, _ = issued[0] - assert subject == "custom_token=yes,language=default" + # The token is bound by its digest and never written into the subject. + assert subject == tmdb_config_subject("x", None) + assert "sha256:" in subject and tmdb_config_subject("y", None) != subject async def test_subject_says_no_custom_token_when_none_given(tmp_path): @@ -146,7 +148,10 @@ async def test_subject_says_no_custom_token_when_none_given(tmp_path): session._do_tmdb_config({}) _, subject, _, _ = issued[0] - assert subject == "custom_token=no,language=default" + # Nothing given means both unchanged — distinct from clearing either. + assert subject == tmdb_config_subject(None, None) + assert subject != tmdb_config_subject("", None) + assert subject != tmdb_config_subject(None, "") async def test_subject_reflects_a_configured_language(tmp_path): @@ -158,7 +163,7 @@ async def test_subject_reflects_a_configured_language(tmp_path): session._do_tmdb_config({"language": "fr-FR"}) _, subject, payload, _ = issued[0] - assert subject == "custom_token=no,language=fr-FR" + assert subject == tmdb_config_subject(None, "fr-FR") assert payload["language"] == "fr-FR" diff --git a/packages/meshbay-node/tests/test_webrtc_transport.py b/packages/meshbay-node/tests/test_webrtc_transport.py index 990b1da..7a6f517 100644 --- a/packages/meshbay-node/tests/test_webrtc_transport.py +++ b/packages/meshbay-node/tests/test_webrtc_transport.py @@ -34,13 +34,12 @@ from meshbay_common.adminop import ( OP_INVITE_CREATE, OP_INVITE_LINK_CREATE, admin_transcript, + invite_create_subject, ) from meshbay_common.crypto import ( generate_gek, pk_to_b64, - unwrap_gek, unwrap_gek_aes, - wrap_gek, wrap_gek_aes, ) from meshbay_common.groupbox import PURPOSE_ACK, PURPOSE_INDEX, unseal @@ -1337,7 +1336,8 @@ async def test_invite_then_join_delivers_the_gek(sk_node, sk_hub, gek, shared_di challenge_msg = await asyncio.wait_for(q_admin.get(), timeout=5.0) assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE assert challenge_msg["op"] == OP_INVITE_CREATE - assert challenge_msg["subject"] == "user-002" + # The name the invitation records is signed with the account it is for. + assert challenge_msg["subject"] == invite_create_subject("user-002", "bob") ch_admin.send(_pack({ "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION, @@ -1570,7 +1570,7 @@ async def test_gek_bundle_fetch_during_handshake(sk_node, sk_hub, gek, shared_di await bundle_store.open() # Pre-populate a bundle for user-001 in group "g" - bundle = wrap_gek(gek, pk_x_raw) + bundle = wrap_gek_aes(gek, pk_x_raw) await bundle_store.store("g", "user-001", bundle["pk_eph_b64"], bundle["nonce_b64"], bundle["wrapped_b64"]) @@ -1631,7 +1631,7 @@ async def test_gek_bundle_fetch_during_handshake(sk_node, sk_hub, gek, shared_di assert bundle_resp["found"] is True # Step 3: Unwrap GEK and compute HMAC proof - recovered_gek = unwrap_gek(bundle_resp, sk_x_raw, pk_x_raw) + recovered_gek = unwrap_gek_aes(bundle_resp, sk_x_raw, pk_x_raw) assert recovered_gek == gek nonce_s = base64.b64decode(msg["nonce"]) diff --git a/packages/meshbay-node/tests/transfer_probe.py b/packages/meshbay-node/tests/transfer_probe.py index 7aa4902..6666960 100755 --- a/packages/meshbay-node/tests/transfer_probe.py +++ b/packages/meshbay-node/tests/transfer_probe.py @@ -274,7 +274,14 @@ async def operator_checks(client, ack, group, node_id) -> int: import re as _re failures = 0 - member_cap = int((ack.get("transfer_limits") or {}).get("download") or 0) + # This member's own cap, as the node states it on every `transfer_state`. + probe_tr, probe_state = await _open_transfer(client) + member_cap = int(probe_state.get("cap") or 0) + client.send({"type": "transfer_close", "v": "0.1", "tr": probe_tr, "reason": "done"}) + while True: # closed before anything below counts what is in use + reply = await client.recv_type("transfer_state", timeout=15) + if reply.get("tr") == probe_tr and reply.get("state") == "closed": + break # The queue has to be held by the *node* cap, not by this member's own. # @@ -446,16 +453,6 @@ async def probe(args) -> int: ack = await alice.connect(http, group["id"], node_id) - limits = ack.get("transfer_limits") - if limits is None: - print("This node does not hand out transfer slots — it predates " - "them, or the handshake ack lost the field. Nothing below " - "can be measured.") - await alice.close() - return 1 - cap = int(limits.get("download") or 0) - print(f"node reports this member may run {cap} download(s) at once\n") - if args.operator: return await operator_checks(alice, ack, group, node_id) @@ -472,9 +469,17 @@ async def probe(args) -> int: # held. So the first reply is read for what the node says is already # in use, and the run stops rather than measuring against a moving # floor. - want = args.want or (cap + 2) opened = [await _open_transfer(alice)] first = opened[0][1] + # This member's cap, as the node states it on every `transfer_state`. + cap = int(first.get("cap") or 0) + if not cap: + print("This node does not state a transfer cap — nothing below can " + "be measured.") + await alice.close() + return 1 + print(f"node reports this member may run {cap} download(s) at once\n") + want = args.want or (cap + 2) if first.get("state") != "granted" or first.get("used", 1) != 1: print(f"this node is not idle: it reports {first.get('used')} of " f"{first.get('cap')} slots already used by this member, and " |