diff options
Diffstat (limited to 'packages')
96 files changed, 2461 insertions, 2153 deletions
diff --git a/packages/meshbay-client/src/main.js b/packages/meshbay-client/src/main.js index 9116d1a..f8bb3f4 100644 --- a/packages/meshbay-client/src/main.js +++ b/packages/meshbay-client/src/main.js @@ -850,16 +850,6 @@ function registerBridge() { return out; }); - ipcMain.handle('hub:probe', async (_e, url) => { - const target = String(url || config.hubBase || '').replace(/\/+$/, ''); - if (!target) return null; - try { - const r = await fetch(`${target}/v1/hub/version`, - { signal: AbortSignal.timeout(10000) }); - return r.ok ? await r.json() : null; - } catch { return null; } - }); - // Hide, never close: `window-all-closed` quits the app, and closing here would // make "minimise to tray" mean "exit". A hidden window keeps the session, the // transfers and the node connection exactly as they were. diff --git a/packages/meshbay-common/pyproject.toml b/packages/meshbay-common/pyproject.toml index a454521..9f43cf2 100644 --- a/packages/meshbay-common/pyproject.toml +++ b/packages/meshbay-common/pyproject.toml @@ -12,7 +12,6 @@ dependencies = [ "PyJWT>=2.9", "blake3>=1.0", "msgpack>=1.1", - "zstandard>=0.23", ] [project.optional-dependencies] diff --git a/packages/meshbay-common/src/meshbay_common/__init__.py b/packages/meshbay-common/src/meshbay_common/__init__.py index 842c52d..fbfdd4d 100644 --- a/packages/meshbay-common/src/meshbay_common/__init__.py +++ b/packages/meshbay-common/src/meshbay_common/__init__.py @@ -249,5 +249,16 @@ __version__ = "0.16.0" # API. A pre-4.0 client presents the session token and is refused at the # handshake — there is no compatibility branch, because leaving one would keep # the disclosure reachable on every node. So the floor moves with it. -MNP_VERSION = "4.0" +# +# 5.0 (2026-09-28) is a MAJOR — four signed operations now sign everything they +# do. `root_add` signed its path and not whether every member may write there; +# `group_attach` signed a group's name and not the directory it exposes; +# `invite_create` did not sign the name it records; `tmdb_config` did not bind +# the token (it now names its SHA-256). Each subject is canonical JSON of every +# value the node acts on (`adminop.structured_subject`). A 4.x client refuses to +# sign the new subjects, and a 5.0 client the old ones — so those four fail, with +# a refusal, across the break. The break is confined to them, so the floor stays +# at 4.0: everything else a 4.x peer does still works, and nothing is left +# unsigned on either side — no node accepts the old subjects. +MNP_VERSION = "5.0" MHP_VERSION = "0.1" diff --git a/packages/meshbay-common/src/meshbay_common/adminop.py b/packages/meshbay-common/src/meshbay_common/adminop.py index c679718..4379519 100644 --- a/packages/meshbay-common/src/meshbay_common/adminop.py +++ b/packages/meshbay-common/src/meshbay_common/adminop.py @@ -30,6 +30,9 @@ fields it received, the node from the state it stored. They are compared by producing the same bytes, never by trusting a value off the wire. """ +import hashlib +import json + ADMIN_TRANSCRIPT_PREFIX = b"meshbay:admin:v1" # Operations that require node-operator authority. @@ -135,6 +138,46 @@ OP_GROUP_DETACH = "group_detach" ADMIN_CHALLENGE_TTL = 120 # seconds +def structured_subject(fields: dict) -> str: + """ + The subject of an operation whose effect is more than one value. + + Every value the executor acts on is in here, because the signature covers the + subject and nothing else of the request: a root's path alone left whether + every member may write there unsigned. Canonical JSON — sorted keys, no + whitespace, UTF-8 — so `null`, `""` and a value stay distinct, and the + browser's `adminSubject` (static/crypto.js) produces the same bytes. + """ + return json.dumps(fields, sort_keys=True, separators=(",", ":"), ensure_ascii=False) + + +def secret_digest(value: str | None) -> str | None: + """A secret named in a subject without being written there: `None` (leave it + unchanged) and `""` (clear it) as themselves, anything else as its SHA-256.""" + if not value: + return value + return "sha256:" + hashlib.sha256(value.encode()).hexdigest() + + +def root_add_subject(path: str, name: str, kind: str, writable: bool, + removable: bool) -> str: + return structured_subject({"path": path, "name": name, "kind": kind, + "writable": writable, "removable": removable}) + + +def group_attach_subject(name: str, shared_dir: str, writable: bool) -> str: + return structured_subject({"name": name, "shared_dir": shared_dir, + "writable": writable}) + + +def invite_create_subject(user_id: str, username: str) -> str: + return structured_subject({"user_id": user_id, "username": username}) + + +def tmdb_config_subject(token: str | None, language: str | None) -> str: + return structured_subject({"token": secret_digest(token), "language": language}) + + def admin_transcript( op: str, node_pk_b64: str, diff --git a/packages/meshbay-common/src/meshbay_common/crypto.py b/packages/meshbay-common/src/meshbay_common/crypto.py index e537500..8fb80f8 100644 --- a/packages/meshbay-common/src/meshbay_common/crypto.py +++ b/packages/meshbay-common/src/meshbay_common/crypto.py @@ -41,25 +41,6 @@ def generate_gek() -> bytes: """Generate a fresh 256-bit Group Encryption Key.""" return ChaCha20Poly1305.generate_key() -def chunk_key(gek: bytes, file_hash: bytes, chunk_index: int) -> bytes: - """Derive a per-chunk encryption key from the GEK (deterministic).""" - return HKDF( - algorithm=hashes.SHA256(), - length=32, - salt=None, - info=b"file:" + file_hash + b":chunk:" + chunk_index.to_bytes(4, "big"), - ).derive(gek) - -def encrypt_chunk(key: bytes, plaintext: bytes) -> tuple[bytes, bytes]: - """Encrypt plaintext with ChaCha20-Poly1305. Returns (nonce, ciphertext).""" - nonce = os.urandom(12) - ct = ChaCha20Poly1305(key).encrypt(nonce, plaintext, None) - return nonce, ct - -def decrypt_chunk(key: bytes, nonce: bytes, ciphertext: bytes) -> bytes: - """Decrypt ciphertext. Raises InvalidTag on authentication failure.""" - return ChaCha20Poly1305(key).decrypt(nonce, ciphertext, None) - def file_hash(path_or_bytes) -> bytes: """Compute blake3 hash of a file (bytes or path-like).""" if isinstance(path_or_bytes, (str, bytes)) and not isinstance(path_or_bytes, bytes): @@ -73,58 +54,6 @@ def file_hash(path_or_bytes) -> bytes: # ── GEK wrapping (ECIES-like) ───────────────────────────────────────────────── -GEK_WRAP_INFO = b"meshbay:gek_wrap:v1" - -def wrap_gek(gek: bytes, pk_recipient: bytes) -> dict: - """ - Wrap a GEK for a recipient using ephemeral X25519 + HKDF + ChaCha20-Poly1305. - - Protocol: - 1. Generate ephemeral (sk_eph, pk_eph) - 2. shared = X25519(sk_eph, pk_recipient) - 3. wrap_key = HKDF(shared, salt=pk_eph, info=GEK_WRAP_INFO) - 4. wrapped = ChaCha20-Poly1305(wrap_key).encrypt(nonce, gek, aad=pk_recipient) - - The hub stores {pk_eph, nonce, wrapped} — opaque, cannot decrypt. - """ - sk_eph = X25519PrivateKey.generate() - pk_eph_raw = pk_to_raw(sk_eph.public_key()) - - shared = sk_eph.exchange(X25519PublicKey.from_public_bytes(pk_recipient)) - wrap_key = HKDF( - algorithm=hashes.SHA256(), length=32, - salt=pk_eph_raw, info=GEK_WRAP_INFO, - ).derive(shared) - - nonce = os.urandom(12) - wrapped = ChaCha20Poly1305(wrap_key).encrypt(nonce, gek, pk_recipient) - - return { - "pk_eph_b64": base64.b64encode(pk_eph_raw).decode(), - "nonce_b64": base64.b64encode(nonce).decode(), - "wrapped_b64": base64.b64encode(wrapped).decode(), - } - -def unwrap_gek(bundle: dict, sk_recipient: bytes, pk_recipient: bytes) -> bytes: - """ - Unwrap a GEK bundle using the recipient's X25519 private key. - Raises InvalidTag if the key is wrong or the bundle was tampered. - """ - pk_eph_raw = base64.b64decode(bundle["pk_eph_b64"]) - nonce = base64.b64decode(bundle["nonce_b64"]) - wrapped = base64.b64decode(bundle["wrapped_b64"]) - - shared = X25519PrivateKey.from_private_bytes(sk_recipient).exchange( - X25519PublicKey.from_public_bytes(pk_eph_raw) - ) - wrap_key = HKDF( - algorithm=hashes.SHA256(), length=32, - salt=pk_eph_raw, info=GEK_WRAP_INFO, - ).derive(shared) - - return ChaCha20Poly1305(wrap_key).decrypt(nonce, wrapped, pk_recipient) - - GEK_WRAP_INFO_AES = b"meshbay:gek_wrap:v1:aes" def wrap_gek_aes(gek: bytes, pk_recipient: bytes) -> dict: @@ -220,24 +149,3 @@ def decrypt_keystore(iv: bytes, ciphertext: bytes, tag: bytes, key: bytes) -> by """Decrypt keystore blob. Raises on authentication failure.""" dec = Cipher(algorithms.AES(key), modes.GCM(iv, tag)).decryptor() return dec.update(ciphertext) + dec.finalize() - -# ── Chunk signing ───────────────────────────────────────────────────────────── - -def sign_chunk(sk_node: Ed25519PrivateKey, chunk_index: int, - nonce: bytes, ct_hash: bytes) -> bytes: - """ - Sign chunk metadata. Payload: chunk_index || nonce || ct_hash. - - Despite the name, this no longer signs file chunks — Phase 9.15 dropped per-chunk - signatures on the WebRTC path and 2026-09-03 dropped the QUIC copy that had been - left behind. Its one caller is `GroupIndex.serialize()`, which signs a whole index - envelope under the pseudo-index `INDEX_CHUNK`. - """ - payload = chunk_index.to_bytes(4, "big") + nonce + ct_hash - return sk_node.sign(payload) - -def verify_chunk_signature(pk_node: Ed25519PublicKey, chunk_index: int, - nonce: bytes, ct_hash: bytes, signature: bytes) -> None: - """Verify chunk signature. Raises InvalidSignature on failure.""" - payload = chunk_index.to_bytes(4, "big") + nonce + ct_hash - pk_node.verify(signature, payload) diff --git a/packages/meshbay-common/src/meshbay_common/groupbox.py b/packages/meshbay-common/src/meshbay_common/groupbox.py index 3b45dba..260e362 100644 --- a/packages/meshbay-common/src/meshbay_common/groupbox.py +++ b/packages/meshbay-common/src/meshbay_common/groupbox.py @@ -32,10 +32,10 @@ it is an oversight. The node holds the GEK for its own group, so unlike the inde this direction seals *towards* the node: it opens the payload before it writes anything to disk, and refuses a chunk that does not open rather than guessing. -Purpose separation is deliberate. `GroupIndex.serialize()` reuses -`chunk_key_aes(gek, file_hash, chunk_index)` with a pseudo-file ("the index as chunk -0 of a virtual index file"), which borrows a file's key space for something that is -not a file. Each purpose here derives its own subkey instead. +Purpose separation is deliberate. The alternative — reusing +`chunk_key_aes(gek, file_hash, chunk_index)` with a pseudo-file, "the index as chunk +0 of a virtual index file" — borrows a file's key space for something that is not a +file. Each purpose here derives its own subkey instead. """ from __future__ import annotations diff --git a/packages/meshbay-common/src/meshbay_common/handshake.py b/packages/meshbay-common/src/meshbay_common/handshake.py index 992fe8a..c3d5b94 100644 --- a/packages/meshbay-common/src/meshbay_common/handshake.py +++ b/packages/meshbay-common/src/meshbay_common/handshake.py @@ -89,6 +89,8 @@ CHALLENGE_PREFIX = b"meshbay:mnp:challenge:v1" # hub session token (MNP_VERSION note). A pre-4.0 peer presents the session # token, which this node now refuses — so the floor moves to 4.0 rather than # leaving a branch that would keep a hub credential reachable by every node. +# 5.0 (2026-09-28) does not move it: the break is confined to four signed +# operations, which a peer across it refuses to sign (MNP_VERSION note). MNP_MIN_SUPPORTED = "4.0" ROLE_CLIENT = "client" diff --git a/packages/meshbay-common/src/meshbay_common/paths.py b/packages/meshbay-common/src/meshbay_common/paths.py index 45bf34e..b673649 100644 --- a/packages/meshbay-common/src/meshbay_common/paths.py +++ b/packages/meshbay-common/src/meshbay_common/paths.py @@ -137,7 +137,9 @@ def sanitize_for_download(name: str, *, replacement: str = "_") -> str: For the client saving a file, never for the node storing one. Returns the name unchanged when it is already portable, so the common case is identity - and the caller can tell whether it renamed anything by comparing. + and the caller can tell whether it renamed anything by comparing. The + browser's copy is `static/portable-name.js`, held to this one by + `test_portable_name_parity.py`. """ if is_portable_name(name): return name diff --git a/packages/meshbay-common/src/meshbay_common/protocol.py b/packages/meshbay-common/src/meshbay_common/protocol.py index 6b929cd..af2d93b 100644 --- a/packages/meshbay-common/src/meshbay_common/protocol.py +++ b/packages/meshbay-common/src/meshbay_common/protocol.py @@ -30,7 +30,6 @@ number — it is a label chosen by the peer, and the only thing it decides is which local promise a reply belongs to. """ -import os from dataclasses import dataclass, field # The wire versions live in meshbay_common/__init__.py — one source, because a @@ -75,7 +74,6 @@ class MNP: # plaintext form on the wire (`chatbox.py`, docs/MESHBAY_DESIGN.md §4.5); # `format` distinguishes a *stored* pre-2.0 row, which is still served. CHAT_MESSAGE = "chat_msg" # one chat message, sealed and signed - CHAT_ATTACHMENT = "chat_attach" # attachment metadata CHAT_HISTORY = "chat_hist" # request message history (newest, or before a cursor) CHAT_HISTORY_RESPONSE = "chat_hist_resp" # history response with messages # Link unfurl: the node fetches a URL a member pasted and returns an @@ -123,7 +121,6 @@ class MNP: STREAM_END = "stream_end" # node signals end of stream STREAM_MORE = "stream_more" # client → node: room for N more segments STREAM_STOP = "stream_stop" # client → node: nobody is watching any more - EPHEMERAL_STREAM = "ephemeral_stream" # reserved — mobile live push HANDSHAKE_CHALLENGE = "handshake_challenge" # node → client: GEK proof nonce HANDSHAKE_RESPONSE = "handshake_response" # client → node: HMAC(GEK, nonce) ADMIN_CHALLENGE = "admin_challenge" # node → client: Ed25519 sign challenge @@ -318,9 +315,8 @@ class IndexEntry: def index_entry_wire(e: IndexEntry) -> dict: """ - The wire-dict shape used by INDEX_SYNC/INDEX_DELTA hand-built messages - (as opposed to GroupIndex.serialize()'s asdict() encoding of the whole - index). Centralized so the three call sites that build these + The wire-dict shape used by INDEX_SYNC/INDEX_DELTA hand-built messages. + Centralized so the three call sites that build these (webrtc/files.py's _do_index_sync, daemon._broadcast_index_change's two branches) can't drift from each other as fields are added. """ @@ -369,8 +365,7 @@ class IndexDelta: # under a key derived from the GEK, which only group members hold, and since C3 the # node authenticates itself in the handshake and is pinned by the client. A # signature per chunk re-proved, once per megabyte, what the session established -# once. (`sign_chunk` still exists in `crypto.py` — it signs the serialized index -# envelope, which is a different artifact; see `indexer/group_index.py`.) +# once. def chunk_ciphertext( @@ -460,11 +455,6 @@ def file_chunk_plaintext( UPLOAD_ID_LEN = 16 # 128 bits of client-chosen correlation, hex on the wire -def new_upload_id() -> str: - """A fresh correlation id for one upload.""" - return os.urandom(UPLOAD_ID_LEN).hex() - - def file_upload_wire( gek: bytes, group_id: str, diff --git a/packages/meshbay-common/src/meshbay_common/webcrypto.py b/packages/meshbay-common/src/meshbay_common/webcrypto.py index 3bd5ae6..7119e2e 100644 --- a/packages/meshbay-common/src/meshbay_common/webcrypto.py +++ b/packages/meshbay-common/src/meshbay_common/webcrypto.py @@ -1,28 +1,21 @@ """ -MeshBay — AES-256-GCM cipher variant for browser-accessible groups. +MeshBay — the content cipher: AES-256-GCM, per-chunk keys derived from the GEK. -The ChaCha20-Poly1305 GEK used in MNP (TCP+TLS and QUIC transport) -is NOT available in the WebCrypto API. For groups whose content must -be decryptable by a web browser (using SubtleCrypto), an AES-256-GCM -variant is used instead. - -The GEK wrapping (X25519 + HKDF) is identical — only the content -cipher changes. The hub stores and distributes GEK bundles the same way. - -Cipher selection is declared per-group in the hub registry: - "cipher": "chacha20-poly1305" (default, native clients) - "cipher": "aes-256-gcm" (browser-compatible groups) +AES-GCM because it is what WebCrypto offers, and one cipher serves every client: +the browser, the desktop client (the same engine) and the Python side here. Python side (this module): - encrypt_chunk_aes / decrypt_chunk_aes + chunk_key_aes / encrypt_chunk_aes / decrypt_chunk_aes JavaScript side (in static/crypto.js): - Uses SubtleCrypto.importKey + SubtleCrypto.decrypt with AES-GCM. + SubtleCrypto.importKey + SubtleCrypto.decrypt with AES-GCM. -Key derivation for AES variant — same HKDF info string with suffix: +Key derivation: info = b"file:" + file_hash + b":chunk:" + chunk_index + b":aes" -This ensures AES and ChaCha20 keys are always distinct even from the same GEK. +The `:aes` suffix dates from a ChaCha20-Poly1305 variant derived from the same +GEK without it, which nothing used and which is gone. It stays: it is part of +every chunk key in existence, and changing it would change them all. """ import os @@ -33,7 +26,7 @@ from cryptography.hazmat.primitives.kdf.hkdf import HKDF def chunk_key_aes(gek: bytes, file_hash: bytes, chunk_index: int) -> bytes: - """Derive a per-chunk AES-256 key. Distinct from ChaCha20 key.""" + """Derive a per-chunk AES-256 key from the GEK.""" return HKDF( algorithm=hashes.SHA256(), length=32, salt=None, info=b"file:" + file_hash + b":chunk:" + chunk_index.to_bytes(4, "big") + b":aes", diff --git a/packages/meshbay-common/tests/test_admin_subject_parity.py b/packages/meshbay-common/tests/test_admin_subject_parity.py new file mode 100644 index 0000000..7a229a9 --- /dev/null +++ b/packages/meshbay-common/tests/test_admin_subject_parity.py @@ -0,0 +1,145 @@ +""" +The subjects of multi-value admin operations are byte-identical in the browser and +in Python. + +The subject is what the operator's signature covers of a request, and each side +builds it on its own — the node from the request it stored, the client from what +the person asked for. A one-byte disagreement does not weaken anything (the client +refuses to sign), but it makes the operation impossible from a browser, and nothing +else in the suite crosses this boundary. + +Skipped when node is unavailable; that is a coverage gap, not a pass. +""" + +import json +import shutil +import subprocess +from pathlib import Path + +import pytest +from meshbay_common.adminop import ( + group_attach_subject, + invite_create_subject, + root_add_subject, + secret_digest, + structured_subject, + tmdb_config_subject, +) + +CRYPTO_JS = (Path(__file__).resolve().parents[2] + / "meshbay-hub" / "src" / "meshbay_hub" / "static" / "crypto.js") + +pytestmark = pytest.mark.skipif( + shutil.which("node") is None or not CRYPTO_JS.exists(), + reason="node or crypto.js unavailable — parity cannot be checked", +) + +ROOT_ADD = [ + ("/srv/Films", "", "generic", False, False), + ("/srv/Films", "Films", "video", True, True), + ("C:\\Users\\me\\Share", "Partagé", "photo", True, False), + ('/srv/a "quoted", odd:name|x', "名前", "audio", False, True), + ("/srv/tab\there\nnewline\x01ctl", "é", "generic", True, False), +] +GROUP_ATTACH = [ + ("photos", "/srv/photos", True), + ("famille-été", "/mnt/disque externe/Photos", False), +] +INVITE_CREATE = [ + ("0f8fad5b-d9cb-469f-a165-70867728950e", ""), + ("0f8fad5b-d9cb-469f-a165-70867728950e", "Élodie \"E\" 🙂"), +] +TMDB_CONFIG = [ + (None, None), ("", None), (None, ""), ("", ""), + ("eyJhbGciOiJIUzI1NiJ9.token", "fr-FR"), + ("abc", "keep"), +] + +_HARNESS = r""" +const fs = require('fs'); +globalThis.window = {}; +const src = fs.readFileSync(process.argv[2], 'utf8'); +const M = new Function(src + '\nreturn { rootAddSubject, groupAttachSubject, ' + + 'inviteCreateSubject, tmdbConfigSubject };')(); +const v = JSON.parse(fs.readFileSync(process.argv[3], 'utf8')); +(async () => { + const out = { + root_add: v.root_add.map((a) => M.rootAddSubject(...a)), + group_attach: v.group_attach.map((a) => M.groupAttachSubject(...a)), + invite_create: v.invite_create.map((a) => M.inviteCreateSubject(...a)), + tmdb_config: [], + }; + for (const a of v.tmdb_config) out.tmdb_config.push(await M.tmdbConfigSubject(...a)); + process.stdout.write(JSON.stringify(out)); +})(); +""" + + +@pytest.fixture(scope="module") +def js(tmp_path_factory): + d = tmp_path_factory.mktemp("subject-parity") + (d / "harness.js").write_text(_HARNESS, encoding="utf-8") + (d / "vectors.json").write_text(json.dumps({ + "root_add": ROOT_ADD, "group_attach": GROUP_ATTACH, + "invite_create": INVITE_CREATE, "tmdb_config": TMDB_CONFIG, + }), encoding="utf-8") + proc = subprocess.run( + ["node", str(d / "harness.js"), str(CRYPTO_JS), str(d / "vectors.json")], + capture_output=True, text=True, encoding="utf-8", timeout=60) + if proc.returncode != 0: + pytest.fail(f"node harness failed:\n{proc.stderr}") + return json.loads(proc.stdout) + + +def _bytes(s: str) -> bytes: + return s.encode("utf-8") + + +@pytest.mark.parametrize("i,args", list(enumerate(ROOT_ADD))) +def test_root_add_subject_parity(i, args, js): + assert _bytes(js["root_add"][i]) == _bytes(root_add_subject(*args)) + + +@pytest.mark.parametrize("i,args", list(enumerate(GROUP_ATTACH))) +def test_group_attach_subject_parity(i, args, js): + assert _bytes(js["group_attach"][i]) == _bytes(group_attach_subject(*args)) + + +@pytest.mark.parametrize("i,args", list(enumerate(INVITE_CREATE))) +def test_invite_create_subject_parity(i, args, js): + assert _bytes(js["invite_create"][i]) == _bytes(invite_create_subject(*args)) + + +@pytest.mark.parametrize("i,args", list(enumerate(TMDB_CONFIG))) +def test_tmdb_config_subject_parity(i, args, js): + assert _bytes(js["tmdb_config"][i]) == _bytes(tmdb_config_subject(*args)) + + +def test_every_value_changes_the_subject(): + base = ("/srv/Films", "Films", "video", False, False) + variants = {root_add_subject(*base)} + for i, other in enumerate(("/srv/Other", "Other", "audio", True, True)): + args = list(base) + args[i] = other + variants.add(root_add_subject(*args)) + assert len(variants) == 6 + + +def test_unchanged_cleared_and_set_are_three_subjects(): + assert len({tmdb_config_subject(None, None), tmdb_config_subject("", None), + tmdb_config_subject("t", None)}) == 3 + assert len({tmdb_config_subject(None, None), tmdb_config_subject(None, ""), + tmdb_config_subject(None, "fr-FR")}) == 3 + + +def test_the_token_is_never_written_into_the_subject(): + token = "eyJhbGciOiJIUzI1NiJ9.a-real-looking-secret" + assert token not in tmdb_config_subject(token, "fr-FR") + assert secret_digest(token).startswith("sha256:") + + +def test_a_crafted_field_cannot_impersonate_another(): + # Under a naive "path|name" join these two would collide. + a = structured_subject({"path": "/a|name=b", "name": ""}) + b = structured_subject({"path": "/a", "name": "b"}) + assert a != b diff --git a/packages/meshbay-common/tests/test_groupbox.py b/packages/meshbay-common/tests/test_groupbox.py index 65d8ce6..6e687ea 100644 --- a/packages/meshbay-common/tests/test_groupbox.py +++ b/packages/meshbay-common/tests/test_groupbox.py @@ -45,8 +45,8 @@ def test_purposes_are_separate_key_spaces(gek): The reason there are two info strings rather than one key reused. An ack sealed under the index subkey would otherwise be openable by anything - holding the index subkey, which is the confusion `GroupIndex.serialize()`'s - "the index as chunk 0 of a virtual index file" creates for chunk keys. + holding the index subkey — the confusion that treating "the index as chunk 0 of + a virtual index file" would create for chunk keys. """ assert group_key(gek, PURPOSE_INDEX) != group_key(gek, PURPOSE_ACK) sealed = seal(gek, PURPOSE_INDEX, "index_sync", "g1", PAYLOAD) diff --git a/packages/meshbay-common/tests/test_portable_name_parity.py b/packages/meshbay-common/tests/test_portable_name_parity.py new file mode 100644 index 0000000..3a491a7 --- /dev/null +++ b/packages/meshbay-common/tests/test_portable_name_parity.py @@ -0,0 +1,78 @@ +""" +The browser makes a name writable everywhere exactly as Python does. + +`static/portable-name.js` renames a file at the moment a member saves it; +`meshbay_common.paths.sanitize_for_download` is the same rule in Python. Two +copies of a rule that differ decide differently which files get renamed, so the +real module runs under node here against the real Python. + +Skipped when node is unavailable; that is a coverage gap, not a pass. +""" + +import json +import shutil +import subprocess +from pathlib import Path + +import pytest +from meshbay_common.paths import is_portable_name, sanitize_for_download + +PORTABLE_JS = (Path(__file__).resolve().parents[2] + / "meshbay-hub" / "src" / "meshbay_hub" / "static" / "portable-name.js") + +pytestmark = pytest.mark.skipif( + shutil.which("node") is None or not PORTABLE_JS.exists(), + reason="node or portable-name.js unavailable — parity cannot be checked", +) + +NAMES = [ + "plain.txt", "Réunion 12:30.pdf", 'a<b>c:d"e/f\\g|h?i*j.txt', "tab\tnew\nline", + "ends with dot.", "ends with space ", "trailing . . ", "CON", "con.txt", "aux.tar.gz", + "COM1", "com10.txt", "LPT9.log", "nul.", ".", "..", "", " ", ".bashrc", "...", + "名前:ファイル.mkv", "emoji 🙂?.png", "\x01\x1f.bin", "prn .txt", "Con", + "a.b.c", "COM1.", "normal name (2).mp4", +] + +_HARNESS = r""" +const fs = require('fs'); +const src = fs.readFileSync(process.argv[2], 'utf8').replace(/^export /gm, ''); +const M = new Function(src + '\nreturn { portableName, portablePath };')(); +const names = JSON.parse(fs.readFileSync(process.argv[3], 'utf8')); +process.stdout.write(JSON.stringify({ + names: names.map((n) => M.portableName(n)), + path: M.portablePath('Top:Folder/sub//aux.txt'), +})); +""" + + +@pytest.fixture(scope="module") +def js(tmp_path_factory): + d = tmp_path_factory.mktemp("portable") + (d / "harness.js").write_text(_HARNESS, encoding="utf-8") + (d / "names.json").write_text(json.dumps(NAMES), encoding="utf-8") + proc = subprocess.run( + ["node", str(d / "harness.js"), str(PORTABLE_JS), str(d / "names.json")], + capture_output=True, text=True, encoding="utf-8", timeout=60) + if proc.returncode != 0: + pytest.fail(f"node harness failed:\n{proc.stderr}") + return json.loads(proc.stdout) + + +@pytest.mark.parametrize("i,name", list(enumerate(NAMES))) +def test_the_browser_renames_as_python_does(i, name, js): + assert js["names"][i] == sanitize_for_download(name) + + +@pytest.mark.parametrize("name", NAMES) +def test_what_comes_out_can_be_written_everywhere(name): + out = sanitize_for_download(name) + assert is_portable_name(out), (name, out) + assert sanitize_for_download(out) == out + + +def test_a_portable_name_is_left_alone(): + assert sanitize_for_download("plain.txt") == "plain.txt" + + +def test_a_path_is_made_portable_segment_by_segment(js): + assert js["path"] == "Top_Folder/sub//aux_.txt" diff --git a/packages/meshbay-common/tests/test_webcrypto.py b/packages/meshbay-common/tests/test_webcrypto.py index fbffdc1..ffe3f36 100644 --- a/packages/meshbay-common/tests/test_webcrypto.py +++ b/packages/meshbay-common/tests/test_webcrypto.py @@ -17,17 +17,6 @@ def test_aes_roundtrip(): assert decrypt_chunk_aes(key, nonce, ct) == data -def test_aes_key_distinct_from_chacha_key(): - """AES and ChaCha20 keys for the same chunk must differ.""" - from meshbay_common.crypto import chunk_key as chacha_key - gek = generate_gek() - data = os.urandom(100) - fh = blake3.blake3(data).digest() - aes_k = chunk_key_aes(gek, fh, 0) - chacha_k = chacha_key(gek, fh, 0) - assert aes_k != chacha_k - - def test_aes_wrong_key_rejected(): gek = generate_gek() data = b"private content" @@ -76,19 +65,3 @@ def test_aes_gek_wrap_wrong_key_rejected(): bundle = wrap_gek_aes(gek, pk_to_raw(sk_a.public_key())) with pytest.raises(Exception): unwrap_gek_aes(bundle, sk_to_raw(sk_b), pk_to_raw(sk_b.public_key())) - - -def test_aes_gek_wrap_differs_from_chacha_wrap(): - from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey - from meshbay_common.crypto import ( - pk_to_raw, - wrap_gek, - wrap_gek_aes, - ) - gek = generate_gek() - sk = X25519PrivateKey.generate() - pk_raw = pk_to_raw(sk.public_key()) - - bundle_aes = wrap_gek_aes(gek, pk_raw) - bundle_chacha = wrap_gek(gek, pk_raw) - assert bundle_aes["wrapped_b64"] != bundle_chacha["wrapped_b64"] diff --git a/packages/meshbay-hub/src/meshbay_hub/api/admin.py b/packages/meshbay-hub/src/meshbay_hub/api/admin.py index 381378c..dbda197 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/admin.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/admin.py @@ -43,6 +43,8 @@ class SettingsPatchRequest(BaseModel): login: dict[str, int] | None = None # Session lifetime, in hours, each optional. session: dict[str, int] | None = None + # Content reports: who may report, how often, what a report leads to. + reports: dict[str, int] | None = None # ── Instance settings ──────────────────────────────────────────────────────── @@ -63,6 +65,9 @@ async def _settings_payload(db: AsyncSession) -> dict: "session": await hub_settings.session_limits(db), "session_defaults": dict(hub_settings.SESSION_DEFAULTS), "session_bounds": {k: list(v) for k, v in hub_settings.SESSION_BOUNDS.items()}, + "reports": await hub_settings.report_limits(db), + "reports_defaults": dict(hub_settings.REPORT_DEFAULTS), + "reports_bounds": {k: list(v) for k, v in hub_settings.REPORT_BOUNDS.items()}, } @@ -162,6 +167,26 @@ async def admin_patch_settings( )) await db.commit() + if body.reports: + unknown = sorted(set(body.reports) - set(hub_settings.REPORT_KEYS)) + if unknown: + raise HTTPException( + status_code=422, detail=f"Unknown report setting(s): {unknown}") + changed = [] + for key, value in body.reports.items(): + clamped = hub_settings.clamp_report_value(key, value) + await hub_settings.set_raw(db, f"reports.{key}", str(clamped)) + changed.append(f"{key}={clamped}") + log.info("Report policy changed by %s: %s", + current_user.username, ", ".join(changed)) + db.add(IPLog( + user_id=current_user.id, + event="admin_reports_update", + ip_address="admin", + detail=", ".join(changed)[:255], + )) + await db.commit() + return await _settings_payload(db) diff --git a/packages/meshbay-hub/src/meshbay_hub/api/deps.py b/packages/meshbay-hub/src/meshbay_hub/api/deps.py index 42f4101..2fd8b99 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/deps.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/deps.py @@ -90,6 +90,17 @@ async def require_user_scope( return current_user +async def require_node_scope( + payload: dict = Depends(_decode_token), + current_user: User = Depends(get_current_user), +) -> User: + """Only a node daemon's own token — for what a node fetches on its own behalf.""" + if payload.get("scope") != "node": + raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, + detail="Node token required") + return current_user + + def user_is_admin(user: User) -> bool: """Admin by DB role or by the config allow-list. Use inside a handler that already depends on `require_moderator` but has to draw the admin line for diff --git a/packages/meshbay-hub/src/meshbay_hub/api/groups.py b/packages/meshbay-hub/src/meshbay_hub/api/groups.py index 3a11345..df2f336 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/groups.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/groups.py @@ -20,16 +20,11 @@ from meshbay_hub.db.models import ( Group, GroupMember, IPLog, - SwarmSource, User, ) router = APIRouter(prefix="/v1/groups", tags=["groups"]) -# Swarm endpoints live at /v1/swarm/*. They were previously declared on the groups -# router with a full path, which mounted them at /v1/groups/v1/swarm/* (H7). -swarm_router = APIRouter(prefix="/v1/swarm", tags=["swarm"]) - @router.get("/mine") async def my_groups( @@ -211,113 +206,6 @@ async def list_public_groups( return {"groups": groups, "total": len(groups)} -# ── Swarm (content replication) ─────────────────────────────────────────────── - -class SwarmRegisterRequest(BaseModel): - content_hash: str # blake3 hex - endpoint: str # "<scheme>:<port>" — a port on the caller, never a host - - -# A transport and a port, and deliberately no host. The field used to be free -# text documented as "ip:port", so a caller could name *someone else's* -# address as a source; nothing dials a swarm source today, which is the only -# reason that was not already a reflection primitive. A reader learns where a -# node is from the node record, which is stamped with the address the announce -# actually came from — so a host here would be a second, weaker, answer to a -# question already settled elsewhere. -_SWARM_ENDPOINT = re.compile(r"^(webrtc|quic):([0-9]{1,5})$") - -# One account, this many public hashes. Rows are keyed (hash, account) with no -# cap, so a loop of invented hashes was unbounded storage growth on a hub -# shared with everyone else. A public library far larger than this is a real -# thing — but it is one a hub operator should be asked about, not something a -# client establishes by writing rows. -MAX_SWARM_HASHES_PER_ACCOUNT = 10_000 - - -@swarm_router.post("/register", status_code=201) -@limiter.limit("120/minute") -async def swarm_register( - body: SwarmRegisterRequest, - request: Request, - current_user: User = Depends(get_current_user), - db: AsyncSession = Depends(get_db), -): - """ - Node registers itself as a source for a PUBLIC content hash. - - Finding H7: the node registered hashes for every group it hosted, private ones - included, and this route was mounted at /v1/groups/v1/swarm/register — so the - node's calls 404'd and the leak was masked by a routing bug rather than - prevented. Nodes now filter by group visibility before calling, and the path is - correct, so the filter has to be right. - - Availability: the endpoint is a port, not an address, and the number of - hashes one account may claim is bounded. See the two constants above. - """ - from meshbay_hub.csam import check_content_hash - if check_content_hash(body.content_hash): - raise HTTPException(status_code=451, detail="Content blocked") - - m = _SWARM_ENDPOINT.match(body.endpoint or "") - if not m or not (0 < int(m.group(2)) < 65536): - raise HTTPException( - status_code=422, - detail="endpoint must be '<webrtc|quic>:<port>' — a port on the " - "registering node, not an address") - - from datetime import datetime - existing = await db.get(SwarmSource, (body.content_hash, current_user.id)) - now = datetime.now(UTC) - if existing: - existing.endpoint = body.endpoint - existing.last_seen = now - else: - held = (await db.execute( - select(func.count()).select_from(SwarmSource) - .where(SwarmSource.node_id == current_user.id))).scalar() or 0 - if held >= MAX_SWARM_HASHES_PER_ACCOUNT: - raise HTTPException( - status_code=429, - detail="This account already claims the maximum number of " - "public content hashes") - db.add(SwarmSource( - content_hash=body.content_hash, - node_id=current_user.id, - endpoint=body.endpoint, - )) - await db.commit() - return {"status": "registered", "hash": body.content_hash} - - -@swarm_router.get("/{content_hash}") -async def swarm_sources( - content_hash: str, - current_user: User = Depends(get_current_user), - db: AsyncSession = Depends(get_db), -): - """ - Return nodes that can serve a content hash. - - Authenticated (H7): an open endpoint lets anyone probe whether a given file - exists anywhere in the network and which node holds it. - """ - from datetime import datetime, timedelta - cutoff = datetime.now(UTC) - timedelta(minutes=30) - result = await db.execute( - select(SwarmSource) - .where( - SwarmSource.content_hash == content_hash, - SwarmSource.last_seen > cutoff, - ) - ) - sources = result.scalars().all() - return { - "hash": content_hash, - "sources": [{"node_id": s.node_id, "endpoint": s.endpoint} for s in sources], - } - - @router.get("/{group_id}/members") async def group_members( group_id: str, diff --git a/packages/meshbay-hub/src/meshbay_hub/api/moderation.py b/packages/meshbay-hub/src/meshbay_hub/api/moderation.py index 35c688c..038f310 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/moderation.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/moderation.py @@ -1,59 +1,88 @@ """ -MeshBay Hub — moderation endpoints. +MeshBay Hub — moderation endpoints (docs/MESHBAY_DESIGN.md §7.5). -Reporting flow: - POST /v1/reports — report a content hash (sign-in required) +Reporting: + POST /v1/reports — a member of a public group reports a file they saw there. - Thresholds (counted as DISTINCT reporting accounts, not raw rows): - < AUTO_BLOCK_THRESHOLD distinct reporters → logged - >= AUTO_BLOCK_THRESHOLD distinct reporters → hash added to the blocklist + Who may: a person's account (never a node's token), old enough + (`reports.min_account_age_hours`), an active member of that public group, + within a daily allowance (`reports.daily_per_account`) as well as the per-address + rate limit. One report per account per hash. + + What it leads to: once `reports.review_threshold` distinct accounts have + reported a hash, it is queued for an administrator (`content_reviews`), who is + notified and blocks or dismisses it. With `reports.auto_block` on, it is blocked + at once instead — the instance's choice, off by default, because a handful of + accounts made for the purpose would then be enough to take a file down. The flow only runs while the hub brokers public content: with public groups - switched off instance-wide there is nothing here to serve a reported hash from, - so it is refused rather than left open as an unauthenticated write surface. + switched off there is nothing here to report. Admin endpoints: - GET /v1/admin/blocklist — list blocked hashes - POST /v1/admin/blocklist — manually add a hash - DELETE /v1/admin/blocklist/{hash} — remove a hash + GET /v1/admin/reports — hashes waiting for a decision + POST /v1/admin/reports/{hash}/block — block it, and tell the nodes + POST /v1/admin/reports/{hash}/dismiss — close it without blocking + GET /v1/admin/blocklist — list blocked hashes + POST /v1/admin/blocklist — manually add a hash + DELETE /v1/admin/blocklist/{hash} — remove a hash Node integration: - GET /v1/blocklist/check?hash=<blake3> — check if a hash is blocked - GET /v1/blocklist — full blocklist (for node sync) + GET /v1/blocklist?after=<hash> — the list, paged, for a node's own token + WebSocket `blocklist_update` — additions and removals, pushed to nodes + hosting a public group """ import logging +from collections import Counter +from datetime import UTC, datetime, timedelta +from typing import Literal from fastapi import APIRouter, Depends, HTTPException, Query, Request -from pydantic import BaseModel +from pydantic import BaseModel, Field from sqlalchemy import func, select from sqlalchemy.ext.asyncio import AsyncSession from meshbay_hub import hub_settings -from meshbay_hub.api.deps import get_current_user, require_admin +from meshbay_hub.api.deps import ( + require_admin, + require_moderator, + require_node_scope, + require_user_scope, + user_is_admin, +) from meshbay_hub.api.middleware import limiter from meshbay_hub.api.netutil import client_ip +from meshbay_hub.api.revocation import broadcast_blocklist_update from meshbay_hub.db.engine import get_db -from meshbay_hub.db.models import ContentBlocklist, ContentReport, User +from meshbay_hub.db.models import ( + ContentBlocklist, + ContentReport, + ContentReview, + Group, + GroupMember, + User, +) log = logging.getLogger(__name__) router = APIRouter(tags=["moderation"]) -# Distinct reporting accounts before a hash is auto-blocked. Kept low for a -# responsive community signal, but note it is only as strong as account -# creation: while a bot can register freely (see the reCAPTCHA gap), the real -# control is the admin reviewing `GET /v1/admin/blocklist` and the audit log. -AUTO_BLOCK_THRESHOLD = 3 +# One answer for every reason a report is not accepted from this account for this +# group, so the endpoint does not tell anyone which groups exist or who is in them. +_NOT_YOURS = "You can report a file only in a public group you are a member of." + + +def _is_hash(value: str) -> bool: + return len(value) == 64 and all(c in "0123456789abcdef" for c in value) # ── Models ──────────────────────────────────────────────────────────────────── class ReportRequest(BaseModel): - content_hash: str # blake3 hex (64 chars) - group_id: str | None = None - reason: str = "illegal" - detail: str | None = None + content_hash: str = Field(max_length=64) # blake3 hex (64 chars) + group_id: str = Field(max_length=36) + reason: Literal["illegal", "spam", "copyright", "other"] = "illegal" + detail: str | None = Field(default=None, max_length=256) class BlocklistAddRequest(BaseModel): @@ -61,124 +90,150 @@ class BlocklistAddRequest(BaseModel): reason: str -# ── Public endpoints ────────────────────────────────────────────────────────── +# ── Reporting ───────────────────────────────────────────────────────────────── @router.post("/v1/reports", status_code=201) @limiter.limit("10/hour") async def report_content( body: ReportRequest, request: Request, - current_user: User = Depends(get_current_user), + current_user: User = Depends(require_user_scope), db: AsyncSession = Depends(get_db), ): """ - Report a public content hash for moderation. + Report a file of a public group, as a member of that group. - Sign-in is required. It used to be anonymous, which made it a censorship - primitive: two unauthenticated POSTs naming any blake3 id auto-added it to - the blocklist that nodes enforce, network-wide, with manual admin removal the - only undo. The threshold now counts *distinct reporting accounts*, one vote - per account per hash. - - Refused entirely when the hub has public groups switched off: nothing here - brokers public content then, nothing syncs the blocklist, and an open write - endpoint would only be abuse surface. + Every bound here answers what a report costs someone else: a file taken out + of a group everyone else uses, and an administrator's time. So a report takes + a person's account (a node's token is refused), one that has existed for a + while, membership of the public group the file was seen in, and a daily + allowance per account besides the rate limit per address — an address is one + of thousands a subscriber holds. It never blocks anything by itself unless the + instance chose automatic blocking: past the threshold, an administrator + decides. """ if not await hub_settings.public_groups_allowed(db): raise HTTPException( status_code=403, detail="This hub does not broker public content, so there is nothing to report here.") - - if len(body.content_hash) != 64 or not all(c in "0123456789abcdef" for c in body.content_hash): + if not _is_hash(body.content_hash): raise HTTPException(status_code=422, detail="content_hash must be 64 hex chars (blake3)") + limits = await hub_settings.report_limits(db) + now = datetime.now(UTC) + + created = current_user.created_at + if created is not None and created.tzinfo is None: + created = created.replace(tzinfo=UTC) + if created is not None and \ + now - created < timedelta(hours=limits["min_account_age_hours"]): + raise HTTPException(status_code=403, + detail="This account is too new to report content yet.") + + group = await db.get(Group, body.group_id) + member = await db.scalar(select(GroupMember.user_id).where( + GroupMember.group_id == body.group_id, + GroupMember.user_id == current_user.id)) + if group is None or group.visibility != "public" or group.status != "active" \ + or member is None: + raise HTTPException(status_code=403, detail=_NOT_YOURS) + + today = await db.scalar(select(func.count(ContentReport.id)).where( + ContentReport.reporter_id == current_user.id, + ContentReport.reported_at > now - timedelta(days=1))) or 0 + if today >= limits["daily_per_account"]: + raise HTTPException(status_code=429, + detail="You have reached today's number of reports.") + # One vote per account per hash — a single reporter must not be able to walk # the threshold up on their own by posting repeatedly. already = await db.scalar( select(ContentReport.id).where( ContentReport.content_hash == body.content_hash, ContentReport.reporter_id == current_user.id)) + if already: + return {"status": "already_reported"} - if not already: - db.add(ContentReport( - content_hash=body.content_hash, - reporter_id=current_user.id, - group_id=body.group_id, - reason=body.reason, - detail=body.detail, - ip_address=client_ip(request), - )) - await db.flush() + db.add(ContentReport( + content_hash=body.content_hash, + reporter_id=current_user.id, + group_id=body.group_id, + reason=body.reason, + detail=body.detail, + ip_address=client_ip(request), + )) + await db.flush() distinct_reporters = await db.scalar( select(func.count(func.distinct(ContentReport.reporter_id))) .where(ContentReport.content_hash == body.content_hash)) or 0 - action = "already_reported" if already else "logged" - if distinct_reporters >= AUTO_BLOCK_THRESHOLD: - existing = await db.get(ContentBlocklist, body.content_hash) - if not existing: - db.add(ContentBlocklist( - content_hash=body.content_hash, - reason=f"auto:{body.reason}", - added_by="auto", - )) - action = "auto_blocked" + blocked_now = False + if distinct_reporters >= limits["review_threshold"] \ + and await db.get(ContentBlocklist, body.content_hash) is None: + review = await db.get(ContentReview, body.content_hash) + if review is not None and review.status == "dismissed": + pass # an administrator's decision stands; more reports do not reopen it + elif limits["auto_block"]: + db.add(ContentBlocklist(content_hash=body.content_hash, + reason=f"auto:{body.reason}", added_by="auto")) + if review is None: + db.add(ContentReview(content_hash=body.content_hash, status="blocked", + decided_at=now, decided_by="auto")) + else: + review.status, review.decided_at, review.decided_by = "blocked", now, "auto" + blocked_now = True log.warning("Content auto-blocked after %d distinct reporters: %s", distinct_reporters, body.content_hash[:16]) - + elif review is None: + db.add(ContentReview(content_hash=body.content_hash, status="pending")) + await _notify_admins(db, body.content_hash) + log.warning("Content queued for review after %d distinct reporters: %s", + distinct_reporters, body.content_hash[:16]) await db.commit() - return { - "status": action, - "content_hash": body.content_hash, - "report_count": distinct_reporters, - "threshold": AUTO_BLOCK_THRESHOLD, - } - + if blocked_now: + await broadcast_blocklist_update(db, add=[body.content_hash]) + # The same answer whatever happened next: a reporter is not told how close a + # file is to review, which is a count to aim at. + return {"status": "logged"} -@router.get("/v1/blocklist/check") -@limiter.limit("120/minute") -async def check_blocklist( - hash: str, - request: Request, - db: AsyncSession = Depends(get_db), -): - """Check if a single hash is blocked. Used by nodes before serving public content. - Unauthenticated, because a node consults it before serving public content - and does so on its own behalf. That makes the shape check worth having: - without it any string of any length became a primary-key lookup. - """ - if len(hash) != 64 or not all(c in "0123456789abcdef" for c in hash): - raise HTTPException(status_code=422, detail="hash must be 64 hex chars (blake3)") - blocked = await db.get(ContentBlocklist, hash) - return { - "blocked": blocked is not None, - "hash": hash, - "reason": blocked.reason if blocked else None, - } +async def _notify_admins(db: AsyncSession, content_hash: str) -> None: + from meshbay_hub.api.notifications import create_notification + admins = [u for u in (await db.execute(select(User).where( + User.status == "active"))).scalars().all() if user_is_admin(u)] + for admin in admins: + await create_notification( + db, admin.id, "content_review", + "Reported content is waiting for a decision", + detail=content_hash[:16], link="#/admin", aggregate=False) @router.get("/v1/blocklist") async def get_blocklist( + current_node: User = Depends(require_node_scope), db: AsyncSession = Depends(get_db), - # Bounded, like every other list. This one takes no authentication — a - # node syncs it at startup — and had no ceiling at all, so any stranger - # could ask for the table in one query, repeatedly. 10 000 is what a node - # asks for, so it is the default and also the most anyone may have. + after: str = Query(default="", max_length=64), limit: int = Query(default=10000, ge=1, le=10000), ): """ - Return the full blocklist. Nodes sync this on startup. - Returns hashes only (not reasons) to minimize data exposure. + The content blocklist, a page at a time, for a node hosting a public group. + + Hashes only, never the reasons. Ordered by hash so `after` (the last hash of + the previous page) is a stable cursor: a list longer than one page used to + be cut at 10 000 with no way to ask for the rest, and the node applying it + silently served everything past the cut. A node's own token, because this + is what a node fetches on its own behalf and nothing else asks for it. """ result = await db.execute( select(ContentBlocklist.content_hash) - .order_by(ContentBlocklist.added_at.desc()) + .where(ContentBlocklist.content_hash > after) + .order_by(ContentBlocklist.content_hash) .limit(limit) ) hashes = [row[0] for row in result.fetchall()] - return {"count": len(hashes), "hashes": hashes} + return {"hashes": hashes, + "next": hashes[-1] if len(hashes) == limit else None} # ── Admin endpoints ─────────────────────────────────────────────────────────── @@ -214,16 +269,19 @@ async def admin_add_blocklist( current_user: User = Depends(require_admin), db: AsyncSession = Depends(get_db), ): + if not _is_hash(body.content_hash): + raise HTTPException(status_code=422, detail="content_hash must be 64 hex chars (blake3)") existing = await db.get(ContentBlocklist, body.content_hash) if existing: raise HTTPException(status_code=409, detail="Hash already blocked") db.add(ContentBlocklist( content_hash=body.content_hash, - reason=body.reason, + reason=body.reason[:64], added_by=current_user.username, )) await db.commit() + await broadcast_blocklist_update(db, add=[body.content_hash]) return {"status": "blocked", "hash": body.content_hash} @@ -238,5 +296,73 @@ async def admin_remove_blocklist( raise HTTPException(status_code=404, detail="Hash not in blocklist") await db.delete(entry) await db.commit() + await broadcast_blocklist_update(db, remove=[content_hash]) return {"status": "unblocked", "hash": content_hash} + +# ── Review queue ────────────────────────────────────────────────────────────── + +@router.get("/v1/admin/reports") +async def admin_list_reports( + current_user: User = Depends(require_moderator), + db: AsyncSession = Depends(get_db), + limit: int = Query(default=100, ge=1, le=500), +): + """Hashes waiting for a decision, oldest first, with what was said about them.""" + reviews = (await db.execute( + select(ContentReview).where(ContentReview.status == "pending") + .order_by(ContentReview.opened_at).limit(limit))).scalars().all() + out = [] + for r in reviews: + reports = (await db.execute(select(ContentReport).where( + ContentReport.content_hash == r.content_hash))).scalars().all() + group_ids = sorted({x.group_id for x in reports if x.group_id}) + names = dict((await db.execute(select(Group.id, Group.name).where( + Group.id.in_(group_ids)))).all()) if group_ids else {} + out.append({ + "hash": r.content_hash, + "opened_at": r.opened_at.isoformat(), + "reporters": len({x.reporter_id for x in reports}), + "reasons": dict(Counter(x.reason for x in reports)), + "details": [x.detail for x in reports if x.detail][:10], + "groups": [{"id": g, "name": names.get(g, "")} for g in group_ids], + }) + return {"reports": out} + + +async def _decide(db: AsyncSession, content_hash: str, status: str, by: str) -> ContentReview: + review = await db.get(ContentReview, content_hash) + if review is None or review.status != "pending": + raise HTTPException(status_code=404, detail="Nothing waiting for this hash") + review.status, review.decided_at, review.decided_by = status, datetime.now(UTC), by + return review + + +@router.post("/v1/admin/reports/{content_hash}/block") +async def admin_block_reported( + content_hash: str, + current_user: User = Depends(require_admin), + db: AsyncSession = Depends(get_db), +): + await _decide(db, content_hash, "blocked", current_user.username) + reasons = Counter((await db.execute(select(ContentReport.reason).where( + ContentReport.content_hash == content_hash))).scalars().all()) + if await db.get(ContentBlocklist, content_hash) is None: + db.add(ContentBlocklist( + content_hash=content_hash, + reason=f"reported:{reasons.most_common(1)[0][0] if reasons else 'other'}", + added_by=current_user.username)) + await db.commit() + await broadcast_blocklist_update(db, add=[content_hash]) + return {"status": "blocked", "hash": content_hash} + + +@router.post("/v1/admin/reports/{content_hash}/dismiss") +async def admin_dismiss_reported( + content_hash: str, + current_user: User = Depends(require_admin), + db: AsyncSession = Depends(get_db), +): + await _decide(db, content_hash, "dismissed", current_user.username) + await db.commit() + return {"status": "dismissed", "hash": content_hash} diff --git a/packages/meshbay-hub/src/meshbay_hub/api/relay.py b/packages/meshbay-hub/src/meshbay_hub/api/relay.py deleted file mode 100644 index 7bb3f66..0000000 --- a/packages/meshbay-hub/src/meshbay_hub/api/relay.py +++ /dev/null @@ -1,166 +0,0 @@ -""" -MeshBay Hub — Mesh Relay registration protocol (5.3). - -Community-operated TURN relays register with hubs. -Nodes query the hub for available relays when UDP hole punching fails. - -Relay registration: - POST /v1/relays/register — relay announces itself (signed JWT) - GET /v1/relays — list active relays (for nodes) - -Relay authentication: relay generates an Ed25519 keypair at install time. An -admin approves the public key, and every register call carries an Ed25519 -signature over "meshbay:relay_register:<relay_id>:<endpoint>:<timestamp>" — -the same proof-of-possession shape as /v1/nodes/announce. - -Relay is responsible for E2E encrypted QUIC traffic only (it cannot -read the application-layer content, only forward UDP packets). -""" - -import base64 -import logging -import time - -from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey -from fastapi import APIRouter, Depends, HTTPException -from pydantic import BaseModel -from sqlalchemy.ext.asyncio import AsyncSession - -from meshbay_hub.api.deps import require_admin -from meshbay_hub.db.engine import get_db -from meshbay_hub.db.models import User - -log = logging.getLogger(__name__) - -# **Closed, the same way and for a similar reason as federation.** Nothing in the -# tree calls these routes — no node asks for a relay, no client offers one — and -# §11.1 measured two ISPs with no TURN relay needed. Two of the three take no -# account and answer anyone who can reach the hub, so a registry nothing uses -# was an unauthenticated surface kept for its own sake. A constant, not a -# setting: re-opening it means building the node side first, then flipping this. -RELAYS_ENABLED = False - - -def _relays_open() -> None: - """Refuse every route on this router while the registry is closed. - - On the router rather than in each handler, so a route added later is closed - before anybody remembers to write the check (C6). - """ - if not RELAYS_ENABLED: - raise HTTPException(status_code=503, - detail="The relay registry is not enabled on this hub") - - -router = APIRouter(prefix="/v1/relays", tags=["relay"], - dependencies=[Depends(_relays_open)]) - -# In-memory relay registry (production: DB table) -_relays: dict[str, dict] = {} # relay_id → {endpoint, pk, last_seen, capacity} - - -# ── Models ──────────────────────────────────────────────────────────────────── - -class RelayRegisterRequest(BaseModel): - """Relay self-registers, proving possession of its approved key.""" - relay_id: str - endpoint: str # "ip:port" (UDP) - pk_relay: str # base64 Ed25519 public key - capacity: int = 100 # max concurrent connections - timestamp: int | None = None # unix seconds - signature: str | None = None # base64 Ed25519 over the register message - - -class RelayAdminApproveRequest(BaseModel): - relay_id: str - pk_relay: str # admin approves by registering the relay's public key - - -# ── Relay endpoints ─────────────────────────────────────────────────────────── - -REGISTER_TIMESTAMP_WINDOW = 300 # seconds either side, as /v1/nodes/announce - - -@router.post("/register", status_code=201) -async def relay_register( - body: RelayRegisterRequest, - db: AsyncSession = Depends(get_db), -): - """ - Relay announces itself. Must be pre-approved by a hub admin, and must prove - it holds the private key that approval registered. - - This endpoint has no `Depends` on an account on purpose — a relay is not a - user — but it had no proof of anything either: it compared `pk_relay` - against the approved value, which is a **public** key, so anyone who could - read it could rewrite where the hub tells nodes to send relayed traffic. - The module docstring said "signs keepalive JWTs" and nothing verified a - signature; `jwt` was imported and never used. A key is not a password, and - the fix is the proof-of-possession pattern already used by - /v1/nodes/announce and /v1/nodes/auth. - """ - approved = _relays.get(body.relay_id) - if not approved or approved.get("pk") != body.pk_relay: - raise HTTPException(status_code=403, - detail="Relay not approved — ask the hub admin to " - "run POST /v1/relays/approve") - - if body.timestamp is None or not body.signature: - raise HTTPException( - status_code=400, - detail="register requires timestamp and signature (proof of possession)") - if abs(int(time.time()) - body.timestamp) > REGISTER_TIMESTAMP_WINDOW: - raise HTTPException(status_code=401, detail="Timestamp too old or too far ahead") - - message = (f"meshbay:relay_register:{body.relay_id}:" - f"{body.endpoint}:{body.timestamp}").encode() - try: - pk = Ed25519PublicKey.from_public_bytes(base64.b64decode(body.pk_relay)) - pk.verify(base64.b64decode(body.signature), message) - except Exception: - log.warning("Relay %s failed proof of possession", body.relay_id[:8]) - raise HTTPException(status_code=401, detail="Invalid relay key proof of possession") - - _relays[body.relay_id].update({ - "endpoint": body.endpoint, - "capacity": body.capacity, - "last_seen": int(time.time()), - "active": True, - }) - log.info("Relay registered: %s at %s", body.relay_id[:8], body.endpoint) - return {"status": "registered", "relay_id": body.relay_id} - - -@router.get("") -async def list_relays(): - """ - List active Mesh Relays. Called by nodes when UDP hole punching fails. - Returns only active relays (seen in the last 5 minutes). - """ - cutoff = int(time.time()) - 300 - active = [ - { - "relay_id": rid, - "endpoint": r["endpoint"], - "capacity": r["capacity"], - } - for rid, r in _relays.items() - if r.get("active") and r.get("last_seen", 0) > cutoff - ] - return {"relays": active, "count": len(active)} - - -@router.post("/approve", status_code=201) -async def admin_approve_relay( - body: RelayAdminApproveRequest, - current_user: User = Depends(require_admin), -): - """Admin: pre-approve a relay by registering its public key.""" - _relays[body.relay_id] = { - "pk": body.pk_relay, - "approved_by": current_user.username, - "approved_at": int(time.time()), - "active": False, # becomes True after first register call - } - log.info("Relay approved by %s: %s", current_user.username, body.relay_id[:8]) - return {"status": "approved", "relay_id": body.relay_id} diff --git a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py index 2c0b8db..fe9edf3 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py @@ -165,6 +165,35 @@ async def broadcast_revocation(token: str) -> int: return sent +async def broadcast_blocklist_update(db: AsyncSession, *, add: list[str] = (), + remove: list[str] = ()) -> int: + """ + Tell every connected node that hosts a public group what changed on the + content blocklist. Returns how many were told. + + Only those nodes: the list names public content, and a node hosting only + private groups has nothing to apply it to. A node that is offline now syncs + the whole list when it next connects (`GET /v1/blocklist`), so a missed push + is only late, never lost. + """ + if not add and not remove: + return 0 + public = set((await db.execute( + select(Group.id).where(Group.visibility == "public"))).scalars().all()) + payload = json.dumps({"type": "blocklist_update", + "add": list(add), "remove": list(remove)}) + sent = 0 + for node_id, ws in list(_connected_nodes.items()): + if not public.intersection(_node_groups.get(node_id, [])): + continue + try: + await ws.send_text(payload) + sent += 1 + except Exception: + _connected_nodes.pop(node_id, None) + return sent + + def _sign_revocation(target: str, target_id: str, reason: str) -> str: """Issue a signed revocation token (JWT EdDSA).""" from meshbay_hub.auth import _hub_id, _hub_sk_pem diff --git a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py index fc40204..6c9699b 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py @@ -171,7 +171,7 @@ async def webrtc_offer( Finding H4: it also ignored group status, so "suspend a group" did not stop new connections from being brokered to nodes hosting it. """ - from meshbay_hub.api.revocation import _connected_nodes, _node_groups + from meshbay_hub.api.revocation import _connected_nodes if len(body.sdp) > MAX_SDP_BYTES: raise HTTPException(status_code=413, detail="SDP too large") diff --git a/packages/meshbay-hub/src/meshbay_hub/api/users.py b/packages/meshbay-hub/src/meshbay_hub/api/users.py index 7046c2f..3e996a7 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/users.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/users.py @@ -41,7 +41,6 @@ from meshbay_hub.db.models import ( Node, Notification, RefreshToken, - SwarmSource, User, UserDevice, UserPreference, @@ -1379,14 +1378,13 @@ async def erase_account(db: AsyncSession, user: User, owned_groups: str = "refus the person it is about. Gone: credentials, email, node key, group memberships, notifications, refresh - tokens, node registrations, device keys, public-swarm sources. The username + tokens, node registrations, device keys. The username is released. Device keys go even though the desktop client keeps its private half: left behind, the key still belongs to this tombstone, so an account created later from the same installation is refused that device ("belongs to another - account"). Swarm sources are keyed by the *user* id and carry the node's - ip:port. + account"). Kept: the row itself, emptied, and the IP log that points at it. Those logs exist for one year to answer legal requests, and a log that cannot say whose @@ -1419,7 +1417,6 @@ async def erase_account(db: AsyncSession, user: User, owned_groups: str = "refus await db.execute(delete(RefreshToken).where(RefreshToken.user_id == user.id)) await db.execute(delete(Node).where(Node.user_id == user.id)) await db.execute(delete(UserDevice).where(UserDevice.user_id == user.id)) - await db.execute(delete(SwarmSource).where(SwarmSource.node_id == user.id)) await db.execute(delete(EmailVerification).where(EmailVerification.user_id == user.id)) # Links this account issued for a group it no longer owns; the ones for its # own groups went with them above. A used link keeps pointing at the diff --git a/packages/meshbay-hub/src/meshbay_hub/app.py b/packages/meshbay-hub/src/meshbay_hub/app.py index 7ccf598..16a9d4a 100644 --- a/packages/meshbay-hub/src/meshbay_hub/app.py +++ b/packages/meshbay-hub/src/meshbay_hub/app.py @@ -21,7 +21,6 @@ from meshbay_hub.api.admin import router as admin_router from meshbay_hub.api.deps import set_admin_usernames from meshbay_hub.api.federation import router as federation_router from meshbay_hub.api.groups import router as groups_router -from meshbay_hub.api.groups import swarm_router from meshbay_hub.api.health import router as health_router from meshbay_hub.api.hub import router as hub_router from meshbay_hub.api.hub import set_config as hub_set_config @@ -31,7 +30,6 @@ from meshbay_hub.api.middleware import limiter from meshbay_hub.api.moderation import router as moderation_router from meshbay_hub.api.nodes import router as nodes_router from meshbay_hub.api.notifications import router as notifications_router -from meshbay_hub.api.relay import router as relay_router from meshbay_hub.api.revocation import router as revocation_router from meshbay_hub.api.signaling import router as signaling_router from meshbay_hub.api.users import router as users_router @@ -41,7 +39,6 @@ from meshbay_hub.api.webapp import configure as webapp_configure from meshbay_hub.api.webapp import router as webapp_router from meshbay_hub.auth import generate_hub_keypair, load_hub_keypair from meshbay_hub.config import HubConfig -from meshbay_hub.csam import csam_router from meshbay_hub.db.engine import close_db, init_db @@ -129,9 +126,6 @@ def create_app(cfg: HubConfig | None = None) -> FastAPI: if cfg.identity.admin_usernames: await _sync_admin_roles(cfg.identity.admin_usernames) - from meshbay_hub.csam import get_csam_checker - get_csam_checker().load() - from meshbay_hub.db.engine import get_session_factory from meshbay_hub.tasks.cleanup import cleanup_loop cleanup_task = asyncio.create_task(cleanup_loop(get_session_factory())) @@ -203,13 +197,10 @@ def create_app(cfg: HubConfig | None = None) -> FastAPI: app.include_router(groups_router) app.include_router(invite_links_router) app.include_router(invite_links_redeem_router) - app.include_router(swarm_router) app.include_router(revocation_router) app.include_router(moderation_router) app.include_router(federation_router) - app.include_router(csam_router) app.include_router(health_router) - app.include_router(relay_router) app.include_router(signaling_router) app.include_router(admin_router) app.include_router(notifications_router) diff --git a/packages/meshbay-hub/src/meshbay_hub/csam.py b/packages/meshbay-hub/src/meshbay_hub/csam.py deleted file mode 100644 index b8e8d04..0000000 --- a/packages/meshbay-hub/src/meshbay_hub/csam.py +++ /dev/null @@ -1,158 +0,0 @@ -""" -MeshBay Hub — CSAM hash matching. - -Checks public content hashes against known CSAM (Child Sexual Abuse Material) -hash databases before allowing content to be registered or served publicly. - -Production integration: - - NCMEC (National Center for Missing & Exploited Children): PhotoDNA hash database - Access requires formal application: https://www.missingkids.org/gethelpnow/cybertipline - - IWF (Internet Watch Foundation): URL and hash list (UK-based) - Access via IWF membership: https://www.iwf.org.uk/our-technology/our-products/hash-list/ - -This module provides: - 1. A local CSAM hash database (SQLite file, populated from official sources) - 2. A check function used before content registration - 3. An admin endpoint to update the hash list - -IMPORTANT: Never log matched hashes or file contents. CSAM detection -must be reported to NCMEC (US law) or relevant authority immediately. -""" - -import logging -from pathlib import Path - -from fastapi import APIRouter, Depends, HTTPException - -from meshbay_hub.api.deps import require_admin -from meshbay_hub.db.models import User - -log = logging.getLogger(__name__) - -# Default path for the CSAM hash database (blake3 hex hashes, one per line) -DEFAULT_CSAM_DB_PATH = Path("/var/lib/meshbay/hub/csam_hashes.txt") - - -class CSAMChecker: - """ - Checks content hashes against a known CSAM hash database. - - Usage: - checker = CSAMChecker() - checker.load() - if checker.is_known_csam(blake3_hex): - # refuse to serve, report to authority - pass - """ - - def __init__(self, db_path: Path = DEFAULT_CSAM_DB_PATH): - self._db_path = db_path - self._hashes: set[str] = set() - self._loaded = False - - def load(self, db_path: Path | None = None) -> int: - """ - Load CSAM hashes from the hash database file. - Returns the number of hashes loaded. - - File format: one blake3 hex hash per line (64 chars), comments with #. - """ - path = db_path or self._db_path - if not path.exists(): - log.warning("CSAM hash database not found: %s. " - "Contact NCMEC (US) or IWF (EU) for access.", path) - self._loaded = True - return 0 - - count = 0 - with open(path) as f: - for line in f: - line = line.strip() - if line and not line.startswith("#") and len(line) == 64: - self._hashes.add(line.lower()) - count += 1 - - self._loaded = True - log.info("CSAM hash database loaded: %d hashes from %s", count, path) - return count - - def is_known_csam(self, content_hash_hex: str) -> bool: - """ - Return True if the hash matches a known CSAM hash. - NEVER logs the hash or any file information. - """ - if not self._loaded: - self.load() - return content_hash_hex.lower() in self._hashes - - @property - def hash_count(self) -> int: - return len(self._hashes) - - def add_hash(self, hash_hex: str) -> None: - """Add a hash to the in-memory set (and optionally persist).""" - self._hashes.add(hash_hex.lower()) - - def update_from_file(self, new_db_path: Path) -> int: - """Hot-reload from a new hash database file.""" - old_count = len(self._hashes) - self._hashes.clear() - count = self.load(new_db_path) - log.info("CSAM database updated: %d → %d hashes", old_count, count) - return count - - -# Module-level singleton (initialised in hub lifespan) -_checker = CSAMChecker() - - -def get_csam_checker() -> CSAMChecker: - return _checker - - -def check_content_hash(blake3_hex: str) -> bool: - """ - Check a content hash against the CSAM database. - Returns True if the content is KNOWN CSAM — block immediately. - - Callers MUST: - 1. Refuse to serve the content - 2. Log the event (without the hash) for legal audit purposes - 3. Report to NCMEC CyberTipline if operating in the US: - https://www.missingkids.org/gethelpnow/cybertipline - """ - return _checker.is_known_csam(blake3_hex) - - -# ── Hub API integration ─────────────────────────────────────────────────────── - -csam_router = APIRouter(prefix="/v1/admin/csam", tags=["csam"]) - - -@csam_router.get("/status") -async def csam_status(current_user: User = Depends(require_admin)): - """Return CSAM checker status (hash count, database path).""" - return { - "hash_count": _checker.hash_count, - "db_path": str(_checker._db_path), - "loaded": _checker._loaded, - "note": "Contact NCMEC or IWF for hash database access.", - } - - -@csam_router.post("/check") -async def check_hash( - body: dict, - current_user: User = Depends(require_admin), -): - """ - Check a single hash. Admin use only. - Returns True/False WITHOUT logging the hash (legal requirement). - """ - hash_hex = body.get("hash", "") - if len(hash_hex) != 64: - raise HTTPException(status_code=422, detail="hash must be 64 hex chars") - matched = check_content_hash(hash_hex) - # Do NOT log whether a match was found — only log the check attempt - log.info("CSAM check performed by admin %s", current_user.username) - return {"matched": matched} diff --git a/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/e6f7a8b9c0d1_drop_swarm_sources.py b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/e6f7a8b9c0d1_drop_swarm_sources.py new file mode 100644 index 0000000..0e61390 --- /dev/null +++ b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/e6f7a8b9c0d1_drop_swarm_sources.py @@ -0,0 +1,38 @@ +"""drop the public-content swarm table + +Nodes registered the hashes of their public groups here and nothing ever read +them back: the swarm was written and never used. + +Revision ID: e6f7a8b9c0d1 +Revises: c4d5e6f7a8b9 +""" + +from collections.abc import Sequence + +import sqlalchemy as sa +from alembic import op + +revision: str = "e6f7a8b9c0d1" +down_revision: str | Sequence[str] | None = "c4d5e6f7a8b9" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.drop_index("ix_swarm_hash", table_name="swarm_sources") + op.drop_table("swarm_sources") + + +def downgrade() -> None: + op.create_table( + "swarm_sources", + sa.Column("content_hash", sa.String(64), nullable=False), + sa.Column("node_id", sa.String(36), nullable=False), + sa.Column("endpoint", sa.String(128), nullable=False), + sa.Column("registered_at", sa.DateTime(timezone=True), nullable=False, + server_default=sa.func.now()), + sa.Column("last_seen", sa.DateTime(timezone=True), nullable=False, + server_default=sa.func.now()), + sa.PrimaryKeyConstraint("content_hash", "node_id"), + ) + op.create_index("ix_swarm_hash", "swarm_sources", ["content_hash"]) diff --git a/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/f7a8b9c0d1e2_add_content_reviews.py b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/f7a8b9c0d1e2_add_content_reviews.py new file mode 100644 index 0000000..18466b8 --- /dev/null +++ b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/f7a8b9c0d1e2_add_content_reviews.py @@ -0,0 +1,35 @@ +"""content reports wait for an administrator + +A hash reported by enough distinct accounts is queued for review instead of +being blocked on the spot, unless the instance chose automatic blocking. + +Revision ID: f7a8b9c0d1e2 +Revises: e6f7a8b9c0d1 +""" + +from collections.abc import Sequence + +import sqlalchemy as sa +from alembic import op + +revision: str = "f7a8b9c0d1e2" +down_revision: str | Sequence[str] | None = "e6f7a8b9c0d1" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.create_table( + "content_reviews", + sa.Column("content_hash", sa.String(64), nullable=False), + sa.Column("status", sa.String(16), nullable=False), + sa.Column("opened_at", sa.DateTime(timezone=True), nullable=False, + server_default=sa.func.now()), + sa.Column("decided_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("decided_by", sa.String(64), nullable=True), + sa.PrimaryKeyConstraint("content_hash"), + ) + + +def downgrade() -> None: + op.drop_table("content_reviews") diff --git a/packages/meshbay-hub/src/meshbay_hub/db/models.py b/packages/meshbay-hub/src/meshbay_hub/db/models.py index 09eb437..5b140f1 100644 --- a/packages/meshbay-hub/src/meshbay_hub/db/models.py +++ b/packages/meshbay-hub/src/meshbay_hub/db/models.py @@ -266,22 +266,6 @@ class UserDevice(Base): __table_args__ = (Index("ix_user_devices_user", "user_id"),) -class SwarmSource(Base): - """ - Tracks which nodes can serve a given content hash (public swarm). - Hub maintains this for load-balanced public content delivery. - """ - __tablename__ = "swarm_sources" - - content_hash: Mapped[str] = mapped_column(String(64), primary_key=True) - node_id: Mapped[str] = mapped_column(String(36), primary_key=True) - endpoint: Mapped[str] = mapped_column(String(128), nullable=False) - registered_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_now) - last_seen: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_now) - - __table_args__ = (Index("ix_swarm_hash", "content_hash"),) - - class Notification(Base): __tablename__ = "notifications" @@ -341,6 +325,24 @@ class ContentBlocklist(Base): added_by: Mapped[str | None] = mapped_column(String(64)) # "auto" or admin username +class ContentReview(Base): + """ + A reported hash waiting for — or given — an administrator's decision. + + Opened when enough distinct accounts have reported it; `status` is + `pending`, then `blocked` or `dismissed`. A dismissed hash stays dismissed: + more reports of it do not reopen it, and an administrator can still block + it from the blocklist. The reports themselves stay in `content_reports`. + """ + __tablename__ = "content_reviews" + + content_hash: Mapped[str] = mapped_column(String(64), primary_key=True) + status: Mapped[str] = mapped_column(String(16), default="pending") + opened_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_now) + decided_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) + decided_by: Mapped[str | None] = mapped_column(String(64)) + + class UserPreference(Base): __tablename__ = "user_preferences" diff --git a/packages/meshbay-hub/src/meshbay_hub/hub_settings.py b/packages/meshbay-hub/src/meshbay_hub/hub_settings.py index e01bfd2..d14cc0c 100644 --- a/packages/meshbay-hub/src/meshbay_hub/hub_settings.py +++ b/packages/meshbay-hub/src/meshbay_hub/hub_settings.py @@ -158,6 +158,42 @@ async def session_limits(db: AsyncSession) -> dict[str, int]: for k in SESSION_KEYS} +# ── Content reports ────────────────────────────────────────────────────────── +# +# Who may report public content, how often, and what a report leads to +# (docs/MESHBAY_DESIGN.md §7.5). `auto_block` is 0 or 1: off, a hash that +# reaches the threshold waits for an administrator; on, it is blocked at once — +# three accounts made for the purpose would then be enough to take a file down. + +REPORT_KEYS = ("min_account_age_hours", "daily_per_account", + "review_threshold", "auto_block") + +REPORT_DEFAULTS: dict[str, int] = { + "min_account_age_hours": 24, + "daily_per_account": 20, + "review_threshold": 3, + "auto_block": 0, +} + +REPORT_BOUNDS: dict[str, tuple[int, int]] = { + "min_account_age_hours": (0, 720), # 30 days + "daily_per_account": (1, 1_000), + "review_threshold": (1, 100), + "auto_block": (0, 1), +} + + +def clamp_report_value(key: str, value: int) -> int: + low, high = REPORT_BOUNDS[key] + return max(low, min(high, int(value))) + + +async def report_limits(db: AsyncSession) -> dict[str, int]: + return {k: clamp_report_value( + k, await get_int(db, f"reports.{k}", REPORT_DEFAULTS[k])) + for k in REPORT_KEYS} + + async def get_raw(db: AsyncSession, key: str) -> str | None: row = await db.get(HubSetting, key) return row.value if row else None diff --git a/packages/meshbay-hub/src/meshbay_hub/static/admin-page.js b/packages/meshbay-hub/src/meshbay_hub/static/admin-page.js index 8cbfb5a..9c32475 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/admin-page.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/admin-page.js @@ -28,6 +28,8 @@ export function AdminPage({ token, role }) { const [logEvent, setLogEvent] = useState(''); const [logOffset, setLogOffset] = useState(0); const [blocklist, setBlocklist] = useState([]); + const [reports, setReports] = useState([]); + const [reportsDraft, setReportsDraft] = useState(null); const [nodes, setNodes] = useState([]); const [detailUser, setDetailUser] = useState(null); const [error, setError] = useState(''); @@ -48,6 +50,7 @@ export function AdminPage({ token, role }) { setMailDraft({ ...data.mail }); setLoginDraft({ ...data.login }); setSessionDraft({ ...data.session }); + setReportsDraft({ ...data.reports }); } catch (e) { setError(e.message); } try { setMailStatus(await hubFetch('/v1/admin/mail', { token })); @@ -68,6 +71,7 @@ export function AdminPage({ token, role }) { setMailDraft({ ...data.mail }); setLoginDraft({ ...data.login }); setSessionDraft({ ...data.session }); + setReportsDraft({ ...data.reports }); if (patch.mail) { try { setMailStatus(await hubFetch('/v1/admin/mail', { token })); @@ -102,6 +106,23 @@ export function AdminPage({ token, role }) { } catch (e) { setError(e.message); } }, [token]); + const loadReports = useCallback(async () => { + try { + const data = await hubFetch('/v1/admin/reports', { token }); + setReports(data.reports); + } catch (e) { setError(e.message); } + }, [token]); + + // Block or dismiss one reported hash. Blocking names it on the list every + // node hosting a public group applies; dismissing closes it for good. + const decideReport = useCallback(async (hash, verdict) => { + if (verdict === 'block' && !await ask(t('admin.report_block_confirm'))) return; + try { + await hubFetch(`/v1/admin/reports/${hash}/${verdict}`, { method: 'POST', token }); + loadReports(); + } catch (e) { setError(e.message); } + }, [token]); + const loadBlocklist = useCallback(async () => { try { const data = await hubFetch('/v1/admin/blocklist', { token }); @@ -124,6 +145,7 @@ export function AdminPage({ token, role }) { .then(d => setNodes(d.nodes || [])).catch(e => setError(e.message)); } else if (tab === 'logs') { setLogOffset(0); loadLogs(logEvent, 0); } + else if (tab === 'reports') loadReports(); else if (tab === 'blocklist') loadBlocklist(); }, [tab]); @@ -205,13 +227,16 @@ const LOGIN_FIELDS = ['max_failures', 'lockout_minutes']; const SESSION_FIELDS = ['browser_idle_hours', 'refresh_idle_hours', 'max_hours']; +const REPORT_FIELDS = ['min_account_age_hours', 'daily_per_account', 'review_threshold', + 'auto_block']; + // Only what changed, and only what is a number: an empty field is someone // mid-edit, not a request to set zero. const changedNumbers = (fields, draft, stored) => Object.fromEntries(fields .filter(k => draft[k] !== '' && draft[k] !== null && Number(draft[k]) !== stored[k]) .map(k => [k, Number(draft[k])])); -const TABS = ['general', 'stats', 'users', 'groups', 'nodes', 'logs', 'blocklist']; +const TABS = ['general', 'stats', 'users', 'groups', 'nodes', 'logs', 'reports', 'blocklist']; const canEditSettings = role === 'admin'; return html` @@ -304,6 +329,37 @@ const TABS = ['general', 'stats', 'users', 'groups', 'nodes', 'logs', 'blocklist `} </div> + ${settings.reports && html` + <div class="settings-section"> + <h3 class="settings-heading">${t('admin.reports_heading')}</h3> + <p class="settings-hint">${t('admin.reports_hint')}</p> + + ${reportsDraft && REPORT_FIELDS.map(key => html` + <div class="settings-row" key=${key}> + <span class="settings-label">${t('admin.reports_' + key)}</span> + <input type="number" class="settings-number" + min=${(settings.reports_bounds?.[key] || [0])[0]} + max=${(settings.reports_bounds?.[key] || [0, 0])[1]} + value=${reportsDraft[key]} + disabled=${!canEditSettings || settingsSaving} + onInput=${e => setReportsDraft(d => ({ ...d, [key]: e.target.value }))} /> + </div> + `)} + + ${canEditSettings && reportsDraft && html` + <div class="settings-row"> + <button class="btn" disabled=${settingsSaving} + onClick=${() => saveSettings({ + reports: changedNumbers(REPORT_FIELDS, reportsDraft, settings.reports), + })}>${t('admin.reports_save')}</button> + <button class="btn btn-secondary" disabled=${settingsSaving} + onClick=${() => setReportsDraft({ ...settings.reports_defaults })} + >${t('admin.mail_reset_defaults')}</button> + </div> + `} + </div> + `} + <div class="settings-section"> <h3 class="settings-heading">${t('admin.session_heading')}</h3> <p class="settings-hint">${t('admin.session_hint')}</p> @@ -546,6 +602,40 @@ const TABS = ['general', 'stats', 'users', 'groups', 'nodes', 'logs', 'blocklist `} `} + ${tab === 'reports' && html` + <table class="admin-table"> + <thead><tr> + <th>${t('admin.col_hash')}</th> + <th>${t('admin.col_reporters')}</th> + <th>${t('admin.col_reason')}</th> + <th>${t('admin.col_groups')}</th> + <th>${t('admin.col_date')}</th> + <th>${t('admin.col_actions')}</th> + </tr></thead> + <tbody> + ${reports.length === 0 && html`<tr><td colspan="6" class="admin-empty">${t('admin.no_reports')}</td></tr>`} + ${reports.map(r => html` + <tr key=${r.hash}> + <td style="font-family:monospace;font-size:0.8em">${r.hash.slice(0, 16)}...</td> + <td>${r.reporters}</td> + <td> + ${Object.entries(r.reasons).map(([k, n]) => `${t('report.reason_' + k)} (${n})`).join(', ')} + ${r.details.map(d => html`<div class="settings-hint">${d}</div>`)} + </td> + <td>${r.groups.map(g => g.name || g.id.slice(0, 8)).join(', ')}</td> + <td>${new Date(r.opened_at).toLocaleDateString()}</td> + <td> + ${role === 'admin' && html` + <button class="admin-btn" onClick=${() => decideReport(r.hash, 'block')}>${t('admin.btn_block')}</button> + <button class="admin-btn" onClick=${() => decideReport(r.hash, 'dismiss')}>${t('admin.btn_dismiss')}</button> + `} + </td> + </tr> + `)} + </tbody> + </table> + `} + ${tab === 'blocklist' && html` <${BlocklistForm} onAdd=${addToBlocklist} /> <table class="admin-table"> diff --git a/packages/meshbay-hub/src/meshbay_hub/static/app.js b/packages/meshbay-hub/src/meshbay_hub/static/app.js index 3101be7..875b1aa 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/app.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/app.js @@ -3,7 +3,7 @@ import { clearPending, loadPending } from './invite-link.js'; import { html, render, useState, useEffect, useLayoutEffect, useCallback, useRef, - createContext, useContext, + createContext, } from './vendor/htm-preact.js'; import { t, getLocale, setLocale, initLocale, LOCALES } from './i18n.js'; import { ZipStream, entriesUnder } from './zipstream.js'; @@ -72,7 +72,6 @@ function useRoute() { // ── Context ────────────────────────────────────────────────────────────────── const AuthContext = createContext(null); -function useAuth() { return useContext(AuthContext); } // The M of the wordmark is a picture; the rest is text. Resolved from this // module's own URL so the hub's fingerprinted path and the application's @@ -304,6 +303,7 @@ function TransferRow({ it }) { </button> `} </div> + ${it.note && html`<div class="transfer-meta transfer-note">${it.note}</div>`} ${it.status === 'preparing' ? html` ${/* Not a progress bar at 0%: nothing is wrong and nothing is @@ -645,10 +645,6 @@ function LazyCreateGroupPage(props) { return html`<${_CreateGroupPage} ...${props} />`; } -// ── Settings Page ─────────────────────────────────────────────────────────── - -const THEME_OPTIONS = ['light', 'dark', 'system']; - // ── Profile Page ──────────────────────────────────────────────────────────── // // ── Lazy-loaded Admin page (admin/moderator only) ───────────────────────── diff --git a/packages/meshbay-hub/src/meshbay_hub/static/crypto.js b/packages/meshbay-hub/src/meshbay_hub/static/crypto.js index a3680ce..0ca8ee5 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/crypto.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/crypto.js @@ -1,20 +1,12 @@ /** - * MeshBay Browser Crypto — AES-256-GCM private group decryption. - * Uses WebCrypto SubtleCrypto API (available in all modern browsers). - * - * Handles groups with cipher="aes-256-gcm" (browser-accessible groups). - * ChaCha20-Poly1305 groups (cipher="chacha20-poly1305") require the - * native client (node) for decryption — not supported in browser. + * MeshBay Browser Crypto — AES-256-GCM, through WebCrypto's SubtleCrypto API. + * The one content cipher, for every client (meshbay_common/webcrypto.py). * * Usage: * const gek = await importGEK(gekB64); * const plaintext = await decryptChunkBin(gek, fileHashHex, chunkIndex, nonce, ct); */ -const CIPHER_INFO_PREFIX = new TextEncoder().encode('file:'); -const CIPHER_INFO_SUFFIX_AES = new TextEncoder().encode(':aes'); - - // ── Key derivation ──────────────────────────────────────────────────────────── /** @@ -284,38 +276,7 @@ async function verifyChatSignature(deviceRaw, groupId, epoch, nonce, ct, sig) { } -// ── GEK generation + ECIES wrapping ────────────────────────────────────────── - -function generateGEK() { - return crypto.getRandomValues(new Uint8Array(32)); -} - -async function wrapGEK(gek, pkXRaw) { - const skEph = await crypto.subtle.generateKey({ name: 'X25519' }, true, ['deriveBits']); - const pkEphRaw = new Uint8Array(await crypto.subtle.exportKey('raw', skEph.publicKey)); - - const pkRecip = await crypto.subtle.importKey('raw', pkXRaw, { name: 'X25519' }, false, []); - const sharedBits = await crypto.subtle.deriveBits( - { name: 'X25519', public: pkRecip }, skEph.privateKey, 256); - - const sharedKey = await crypto.subtle.importKey( - 'raw', sharedBits, 'HKDF', false, ['deriveKey']); - const wrapKey = await crypto.subtle.deriveKey( - { name: 'HKDF', hash: 'SHA-256', salt: pkEphRaw, - info: new TextEncoder().encode('meshbay:gek_wrap:v1:aes') }, - sharedKey, - { name: 'AES-GCM', length: 256 }, false, ['encrypt']); - - const nonce = crypto.getRandomValues(new Uint8Array(12)); - const ct = await crypto.subtle.encrypt( - { name: 'AES-GCM', iv: nonce, additionalData: pkXRaw }, wrapKey, gek); - - return { - pk_eph_b64: btoa(String.fromCharCode(...pkEphRaw)), - nonce_b64: btoa(String.fromCharCode(...nonce)), - wrapped_b64: btoa(String.fromCharCode(...new Uint8Array(ct))), - }; -} +// ── GEK unwrapping (ECIES) ───────────────────────────────────────────────────── async function unwrapGEK(bundle, skXPkcs8, pkXRaw) { const pkEphRaw = b64decode(bundle.pk_eph_b64); @@ -342,18 +303,45 @@ async function unwrapGEK(bundle, skXPkcs8, pkXRaw) { return new Uint8Array(plain); } -// ── Chunk encryption (for upload) ──────────────────────────────────────────── +function b64encode(bytes) { + return btoa(String.fromCharCode(...bytes)); +} -async function encryptChunk(gek, fileHashHex, chunkIndex, plaintext) { - const chunkKey = await deriveChunkKey(gek, fileHashHex, chunkIndex); - const nonce = crypto.getRandomValues(new Uint8Array(12)); - const ct = await crypto.subtle.encrypt( - { name: 'AES-GCM', iv: nonce }, chunkKey, plaintext); - return { nonce, ct: new Uint8Array(ct) }; +// ── Admin operation subjects ──────────────────────────────────────────────── +// Mirrors meshbay_common/adminop.py. The subject is what the signature covers of +// a request, so an operation whose effect is several values names them all. +// Canonical JSON — sorted keys, no whitespace — so both sides build the same +// bytes, and `null`, `""` and a value stay distinct. + +function adminSubject(fields) { + const sorted = {}; + for (const k of Object.keys(fields).sort()) sorted[k] = fields[k]; + return JSON.stringify(sorted); } -function b64encode(bytes) { - return btoa(String.fromCharCode(...bytes)); +// A secret named without being written: null (unchanged) and '' (clear) as +// themselves, anything else as its SHA-256. +async function secretDigest(value) { + if (!value) return value; + const d = await crypto.subtle.digest('SHA-256', new TextEncoder().encode(value)); + return 'sha256:' + Array.from(new Uint8Array(d)) + .map((b) => b.toString(16).padStart(2, '0')).join(''); +} + +function rootAddSubject(path, name, kind, writable, removable) { + return adminSubject({ path, name, kind, writable, removable }); +} + +function groupAttachSubject(name, sharedDir, writable) { + return adminSubject({ name, shared_dir: sharedDir, writable }); +} + +function inviteCreateSubject(userId, username) { + return adminSubject({ user_id: userId, username }); +} + +async function tmdbConfigSubject(token, language) { + return adminSubject({ token: await secretDigest(token), language }); } // ── Admin operation transcript ─────────────────────────────────────────────── @@ -583,8 +571,9 @@ async function verifyNodeSignature(nodePkB64, sigB64, transcript) { window.MeshBayCrypto = { importGEK, deriveChunkKey, decryptChunkBin, openGroup, sealGroup, - generateGEK, wrapGEK, unwrapGEK, encryptChunk, b64encode, b64decode, - adminTranscript, handshakeTranscript, handshakeProof, webrtcBinding, + unwrapGEK, b64encode, b64decode, + adminTranscript, adminSubject, rootAddSubject, groupAttachSubject, + inviteCreateSubject, tmdbConfigSubject, handshakeTranscript, handshakeProof, webrtcBinding, challengeTranscript, joinTranscript, verifyNodeSignature, constantTimeEqual, deviceRequestTranscript, deviceAddTranscript, deviceHelloTranscript, deviceCodeHash, diff --git a/packages/meshbay-hub/src/meshbay_hub/static/downloads.js b/packages/meshbay-hub/src/meshbay_hub/static/downloads.js index 9c23d3c..50bca46 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/downloads.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/downloads.js @@ -193,13 +193,6 @@ export async function openTarget(filename) { }; } -/** - * Below this, a download with no granted folder and no service worker is - * collected in memory and handed to the browser. Above it that would mean - * holding gigabytes in a tab, so it is worth one Save As dialog instead. - */ -export const BLOB_LIMIT = 512 * 1024 * 1024; - // ── Streaming to disk without the File System Access API ──────────────────── const SW_PATH = '/sw.js'; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/file-utils.js b/packages/meshbay-hub/src/meshbay_hub/static/file-utils.js index bdfba6e..255c162 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/file-utils.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/file-utils.js @@ -3,6 +3,7 @@ import * as platform from './platform.js'; import { t } from './i18n.js'; import { ask } from './ask.js'; import { ZipStream, entriesUnder } from './zipstream.js'; +import { portableName, portablePath } from './portable-name.js'; const FILE_ICONS = { video: '\u{1F3AC}', audio: '\u{1F3B5}', image: '\u{1F5BC}', @@ -428,6 +429,9 @@ async function pipelinedDownload(transport, gekKey, fileId, totalChunks, onChunk */ async function downloadEntry(transfers, transport, gek, entry) { const totalChunks = Math.ceil(entry.size / CHUNK_SIZE); + // Saved under a name every platform can write, and the row says so when that + // is not the node's (portable-name.js, docs/MESHBAY_DESIGN.md §10). + const saveName = portableName(entry.name); const openRef = { url: null }; let target = null; // The in-memory fallback's accumulator, held out here so a pause does not @@ -435,14 +439,15 @@ async function downloadEntry(transfers, transport, gek, entry) { const memoryChunks = new Array(totalChunks); transfers.start({ - kind: 'download', name: entry.name, total: entry.size, transport, + kind: 'download', name: saveName, total: entry.size, transport, + note: saveName !== entry.name ? t('transfers.renamed', { name: entry.name }) : '', // The row exists from the click. Opening a target is what takes the time — // the streamed path waits for the worker (twice), a Save As dialog waits // for a person — and doing it before the row meant three clicks produced no // panel at all and then several rows at once. prepare: async () => { - target = await _openTargetInTurn(entry.name, entry.size); + target = await _openTargetInTurn(saveName, entry.size); // Dismissed: nothing was started, so nothing is left on screen. if (target === false) return false; // `pausable` travels with the target, because only the target knows. The @@ -489,7 +494,7 @@ async function downloadEntry(transfers, transport, gek, entry) { transport, gek, entry.id, totalChunks, onChunk, null, signal, lease && lease.tr, from, memoryChunks); const blob = new Blob(chunks); - _saveBlob(blob, entry.name); + _saveBlob(blob, saveName); openRef.url = URL.createObjectURL(blob); } }, @@ -523,7 +528,10 @@ async function downloadDirectory(transfers, transport, gek, entries, dir, { setE return; } const totalBytes = files.reduce((n, f) => n + (f.entry.size || 0), 0); - const suggested = (dir.split('/').pop() || 'files') + '.zip'; + const suggested = portableName(dir.split('/').pop() || 'files') + '.zip'; + // Every name in the archive is made writable everywhere, or a Windows + // extraction refuses it; the row says how many changed. + const renamed = files.filter(f => portablePath(f.name) !== f.name).length; // Checked here rather than by disabling the button: Files zips a whole // multi-directory selection in one click (`for (const d of selectedDirs)`), @@ -548,6 +556,7 @@ async function downloadDirectory(transfers, transport, gek, entries, dir, { setE transfers.start({ kind: 'download', name: suggested, total: totalBytes, transport, + note: renamed ? t('transfers.renamed_n', { n: renamed }) : '', // Same order as downloadEntry: the row first, then the target, then the // slot. A folder of forty files is exactly where the wait is longest. @@ -584,7 +593,7 @@ async function downloadDirectory(transfers, transport, gek, entries, dir, { setE }); for (const { entry, name } of files) { - await zip.begin(name, entry.size, + await zip.begin(portablePath(name), entry.size, new Date((entry.added_at || 0) * 1000)); // A zero-byte file has no chunk to ask for; the header and an empty // descriptor are the whole entry. diff --git a/packages/meshbay-hub/src/meshbay_hub/static/files-app.js b/packages/meshbay-hub/src/meshbay_hub/static/files-app.js index ac978e2..cda34b5 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/files-app.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/files-app.js @@ -2,7 +2,7 @@ import { html, useState, useEffect, useRef, useCallback, } from './vendor/htm-preact.js'; import { t } from './i18n.js'; -import { ask } from './ask.js'; +import { ask, tell } from './ask.js'; import { Icon } from './icon.js'; import { entriesUnder } from './zipstream.js'; import { transfers } from './transfers.js'; @@ -12,6 +12,7 @@ import { } from './file-utils.js'; import { useStickyBand } from './sticky.js'; import { Menu, useMenu } from './menu.js'; +import { askReport } from './report.js'; // ── Files ──────────────────────────────────────────────────────────────────── // @@ -140,7 +141,7 @@ function FilesPanel({ groupId, transportRef, gekRef, status, entries, nodeDirs, nodeRoots, setEntries, setNodeDirs, setNodeRoots, applyIndex, isNodeAdmin, operatorPaired, userId, setError, onPreview, - showGroup, readOnly, getTransport, onRefreshIndex, showRefresh, + showGroup, readOnly, getTransport, onRefreshIndex, showRefresh, onReport, }) { const [selected, setSelected] = useState(() => new Set()); const [sortKey, setSortKey] = useState('name'); @@ -668,6 +669,20 @@ function FilesPanel({ ? [[], [key.slice(4)]] : [entries.filter(x => x.id === key), []]; const items = actionsFor(files, dirs, selected.has(key)).filter(a => !a.disabled); + // One file at a time, and from the menu only: a report is about a file + // somebody looked at, not a batch action for a toolbar. + const one = files.length === 1 && dirs.length === 0 ? files[0] : null; + if (onReport && one) { + items.push({ key: 'report', icon: 'shield', label: t('report.action'), + onSelect: async () => { + const answer = await askReport(one.name); + if (!answer) return; + try { + await onReport(one.id, answer.reason, answer.detail); + await tell(t('report.sent')); + } catch (err) { setError(err.message); } + } }); + } if (items.length) openAt(e, items); }; @@ -893,7 +908,6 @@ function FilesPanel({ // ── File Preview (text, images) ───────────────────────────────────────── -const TEXT_EXTS = /\.(txt|md|json|csv|log|xml|yaml|yml|ini|conf|py|js|html|css|sh|c|h|java|rs|go|rb|toml)$/i; const IMAGE_EXTS = /\.(jpg|jpeg|png|gif|webp|svg|bmp|ico)$/i; function FilePreview({ entry, transportRef, gekRef, onClose, onDownload }) { diff --git a/packages/meshbay-hub/src/meshbay_hub/static/group-page.js b/packages/meshbay-hub/src/meshbay_hub/static/group-page.js index 7c9b2f0..b6f3d37 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/group-page.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/group-page.js @@ -834,8 +834,18 @@ function GroupPage({ groupId, group, token, username, userId, userPrefs, }), [perAppDirectories, chatDirectory, chatLinkPreview, tmdbConfig, musicbrainzConfig]); + // Reporting a file is offered in a public group only: that is the only place + // the hub's moderation reaches (docs/MESHBAY_DESIGN.md §7.5), and the hub + // refuses a report from anyone who is not a member of the group named here. + const isPublic = !!(group && group.visibility === 'public'); + const reportContent = useCallback((contentHash, reason, detail) => hubFetch( + '/v1/reports', { method: 'POST', token, + body: { content_hash: contentHash, group_id: groupId, reason, detail } }), + [groupId, token]); + const commonProps = { groupId, transportRef, gekRef, status, username, deviceReady, + onReport: isPublic ? reportContent : null, entries, availableEntries, nodeDirs, nodeRoots, setEntries, setNodeDirs, setNodeRoots, applyIndex, isNodeAdmin, operatorPaired, attachRoot, attachDir, userId, setError, onPreview, diff --git a/packages/meshbay-hub/src/meshbay_hub/static/i18n.js b/packages/meshbay-hub/src/meshbay_hub/static/i18n.js index 56d35f7..bc78266 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/i18n.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/i18n.js @@ -131,10 +131,6 @@ export function setLocale(code) { return true; } -export function addLocale(code, strings) { - _strings[code] = strings; -} - function _pluralRules(locale) { if (!_plurals[locale]) _plurals[locale] = new Intl.PluralRules(locale); return _plurals[locale]; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js b/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js index 879f56f..33b1cf2 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js @@ -256,15 +256,9 @@ async function deriveRecoveryKey(R, username) { // ── Bundle encryption ───────────────────────────────────────────────────────── /** - * Encrypt the keypair bundle with the password-derived AES key. - * Bundle format: JSON { skEd: base64(pkcs8), skX: base64(pkcs8) } + * Encrypt the keypair bundle with a bundle key derived at sign-in. Always + * writes v2. Bundle format: JSON { skEd: base64(pkcs8), skX: base64(pkcs8) } */ -async function encryptBundle(skEdRaw, skXRaw, password, username) { - const aesKey = await deriveEncryptionKey(password, username); - return encryptBundleWithKey(skEdRaw, skXRaw, aesKey); -} - -/** Same, when the key was already derived at sign-in. Always writes v2. */ async function encryptBundleWithKey(skEdRaw, skXRaw, aesKey) { const nonce = crypto.getRandomValues(new Uint8Array(12)); const data = new TextEncoder().encode(JSON.stringify({ @@ -289,16 +283,6 @@ function bundleVersion(bundleB64) { } catch { return 1; } } -/** - * Decrypt a keypair bundle. Throws if password is wrong. - */ -async function decryptBundle(bundleB64, password, username) { - const key = bundleVersion(bundleB64) === 2 - ? await deriveEncryptionKey(password, username) - : await deriveEncryptionKeyV1(password, username); - return decryptBundleWithKey(bundleB64, key); -} - // ── Registration ────────────────────────────────────────────────────────────── /** diff --git a/packages/meshbay-hub/src/meshbay_hub/static/locales/de.js b/packages/meshbay-hub/src/meshbay_hub/static/locales/de.js index 0212e53..a83f44b 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/locales/de.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/locales/de.js @@ -1194,4 +1194,32 @@ export default { // In-page confirm/alert (ask.js) 'dialog.ok': 'OK', 'dialog.cancel': 'Abbrechen', + + // Content reports (report.js, admin-page.js) + 'report.action': "Melden", + 'report.title': "Diese Datei melden", + 'report.hint': "Ein Administrator dieses Hubs wird sie prüfen.", + 'report.reason_illegal': "Illegaler Inhalt", + 'report.reason_spam': "Spam", + 'report.reason_copyright': "Urheberrechtsverletzung", + 'report.reason_other': "Sonstiges", + 'report.detail_placeholder': "Details (optional)", + 'report.send': "Meldung senden", + 'report.sent': "Danke. Ihre Meldung wurde erfasst.", + 'admin.tab_reports': "Meldungen", + 'admin.no_reports': "Nichts wartet auf eine Entscheidung", + 'admin.col_reporters': "Meldende", + 'admin.btn_dismiss': "Verwerfen", + 'admin.report_block_confirm': "Diese Datei sperren? Jeder Knoten mit einer öffentlichen Gruppe stellt sie dort nicht mehr bereit.", + 'admin.reports_heading': "Inhaltsmeldungen", + 'admin.reports_hint': "Wer eine Datei in einer öffentlichen Gruppe melden darf, wie oft, und was geschieht, wenn genug Mitglieder es getan haben.", + 'admin.reports_min_account_age_hours': "Mindestalter des Kontos (Stunden)", + 'admin.reports_daily_per_account': "Meldungen pro Konto und Tag", + 'admin.reports_review_threshold': "Mitglieder bis zur Prüfung", + 'admin.reports_auto_block': "Ohne Prüfung sperren (1 = ja, 0 = nein)", + 'admin.reports_save': "Meldeeinstellungen speichern", + + // Names made writable everywhere when saved (portable-name.js) + 'transfers.renamed': "Umbenannt: „{name}“ ist nicht auf jedem System ein gültiger Name", + 'transfers.renamed_n': "Namen geändert, damit sie auf jedem System gültig sind: {n}", }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/locales/en.js b/packages/meshbay-hub/src/meshbay_hub/static/locales/en.js index ce6012c..46c5094 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/locales/en.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/locales/en.js @@ -1175,4 +1175,32 @@ export default { // In-page confirm/alert (ask.js) 'dialog.ok': 'OK', 'dialog.cancel': 'Cancel', + + // Content reports (report.js, admin-page.js) + 'report.action': "Report", + 'report.title': "Report this file", + 'report.hint': "An administrator of this hub will review it.", + 'report.reason_illegal': "Illegal content", + 'report.reason_spam': "Spam", + 'report.reason_copyright': "Copyright infringement", + 'report.reason_other': "Other", + 'report.detail_placeholder': "Details (optional)", + 'report.send': "Send report", + 'report.sent': "Thank you. Your report has been recorded.", + 'admin.tab_reports': "Reports", + 'admin.no_reports': "Nothing is waiting for a decision", + 'admin.col_reporters': "Reporters", + 'admin.btn_dismiss': "Dismiss", + 'admin.report_block_confirm': "Block this file? Every node hosting a public group will stop serving it there.", + 'admin.reports_heading': "Content reports", + 'admin.reports_hint': "Who may report a file in a public group, how often, and what happens once enough members have.", + 'admin.reports_min_account_age_hours': "Minimum account age (hours)", + 'admin.reports_daily_per_account': "Reports per account per day", + 'admin.reports_review_threshold': "Members before review", + 'admin.reports_auto_block': "Block without review (1 = yes, 0 = no)", + 'admin.reports_save': "Save report settings", + + // Names made writable everywhere when saved (portable-name.js) + 'transfers.renamed': "Renamed: “{name}” is not a valid name on every system", + 'transfers.renamed_n': "Names changed to be valid on every system: {n}", }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/locales/es.js b/packages/meshbay-hub/src/meshbay_hub/static/locales/es.js index 3ee45ac..0d0e865 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/locales/es.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/locales/es.js @@ -1188,4 +1188,32 @@ export default { // In-page confirm/alert (ask.js) 'dialog.ok': 'Aceptar', 'dialog.cancel': 'Cancelar', + + // Content reports (report.js, admin-page.js) + 'report.action': "Denunciar", + 'report.title': "Denunciar este archivo", + 'report.hint': "Un administrador de este hub lo revisará.", + 'report.reason_illegal': "Contenido ilegal", + 'report.reason_spam': "Spam", + 'report.reason_copyright': "Infracción de derechos de autor", + 'report.reason_other': "Otro", + 'report.detail_placeholder': "Detalles (opcional)", + 'report.send': "Enviar denuncia", + 'report.sent': "Gracias. Su denuncia ha quedado registrada.", + 'admin.tab_reports': "Denuncias", + 'admin.no_reports': "Nada espera una decisión", + 'admin.col_reporters': "Denunciantes", + 'admin.btn_dismiss': "Descartar", + 'admin.report_block_confirm': "¿Bloquear este archivo? Todos los nodos que alojan un grupo público dejarán de servirlo allí.", + 'admin.reports_heading': "Denuncias de contenido", + 'admin.reports_hint': "Quién puede denunciar un archivo en un grupo público, con qué frecuencia y qué ocurre cuando suficientes miembros lo han hecho.", + 'admin.reports_min_account_age_hours': "Antigüedad mínima de la cuenta (horas)", + 'admin.reports_daily_per_account': "Denuncias por cuenta y día", + 'admin.reports_review_threshold': "Miembros antes de la revisión", + 'admin.reports_auto_block': "Bloquear sin revisión (1 = sí, 0 = no)", + 'admin.reports_save': "Guardar ajustes de denuncias", + + // Names made writable everywhere when saved (portable-name.js) + 'transfers.renamed': "Renombrado: «{name}» no es un nombre válido en todos los sistemas", + 'transfers.renamed_n': "Nombres cambiados para ser válidos en todos los sistemas: {n}", }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/locales/fr.js b/packages/meshbay-hub/src/meshbay_hub/static/locales/fr.js index 0ee20e3..3c86cd7 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/locales/fr.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/locales/fr.js @@ -1203,4 +1203,32 @@ export default { // In-page confirm/alert (ask.js) 'dialog.ok': 'OK', 'dialog.cancel': 'Annuler', + + // Content reports (report.js, admin-page.js) + 'report.action': "Signaler", + 'report.title': "Signaler ce fichier", + 'report.hint': "Un administrateur de ce hub l’examinera.", + 'report.reason_illegal': "Contenu illégal", + 'report.reason_spam': "Spam", + 'report.reason_copyright': "Atteinte au droit d’auteur", + 'report.reason_other': "Autre", + 'report.detail_placeholder': "Précisions (facultatif)", + 'report.send': "Envoyer le signalement", + 'report.sent': "Merci. Votre signalement a été enregistré.", + 'admin.tab_reports': "Signalements", + 'admin.no_reports': "Rien n’attend de décision", + 'admin.col_reporters': "Signalements reçus", + 'admin.btn_dismiss': "Écarter", + 'admin.report_block_confirm': "Bloquer ce fichier ? Chaque nœud qui héberge un groupe public cessera de le servir.", + 'admin.reports_heading': "Signalement de contenu", + 'admin.reports_hint': "Qui peut signaler un fichier dans un groupe public, à quelle fréquence, et ce qui se passe quand assez de membres l’ont fait.", + 'admin.reports_min_account_age_hours': "Âge minimal du compte (heures)", + 'admin.reports_daily_per_account': "Signalements par compte et par jour", + 'admin.reports_review_threshold': "Membres avant examen", + 'admin.reports_auto_block': "Bloquer sans examen (1 = oui, 0 = non)", + 'admin.reports_save': "Enregistrer les réglages de signalement", + + // Names made writable everywhere when saved (portable-name.js) + 'transfers.renamed': "Renommé : « {name} » n’est pas un nom valide sur tous les systèmes", + 'transfers.renamed_n': "Noms modifiés pour être valides sur tous les systèmes : {n}", }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/locales/it.js b/packages/meshbay-hub/src/meshbay_hub/static/locales/it.js index f1631cd..3a74b4d 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/locales/it.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/locales/it.js @@ -1202,4 +1202,32 @@ export default { // In-page confirm/alert (ask.js) 'dialog.ok': 'OK', 'dialog.cancel': 'Annulla', + + // Content reports (report.js, admin-page.js) + 'report.action': "Segnala", + 'report.title': "Segnala questo file", + 'report.hint': "Un amministratore di questo hub lo esaminerà.", + 'report.reason_illegal': "Contenuto illegale", + 'report.reason_spam': "Spam", + 'report.reason_copyright': "Violazione del diritto d’autore", + 'report.reason_other': "Altro", + 'report.detail_placeholder': "Dettagli (facoltativo)", + 'report.send': "Invia segnalazione", + 'report.sent': "Grazie. La tua segnalazione è stata registrata.", + 'admin.tab_reports': "Segnalazioni", + 'admin.no_reports': "Nulla attende una decisione", + 'admin.col_reporters': "Segnalanti", + 'admin.btn_dismiss': "Archivia", + 'admin.report_block_confirm': "Bloccare questo file? Ogni nodo che ospita un gruppo pubblico smetterà di servirlo lì.", + 'admin.reports_heading': "Segnalazioni di contenuti", + 'admin.reports_hint': "Chi può segnalare un file in un gruppo pubblico, quanto spesso e cosa succede quando abbastanza membri lo hanno fatto.", + 'admin.reports_min_account_age_hours': "Età minima dell’account (ore)", + 'admin.reports_daily_per_account': "Segnalazioni per account al giorno", + 'admin.reports_review_threshold': "Membri prima dell’esame", + 'admin.reports_auto_block': "Blocca senza esame (1 = sì, 0 = no)", + 'admin.reports_save': "Salva impostazioni segnalazioni", + + // Names made writable everywhere when saved (portable-name.js) + 'transfers.renamed': "Rinominato: «{name}» non è un nome valido su tutti i sistemi", + 'transfers.renamed_n': "Nomi modificati per essere validi su tutti i sistemi: {n}", }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/locales/ja.js b/packages/meshbay-hub/src/meshbay_hub/static/locales/ja.js index 94e8f75..e0c38de 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/locales/ja.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/locales/ja.js @@ -1186,4 +1186,32 @@ export default { // In-page confirm/alert (ask.js) 'dialog.ok': 'OK', 'dialog.cancel': 'キャンセル', + + // Content reports (report.js, admin-page.js) + 'report.action': "報告", + 'report.title': "このファイルを報告", + 'report.hint': "このハブの管理者が確認します。", + 'report.reason_illegal': "違法なコンテンツ", + 'report.reason_spam': "スパム", + 'report.reason_copyright': "著作権侵害", + 'report.reason_other': "その他", + 'report.detail_placeholder': "詳細(任意)", + 'report.send': "報告を送信", + 'report.sent': "ありがとうございます。報告を受け付けました。", + 'admin.tab_reports': "報告", + 'admin.no_reports': "判断待ちの項目はありません", + 'admin.col_reporters': "報告者数", + 'admin.btn_dismiss': "却下", + 'admin.report_block_confirm': "このファイルをブロックしますか?公開グループをホストするすべてのノードが、そこでの提供を停止します。", + 'admin.reports_heading': "コンテンツの報告", + 'admin.reports_hint': "公開グループのファイルを誰が、どのくらいの頻度で報告できるか、そして十分な数のメンバーが報告したときに何が起こるか。", + 'admin.reports_min_account_age_hours': "アカウントの最低経過時間(時間)", + 'admin.reports_daily_per_account': "1アカウントあたり1日の報告数", + 'admin.reports_review_threshold': "確認までのメンバー数", + 'admin.reports_auto_block': "確認せずにブロック(1 = はい、0 = いいえ)", + 'admin.reports_save': "報告の設定を保存", + + // Names made writable everywhere when saved (portable-name.js) + 'transfers.renamed': "名前を変更しました:「{name}」はすべてのシステムで有効な名前ではありません", + 'transfers.renamed_n': "すべてのシステムで有効になるよう変更した名前:{n}", }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/locales/nl.js b/packages/meshbay-hub/src/meshbay_hub/static/locales/nl.js index ef97f1f..83b9eed 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/locales/nl.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/locales/nl.js @@ -1204,4 +1204,32 @@ export default { // In-page confirm/alert (ask.js) 'dialog.ok': 'OK', 'dialog.cancel': 'Annuleren', + + // Content reports (report.js, admin-page.js) + 'report.action': "Melden", + 'report.title': "Dit bestand melden", + 'report.hint': "Een beheerder van deze hub bekijkt het.", + 'report.reason_illegal': "Illegale inhoud", + 'report.reason_spam': "Spam", + 'report.reason_copyright': "Inbreuk op auteursrecht", + 'report.reason_other': "Overig", + 'report.detail_placeholder': "Details (optioneel)", + 'report.send': "Melding versturen", + 'report.sent': "Dank u. Uw melding is vastgelegd.", + 'admin.tab_reports': "Meldingen", + 'admin.no_reports': "Niets wacht op een beslissing", + 'admin.col_reporters': "Melders", + 'admin.btn_dismiss': "Afwijzen", + 'admin.report_block_confirm': "Dit bestand blokkeren? Elke node met een openbare groep stopt met het daar aanbieden.", + 'admin.reports_heading': "Inhoudsmeldingen", + 'admin.reports_hint': "Wie een bestand in een openbare groep mag melden, hoe vaak, en wat er gebeurt als genoeg leden dat hebben gedaan.", + 'admin.reports_min_account_age_hours': "Minimale leeftijd van het account (uren)", + 'admin.reports_daily_per_account': "Meldingen per account per dag", + 'admin.reports_review_threshold': "Leden vóór beoordeling", + 'admin.reports_auto_block': "Blokkeren zonder beoordeling (1 = ja, 0 = nee)", + 'admin.reports_save': "Meldingsinstellingen opslaan", + + // Names made writable everywhere when saved (portable-name.js) + 'transfers.renamed': "Hernoemd: ‘{name}’ is niet op elk systeem een geldige naam", + 'transfers.renamed_n': "Namen aangepast zodat ze op elk systeem geldig zijn: {n}", }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/locales/pl.js b/packages/meshbay-hub/src/meshbay_hub/static/locales/pl.js index 5a2257b..3edba2d 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/locales/pl.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/locales/pl.js @@ -1230,4 +1230,32 @@ export default { // In-page confirm/alert (ask.js) 'dialog.ok': 'OK', 'dialog.cancel': 'Anuluj', + + // Content reports (report.js, admin-page.js) + 'report.action': "Zgłoś", + 'report.title': "Zgłoś ten plik", + 'report.hint': "Administrator tego huba go sprawdzi.", + 'report.reason_illegal': "Treść nielegalna", + 'report.reason_spam': "Spam", + 'report.reason_copyright': "Naruszenie praw autorskich", + 'report.reason_other': "Inne", + 'report.detail_placeholder': "Szczegóły (opcjonalnie)", + 'report.send': "Wyślij zgłoszenie", + 'report.sent': "Dziękujemy. Zgłoszenie zostało zapisane.", + 'admin.tab_reports': "Zgłoszenia", + 'admin.no_reports': "Nic nie czeka na decyzję", + 'admin.col_reporters': "Zgłaszający", + 'admin.btn_dismiss': "Odrzuć", + 'admin.report_block_confirm': "Zablokować ten plik? Każdy węzeł hostujący grupę publiczną przestanie go tam udostępniać.", + 'admin.reports_heading': "Zgłoszenia treści", + 'admin.reports_hint': "Kto może zgłosić plik w grupie publicznej, jak często i co się dzieje, gdy zrobi to wystarczająco wielu członków.", + 'admin.reports_min_account_age_hours': "Minimalny wiek konta (godziny)", + 'admin.reports_daily_per_account': "Zgłoszenia na konto dziennie", + 'admin.reports_review_threshold': "Członkowie przed oceną", + 'admin.reports_auto_block': "Blokuj bez oceny (1 = tak, 0 = nie)", + 'admin.reports_save': "Zapisz ustawienia zgłoszeń", + + // Names made writable everywhere when saved (portable-name.js) + 'transfers.renamed': "Zmieniono nazwę: „{name}” nie jest prawidłową nazwą w każdym systemie", + 'transfers.renamed_n': "Nazwy zmienione, by były prawidłowe w każdym systemie: {n}", }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/locales/pt-BR.js b/packages/meshbay-hub/src/meshbay_hub/static/locales/pt-BR.js index 65d2102..3f44570 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/locales/pt-BR.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/locales/pt-BR.js @@ -1189,4 +1189,32 @@ export default { // In-page confirm/alert (ask.js) 'dialog.ok': 'OK', 'dialog.cancel': 'Cancelar', + + // Content reports (report.js, admin-page.js) + 'report.action': "Denunciar", + 'report.title': "Denunciar este arquivo", + 'report.hint': "Um administrador deste hub vai analisá-lo.", + 'report.reason_illegal': "Conteúdo ilegal", + 'report.reason_spam': "Spam", + 'report.reason_copyright': "Violação de direitos autorais", + 'report.reason_other': "Outro", + 'report.detail_placeholder': "Detalhes (opcional)", + 'report.send': "Enviar denúncia", + 'report.sent': "Obrigado. Sua denúncia foi registrada.", + 'admin.tab_reports': "Denúncias", + 'admin.no_reports': "Nada aguarda decisão", + 'admin.col_reporters': "Denunciantes", + 'admin.btn_dismiss': "Descartar", + 'admin.report_block_confirm': "Bloquear este arquivo? Todos os nós que hospedam um grupo público deixarão de servi-lo ali.", + 'admin.reports_heading': "Denúncias de conteúdo", + 'admin.reports_hint': "Quem pode denunciar um arquivo em um grupo público, com que frequência e o que acontece quando membros suficientes o fizeram.", + 'admin.reports_min_account_age_hours': "Idade mínima da conta (horas)", + 'admin.reports_daily_per_account': "Denúncias por conta por dia", + 'admin.reports_review_threshold': "Membros antes da análise", + 'admin.reports_auto_block': "Bloquear sem análise (1 = sim, 0 = não)", + 'admin.reports_save': "Salvar configurações de denúncias", + + // Names made writable everywhere when saved (portable-name.js) + 'transfers.renamed': "Renomeado: “{name}” não é um nome válido em todos os sistemas", + 'transfers.renamed_n': "Nomes alterados para serem válidos em todos os sistemas: {n}", }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/locales/zh-CN.js b/packages/meshbay-hub/src/meshbay_hub/static/locales/zh-CN.js index 49ce189..1168460 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/locales/zh-CN.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/locales/zh-CN.js @@ -1175,4 +1175,32 @@ export default { // In-page confirm/alert (ask.js) 'dialog.ok': '确定', 'dialog.cancel': '取消', + + // Content reports (report.js, admin-page.js) + 'report.action': "举报", + 'report.title': "举报此文件", + 'report.hint': "本中心的管理员会进行审核。", + 'report.reason_illegal': "违法内容", + 'report.reason_spam': "垃圾信息", + 'report.reason_copyright': "侵犯版权", + 'report.reason_other': "其他", + 'report.detail_placeholder': "详细说明(可选)", + 'report.send': "提交举报", + 'report.sent': "谢谢,您的举报已记录。", + 'admin.tab_reports': "举报", + 'admin.no_reports': "暂无待处理事项", + 'admin.col_reporters': "举报人数", + 'admin.btn_dismiss': "驳回", + 'admin.report_block_confirm': "要屏蔽此文件吗?所有托管公开群组的节点都将停止在那里提供它。", + 'admin.reports_heading': "内容举报", + 'admin.reports_hint': "谁可以举报公开群组中的文件、举报频率,以及足够多的成员举报后会发生什么。", + 'admin.reports_min_account_age_hours': "账户最短注册时长(小时)", + 'admin.reports_daily_per_account': "每个账户每天的举报数", + 'admin.reports_review_threshold': "进入审核所需人数", + 'admin.reports_auto_block': "无需审核直接屏蔽(1 = 是,0 = 否)", + 'admin.reports_save': "保存举报设置", + + // Names made writable everywhere when saved (portable-name.js) + 'transfers.renamed': "已重命名:“{name}”并非在所有系统上都是有效的名称", + 'transfers.renamed_n': "为在所有系统上有效而修改的名称:{n}", }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/portable-name.js b/packages/meshbay-hub/src/meshbay_hub/static/portable-name.js new file mode 100644 index 0000000..bb2438c --- /dev/null +++ b/packages/meshbay-hub/src/meshbay_hub/static/portable-name.js @@ -0,0 +1,51 @@ +/** + * A name that can be written on every platform a member saves to + * (docs/MESHBAY_DESIGN.md §10). + * + * A node serves a file under the name its own disk gave it, and a name that is + * fine on ext4 can be impossible on Windows or on an exFAT drive: reserved + * characters, a trailing dot or space, `CON` or `aux.txt`. The node never + * rewrites a name — it is the string that opens the file — so the client does, + * at the moment it saves, and says so. + * + * The same rules as `meshbay_common.paths.sanitize_for_download`, and held to + * them by `test_portable_name_parity.py`: two copies of a rule that differ + * decide differently which files get renamed. + */ + +const WINDOWS_RESERVED = new Set([ + 'CON', 'PRN', 'AUX', 'NUL', + ...[1, 2, 3, 4, 5, 6, 7, 8, 9].map((i) => `COM${i}`), + ...[1, 2, 3, 4, 5, 6, 7, 8, 9].map((i) => `LPT${i}`), +]); + +const RESERVED_CHARS = new Set('<>:"/\\|?*'); + +const reserved = (c) => RESERVED_CHARS.has(c) || c.charCodeAt(0) < 32; + +function isPortable(name) { + if (!name || name === '.' || name === '..') return false; + for (const c of name) if (reserved(c)) return false; + if (name.endsWith(' ') || name.endsWith('.')) return false; + return !WINDOWS_RESERVED.has(name.split('.', 1)[0].toUpperCase()); +} + +/** `name` if it can be written everywhere, otherwise the nearest name that can. */ +export function portableName(name, replacement = '_') { + const text = String(name ?? ''); + if (isPortable(text)) return text; + let out = Array.from(text).map((c) => (reserved(c) ? replacement : c)).join(''); + out = out.replace(/[ .]+$/, ''); + const dot = out.indexOf('.'); + let stem = dot === -1 ? out : out.slice(0, dot); + const rest = dot === -1 ? '' : out.slice(dot); + if (WINDOWS_RESERVED.has(stem.toUpperCase())) stem += replacement; + out = stem + rest; + return out || 'unnamed'; +} + +/** Every segment of a relative path made portable, keeping `/` between them. */ +export function portablePath(path) { + return String(path ?? '').split('/') + .map((segment) => (segment ? portableName(segment) : segment)).join('/'); +} diff --git a/packages/meshbay-hub/src/meshbay_hub/static/report.js b/packages/meshbay-hub/src/meshbay_hub/static/report.js new file mode 100644 index 0000000..f433780 --- /dev/null +++ b/packages/meshbay-hub/src/meshbay_hub/static/report.js @@ -0,0 +1,70 @@ +import { html, render, useEffect, useRef, useState } from './vendor/htm-preact.js'; +import { t } from './i18n.js'; + +/** + * Ask why a file is being reported, drawn by the page like `ask.js`. + * + * Resolves `{ reason, detail }`, or null when cancelled. The reasons are the + * hub's own closed list (api/moderation.py `ReportRequest`), and the detail is + * bounded to what the hub keeps (256 characters). + */ + +const REASONS = ['illegal', 'spam', 'copyright', 'other']; + +function ReportDialog({ name, onDone }) { + const [reason, setReason] = useState('illegal'); + const [detail, setDetail] = useState(''); + const firstRef = useRef(null); + useEffect(() => { if (firstRef.current) firstRef.current.focus(); }, []); + + return html` + <div class="video-overlay" onClick=${(e) => { + if (e.target.classList.contains('video-overlay')) onDone(null); + }}> + <form class="music-detail playlist-modal" role="dialog" aria-modal="true" + onKeyDown=${(e) => { if (e.key === 'Escape') { e.preventDefault(); onDone(null); } }} + onSubmit=${(e) => { + e.preventDefault(); + onDone({ reason, detail: detail.trim() || null }); + }}> + <div class="playlist-modal-body"> + <div class="ask-message"><strong>${t('report.title')}</strong></div> + <div class="ask-message file-name">${name}</div> + <p class="settings-hint">${t('report.hint')}</p> + <select ref=${firstRef} value=${reason} + onChange=${(e) => setReason(e.target.value)}> + ${REASONS.map((r) => html` + <option key=${r} value=${r}>${t('report.reason_' + r)}</option>`)} + </select> + <textarea maxlength="256" placeholder=${t('report.detail_placeholder')} + value=${detail} onInput=${(e) => setDetail(e.target.value)} /> + <div class="playlist-modal-actions"> + <button type="button" class="tb-btn" onClick=${() => onDone(null)}> + ${t('dialog.cancel')}</button> + <button type="submit" class="admin-btn">${t('report.send')}</button> + </div> + </div> + </form> + </div> + `; +} + +export function askReport(name) { + return new Promise((resolve) => { + const host = document.createElement('div'); + document.body.appendChild(host); + const previous = document.activeElement; + let settled = false; + const onDone = (value) => { + if (settled) return; + settled = true; + render(null, host); + host.remove(); + if (previous && previous.isConnected && typeof previous.focus === 'function') { + previous.focus(); + } + resolve(value); + }; + render(html`<${ReportDialog} name=${String(name)} onDone=${onDone} />`, host); + }); +} diff --git a/packages/meshbay-hub/src/meshbay_hub/static/style.css b/packages/meshbay-hub/src/meshbay_hub/static/style.css index 8fcb9cb..fc91596 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/style.css +++ b/packages/meshbay-hub/src/meshbay_hub/static/style.css @@ -2588,6 +2588,7 @@ a.transfer-name { margin-top: 3px; } .transfer-failed { color: var(--error); } +.transfer-note { justify-content: flex-start; font-style: italic; } /* ── File selection ──────────────────────────────────────────────────────── */ @@ -5351,7 +5352,9 @@ h2 .gn-owner, h3 .gn-owner { font-size: 0.55em; } .playlist-modal { max-width: 420px; } .playlist-modal-body { padding: 16px; display: flex; flex-direction: column; gap: 12px; } -.playlist-modal-body input { +.playlist-modal-body input, +.playlist-modal-body select, +.playlist-modal-body textarea { width: 100%; padding: 9px 12px; border: 1px solid var(--border); @@ -5360,7 +5363,10 @@ h2 .gn-owner, h3 .gn-owner { font-size: 0.55em; } color: var(--text); font: inherit; } -.playlist-modal-body input:focus { +.playlist-modal-body textarea { resize: vertical; min-height: 4.5em; } +.playlist-modal-body input:focus, +.playlist-modal-body select:focus, +.playlist-modal-body textarea:focus { outline: none; border-color: var(--border-focus); } diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transfers.js b/packages/meshbay-hub/src/meshbay_hub/static/transfers.js index 9a1455a..0d012ed 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/transfers.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/transfers.js @@ -29,13 +29,6 @@ function _live(status) { || status === 'paused'; } -/** Raised by `run` when it stopped because the transfer was paused. */ -function _pausedError() { - const err = new Error('Paused'); - err.name = 'PausedError'; - return err; -} - function _abortError() { const err = new Error('Cancelled'); err.name = 'AbortError'; @@ -80,6 +73,7 @@ export class TransferStore { queuedByOwnLimit: Boolean( it.lease && it.lease.cap && it.lease.used >= it.lease.cap), error: it.error || '', + note: it.note || '', speed: this._speed(it), // The ETA is drawn only once the window holds a few seconds of real // measurement -- see etaSeconds. @@ -141,10 +135,14 @@ export class TransferStore { * there is somewhere to write — see file-utils.js's downloadEntry. */ start({ kind, name, total = 0, transport = null, run, open = null, - lease = null, prepare = null, makeLease = null, pausable = false }) { + lease = null, prepare = null, makeLease = null, pausable = false, + note = '' }) { const item = { id: _nextId++, kind, name, total, transport, open, lease, + // One line said under the name for the life of the row — that a name was + // changed to be written here, for instance. + note, done: 0, // A transfer that has to wait for a slot starts as 'queued', not // 'running'. Two different things are true of it — nothing is moving, and diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport-admin.js b/packages/meshbay-hub/src/meshbay_hub/static/transport-admin.js index 03a0f33..c7c3f47 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/transport-admin.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/transport-admin.js @@ -263,7 +263,10 @@ extendTransport(class { }); if (msg.type === 'error') throw new Error(msg.detail); if (msg.type === 'admin_challenge') { - return this._authorizeAdminOp(msg, 'root_add', path, signFn); + // Everything the node will act on is in the subject, `writable` included. + const subject = window.MeshBayCrypto.rootAddSubject( + path, name || '', kind || 'generic', !!writable, !!removable); + return this._authorizeAdminOp(msg, 'root_add', subject, signFn); } return msg; } @@ -369,11 +372,13 @@ extendTransport(class { async attachGroup(name, sharedDir, uploadDir, signFn) { const msg = await this._sendAndWait({ type: 'group_attach', v: '0.1', - name, shared_dir: sharedDir, upload_dir: uploadDir || '', + name, shared_dir: sharedDir, upload_dir: uploadDir || '', writable: true, }); if (msg.type === 'error') throw new Error(msg.detail); if (msg.type === 'admin_challenge') { - return this._authorizeAdminOp(msg, 'group_attach', name, signFn); + // The directory being exposed is signed, not only the group's name. + const subject = window.MeshBayCrypto.groupAttachSubject(name, sharedDir, true); + return this._authorizeAdminOp(msg, 'group_attach', subject, signFn); } return msg; } @@ -416,7 +421,10 @@ extendTransport(class { }); if (msg.type === 'error') throw new Error(msg.detail); if (msg.type === 'admin_challenge') { - return this._authorizeAdminOp(msg, 'invite_create', userId, signFn); + // The node keeps the first 64 code points of the name, as Python slices. + const name = Array.from(username || '').slice(0, 64).join(''); + const subject = window.MeshBayCrypto.inviteCreateSubject(userId, name); + return this._authorizeAdminOp(msg, 'invite_create', subject, signFn); } return msg; } diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport-chat.js b/packages/meshbay-hub/src/meshbay_hub/static/transport-chat.js index d6d733f..270c76e 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/transport-chat.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/transport-chat.js @@ -212,7 +212,6 @@ extendTransport(class { if (msg.epoch) this.chatEpoch = msg.epoch; this._chatKeys = null; this._chatKeysInFlight = null; - if (this._onChatEpoch) this._onChatEpoch(this.chatEpoch); } /** diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport-devices.js b/packages/meshbay-hub/src/meshbay_hub/static/transport-devices.js index 4291cf8..199b6a2 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/transport-devices.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/transport-devices.js @@ -268,6 +268,10 @@ extendTransport(class { * The counterpart of storeKeypairBundle: turning the setting off has to remove * what is already stored, not merely stop adding to it — otherwise the blob * stays on every node the account has ever joined (C4). + * + * Nothing calls this yet, on purpose: it is reserved for `device_policy` + * (docs/MESHBAY_DESIGN.md §3.7, O3). Offered alone, it would strand the next + * browser that signs in to this node. */ async deleteKeypairBundle() { const msg = await this._sendAndWait({ diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport-media.js b/packages/meshbay-hub/src/meshbay_hub/static/transport-media.js index f94eb40..cf53c08 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/transport-media.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/transport-media.js @@ -106,18 +106,17 @@ extendTransport(class { * `language`, to leave whatever is stored unchanged. */ async setTmdbConfig(token, language, signFn) { + const tok = token === undefined ? null : token; + const lang = language === undefined ? null : language; const msg = await this._sendAndWait({ - type: 'tmdb_config', v: '0.7', - token: token === undefined ? null : token, - language: language === undefined ? null : language, + type: 'tmdb_config', v: '0.7', token: tok, language: lang, }); if (msg.type === 'error') throw new Error(msg.detail); if (msg.type === 'admin_challenge') { - // Must match the node's subject byte-for-byte (apps/video_meta.py - // _do_tmdb_config) — the token itself is never part of the subject - // (it would end up in the audit log in plaintext), only whether one - // was supplied. The language is not a secret, so it appears as-is. - const subject = `custom_token=${token ? 'yes' : 'no'},language=${language || 'default'}`; + // Must match the node's subject byte for byte (apps/video_meta.py + // _do_tmdb_config). The token is named by its SHA-256, never written: + // the subject ends up in the audit log. + const subject = await window.MeshBayCrypto.tmdbConfigSubject(tok, lang); return this._authorizeAdminOp(msg, 'tmdb_config', subject, signFn); } return msg; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport.js b/packages/meshbay-hub/src/meshbay_hub/static/transport.js index 546d01b..98d010e 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/transport.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/transport.js @@ -305,8 +305,10 @@ window.addEventListener('hashchange', () => { // The `v: '0.1'` on every other message in this file is the historical value // and is read by nothing; it is left alone deliberately. The range is // negotiated once, at the start, not restated per message. -const MNP_V = '4.0'; -// Raised with it: 4.0 is a flag day. A member now presents a short-lived +const MNP_V = '5.0'; +// Not raised with 5.0 (see meshbay_common/__init__.py): the break is confined to +// four signed operations, which a peer on the other side of it refuses to sign. +// Set at 4.0, a flag day. A member now presents a short-lived // MNP-audience token in the handshake, not its hub session token — a node // older than 4.0 expected the session token, and one newer refuses it, so the // two cannot authenticate across the break. This is the C6 rule: no @@ -498,11 +500,6 @@ class MeshBayTransport { this._inFlightUploads = new Set(); // tr → Lease. A transfer's slot on the node, from the client's side. this._leases = new Map(); - // Set from the handshake ack: a node that answers with `transfer_limits` - // speaks transfer slots. Used instead of a timeout, because "no answer - // yet" and "this node will never answer" are indistinguishable in time and - // guessing wrong either stalls every download or defeats the cap. - this._transferLimits = null; // Set once close() runs — stops the automatic reconnect from firing on a // connection the caller tore down on purpose (leaving the group, page // unload), which would otherwise race back in right as everything else @@ -571,18 +568,11 @@ class MeshBayTransport { set onIndexDelta(fn) { this._onIndexDelta = fn; } set onRootsChanged(fn) { this._onRootsChanged = fn; } - /** The MNP version the connected node declared, or '' before a handshake. */ - get nodeVersion() { return this._nodeVersion || ''; } - - /** This member's own caps in this group, or null when the node said nothing. */ - get transferLimits() { return this._transferLimits; } - set onAppsEnabled(fn) { this._onAppsEnabled = fn; } set onAppDirectories(fn) { this._onAppDirectories = fn; } set onChatDirectory(fn) { this._onChatDirectory = fn; } set onChatLinkPreview(fn) { this._onChatLinkPreview = fn; } set onSearchListed(fn) { this._onSearchListed = fn; } - set onChatEpoch(fn) { this._onChatEpoch = fn; } set onTmdbConfig(fn) { this._onTmdbConfig = fn; } set onTmdbEnabled(fn) { this._onTmdbEnabled = fn; } set onMusicbrainzEnabled(fn) { this._onMusicbrainzEnabled = fn; } @@ -941,12 +931,9 @@ class MeshBayTransport { if (reply.type === 'handshake_challenge') { // The node's half of the range. Checked before anything else in this // block, because everything below — the join, the proof, the sealed ack - // — assumes both sides mean the same thing by each message. + // — assumes both sides mean the same thing by each message. Nothing else + // reads the node's version: a peer this admits speaks every message here. _checkNodeVersion(reply); - // Kept for diagnostics only. Nothing branches on it: the range check - // above is what decides whether these two can talk at all, and a peer it - // admits speaks every message in this file. - this._nodeVersion = String(reply.v || ''); if (!window.MeshBayCrypto) { throw new Error('Node requires GEK proof but no crypto available'); } @@ -959,16 +946,14 @@ class MeshBayTransport { // nonce_node ties a join to this connection, so one cannot be lifted onto // another. node_pk is announced here because a first-time member has no // GEK and so cannot complete the handshake that would prove it. From an - // older node it is unverified until the ack below checks it. + // The signature checked next is what proves it here, before the ack. this._nonceNode = window.MeshBayCrypto.b64decode(reply.nonce); this.nodePk = reply.node_pk || null; - // Since MNP 3.4 the node signs its challenge over this connection, so - // node_pk is proved here and not only at the ack — which comes after any - // join. A signature that does not verify is a peer lying about which node - // it is, and is refused. An absent one is an older node: `nodePkProved` - // stays false, and whatever needs the key proved before a code leaves - // (an invitation link names its node) reads that — never a version. - this.nodePkProved = await _challengeProvesNodeKey( + // The node signs its challenge over this connection, so node_pk is proved + // here and not only at the ack — which comes after any join. Every node + // this client can reach signs (the floor is 4.0, and signing is 3.4), so + // a missing signature is refused exactly like a wrong one. + await _challengeProvesNodeKey( reply, groupId || '', this._nonceClient, this._pc.localDescription.sdp, this._rawAnswerSdp); @@ -1064,8 +1049,7 @@ class MeshBayTransport { // (docs/MESHBAY_DESIGN.md §3.4). Otherwise nothing is sent at all — not // even a join without the code, which this node would answer by asking // for one. - const linkRefusal = _linkJoinRefusal(joinNodePk, joinCode, this.nodePk, - this.nodePkProved); + const linkRefusal = _linkJoinRefusal(joinNodePk, joinCode, this.nodePk); if (linkRefusal) { this._joinError = linkRefusal; } else if (!gekRaw && this._sessionKeys && userId) { @@ -1165,7 +1149,6 @@ class MeshBayTransport { delete ack.nonce; delete ack.ct; Object.assign(ack, config); - this._transferLimits = ack.transfer_limits || null; // From the *sealed* part of the ack: a forged epoch would have this // client sealing under a key the group has retired. @@ -2198,35 +2181,29 @@ class MeshBayTransport { * Why a code from an invitation link must not go to this node, or null. * * `link_other_node` is the caller's cue to try the next node the hub listed, - * as for `not_hosted`: the link names one node, and this is not it. An older - * node that cannot prove its key early is refused rather than trusted — it - * cannot have issued a link code anyway. + * as for `not_hosted`: the link names one node, and this is not it. `nodePk` + * has already been proved by the challenge signature, which is required. */ -function _linkJoinRefusal(joinNodePk, joinCode, nodePk, nodePkProved) { +function _linkJoinRefusal(joinNodePk, joinCode, nodePk) { if (!joinNodePk || !joinCode) return null; if (nodePk !== joinNodePk) { const err = new Error('This invitation was issued by another machine hosting this group.'); err.reason = 'link_other_node'; return err; } - if (!nodePkProved) { - const err = new Error('This node is too old to accept invitation links.'); - err.reason = 'link_node_unproved'; - return err; - } return null; } /** - * Whether `handshake_challenge` proves the key it announces (MNP 3.4). + * Check that `handshake_challenge` proves the key it announces; throw if not. * - * True when it carries a signature that verifies over this connection, false - * when it carries none — an older node, which proves its key only at the ack. - * A signature that does not verify is a peer lying about which node it is, and - * throws: that is a refusal, not a node that merely cannot say. + * A node signs whenever it has a channel binding, and one without a binding + * could not complete the handshake anyway (its proof is refused), so a missing + * signature is refused like a wrong one: both are a peer that cannot show it is + * the node it names. There is no "older node" case — the floor is 4.0. */ async function _challengeProvesNodeKey(reply, groupId, nonceClient, offerSdp, answerSdp) { - if (!reply.sig) return false; + if (!reply.sig) throw new Error('Node challenge is not signed — refusing connection'); const C = window.MeshBayCrypto; let ok = false; try { diff --git a/packages/meshbay-hub/src/meshbay_hub/static/video-player.js b/packages/meshbay-hub/src/meshbay_hub/static/video-player.js index 2b274b5..fa10d15 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/video-player.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/video-player.js @@ -270,12 +270,6 @@ function reconnectPlan(playhead, range, ended, castActive) { return { mode: 'seek', at: playhead }; } -function _mseSupported(codec) { - if (!window.MediaSource) return false; - const mime = `video/mp4; codecs="${codec}"`; - return MediaSource.isTypeSupported(mime); -} - /** Seconds as h:mm:ss, or m:ss under an hour. */ function formatClock(seconds) { const s = Math.max(0, Math.floor(seconds || 0)); diff --git a/packages/meshbay-hub/tests/test_account_deletion.py b/packages/meshbay-hub/tests/test_account_deletion.py index 2653f0d..64c2be8 100644 --- a/packages/meshbay-hub/tests/test_account_deletion.py +++ b/packages/meshbay-hub/tests/test_account_deletion.py @@ -130,16 +130,9 @@ def _device_pk() -> str: @pytest.mark.asyncio -async def test_deletion_clears_device_keys_and_swarm_sources(client, db_session): - """ - The privacy statement says every account row goes but the IP log. Swarm - sources are keyed by the *user* id despite the column's name, and carry the - node's transport and port — `webrtc:<port>`, which is what `daemon.py` - actually sends. This asked with `192.0.2.7:4433`, from the days when the - field was free text documented as "ip:port": a shape no node has ever - produced, and one that let a caller name a third party's address. - """ - from meshbay_hub.db.models import SwarmSource, UserDevice +async def test_deletion_clears_device_keys(client, db_session): + """The privacy statement says every account row goes but the IP log.""" + from meshbay_hub.db.models import UserDevice token, password = await _register(client, "devicer_test") headers = {"Authorization": f"Bearer {token}"} @@ -149,15 +142,10 @@ async def test_deletion_clears_device_keys_and_swarm_sources(client, db_session) r = await client.post("/v1/users/devices", headers=headers, json={"pk_auth_ed25519": _device_pk(), "label": "desktop"}) assert r.status_code == 201, r.text - r = await client.post("/v1/swarm/register", headers=headers, - json={"content_hash": "ab" * 32, "endpoint": "webrtc:4433"}) - assert r.status_code == 201, r.text # Present before, or the emptiness asserted below proves nothing. assert (await db_session.execute( select(UserDevice).where(UserDevice.user_id == uid))).scalars().all() - assert (await db_session.execute( - select(SwarmSource).where(SwarmSource.node_id == uid))).scalars().all() r = await client.request("DELETE", "/v1/users/me", headers=headers, json={"auth_key": _auth_key(password, "devicer_test")}) @@ -166,8 +154,6 @@ async def test_deletion_clears_device_keys_and_swarm_sources(client, db_session) db_session.expire_all() assert (await db_session.execute( select(UserDevice).where(UserDevice.user_id == uid))).scalars().all() == [] - assert (await db_session.execute( - select(SwarmSource).where(SwarmSource.node_id == uid))).scalars().all() == [] @pytest.mark.asyncio diff --git a/packages/meshbay-hub/tests/test_availability_between_members.py b/packages/meshbay-hub/tests/test_availability_between_members.py index 24d85a0..e66be31 100644 --- a/packages/meshbay-hub/tests/test_availability_between_members.py +++ b/packages/meshbay-hub/tests/test_availability_between_members.py @@ -202,56 +202,6 @@ async def test_the_notify_budget_is_not_refilled_by_reconnecting(client): rev._notify_window.pop(node_id, None) -# ── A member must not aim other people's traffic ───────────────────────────── - -@pytest.mark.asyncio -async def test_a_swarm_source_cannot_name_someone_elses_address(client): - """ - `endpoint` was free text documented as "ip:port", so an account could - publish a third party's address as a source for any content. Nothing dials - a swarm source today, which is the only reason this was not already the - reflection primitive that `notify_incoming` was fixed for (H6). A port is - all a reader needs: where the node is comes from the node record, which is - stamped with the address its announce arrived from. - """ - user = await _make_user(client, "av_swarm1") - headers = {"Authorization": f"Bearer {user['token']}"} - - for bad in ("192.0.2.7:4433", "evil.example:53", "webrtc:0", "webrtc:70000", - "webrtc:4433 ", "http://example.test"): - r = await client.post("/v1/swarm/register", headers=headers, - json={"content_hash": "ab" * 32, "endpoint": bad}) - assert r.status_code == 422, f"{bad!r} was accepted: {r.text}" - - r = await client.post("/v1/swarm/register", headers=headers, - json={"content_hash": "ab" * 32, "endpoint": "webrtc:19010"}) - assert r.status_code == 201, r.text - - -@pytest.mark.asyncio -async def test_one_account_cannot_fill_the_swarm_table(client, monkeypatch): - """Rows are keyed (hash, account) with no cap — an invented hash each time.""" - import meshbay_hub.api.groups as groups_api - monkeypatch.setattr(groups_api, "MAX_SWARM_HASHES_PER_ACCOUNT", 3) - - user = await _make_user(client, "av_swarm2") - headers = {"Authorization": f"Bearer {user['token']}"} - for i in range(3): - r = await client.post("/v1/swarm/register", headers=headers, - json={"content_hash": f"{i:064x}", - "endpoint": "webrtc:19010"}) - assert r.status_code == 201, r.text - - r = await client.post("/v1/swarm/register", headers=headers, - json={"content_hash": f"{99:064x}", "endpoint": "webrtc:19010"}) - assert r.status_code == 429, r.text - - # Refreshing one already held is not a new claim and must still work. - r = await client.post("/v1/swarm/register", headers=headers, - json={"content_hash": f"{0:064x}", "endpoint": "webrtc:19011"}) - assert r.status_code == 201, r.text - - # ── A member's node must not answer for another's ──────────────────────────── def test_a_node_cannot_answer_an_offer_it_was_never_sent(): @@ -321,75 +271,6 @@ async def test_an_invite_email_says_what_the_hub_knows_not_what_it_is_told( "the sender chose the subject line of a message the hub signs") -# ── A relay is not authenticated by the key it publishes ───────────────────── - -@pytest.mark.asyncio -async def test_a_relay_must_prove_it_holds_the_approved_key(client, monkeypatch): - """ - `relay_register` had no `Depends` and verified nothing: it compared - `pk_relay` against the approved value, which is a **public** key. Anyone - who could read it could rewrite where the hub tells nodes to send relayed - traffic — an unauthenticated write to state other people's machines act - on. The module docstring said the relay "signs keepalive JWTs"; `jwt` was - imported and never used. - """ - from meshbay_hub.api import relay as relay_mod - - # The registry ships closed (`relay.RELAYS_ENABLED`); the proof it demands - # is still what will be wanted the day it opens. - monkeypatch.setattr(relay_mod, "RELAYS_ENABLED", True) - sk = Ed25519PrivateKey.generate() - pk = pk_to_b64(sk.public_key()) - relay_mod._relays["r1"] = {"pk": pk, "active": False} - try: - # The public key alone, which used to be enough. - r = await client.post("/v1/relays/register", json={ - "relay_id": "r1", "endpoint": "198.51.100.9:9999", - "pk_relay": pk, "capacity": 100}) - assert r.status_code == 400, r.text - assert relay_mod._relays["r1"].get("endpoint") is None - - # A signature over someone else's endpoint does not carry either: the - # endpoint is inside the signed message. - ts = int(time.time()) - sig = sk.sign(f"meshbay:relay_register:r1:10.0.0.1:4433:{ts}".encode()) - r = await client.post("/v1/relays/register", json={ - "relay_id": "r1", "endpoint": "198.51.100.9:9999", "pk_relay": pk, - "timestamp": ts, "signature": base64.b64encode(sig).decode()}) - assert r.status_code == 401, r.text - - endpoint = "203.0.113.4:4433" - sig = sk.sign(f"meshbay:relay_register:r1:{endpoint}:{ts}".encode()) - r = await client.post("/v1/relays/register", json={ - "relay_id": "r1", "endpoint": endpoint, "pk_relay": pk, - "timestamp": ts, "signature": base64.b64encode(sig).decode()}) - assert r.status_code == 201, r.text - assert relay_mod._relays["r1"]["endpoint"] == endpoint - finally: - relay_mod._relays.pop("r1", None) - - -@pytest.mark.asyncio -async def test_every_relay_route_is_closed_as_the_hub_ships(client): - """Nothing in the tree uses the registry, and two of its routes take no account. - - A dependency on the router, so a route added later is closed too. The flag - is read as shipped, not set here — a test that closes the gate itself - would keep passing the day somebody opens it. - """ - admin = await _make_user(client, "relayadmin") - from meshbay_hub.api.deps import set_admin_usernames - set_admin_usernames(["relayadmin"]) - auth = {"Authorization": f"Bearer {admin['token']}"} - - for method, path in (("get", "/v1/relays"), - ("post", "/v1/relays/register"), - ("post", "/v1/relays/approve")): - kwargs = {"headers": auth} if method == "get" else {"json": {}, "headers": auth} - r = await getattr(client, method)(path, **kwargs) - assert r.status_code == 503, (path, r.status_code, r.text) - - @pytest.mark.asyncio async def test_a_stranger_who_locks_your_name_does_not_sign_you_out(client, db_session): """AV26. The sign-in lockout is keyed by username, and usernames are public. @@ -442,26 +323,6 @@ async def test_a_stranger_who_locks_your_name_does_not_sign_you_out(client, db_s assert r.status_code == 200, r.text -@pytest.mark.asyncio -async def test_a_captured_relay_registration_is_not_replayable(client, monkeypatch): - """Same reason /v1/nodes/announce bounds its timestamp.""" - from meshbay_hub.api import relay as relay_mod - - monkeypatch.setattr(relay_mod, "RELAYS_ENABLED", True) - sk = Ed25519PrivateKey.generate() - pk = pk_to_b64(sk.public_key()) - relay_mod._relays["r2"] = {"pk": pk, "active": False} - try: - ts = int(time.time()) - 3600 - sig = sk.sign(f"meshbay:relay_register:r2:203.0.113.5:4433:{ts}".encode()) - r = await client.post("/v1/relays/register", json={ - "relay_id": "r2", "endpoint": "203.0.113.5:4433", "pk_relay": pk, - "timestamp": ts, "signature": base64.b64encode(sig).decode()}) - assert r.status_code == 401, r.text - finally: - relay_mod._relays.pop("r2", None) - - # ── Mail: three paths out of the hub, one of them unmetered ────────────────── @pytest.mark.asyncio diff --git a/packages/meshbay-hub/tests/test_challenge_signature_client.py b/packages/meshbay-hub/tests/test_challenge_signature_client.py index b904698..c13c55d 100644 --- a/packages/meshbay-hub/tests/test_challenge_signature_client.py +++ b/packages/meshbay-hub/tests/test_challenge_signature_client.py @@ -85,7 +85,7 @@ def test_the_browser_holds_the_node_to_its_challenge(tmp_path): cases = { "signed over this connection": (case(), "true"), - "an older node, no signature": (case(sig=None), "false"), + "no signature": (case(sig=None), "refused"), "another key announced": (case(pk=sk_other), "refused"), "a relay's fingerprint": (case(answer=os.urandom(32)), "refused"), "a replay under another nonce": (case(nonce=os.urandom(32)), "refused"), diff --git a/packages/meshbay-hub/tests/test_invite_link_client.py b/packages/meshbay-hub/tests/test_invite_link_client.py index aa8c194..3fb8922 100644 --- a/packages/meshbay-hub/tests/test_invite_link_client.py +++ b/packages/meshbay-hub/tests/test_invite_link_client.py @@ -151,14 +151,13 @@ def test_a_link_code_goes_to_the_node_the_link_names_and_no_other(tmp_path): got = _run(tmp_path, fn.group(0) + """ const r = (...a) => { const e = _linkJoinRefusal(...a); return e ? e.reason : null; }; process.stdout.write(JSON.stringify([ - r('KEY', 'K7P2-9WQX', 'KEY', true), - r('KEY', 'K7P2-9WQX', 'OTHER', true), - r('KEY', 'K7P2-9WQX', 'KEY', false), - r(undefined, 'K7P2-9WQX', 'OTHER', false), - r('KEY', null, 'OTHER', false), + r('KEY', 'K7P2-9WQX', 'KEY'), + r('KEY', 'K7P2-9WQX', 'OTHER'), + r(undefined, 'K7P2-9WQX', 'OTHER'), + r('KEY', null, 'OTHER'), ])); """) - assert got == [None, "link_other_node", "link_node_unproved", None, None] + assert got == [None, "link_other_node", None, None] # ── Read from the source ───────────────────────────────────────────────────── diff --git a/packages/meshbay-hub/tests/test_memory_ceiling.py b/packages/meshbay-hub/tests/test_memory_ceiling.py index 9ea6ad3..48aac1a 100644 --- a/packages/meshbay-hub/tests/test_memory_ceiling.py +++ b/packages/meshbay-hub/tests/test_memory_ceiling.py @@ -84,7 +84,6 @@ const platform = {{ bridgeMessage: (e) => String(e), }}; const downloads = {{ - BLOB_LIMIT: 512 * 1024 * 1024, // Called by the refusal to name why the streamed path declined -- absent // from this stub, the error constructor threw TypeError and the test saw the // wrong failure entirely. diff --git a/packages/meshbay-hub/tests/test_moderation.py b/packages/meshbay-hub/tests/test_moderation.py index 3109e29..b599c84 100644 --- a/packages/meshbay-hub/tests/test_moderation.py +++ b/packages/meshbay-hub/tests/test_moderation.py @@ -22,6 +22,20 @@ async def _register_and_login(client, username: str) -> dict: return {"Authorization": f"Bearer {r.json()['access_token']}"} +async def _node_headers(client, username: str) -> dict: + """A node daemon's token for a fresh account — what a node syncs with.""" + from meshbay_hub.auth import issue_access_token + user = await _register_and_login(client, username) + me = (await client.get("/v1/users/me", headers=user)).json() + tok = issue_access_token(me["user_id"], ttl=3600, groups=[], scope="node") + return {"Authorization": f"Bearer {tok}"} + + +async def _blocked(client, h: str) -> bool: + node = await _node_headers(client, f"node_{h[:6]}_{len(h)}") + return h in (await client.get("/v1/blocklist", headers=node)).json()["hashes"] + + @pytest.fixture async def reporter(client): return await _register_and_login(client, "reporter_one") @@ -34,79 +48,198 @@ async def admin_headers(client): return headers +async def _policy(client, admin, **values): + r = await client.patch("/v1/admin/settings", json={"reports": values}, headers=admin) + assert r.status_code == 200, r.text + return r.json()["reports"] + + +async def _public_group(client, db_session, owner, name="commons-mod"): + from datetime import UTC, datetime + + from meshbay_hub.db.models import Group + r = await client.post("/v1/groups", json={"name": name, "visibility": "public", + "join_policy": "open"}, headers=owner) + assert r.status_code == 201, r.text + gid = r.json()["group_id"] + (await db_session.get(Group, gid)).hosted_at = datetime.now(UTC) + await db_session.commit() + return gid + + +async def _members(client, gid, names): + out = [] + for n in names: + h = await _register_and_login(client, n) + assert (await client.post(f"/v1/groups/{gid}/join", headers=h)).status_code == 200 + out.append(h) + return out + + +async def _report(client, headers, gid, h=FAKE_HASH, **extra): + return await client.post("/v1/reports", headers=headers, + json={"content_hash": h, "group_id": gid, + "reason": "illegal", **extra}) + + +@pytest.fixture +async def setting(client, admin_headers): + """Reports allowed from a new account, so tests need not wait a day.""" + await _policy(client, admin_headers, min_account_age_hours=0) + return admin_headers + + @pytest.mark.asyncio async def test_report_requires_auth(client): - # No credentials at all — FastAPI rejects the missing header before the body. r = await client.post("/v1/reports", json={ - "content_hash": FAKE_HASH, "reason": "illegal"}) + "content_hash": FAKE_HASH, "group_id": "g", "reason": "illegal"}) assert r.status_code in (401, 422) - - # A bogus token is a clean 401. r = await client.post("/v1/reports", - json={"content_hash": FAKE_HASH, "reason": "illegal"}, + json={"content_hash": FAKE_HASH, "group_id": "g", + "reason": "illegal"}, headers={"Authorization": "Bearer not-a-real-token"}) assert r.status_code == 401 @pytest.mark.asyncio -async def test_report_content_logged(client, reporter): - r = await client.post("/v1/reports", - json={"content_hash": FAKE_HASH, "reason": "illegal"}, - headers=reporter) - assert r.status_code == 201 - data = r.json() - assert data["report_count"] == 1 - assert data["status"] == "logged" +async def test_a_node_token_cannot_report(client, db_session, setting): + owner = await _register_and_login(client, "owner_nodetok") + gid = await _public_group(client, db_session, owner, "nodetok-grp") + node = await _node_headers(client, "node_reporter") + assert (await _report(client, node, gid)).status_code == 403 + + +@pytest.mark.asyncio +async def test_a_new_account_cannot_report_yet(client, db_session, admin_headers): + owner = await _register_and_login(client, "owner_young") + gid = await _public_group(client, db_session, owner, "young-grp") + [young] = await _members(client, gid, ["young_member"]) + r = await _report(client, young, gid) + assert r.status_code == 403 and "too new" in r.json()["detail"] + + +@pytest.mark.asyncio +async def test_only_a_member_of_that_public_group_may_report(client, db_session, setting): + owner = await _register_and_login(client, "owner_member") + gid = await _public_group(client, db_session, owner, "member-grp") + stranger = await _register_and_login(client, "stranger_one") + private = (await client.post("/v1/groups", json={"name": "priv-mod"}, + headers=owner)).json()["group_id"] + answers = {(await _report(client, stranger, gid)).json()["detail"], + (await _report(client, owner, private)).json()["detail"], + (await _report(client, stranger, "no-such-group")).json()["detail"]} + # One uniform refusal: it must not say which groups exist or who is in them. + assert len(answers) == 1 + assert (await _report(client, owner, gid)).status_code == 201 @pytest.mark.asyncio -async def test_same_reporter_cannot_walk_the_threshold(client, reporter): +async def test_a_report_says_nothing_about_how_close_review_is(client, db_session, setting): + owner = await _register_and_login(client, "owner_quiet") + gid = await _public_group(client, db_session, owner, "quiet-grp") + r = await _report(client, owner, gid) + assert r.status_code == 201 and r.json() == {"status": "logged"} + + +@pytest.mark.asyncio +async def test_same_reporter_cannot_walk_the_threshold(client, db_session, setting): + await _policy(client, setting, review_threshold=1) + owner = await _register_and_login(client, "owner_walk") + gid = await _public_group(client, db_session, owner, "walk-grp") h = "b" * 64 - for _ in range(5): - r = await client.post("/v1/reports", - json={"content_hash": h, "reason": "spam"}, - headers=reporter) - assert r.json()["report_count"] == 1 + [m] = await _members(client, gid, ["walker_one"]) + await _report(client, m, gid, h) + for _ in range(4): + r = await _report(client, m, gid, h) assert r.json()["status"] == "already_reported" - check = await client.get(f"/v1/blocklist/check?hash={h}") - assert check.json()["blocked"] is False - @pytest.mark.asyncio -async def test_auto_block_on_distinct_reporters(client): +async def test_reaching_the_threshold_queues_for_an_administrator(client, db_session, setting): + owner = await _register_and_login(client, "owner_queue") + gid = await _public_group(client, db_session, owner, "queue-grp") h = "c" * 64 - for i in range(3): - headers = await _register_and_login(client, f"reporter_{i}") - r = await client.post("/v1/reports", - json={"content_hash": h, "reason": "illegal"}, - headers=headers) - assert r.json()["status"] == "auto_blocked" - assert r.json()["report_count"] == 3 + for m in await _members(client, gid, ["queue_r0", "queue_r1", "queue_r2"]): + assert (await _report(client, m, gid, h)).status_code == 201 - check = await client.get(f"/v1/blocklist/check?hash={h}") - assert check.json()["blocked"] is True + assert not await _blocked(client, h), "nothing is blocked without a decision" + queue = (await client.get("/v1/admin/reports", headers=setting)).json()["reports"] + [item] = [q for q in queue if q["hash"] == h] + assert item["reporters"] == 3 and item["reasons"] == {"illegal": 3} + assert item["groups"] == [{"id": gid, "name": "queue-grp"}] + notes = (await client.get("/v1/notifications", headers=setting)).json() + assert any(n["kind"] == "content_review" for n in notes["notifications"]) + + r = await client.post(f"/v1/admin/reports/{h}/block", headers=setting) + assert r.status_code == 200 + assert await _blocked(client, h) + queue = (await client.get("/v1/admin/reports", headers=setting)).json()["reports"] + assert h not in [q["hash"] for q in queue] @pytest.mark.asyncio -async def test_reports_refused_when_public_groups_disabled(client, reporter, admin_headers): - await client.patch("/v1/admin/settings", - json={"allow_public_groups": False}, - headers=admin_headers) +async def test_a_dismissed_report_stays_dismissed(client, db_session, setting): + await _policy(client, setting, review_threshold=1) + owner = await _register_and_login(client, "owner_dismiss") + gid = await _public_group(client, db_session, owner, "dismiss-grp") + h = "d" * 64 + [a, b] = await _members(client, gid, ["dismiss_a", "dismiss_b"]) + await _report(client, a, gid, h) + assert (await client.post(f"/v1/admin/reports/{h}/dismiss", + headers=setting)).status_code == 200 + await _report(client, b, gid, h) + queue = (await client.get("/v1/admin/reports", headers=setting)).json()["reports"] + assert h not in [q["hash"] for q in queue] + assert not await _blocked(client, h) - r = await client.post("/v1/reports", - json={"content_hash": "d" * 64, "reason": "illegal"}, - headers=reporter) + +@pytest.mark.asyncio +async def test_automatic_blocking_is_the_instances_choice(client, db_session, setting): + await _policy(client, setting, review_threshold=2, auto_block=1) + owner = await _register_and_login(client, "owner_auto") + gid = await _public_group(client, db_session, owner, "auto-grp") + h = "e" * 64 + for m in await _members(client, gid, ["autoblock_r0", "autoblock_r1"]): + await _report(client, m, gid, h) + assert await _blocked(client, h) + + +@pytest.mark.asyncio +async def test_a_member_has_a_daily_allowance(client, db_session, setting): + await _policy(client, setting, daily_per_account=2) + owner = await _register_and_login(client, "owner_daily") + gid = await _public_group(client, db_session, owner, "daily-grp") + for i in range(2): + assert (await _report(client, owner, gid, f"{i:064x}")).status_code == 201 + assert (await _report(client, owner, gid, f"{9:064x}")).status_code == 429 + + +@pytest.mark.asyncio +async def test_a_report_is_shaped(client, db_session, setting): + owner = await _register_and_login(client, "owner_shape") + gid = await _public_group(client, db_session, owner, "shape-grp") + assert (await _report(client, owner, gid, "not-a-hash")).status_code == 422 + assert (await _report(client, owner, gid, reason="because")).status_code == 422 + assert (await _report(client, owner, gid, detail="x" * 257)).status_code == 422 + + +@pytest.mark.asyncio +async def test_only_an_administrator_decides(client, db_session, setting): + await _policy(client, setting, review_threshold=1) + owner = await _register_and_login(client, "owner_decide") + gid = await _public_group(client, db_session, owner, "decide-grp") + await _report(client, owner, gid, "f" * 64) + r = await client.post(f"/v1/admin/reports/{'f' * 64}/block", headers=owner) assert r.status_code == 403 @pytest.mark.asyncio -async def test_invalid_hash_rejected(client, reporter): - r = await client.post("/v1/reports", - json={"content_hash": "not-a-valid-blake3-hash", - "reason": "test"}, - headers=reporter) - assert r.status_code == 422 +async def test_reports_refused_when_public_groups_disabled(client, db_session, setting): + owner = await _register_and_login(client, "owner_off") + gid = await _public_group(client, db_session, owner, "off-grp") + await client.patch("/v1/admin/settings", json={"allow_public_groups": False}, + headers=setting) + assert (await _report(client, owner, gid)).status_code == 403 @pytest.mark.asyncio @@ -118,14 +251,13 @@ async def test_admin_add_remove_blocklist(client, admin_headers): headers=admin_headers) assert r.status_code == 201 - r = await client.get(f"/v1/blocklist/check?hash={hash4}") - assert r.json()["blocked"] is True + assert await _blocked(client, hash4) r = await client.delete(f"/v1/admin/blocklist/{hash4}", headers=admin_headers) assert r.status_code == 200 - r = await client.get(f"/v1/blocklist/check?hash={hash4}") - assert r.json()["blocked"] is False + node = await _node_headers(client, "node_after_unblock") + assert hash4 not in (await client.get("/v1/blocklist", headers=node)).json()["hashes"] @pytest.mark.asyncio @@ -134,6 +266,80 @@ async def test_full_blocklist(client, admin_headers): await client.post("/v1/admin/blocklist", json={"content_hash": hash5, "reason": "test"}, headers=admin_headers) - r = await client.get("/v1/blocklist") + node = await _node_headers(client, "node_full") + r = await client.get("/v1/blocklist", headers=node) assert r.status_code == 200 assert hash5 in r.json()["hashes"] + + +@pytest.mark.asyncio +async def test_only_a_node_token_reads_the_blocklist(client, reporter): + assert (await client.get("/v1/blocklist")).status_code in (401, 422) + assert (await client.get("/v1/blocklist", headers=reporter)).status_code == 403 + + +@pytest.mark.asyncio +async def test_the_blocklist_pages_past_its_limit(client, admin_headers): + hashes = sorted(f"{i:064x}" for i in range(5)) + for h in hashes: + await client.post("/v1/admin/blocklist", + json={"content_hash": h, "reason": "test"}, + headers=admin_headers) + node = await _node_headers(client, "node_pager") + seen, after = [], "" + while True: + page = (await client.get(f"/v1/blocklist?limit=2&after={after}", + headers=node)).json() + seen += page["hashes"] + if not page["next"]: + break + after = page["next"] + assert [h for h in seen if h in hashes] == hashes + + +@pytest.mark.asyncio +async def test_an_admin_cannot_block_something_that_is_not_a_hash(client, admin_headers): + r = await client.post("/v1/admin/blocklist", + json={"content_hash": "x" * 5000, "reason": "test"}, + headers=admin_headers) + assert r.status_code == 422 + + +@pytest.mark.asyncio +async def test_a_change_is_pushed_to_nodes_hosting_a_public_group_only(db_session): + """A node hosting only private groups has nothing to apply the list to.""" + import json + + import meshbay_hub.api.revocation as rev + from meshbay_hub.db.models import Group, User + + db_session.add(User(id="owner-bl", username="owner_bl", email="x", hub_id="h", + pw_hash=b"x", pw_salt=b"x")) + db_session.add(Group(id="pub-bl", name="pub", admin_id="owner-bl", + visibility="public")) + db_session.add(Group(id="priv-bl", name="priv", admin_id="owner-bl", + visibility="private")) + await db_session.commit() + + class _WS: + def __init__(self): + self.sent = [] + + async def send_text(self, text): + self.sent.append(json.loads(text)) + + public_node, private_node = _WS(), _WS() + saved = dict(rev._connected_nodes), dict(rev._node_groups) + rev._connected_nodes.update({"n-pub": public_node, "n-priv": private_node}) + rev._node_groups.update({"n-pub": ["pub-bl"], "n-priv": ["priv-bl"]}) + try: + sent = await rev.broadcast_blocklist_update(db_session, add=["a" * 64]) + finally: + rev._connected_nodes.clear() + rev._connected_nodes.update(saved[0]) + rev._node_groups.clear() + rev._node_groups.update(saved[1]) + assert sent == 1 + assert public_node.sent == [{"type": "blocklist_update", "add": ["a" * 64], + "remove": []}] + assert private_node.sent == [] diff --git a/packages/meshbay-hub/tests/test_portable_save_names.py b/packages/meshbay-hub/tests/test_portable_save_names.py new file mode 100644 index 0000000..05a4af7 --- /dev/null +++ b/packages/meshbay-hub/tests/test_portable_save_names.py @@ -0,0 +1,42 @@ +""" +Every way a file is saved goes through a name that can be written everywhere +(docs/MESHBAY_DESIGN.md §10, `static/portable-name.js`). + +The rule itself is held to Python's by `test_portable_name_parity.py`. What this +holds is that the save paths use it: a single file under `portableName`, the +entries of a folder's zip and the archive's own name under `portablePath` / +`portableName`, and the node's own name never handed to a save target directly. +Read from the source, because these paths need a browser and a disk to run. +""" + +import re +from pathlib import Path + +STATIC = Path(__file__).resolve().parents[1] / "src" / "meshbay_hub" / "static" +FILE_UTILS = (STATIC / "file-utils.js").read_text(encoding="utf-8") + + +def _body(name: str) -> str: + m = re.search(rf"^async function {name}\(.*?^\}}", FILE_UTILS, re.M | re.S) + assert m, f"file-utils.js no longer has {name}" + return m.group(0) + + +def test_a_single_file_is_saved_under_a_portable_name(): + body = _body("downloadEntry") + assert "portableName(entry.name)" in body + assert "_openTargetInTurn(saveName" in body + assert "_saveBlob(blob, saveName)" in body + assert "_openTargetInTurn(entry.name" not in body + assert "_saveBlob(blob, entry.name)" not in body + + +def test_a_zip_writes_portable_names_and_is_named_portably(): + body = _body("downloadDirectory") + assert "zip.begin(portablePath(name)" in body + assert "portableName(dir.split('/').pop()" in body + + +def test_a_renamed_file_is_said_so_on_its_row(): + assert "t('transfers.renamed'" in _body("downloadEntry") + assert "t('transfers.renamed_n'" in _body("downloadDirectory") diff --git a/packages/meshbay-hub/tests/test_unauthenticated_surface.py b/packages/meshbay-hub/tests/test_unauthenticated_surface.py index f2da42b..563dfa8 100644 --- a/packages/meshbay-hub/tests/test_unauthenticated_surface.py +++ b/packages/meshbay-hub/tests/test_unauthenticated_surface.py @@ -42,10 +42,6 @@ PUBLIC = { ("POST", "/v1/nodes/auth"): "node sign-in — Ed25519 signature over a fresh timestamp", ("WS", "/v1/nodes/ws"): "a node-scoped JWT in the first message, within a timeout", ("GET", "/v1/groups"): "the public directory — empty when public groups are off", - ("GET", "/v1/blocklist"): "nodes sync it on their own behalf — hashes only", - ("GET", "/v1/blocklist/check"): "nodes consult it on their own behalf", - ("GET", "/v1/relays"): "closed: 503 while relay.RELAYS_ENABLED is False", - ("POST", "/v1/relays/register"): "closed; when open, approved key + signature", ("GET", "/mhp/info"): "closed: 503 while federation.FEDERATION_ENABLED is False", ("GET", "/mhp/directory"): "closed; when open, an MHP token", ("POST", "/mhp/directory"): "closed; when open, an MHP token", 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/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 " |