aboutsummaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
Diffstat (limited to 'packages')
-rw-r--r--packages/meshbay-client/scripts/sync-ui.js2
-rw-r--r--packages/meshbay-client/src/keyring.js34
-rw-r--r--packages/meshbay-client/src/main.js24
-rw-r--r--packages/meshbay-common/src/meshbay_common/keyderive.py130
-rw-r--r--packages/meshbay-common/src/meshbay_common/paths.py8
-rw-r--r--packages/meshbay-common/tests/test_keyderive.py57
-rw-r--r--packages/meshbay-common/tests/test_paths.py7
-rw-r--r--packages/meshbay-common/tests/test_portable_name_parity.py1
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/groups.py25
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/revocation.py53
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/signaling.py34
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/users.py112
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/d4e5f6a7b8ca_known_browsers.py32
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/db/models.py19
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/create-group-page.js7
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/csv.js14
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/keyderive.js77
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/node-page.js7
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/portable-name.js7
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/transport-rewrap.js4
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/transport.js53
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/webrtc-test.html265
-rw-r--r--packages/meshbay-hub/tests/test_bundle_key.py48
-rw-r--r--packages/meshbay-hub/tests/test_csv_cell.py32
-rw-r--r--packages/meshbay-hub/tests/test_desktop_keyring.py34
-rw-r--r--packages/meshbay-hub/tests/test_desktop_shell.py23
-rw-r--r--packages/meshbay-hub/tests/test_group_name_checked.py33
-rw-r--r--packages/meshbay-hub/tests/test_hub_work_is_bounded.py51
-rw-r--r--packages/meshbay-hub/tests/test_known_browser.py98
-rw-r--r--packages/meshbay-hub/tests/test_login_lockout.py20
-rw-r--r--packages/meshbay-hub/tests/test_username_case.py24
-rw-r--r--packages/meshbay-node/src/meshbay_node/cli/groups.py11
-rw-r--r--packages/meshbay-node/src/meshbay_node/cli/parser.py3
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py48
-rw-r--r--packages/meshbay-node/src/meshbay_node/hub_client.py3
-rw-r--r--packages/meshbay-node/src/meshbay_node/linkpreview.py157
-rw-r--r--packages/meshbay-node/src/meshbay_node/media_probe.py17
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops/groups.py30
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops/node_toml.py59
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops/roots.py13
-rw-r--r--packages/meshbay-node/src/meshbay_node/roots.py63
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py76
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py19
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py4
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py17
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py16
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py16
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py7
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py40
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py26
-rw-r--r--packages/meshbay-node/src/meshbay_node/ui/app.py1
-rw-r--r--packages/meshbay-node/src/meshbay_node/uploads.py33
-rw-r--r--packages/meshbay-node/tests/golden/cli.json9
-rw-r--r--packages/meshbay-node/tests/test_attach_from_the_hub.py99
-rw-r--r--packages/meshbay-node/tests/test_chat_is_bounded.py35
-rw-r--r--packages/meshbay-node/tests/test_ffprobe_is_bounded.py61
-rw-r--r--packages/meshbay-node/tests/test_linkpreview.py94
-rw-r--r--packages/meshbay-node/tests/test_member_capacity.py151
-rw-r--r--packages/meshbay-node/tests/test_member_errors_are_plain.py22
-rw-r--r--packages/meshbay-node/tests/test_ops.py2
-rw-r--r--packages/meshbay-node/tests/test_partial_uploads.py118
-rw-r--r--packages/meshbay-node/tests/test_revocation_closes_sessions.py53
63 files changed, 1958 insertions, 682 deletions
diff --git a/packages/meshbay-client/scripts/sync-ui.js b/packages/meshbay-client/scripts/sync-ui.js
index d0b6eee..4ebdf35 100644
--- a/packages/meshbay-client/scripts/sync-ui.js
+++ b/packages/meshbay-client/scripts/sync-ui.js
@@ -23,7 +23,7 @@ const DEST = path.resolve(__dirname, '..', 'ui');
const OURS = ['index.html'];
// sw.js has to sit at the root of the scope it serves, which it already does.
-const SKIP = new Set(['webrtc-test.html']);
+const SKIP = new Set();
function copyTree(from, to) {
fs.mkdirSync(to, { recursive: true });
diff --git a/packages/meshbay-client/src/keyring.js b/packages/meshbay-client/src/keyring.js
index 4fb6eac..ec37fd4 100644
--- a/packages/meshbay-client/src/keyring.js
+++ b/packages/meshbay-client/src/keyring.js
@@ -24,6 +24,8 @@ const { transcriptFor } = require('./transcripts.js');
// keyderive.js: the same numbers, or no bundle opens across the two.
const ARGON2 = { memory: 131072, passes: 3, parallelism: 1, tagLength: 32 };
const MAGIC = Buffer.from('MBK3');
+// TRANSITIONAL — the format before MBK3, read once to be replaced (keyderive.js).
+const LEGACY_MAGIC = Buffer.from('MBK2');
const X25519_SPKI = Buffer.from('302a300506032b656e032100', 'hex');
const B32 = 'ABCDEFGHIJKLMNOPQRSTUVWXYZ234567';
@@ -63,6 +65,17 @@ function seal(identity, key, userId, nodePk, pepperVersion) {
return b64(Buffer.concat([MAGIC, Buffer.from([pepperVersion & 0xff]), nonce, ct]));
}
+/** TRANSITIONAL — MBK2: "MBK2" ‖ nonce ‖ AES-GCM under the Argon2 key, no AAD. */
+function openLegacy(bundleB64, key) {
+ const raw = unb64(bundleB64);
+ const nonce = raw.subarray(4, 16);
+ const body = raw.subarray(16, raw.length - 16);
+ const d = crypto.createDecipheriv('aes-256-gcm', key, nonce);
+ d.setAuthTag(raw.subarray(raw.length - 16));
+ const plain = JSON.parse(Buffer.concat([d.update(body), d.final()]).toString());
+ return { ed: plain.skEd, x: plain.skX };
+}
+
function open(bundleB64, key, userId, nodePk) {
const raw = unb64(bundleB64);
if (!raw.subarray(0, 4).equals(MAGIC)) {
@@ -142,7 +155,11 @@ function createKeyring({ load, save, argon2 }) {
const v = pepperVersion || 1;
if (p) { pending.set(userId, { m, v }); return true; }
const s = state();
- s.masters[userId] = { m: b64(m), v };
+ // `legacy` (TRANSITIONAL): the Argon2 key itself, which MBK2 bundles
+ // were sealed under — kept beside `M`, in the same OS-protected store
+ // and for as long, so a node still holding one has it opened and
+ // replaced on the next connection. Remove once no MBK2 bundle is left.
+ s.masters[userId] = { m: b64(m), v, legacy: b64(a) };
save(s);
return true;
},
@@ -150,7 +167,10 @@ function createKeyring({ load, save, argon2 }) {
const p = pending.get(userId);
if (!p) return false;
const s = state();
- s.masters[userId] = { m: b64(p.m), v: p.v };
+ // The legacy key stays the old passphrase's: MBK2 bundles were sealed
+ // under that one, never under the new.
+ const legacy = (s.masters[userId] || {}).legacy;
+ s.masters[userId] = { m: b64(p.m), v: p.v, ...(legacy ? { legacy } : {}) };
save(s);
pending.delete(userId);
return true;
@@ -177,6 +197,16 @@ function createKeyring({ load, save, argon2 }) {
* was entered (a reset on a machine that had never held this identity).
*/
openBundle(userId, nodePk, { bundleEnc, recoveryEnc, recoveryMnemonic, username }) {
+ if (unb64(bundleEnc).subarray(0, 4).equals(LEGACY_MAGIC)) {
+ // TRANSITIONAL. Kept unsealed (`sealedWith: null`), so the next
+ // settle replaces the node's copy with MBK3, or withdraws it when the
+ // account has no browser access.
+ const legacy = (state().masters[userId] || {}).legacy;
+ if (!legacy) throw new Error('no_legacy_key');
+ const id = openLegacy(bundleEnc, unb64(legacy));
+ keep(userId, nodePk, { ...id, sealedWith: null });
+ return publicOf(id);
+ }
const { m } = master(userId);
let id;
try {
diff --git a/packages/meshbay-client/src/main.js b/packages/meshbay-client/src/main.js
index 309d5f6..52ca168 100644
--- a/packages/meshbay-client/src/main.js
+++ b/packages/meshbay-client/src/main.js
@@ -1160,8 +1160,24 @@ function registerBridge() {
return { path: chosen, name: path.basename(chosen) };
});
+ // The Mark-of-the-Web, as a browser leaves on every download: the file came
+ // from somebody else's machine, and Windows decides what that means —
+ // SmartScreen for a program, Protected View for a document. This application
+ // writes its files itself, so nothing else marks them. NTFS only; elsewhere
+ // there is no such stream, and nothing is lost by not having one.
+ function markFromInternet(file) {
+ if (process.platform !== 'win32') return;
+ try {
+ fs.writeFileSync(`${file}:Zone.Identifier`, '[ZoneTransfer]\r\nZoneId=3\r\n');
+ } catch { /* FAT, exFAT, a network share: no alternate data streams */ }
+ }
+
handle('save:begin', async (_e, suggestedName, opts) => {
- const wanted = path.basename(String(suggestedName || 'download'));
+ // Bidirectional controls replaced here as well as in the page
+ // (portable-name.js): "invoice\u202efdp.exe" would be saved, and listed by
+ // the file manager, as "invoiceexe.pdf".
+ const wanted = path.basename(String(suggestedName || 'download'))
+ .replace(/[\u061c\u200e\u200f\u202a-\u202e\u2066-\u2069]/g, '_');
const chosen = chosenDownloadDir();
let target = null;
@@ -1231,6 +1247,7 @@ function registerBridge() {
console.error('[MeshBay] could not finalise download:', err.message);
return false;
}
+ markFromInternet(sink.path);
completedPaths.set(String(id), sink.path);
return true;
});
@@ -2184,7 +2201,10 @@ function registerBridge() {
attachGroup: async (a) => {
const body = { name: aText(a.name, 'the group name'),
shared_dir: aText(a.path, 'the folder', 4096),
- writable: a.writable !== false };
+ writable: a.writable !== false,
+ // The person's choice on the creation form; the node never
+ // takes it from the hub.
+ join_policy: a.joinPolicy === 'open' ? 'open' : 'invite' };
await confirmOrRefuse('native.attach_confirm',
{ name: body.name, path: body.shared_dir });
return ['POST', '/api/groups/attach', body];
diff --git a/packages/meshbay-common/src/meshbay_common/keyderive.py b/packages/meshbay-common/src/meshbay_common/keyderive.py
deleted file mode 100644
index 4b90af3..0000000
--- a/packages/meshbay-common/src/meshbay_common/keyderive.py
+++ /dev/null
@@ -1,130 +0,0 @@
-"""
-MeshBay — Key derivation from username + password.
-
-Allows Ed25519 + X25519 keypairs to be derived deterministically
-from credentials. Same inputs → same keys on any device.
-
-Algorithm: Argon2id (Python CLI / native clients)
- salt = SHA-256("meshbay:v1:" + username)
- seed = Argon2id(password, salt, length=64, ...)
- sk_ed = Ed25519PrivateKey.from_private_bytes(seed[:32])
- sk_x25519 = X25519PrivateKey.from_private_bytes(seed[32:])
-
-Browser alternative (keyderive.js): uses PBKDF2-SHA512 because
-WebCrypto does not support Argon2. The two algorithms produce
-DIFFERENT keys from the same password — a user registered via Python
-CLI and via web browser will have different keypairs.
-
-Resolution: the web client generates RANDOM keypairs on first login
-(WebCrypto, stored encrypted in hub), and uses derive_keys_from_password
-only to encrypt/decrypt the stored keypair bundle. This avoids the
-algorithm mismatch problem entirely.
-
-See keyderive.js for the browser-side implementation.
-"""
-
-import hashlib
-
-from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
-from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey
-from cryptography.hazmat.primitives.kdf.argon2 import Argon2id
-
-# Argon2id parameters — same as keystore (see crypto.py)
-_ITERATIONS = 3
-_MEMORY_COST = 65536 # 64 MB — increase to 262144 for production
-_LANES = 4
-_SEED_LENGTH = 64 # 32 bytes Ed25519 + 32 bytes X25519
-
-
-def _derive_salt(username: str) -> bytes:
- """Deterministic salt: SHA-256 of 'meshbay:v1:<username>'."""
- return hashlib.sha256(f"meshbay:v1:{username}".encode()).digest()
-
-
-def derive_keys_from_password(
- username: str,
- password: str,
-) -> tuple[Ed25519PrivateKey, X25519PrivateKey]:
- """
- Derive Ed25519 + X25519 keypairs deterministically from username + password.
-
- Properties:
- - Same credentials always produce the same keypairs
- - Different usernames produce different keys (even with same password)
- - Password cannot be recovered from the public keys
- - Changing the password invalidates all GEK bundles stored on the hub
-
- Use for:
- - CLI / native node registration (Argon2id available)
- - Recovery of lost keypairs from credentials
-
- Do NOT use for:
- - Web browser registration (use random keypairs + encrypted bundle instead)
- """
- salt = _derive_salt(username)
- kdf = Argon2id(
- salt=salt, length=_SEED_LENGTH,
- iterations=_ITERATIONS, lanes=_LANES, memory_cost=_MEMORY_COST,
- )
- seed = kdf.derive(password.encode())
- return (
- Ed25519PrivateKey.from_private_bytes(seed[:32]),
- X25519PrivateKey.from_private_bytes(seed[32:]),
- )
-
-
-def encrypt_keypair_bundle(
- sk_ed: Ed25519PrivateKey,
- sk_x: X25519PrivateKey,
- password: str,
- username: str,
-) -> bytes:
- """
- Encrypt a keypair bundle with a password-derived key (for hub storage).
- Used by web clients: random keypairs encrypted with password, stored on hub.
- Returns: AES-256-GCM ciphertext (nonce prepended).
- """
- import os
-
- import msgpack
- from cryptography.hazmat.primitives.ciphers.aead import AESGCM
-
- from meshbay_common.crypto import sk_to_raw
-
- # Derive an AES key from the password (different info string from key derivation)
- salt = hashlib.sha256(f"meshbay:bundle:v1:{username}".encode()).digest()
- kdf = Argon2id(salt=salt, length=32, iterations=_ITERATIONS,
- lanes=_LANES, memory_cost=_MEMORY_COST)
- aes_key = kdf.derive(password.encode())
-
- payload = msgpack.packb({
- "sk_ed": sk_to_raw(sk_ed),
- "sk_x": sk_to_raw(sk_x),
- }, use_bin_type=True)
-
- nonce = os.urandom(12)
- ct = AESGCM(aes_key).encrypt(nonce, payload, None)
- return nonce + ct
-
-
-def decrypt_keypair_bundle(
- bundle: bytes,
- password: str,
- username: str,
-) -> tuple[Ed25519PrivateKey, X25519PrivateKey]:
- """Decrypt a keypair bundle. Raises on wrong password."""
- import msgpack
- from cryptography.hazmat.primitives.ciphers.aead import AESGCM
-
- salt = hashlib.sha256(f"meshbay:bundle:v1:{username}".encode()).digest()
- kdf = Argon2id(salt=salt, length=32, iterations=_ITERATIONS,
- lanes=_LANES, memory_cost=_MEMORY_COST)
- aes_key = kdf.derive(password.encode())
-
- nonce, ct = bundle[:12], bundle[12:]
- payload = AESGCM(aes_key).decrypt(nonce, ct, None)
- data = msgpack.unpackb(payload, raw=False)
- return (
- Ed25519PrivateKey.from_private_bytes(data["sk_ed"]),
- X25519PrivateKey.from_private_bytes(data["sk_x"]),
- )
diff --git a/packages/meshbay-common/src/meshbay_common/paths.py b/packages/meshbay-common/src/meshbay_common/paths.py
index b673649..8be4680 100644
--- a/packages/meshbay-common/src/meshbay_common/paths.py
+++ b/packages/meshbay-common/src/meshbay_common/paths.py
@@ -30,8 +30,12 @@ WINDOWS_RESERVED = frozenset({
*(f"LPT{i}" for i in range(1, 10)),
})
-# Reserved on Windows; `/` is reserved everywhere. Control characters go too.
-_RESERVED_CHARS = set('<>:"/\\|?*') | {chr(c) for c in range(32)}
+# Reserved on Windows; `/` is reserved everywhere. Control characters go too,
+# and so do the bidirectional controls: "invoice\u202efdp.exe" displays as
+# "invoiceexe.pdf", and a saved name must say what the file is.
+BIDI_CONTROLS = frozenset("\u061c\u200e\u200f\u202a\u202b\u202c\u202d\u202e"
+ "\u2066\u2067\u2068\u2069")
+_RESERVED_CHARS = set('<>:"/\\|?*') | {chr(c) for c in range(32)} | BIDI_CONTROLS
# Windows without long-path support. A deep media library reaches this.
MAX_PATH_WINDOWS = 260
diff --git a/packages/meshbay-common/tests/test_keyderive.py b/packages/meshbay-common/tests/test_keyderive.py
deleted file mode 100644
index 40b3c71..0000000
--- a/packages/meshbay-common/tests/test_keyderive.py
+++ /dev/null
@@ -1,57 +0,0 @@
-"""Tests for password-based key derivation."""
-
-import pytest
-from meshbay_common.crypto import pk_to_b64
-from meshbay_common.keyderive import (
- decrypt_keypair_bundle,
- derive_keys_from_password,
- encrypt_keypair_bundle,
-)
-
-
-def test_deterministic():
- """Same credentials → same keys."""
- sk_ed1, sk_x1 = derive_keys_from_password("alice", "correct-horse")
- sk_ed2, sk_x2 = derive_keys_from_password("alice", "correct-horse")
- assert pk_to_b64(sk_ed1.public_key()) == pk_to_b64(sk_ed2.public_key())
- assert pk_to_b64(sk_x1.public_key()) == pk_to_b64(sk_x2.public_key())
-
-
-def test_different_users_different_keys():
- sk_ed_a, _ = derive_keys_from_password("alice", "samepassword")
- sk_ed_b, _ = derive_keys_from_password("bob", "samepassword")
- assert pk_to_b64(sk_ed_a.public_key()) != pk_to_b64(sk_ed_b.public_key())
-
-
-def test_different_passwords_different_keys():
- sk_ed1, _ = derive_keys_from_password("alice", "password1")
- sk_ed2, _ = derive_keys_from_password("alice", "password2")
- assert pk_to_b64(sk_ed1.public_key()) != pk_to_b64(sk_ed2.public_key())
-
-
-def test_ed_and_x_keys_independent():
- sk_ed, sk_x = derive_keys_from_password("user", "pass12345")
- from meshbay_common.crypto import sk_to_raw
- assert sk_to_raw(sk_ed) != sk_to_raw(sk_x)
-
-
-def test_bundle_encrypt_decrypt():
- sk_ed, sk_x = derive_keys_from_password("alice", "strongpass!")
- bundle = encrypt_keypair_bundle(sk_ed, sk_x, "password123", "alice")
- sk_ed2, sk_x2 = decrypt_keypair_bundle(bundle, "password123", "alice")
- assert pk_to_b64(sk_ed.public_key()) == pk_to_b64(sk_ed2.public_key())
- assert pk_to_b64(sk_x.public_key()) == pk_to_b64(sk_x2.public_key())
-
-
-def test_bundle_wrong_password_rejected():
- sk_ed, sk_x = derive_keys_from_password("alice", "correctpass")
- bundle = encrypt_keypair_bundle(sk_ed, sk_x, "correctpass", "alice")
- with pytest.raises(Exception):
- decrypt_keypair_bundle(bundle, "wrongpass", "alice")
-
-
-def test_bundle_wrong_username_rejected():
- sk_ed, sk_x = derive_keys_from_password("alice", "pass12345")
- bundle = encrypt_keypair_bundle(sk_ed, sk_x, "pass12345", "alice")
- with pytest.raises(Exception):
- decrypt_keypair_bundle(bundle, "pass12345", "bob") # wrong username salt
diff --git a/packages/meshbay-common/tests/test_paths.py b/packages/meshbay-common/tests/test_paths.py
index 223a2a6..012a32e 100644
--- a/packages/meshbay-common/tests/test_paths.py
+++ b/packages/meshbay-common/tests/test_paths.py
@@ -119,3 +119,10 @@ def test_sanitizing_produces_something_writable():
def test_sanitizing_never_returns_nothing():
assert sanitize_for_download("...") not in ("", None)
assert sanitize_for_download("???") not in ("", None)
+
+
+def test_a_bidi_override_cannot_hide_an_extension():
+ from meshbay_common.paths import portable_name_problem, sanitize_for_download
+ disguised = "invoice‮fdp.exe" # displays as "invoiceexe.pdf"
+ assert sanitize_for_download(disguised) == "invoice_fdp.exe"
+ assert portable_name_problem(disguised)
diff --git a/packages/meshbay-common/tests/test_portable_name_parity.py b/packages/meshbay-common/tests/test_portable_name_parity.py
index 3a491a7..d7df710 100644
--- a/packages/meshbay-common/tests/test_portable_name_parity.py
+++ b/packages/meshbay-common/tests/test_portable_name_parity.py
@@ -26,6 +26,7 @@ pytestmark = pytest.mark.skipif(
)
NAMES = [
+ "invoice\u202efdp.exe", "a\u2066b\u2069.txt", "mark\u200f.txt",
"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", "...",
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/groups.py b/packages/meshbay-hub/src/meshbay_hub/api/groups.py
index 10049b2..fc117bf 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/groups.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/groups.py
@@ -1,6 +1,7 @@
"""Group endpoints — /v1/groups/*"""
import re
+import unicodedata
from datetime import UTC, datetime
from fastapi import APIRouter, Depends, HTTPException, Query, Request
@@ -347,6 +348,25 @@ async def join_group(
"owner_username": owner}
+# The column's width. A longer name was a database error on PostgreSQL (a 500)
+# and silently truncated on SQLite.
+MAX_GROUP_NAME = 128
+# Line breaks and other C0/C1 controls, and the bidirectional overrides that
+# make a name display as something other than what it is. Joiners stay: an
+# emoji family is a ZWJ sequence.
+_BIDI_CONTROLS = frozenset("\u202a\u202b\u202c\u202d\u202e\u2066\u2067\u2068\u2069")
+
+
+def _group_name_problem(name: str) -> str | None:
+ if not name:
+ return "A group needs a name."
+ if len(name) > MAX_GROUP_NAME:
+ return f"A group name is at most {MAX_GROUP_NAME} characters."
+ if any(unicodedata.category(c) == "Cc" or c in _BIDI_CONTROLS for c in name):
+ return "A group name cannot contain control characters."
+ return None
+
+
class GroupCreateRequest(BaseModel):
name: str
visibility: str = "private" # public|private
@@ -426,8 +446,9 @@ async def create_group(
"anyone to be able to join.")
name = body.name.strip()
- if not name:
- raise HTTPException(status_code=422, detail="A group needs a name.")
+ problem = _group_name_problem(name)
+ if problem:
+ raise HTTPException(status_code=422, detail=problem)
# One name per owner, case-insensitively. Two *different* owners may each
# have a "photos" — that is why the check is scoped to `admin_id` and why
# the group's real identity stays its UUID. The DB has a unique index too
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
index 6a6baa7..5a33d77 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
@@ -87,6 +87,16 @@ NOTIFY_WINDOW_SECONDS = 60
_notify_window: dict[str, tuple[float, int]] = {} # node_id → (window start, count)
+# `update_groups` re-reads the node's groups from the database. A node sends one
+# when its configuration is reloaded — an operator attaching a group, an owner's
+# approval arriving — so ten a minute is far past real use, and the same budget
+# rule as chat_notify keeps one node from spending the hub's database for others.
+UPDATE_GROUPS_BURST = 10
+_update_window: dict[str, tuple[float, int]] = {}
+# The groups one node may claim in one message. An operator with fifty groups
+# is a large one.
+MAX_CLAIMED_GROUPS = 1000
+
def forget_node(node_id: str) -> None:
"""Drop everything a disconnected node's socket owned.
@@ -106,25 +116,35 @@ def forget_node(node_id: str) -> None:
_node_users.pop(node_id, None)
-def _notify_budget(node_id: str) -> bool:
- """True if this node may send one more chat_notify now."""
+def _spend(window: dict[str, tuple[float, int]], node_id: str, burst: int) -> bool:
+ """True if this node may send one more message of a budgeted kind now."""
now = time.monotonic()
- if len(_notify_window) > 1000:
+ if len(window) > 1000:
# Swept here rather than on disconnect, which would let a node refill
# its budget by reconnecting — the same token stays valid for an hour.
- for nid, (started, _) in list(_notify_window.items()):
+ for nid, (started, _) in list(window.items()):
if now - started >= NOTIFY_WINDOW_SECONDS:
- _notify_window.pop(nid, None)
- start, count = _notify_window.get(node_id, (now, 0))
+ window.pop(nid, None)
+ start, count = window.get(node_id, (now, 0))
if now - start >= NOTIFY_WINDOW_SECONDS:
start, count = now, 0
- if count >= NOTIFY_BURST:
- _notify_window[node_id] = (start, count)
+ if count >= burst:
+ window[node_id] = (start, count)
return False
- _notify_window[node_id] = (start, count + 1)
+ window[node_id] = (start, count + 1)
return True
+def _notify_budget(node_id: str) -> bool:
+ """True if this node may send one more chat_notify now."""
+ return _spend(_notify_window, node_id, NOTIFY_BURST)
+
+
+def _update_budget(node_id: str) -> bool:
+ """True if this node may send one more update_groups now."""
+ return _spend(_update_window, node_id, UPDATE_GROUPS_BURST)
+
+
async def _mark_hosted(group_ids: list[str]) -> None:
"""Stamp the first time a node announced it hosts each of these groups.
@@ -461,8 +481,8 @@ async def node_websocket(ws: WebSocket):
await _reject(ws, "Node already connected", 4009)
return
- resolved_id, result = await _authorize_node_ws(
- msg["token"], claimed_id, msg.get("group_ids"))
+ claimed = [str(g) for g in (msg.get("group_ids") or [])][:MAX_CLAIMED_GROUPS]
+ resolved_id, result = await _authorize_node_ws(msg["token"], claimed_id, claimed)
if resolved_id is None:
await _reject(ws, result, 4003)
return
@@ -472,7 +492,7 @@ async def node_websocket(ws: WebSocket):
node_id = resolved_id
_connected_nodes[node_id] = ws
_node_groups[node_id] = group_ids
- _node_claims[node_id] = list(msg.get("group_ids") or [])
+ _node_claims[node_id] = claimed
_node_users[node_id] = user_id
await _mark_hosted(group_ids)
log.info("Node WS connected: %s (user=%s, groups=%d)",
@@ -493,13 +513,16 @@ async def node_websocket(ws: WebSocket):
from meshbay_hub.api.signaling import handle_webrtc_answer
handle_webrtc_answer(msg, node_id)
elif msg.get("type") == "update_groups":
+ if not _update_budget(node_id):
+ log.warning("Node %s exceeded its update_groups rate", node_id[:8])
+ continue
# Through the same gate as the registration above. This used to
# assign the message's list verbatim, so the ceiling that makes
# C2 hold at authentication could be stepped over one message
# later: a node had only to reload to claim any group on the hub.
- _node_claims[node_id] = list(msg.get("group_ids") or [])
- new_gids = await resolve_node_groups(
- node_id, user_id, msg.get("group_ids"))
+ claimed = [str(g) for g in (msg.get("group_ids") or [])][:MAX_CLAIMED_GROUPS]
+ _node_claims[node_id] = claimed
+ new_gids = await resolve_node_groups(node_id, user_id, claimed)
_node_groups[node_id] = new_gids
await _mark_hosted(new_gids)
log.info("Node %s updated groups: %d", node_id[:8], len(new_gids))
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py
index 6c9699b..56e0e4b 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py
@@ -20,7 +20,7 @@ import time
import uuid
from fastapi import APIRouter, Depends, HTTPException, Request
-from pydantic import BaseModel
+from pydantic import BaseModel, Field, field_validator
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
@@ -41,9 +41,23 @@ _webrtc_answers: dict[str, asyncio.Future] = {}
_answer_owner: dict[str, str] = {}
+# A browser offers a handful of candidates — a host and a reflexive one per
+# interface — and embeds them in the SDP anyway. The list is relayed to the node
+# as it came, so it is bounded like the SDP beside it.
+MAX_ICE_CANDIDATES = 64
+MAX_ICE_BYTES = 32 * 1024
+
+
class WebRTCOfferRequest(BaseModel):
sdp: str
- ice_candidates: list[dict] = []
+ ice_candidates: list[dict] = Field(default_factory=list, max_length=MAX_ICE_CANDIDATES)
+
+ @field_validator("ice_candidates")
+ @classmethod
+ def _bounded(cls, v: list[dict]) -> list[dict]:
+ if len(json.dumps(v)) > MAX_ICE_BYTES:
+ raise ValueError("ICE candidates too large")
+ return v
class WebRTCOfferResponse(BaseModel):
@@ -176,14 +190,6 @@ async def webrtc_offer(
if len(body.sdp) > MAX_SDP_BYTES:
raise HTTPException(status_code=413, detail="SDP too large")
- # Logged here because this is the moment a browser starts a peer connection,
- # and the address it starts it from is this one — the hub's own view of the
- # TCP connection. Whatever address the peers then discover through STUN is
- # theirs to negotiate and is not what a log should record.
- db.add(IPLog(user_id=current_user.id, event="webrtc_offer",
- ip_address=client_ip(request), detail=node_id[:8]))
- await db.commit()
-
ws = _connected_nodes.get(node_id)
if not ws:
raise HTTPException(status_code=404, detail="Node not connected")
@@ -205,6 +211,14 @@ async def webrtc_offer(
raise HTTPException(status_code=429, detail="Too many connections to this node",
headers={"Retry-After": str(max(1, math.ceil(wait)))})
+ # Logged once the offer is going to a node, not before: the address a peer
+ # connection starts from is the hub's own view of this TCP connection, and an
+ # IP log row is kept a year — written before the checks above, any account
+ # could add rows for any string it named as a node.
+ db.add(IPLog(user_id=current_user.id, event="webrtc_offer",
+ ip_address=client_ip(request), detail=node_id[:8]))
+ await db.commit()
+
peer_id = str(uuid.uuid4())
answer_future: asyncio.Future = asyncio.get_event_loop().create_future()
_webrtc_answers[peer_id] = answer_future
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/users.py b/packages/meshbay-hub/src/meshbay_hub/api/users.py
index 8e780df..9326bfb 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/users.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/users.py
@@ -1,6 +1,7 @@
"""User endpoints — /v1/users/*"""
import base64
+import hashlib
import logging
import re
import secrets
@@ -42,6 +43,7 @@ from meshbay_hub.db.models import (
GroupInviteLink,
GroupMember,
IPLog,
+ KnownBrowser,
Node,
Notification,
RefreshToken,
@@ -161,6 +163,7 @@ class LoginRequest(BaseModel):
username: str
password: str | None = None # legacy (raw password) for migration
auth_key: str | None = None # PBKDF2-derived auth key (new scheme)
+ known_browser: str | None = None # from an earlier sign-in on this browser
class RefreshRequest(BaseModel):
@@ -182,12 +185,18 @@ async def register(
):
eh = hash_email_blind(body.email)
- existing = await db.execute(
- select(User).where(User.username == body.username))
- found = existing.scalar_one_or_none()
+ # Unique regardless of case: invitations and member management name people
+ # by username, and "Alice" beside "alice" is one person to whoever reads it.
+ # Accounts that already differ only by case (made before this) keep their
+ # names; the exact match is the one a retry means.
+ same = (await db.execute(
+ select(User).where(func.lower(User.username) == body.username.lower())
+ )).scalars().all()
+ found = next((u for u in same if u.username == body.username), same[0] if same else None)
if found:
- if found.status == "pending" and found.email_hash == eh:
+ if (found.username == body.username and found.status == "pending"
+ and found.email_hash == eh):
# Same person retrying before validation — resend a code.
# No captcha: the initial registration already passed it.
#
@@ -351,6 +360,58 @@ async def _take_login_attempt(db: AsyncSession, username: str) -> None:
headers={"Retry-After": str(retry_after)})
+def _session_counter(user: User) -> str:
+ """The failure counter for a passphrase re-checked inside an open session.
+
+ Its own, not the sign-in one: a stranger who keeps a name locked at sign-in
+ must not also stop its owner changing their passphrase, deleting their
+ account or registering a device from a session they already hold.
+ """
+ return f"\x00session:{user.id}"
+
+
+async def _browser_counter(db: AsyncSession, username: str,
+ token: str | None) -> tuple[str, "KnownBrowser | None"]:
+ """The failure counter for a sign-in, and the known browser behind it if any.
+
+ A browser that signed in to this account before presents its token and is
+ counted on its own: the username's counter, which anyone can spend, then
+ locks only browsers this account has never used. A token for another
+ account, or none, is the username's counter — the answer is the same either
+ way, so it says nothing about the account (M1).
+ """
+ if token:
+ row = (await db.execute(
+ select(KnownBrowser).join(User, User.id == KnownBrowser.user_id)
+ .where(KnownBrowser.token_hash == _browser_hash(token),
+ User.username == username))).scalar_one_or_none()
+ if row is not None:
+ return f"{username}\x00browser:{row.id}", row
+ return username, None
+
+
+def _browser_hash(token: str) -> str:
+ return hashlib.sha256(f"meshbay:known_browser:{token}".encode()).hexdigest()
+
+
+# How many browsers one account is remembered on. The oldest goes first; a
+# browser forgotten here is only an unknown one again.
+MAX_KNOWN_BROWSERS = 20
+
+
+async def _remember_browser(db: AsyncSession, user: User) -> str:
+ """A new known-browser token for `user`. The caller commits."""
+ raw = secrets.token_urlsafe(32)
+ rows = (await db.execute(
+ select(KnownBrowser.id).where(KnownBrowser.user_id == user.id)
+ .order_by(KnownBrowser.last_used_at.desc()))).scalars().all()
+ stale = rows[MAX_KNOWN_BROWSERS - 1:]
+ if stale:
+ await db.execute(delete(KnownBrowser).where(KnownBrowser.id.in_(stale)))
+ db.add(KnownBrowser(user_id=user.id, token_hash=_browser_hash(raw)))
+ return raw
+
+
async def _prove_passphrase(db: AsyncSession, user: User, auth_key: str) -> None:
"""Refuse with 403 unless `auth_key` is this account's, spending an attempt.
@@ -358,18 +419,18 @@ async def _prove_passphrase(db: AsyncSession, user: User, auth_key: str) -> None
refreshed one, or one lifted from a page, and what it would buy here outlives
the session or reopens the offline search the pepper exists to prevent.
"""
- await _take_login_attempt(db, user.username)
+ await _take_login_attempt(db, _session_counter(user))
if not await verify_password_off_loop(auth_key, user.pw_hash, user.pw_salt,
user.pw_version):
raise HTTPException(status_code=403, detail="Passphrase does not match")
- await login_throttle.clear(db, user.username)
+ await login_throttle.clear(db, _session_counter(user))
async def _login_failed(db: AsyncSession, username: str, ip: str,
- user_id: str | None = None) -> None:
+ user_id: str | None = None, counter: str | None = None) -> None:
"""Record a wrong passphrase and answer 401. Always raises."""
db.add(IPLog(user_id=user_id, event="login_fail", ip_address=ip, detail=username))
- if await login_throttle.is_now_locked(db, username):
+ if await login_throttle.is_now_locked(db, counter or username):
# Once, on the failure that spent the last attempt — so the logs tab
# shows when a name was locked, not every refusal after it.
db.add(IPLog(user_id=user_id, event="login_locked", ip_address=ip,
@@ -431,32 +492,33 @@ async def login(
# Before the account is even looked up: an unknown name spends attempts and
# locks exactly like a real one, so neither answer tells them apart (M1).
- await _take_login_attempt(db, body.username)
+ counter, browser = await _browser_counter(db, body.username, body.known_browser)
+ await _take_login_attempt(db, counter)
result = await db.execute(
select(User).where(User.username == body.username))
user = result.scalar_one_or_none()
if not user:
- await _login_failed(db, body.username, ip)
+ await _login_failed(db, body.username, ip, counter=counter)
if user.pw_version >= 3:
# New scheme: verify auth_key
if not body.auth_key or not await verify_password_off_loop(
body.auth_key, user.pw_hash, user.pw_salt, version=user.pw_version
):
- await _login_failed(db, body.username, ip, user.id)
+ await _login_failed(db, body.username, ip, user.id, counter)
else:
# Legacy scheme: need raw password
if not body.password:
# Nothing was checked, so nothing was guessed.
- await login_throttle.release(db, body.username)
+ await login_throttle.release(db, counter)
await db.commit()
raise HTTPException(status_code=401, detail="auth_upgrade_required")
if not await verify_password_off_loop(
body.password, user.pw_hash, user.pw_salt, version=user.pw_version
):
- await _login_failed(db, body.username, ip, user.id)
+ await _login_failed(db, body.username, ip, user.id, counter)
# Migrate to new scheme if auth_key provided alongside password
if body.auth_key:
new_hash, new_salt = await hash_password_off_loop(body.auth_key)
@@ -471,6 +533,7 @@ async def login(
user.pw_version = 2
# The passphrase was right, whatever the account's status turns out to be.
+ await login_throttle.clear(db, counter)
await login_throttle.clear(db, body.username)
if user.status != "active":
@@ -501,6 +564,11 @@ async def login(
))
db.add(IPLog(user_id=user.id, event="login", ip_address=ip))
pepper = _bundle_pepper(user)
+ if browser is not None:
+ browser.last_used_at = datetime.now(UTC)
+ known = {}
+ else:
+ known = {"known_browser": await _remember_browser(db, user)}
await db.commit()
return {
@@ -509,6 +577,7 @@ async def login(
"token_type": "bearer",
"expires_in": _ttl(),
**pepper,
+ **known,
}
@@ -787,7 +856,8 @@ async def get_current_user_info(
# A passphrase change re-wraps every node's bundle *before* the hub
# accepts the new passphrase, and must not start while the hub would
# then refuse it.
- "passphrase_locked_for": await login_throttle.locked_for(db, current_user.username),
+ "passphrase_locked_for": await login_throttle.locked_for(
+ db, _session_counter(current_user)),
}
@@ -849,13 +919,13 @@ async def update_profile(
raise HTTPException(
status_code=403,
detail="Changing your e-mail requires your passphrase.")
- await _take_login_attempt(db, current_user.username)
+ await _take_login_attempt(db, _session_counter(current_user))
if not await verify_password_off_loop(
body.auth_key, current_user.pw_hash, current_user.pw_salt,
current_user.pw_version):
raise HTTPException(status_code=403,
detail="Passphrase does not match")
- await login_throttle.clear(db, current_user.username)
+ await login_throttle.clear(db, _session_counter(current_user))
# How often one account may point the hub at a *different* address.
# Long, because this is the only path where a signed-in account chooses
@@ -1074,12 +1144,12 @@ async def change_password(
current_user: User = Depends(require_user_scope),
db: AsyncSession = Depends(get_db),
):
- await _take_login_attempt(db, current_user.username)
+ await _take_login_attempt(db, _session_counter(current_user))
if not await verify_password_off_loop(body.old_auth_key, current_user.pw_hash,
current_user.pw_salt, current_user.pw_version):
raise HTTPException(status_code=403,
detail="Current passphrase does not match")
- await login_throttle.clear(db, current_user.username)
+ await login_throttle.clear(db, _session_counter(current_user))
if body.new_auth_key == body.old_auth_key:
raise HTTPException(status_code=400,
detail="New passphrase must differ from the current one")
@@ -1279,6 +1349,7 @@ async def password_reset(
update(RefreshToken).where(RefreshToken.user_id == user.id)
.values(revoked=True))
await db.execute(delete(UserDevice).where(UserDevice.user_id == user.id))
+ await db.execute(delete(KnownBrowser).where(KnownBrowser.user_id == user.id))
# A code sent to the address on file is a stronger proof than a passphrase,
# and it is the way out of a lockout somebody else caused.
await login_throttle.clear(db, user.username)
@@ -1493,6 +1564,7 @@ async def erase_account(db: AsyncSession, user: User, owned_groups: str = "refus
GroupHost.node_id.in_(select(Node.id).where(Node.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(KnownBrowser).where(KnownBrowser.user_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
@@ -1538,11 +1610,11 @@ async def delete_own_account(
borrowed laptop or a session left open. Same value as at sign-in, so the hub
still never sees the passphrase itself.
"""
- await _take_login_attempt(db, current_user.username)
+ await _take_login_attempt(db, _session_counter(current_user))
if not await verify_password_off_loop(body.auth_key, current_user.pw_hash, current_user.pw_salt,
current_user.pw_version):
raise HTTPException(status_code=403, detail="Passphrase does not match")
- await login_throttle.clear(db, current_user.username)
+ await login_throttle.clear(db, _session_counter(current_user))
return await erase_account(db, current_user)
diff --git a/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/d4e5f6a7b8ca_known_browsers.py b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/d4e5f6a7b8ca_known_browsers.py
new file mode 100644
index 0000000..aaad5c7
--- /dev/null
+++ b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/d4e5f6a7b8ca_known_browsers.py
@@ -0,0 +1,32 @@
+"""browsers an account has signed in from, each with its own failure counter
+
+Revision ID: d4e5f6a7b8ca
+Revises: c3d4e5f6a7b9
+"""
+
+from collections.abc import Sequence
+
+import sqlalchemy as sa
+from alembic import op
+
+revision: str = "d4e5f6a7b8ca"
+down_revision: str | Sequence[str] | None = "c3d4e5f6a7b9"
+branch_labels: str | Sequence[str] | None = None
+depends_on: str | Sequence[str] | None = None
+
+
+def upgrade() -> None:
+ op.create_table(
+ "known_browsers",
+ sa.Column("id", sa.String(36), primary_key=True),
+ sa.Column("user_id", sa.String(36), sa.ForeignKey("users.id"), nullable=False),
+ sa.Column("token_hash", sa.String(64), nullable=False, unique=True),
+ sa.Column("created_at", sa.DateTime(timezone=True)),
+ sa.Column("last_used_at", sa.DateTime(timezone=True)),
+ )
+ op.create_index("ix_known_browsers_user_id", "known_browsers", ["user_id"])
+
+
+def downgrade() -> None:
+ op.drop_index("ix_known_browsers_user_id", table_name="known_browsers")
+ op.drop_table("known_browsers")
diff --git a/packages/meshbay-hub/src/meshbay_hub/db/models.py b/packages/meshbay-hub/src/meshbay_hub/db/models.py
index 1e652a6..dbc0f10 100644
--- a/packages/meshbay-hub/src/meshbay_hub/db/models.py
+++ b/packages/meshbay-hub/src/meshbay_hub/db/models.py
@@ -429,6 +429,25 @@ class MailQuota(Base):
last_sent: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
+class KnownBrowser(Base):
+ """A browser this account has signed in from, for the sign-in lockout.
+
+ Its own failure counter, which a stranger cannot spend: the lockout keyed by
+ username alone let anyone who knew a name keep its owner out of every
+ browser, four requests an hour. Only a hash of the token is kept. It is not
+ a credential — a sign-in presenting it still needs the passphrase.
+ """
+
+ __tablename__ = "known_browsers"
+
+ id: Mapped[str] = mapped_column(String(36), primary_key=True, default=_uuid)
+ user_id: Mapped[str] = mapped_column(ForeignKey("users.id"), nullable=False,
+ index=True)
+ token_hash: Mapped[str] = mapped_column(String(64), unique=True, nullable=False)
+ created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_now)
+ last_used_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_now)
+
+
class LoginThrottle(Base):
"""Wrong passphrases per username, for the sign-in lockout (`login_throttle.py`).
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/create-group-page.js b/packages/meshbay-hub/src/meshbay_hub/static/create-group-page.js
index d4c2ab8..ec4a5b6 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/create-group-page.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/create-group-page.js
@@ -77,7 +77,7 @@ function CreateGroupFormSimple({ token, onCreated, allowPublicGroups = true }) {
<div class="form-field">
<label class="form-label">${t('create_group.name')}</label>
<input type="text" placeholder="${t('create_group.name_placeholder')}"
- value=${name} onInput=${e => setName(e.target.value)} required autofocus />
+ value=${name} onInput=${e => setName(e.target.value)} required autofocus maxlength="128" />
</div>
<div class="form-field" style="margin-bottom:0">
@@ -255,6 +255,9 @@ function CreateGroupWizard({ token, username, onCreated, onNodeLinked, allowPubl
name: name.trim(),
path: mainRoot.path,
writable: mainRoot.writable !== false,
+ // How people join is set on the node, from this form — the node does
+ // not take it from the hub.
+ joinPolicy,
};
await platform.node.op('attachGroup', attachBody);
await platform.node.op('reload');
@@ -371,7 +374,7 @@ function CreateGroupWizard({ token, username, onCreated, onNodeLinked, allowPubl
<div class="form-field">
<label class="form-label">${t('create_group.name')}</label>
<input type="text" placeholder="${t('create_group.name_placeholder')}"
- value=${name} onInput=${e => setName(e.target.value)} required autofocus />
+ value=${name} onInput=${e => setName(e.target.value)} required autofocus maxlength="128" />
</div>
<div class="form-field">
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/csv.js b/packages/meshbay-hub/src/meshbay_hub/static/csv.js
new file mode 100644
index 0000000..310bdca
--- /dev/null
+++ b/packages/meshbay-hub/src/meshbay_hub/static/csv.js
@@ -0,0 +1,14 @@
+/**
+ * One CSV cell, quoted when it has to be, and never a formula.
+ *
+ * A spreadsheet runs a cell that starts with `=`, `+`, `-` or `@` (or a tab or a
+ * carriage return in front of one) as a formula. The audit export carries text
+ * a member chose — a refused blob's kind, a file name — so such a cell is given
+ * a leading apostrophe, which spreadsheets read as "this is text", and which is
+ * what other exports do.
+ */
+export function csvCell(value) {
+ let s = value == null ? '' : String(value);
+ if (/^[=+\-@\t\r]/.test(s)) s = `'${s}`;
+ return /[",\n\r]/.test(s) ? `"${s.replace(/"/g, '""')}"` : s;
+}
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js b/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js
index 6bd5896..7bbcac5 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js
@@ -161,6 +161,12 @@ async function deriveBundleSessionKey(password, username, userId, pepperB64, pep
// HKDF keys are non-extractable by specification.
v3: await crypto.subtle.importKey('raw', m, 'HKDF', false, ['deriveKey', 'deriveBits']),
pepperVersion: pepperVersion || 1,
+ // TRANSITIONAL — the key MBK2 bundles were sealed under, which this same
+ // Argon2 run produces anyway. Kept for the session so a node still holding
+ // one has it opened and replaced by MBK3 on the account's next visit,
+ // rather than the member being re-invited. Decrypt only; nothing is sealed
+ // under it. Remove once no MBK2 bundle is left on any node.
+ legacy: await crypto.subtle.importKey('raw', a, { name: 'AES-GCM' }, false, ['decrypt']),
};
}
@@ -321,13 +327,52 @@ async function encryptBundle(skEdRaw, skXRaw, aesKey, { userId, nodePk, pepperVe
return btoa(String.fromCharCode(...out));
}
-/** 'current', or 'retired' for anything written before MBK3. */
+// TRANSITIONAL — the format before MBK3: "MBK2" ‖ nonce (12) ‖ AES-GCM under
+// the passphrase's Argon2 key alone, no associated data. Read once to be
+// replaced; never written.
+const LEGACY_MAGIC = 'MBK2';
+
+/**
+ * 'current'; 'legacy' for MBK2, opened once with the session's legacy key and
+ * replaced; 'retired' for anything older, which is not read at all.
+ */
function bundleFormat(bundleB64) {
try {
- return atob(bundleB64).startsWith(BUNDLE_MAGIC) ? 'current' : 'retired';
+ const head = atob(bundleB64).slice(0, 4);
+ if (head === BUNDLE_MAGIC) return 'current';
+ return head === LEGACY_MAGIC ? 'legacy' : 'retired';
} catch { return 'retired'; }
}
+/** TRANSITIONAL — open an MBK2 bundle (passphrase or recovery copy). */
+async function decryptLegacyBundle(bundleB64, aesKey) {
+ if (bundleFormat(bundleB64) !== 'legacy') throw new Error('not an MBK2 bundle');
+ const raw = _b64bytes(bundleB64);
+ const off = LEGACY_MAGIC.length;
+ const plain = await crypto.subtle.decrypt(
+ { name: 'AES-GCM', iv: raw.slice(off, off + 12) }, aesKey, raw.slice(off + 12));
+ return JSON.parse(new TextDecoder().decode(plain));
+}
+
+/**
+ * TRANSITIONAL — an identity read from an MBK2 bundle, sealed again as MBK3
+ * for the same node (and the recovery copy too, when a recovery key is in
+ * hand), for the caller to store in place of the old one.
+ */
+async function resealLegacyIdentity(keys, sessionKey, recoveryKey, { userId, nodePk }) {
+ const skEd = _b64bytes(keys.skEd);
+ const skX = _b64bytes(keys.skX);
+ const out = {
+ bundleEnc: await encryptBundle(skEd, skX, await nodeBundleKey(sessionKey, nodePk),
+ { userId, nodePk, pepperVersion: sessionKey.pepperVersion }),
+ };
+ if (recoveryKey) {
+ out.bundleEncRecovery = await encryptBundle(skEd, skX, recoveryKey,
+ { userId, nodePk, pepperVersion: 0 });
+ }
+ return out;
+}
+
// ── Registration ──────────────────────────────────────────────────────────────
/**
@@ -426,13 +471,36 @@ async function decryptBundle(bundleB64, aesKey, { userId, nodePk }) {
* decrypts it and returns the keys + encrypted bundle for push to node.
* Otherwise returns bundleKey so the caller can fetch from node during handshake.
*/
+// The token the hub gave this browser at an earlier sign-in, per account. It
+// is not a credential — the passphrase is still asked — but a sign-in that
+// presents it has a failure counter of its own, so a stranger who keeps
+// failing on this account's name locks only browsers it has never used.
+// Kept across sign-outs on purpose: forgetting it would be the lockout again.
+const KNOWN_BROWSERS = 'mb_known_browsers';
+
+function _knownBrowser(username) {
+ try { return (JSON.parse(localStorage.getItem(KNOWN_BROWSERS)) || {})[username] || null; }
+ catch { return null; }
+}
+
+function _rememberBrowser(username, token) {
+ if (!token) return;
+ try {
+ const all = JSON.parse(localStorage.getItem(KNOWN_BROWSERS)) || {};
+ all[username] = token;
+ localStorage.setItem(KNOWN_BROWSERS, JSON.stringify(all));
+ } catch { /* storage refused: this browser stays an unknown one */ }
+}
+
async function loginAndRecover(username, password) {
const authKey = await deriveAuthKey(password, username);
+ const known = _knownBrowser(username);
const resp = await hubCall('/v1/users/login', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
- body: JSON.stringify({ username, auth_key: authKey }),
+ body: JSON.stringify({ username, auth_key: authKey,
+ ...(known ? { known_browser: known } : {}) }),
});
if (!resp.ok) {
@@ -454,6 +522,7 @@ async function loginAndRecover(username, password) {
}
const data = await resp.json();
+ _rememberBrowser(username, data.known_browser);
const result = {
accessToken: data.access_token,
refreshToken: data.refresh_token,
@@ -492,7 +561,7 @@ window.MeshBayKeys = {
// The bundle key (docs/MESHBAY_DESIGN.md §3.1, §3.7): one session key per
// sign-in, one derived key per node, one format.
deriveBundleSessionKey, sessionBundleKey, nodeBundleKey, fetchBundlePepper,
- encryptBundle, decryptBundle, bundleFormat,
+ encryptBundle, decryptBundle, bundleFormat, decryptLegacyBundle, resealLegacyIdentity,
// Account recovery key (docs/MESHBAY_DESIGN.md §3.6).
generateRecoveryKey, deriveRecoveryKey,
};
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/node-page.js b/packages/meshbay-hub/src/meshbay_hub/static/node-page.js
index f940934..41e9567 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/node-page.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/node-page.js
@@ -1,3 +1,4 @@
+import { csvCell } from './csv.js';
import {
html, useState, useEffect, useCallback, useRef,
} from './vendor/htm-preact.js';
@@ -585,17 +586,13 @@ export function NodePage({ groups, token, username }) {
const cols = ['timestamp', 'event', 'user', 'user_id', 'ip',
'group', 'group_id', 'detail'];
- const esc = (v) => {
- const s = v == null ? '' : String(v);
- return /[",\n\r]/.test(s) ? '"' + s.replace(/"/g, '""') + '"' : s;
- };
const lines = [cols.join(',')];
for (const e of rows) {
lines.push([
new Date(e.timestamp * 1000).toISOString(),
e.event, e.username || '', e.user_id || '', e.ip || '',
e.group_name || '', e.group_id || '', e.detail || '',
- ].map(esc).join(','));
+ ].map(csvCell).join(','));
}
const csv = lines.join('\r\n') + '\r\n';
const stamp = new Date().toISOString().slice(0, 19).replace(/[:T]/g, '-');
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/portable-name.js b/packages/meshbay-hub/src/meshbay_hub/static/portable-name.js
index bb2438c..34085ae 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/portable-name.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/portable-name.js
@@ -20,8 +20,13 @@ const WINDOWS_RESERVED = new Set([
]);
const RESERVED_CHARS = new Set('<>:"/\\|?*');
+// The bidirectional controls: "invoice\u202efdp.exe" displays as
+// "invoiceexe.pdf", and a saved name must say what the file is.
+const BIDI_CONTROLS = new Set('\u061c\u200e\u200f\u202a\u202b\u202c\u202d\u202e'
+ + '\u2066\u2067\u2068\u2069');
-const reserved = (c) => RESERVED_CHARS.has(c) || c.charCodeAt(0) < 32;
+const reserved = (c) => RESERVED_CHARS.has(c) || BIDI_CONTROLS.has(c)
+ || c.charCodeAt(0) < 32;
function isPortable(name) {
if (!name || name === '.' || name === '..') return false;
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport-rewrap.js b/packages/meshbay-hub/src/meshbay_hub/static/transport-rewrap.js
index 8bf884c..cdf86b0 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/transport-rewrap.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/transport-rewrap.js
@@ -117,7 +117,9 @@ async function rewrapAllNodes(o) {
anyOk = true;
continue;
}
- if (tp.newNodeBundle) {
+ // An identity read from an MBK2 bundle (TRANSITIONAL) is an existing
+ // one, and is re-sealed below like any other.
+ if (tp.newNodeBundle && !tp.upgradedLegacy) {
// No identity existed on this node — connect just minted one under
// the old key. Don't persist it: the next time this group is opened
// the normal flow creates one under the current key, and storing it
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport.js b/packages/meshbay-hub/src/meshbay_hub/static/transport.js
index 9b86921..f9e1370 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/transport.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/transport.js
@@ -684,6 +684,8 @@ class MeshBayTransport {
/** Set on a first join: the identity created for this node, still to be left with it. */
get newNodeBundle() { return this._newNodeBundle || null; }
+ /** TRANSITIONAL — the identity was read from an MBK2 bundle, not created. */
+ get upgradedLegacy() { return Boolean(this._upgradedLegacy); }
set newNodeBundle(v) { this._newNodeBundle = v; }
/** The recovery-wrapped copy of that same first-join identity, when a recovery key was in hand. */
@@ -765,6 +767,7 @@ class MeshBayTransport {
this._groupId = groupId || '';
this._newNodeBundle = null;
this._newNodeBundleRecovery = null;
+ this._upgradedLegacy = false;
this._joinError = null;
// Per connection, for the same reason the chat keys and the roster are
// dropped further down: the device the *previous* connection identified
@@ -1048,6 +1051,8 @@ class MeshBayTransport {
fresh = await this._settleNativeIdentity(kpResp);
} else if (this._nodeHasBundle && K.bundleFormat(kpResp.bundle_enc) === 'retired') {
throw _retiredBundleError();
+ } else if (this._nodeHasBundle && K.bundleFormat(kpResp.bundle_enc) === 'legacy') {
+ keys = await this._openLegacyBundle(kpResp, sealedFor);
} else if (this._nodeHasBundle) {
try {
keys = await K.decryptBundle(kpResp.bundle_enc,
@@ -1298,6 +1303,46 @@ class MeshBayTransport {
* hangs, the textbox is dead" report. Every exit below names itself.
*/
/**
+ * TRANSITIONAL — an MBK2 bundle, opened with the session's legacy key (or
+ * the recovery copy with the recovery key) and sealed again as MBK3, left
+ * for `settleNodeBundle` to store in its place once the connection is made.
+ *
+ * A session restored from before the legacy key was kept has none: the
+ * passphrase is asked for again (`no_keys`) rather than the identity being
+ * declared lost. A legacy key that does not open it — a bundle sealed under
+ * an older passphrase — is what a current bundle that does not open is: the
+ * caller goes on to a first join.
+ */
+ async _openLegacyBundle(kpResp, sealedFor) {
+ const K = window.MeshBayKeys;
+ let keys = null;
+ if (this._bundleKey.legacy) {
+ try { keys = await K.decryptLegacyBundle(kpResp.bundle_enc, this._bundleKey.legacy); }
+ catch { /* sealed under another passphrase */ }
+ }
+ if (!keys && this._recoveryKey && kpResp.bundle_enc_recovery
+ && K.bundleFormat(kpResp.bundle_enc_recovery) === 'legacy') {
+ try {
+ keys = await K.decryptLegacyBundle(kpResp.bundle_enc_recovery, this._recoveryKey);
+ this._recoveredFromRecovery = true;
+ } catch { /* not this recovery key */ }
+ }
+ if (!keys) {
+ if (!this._bundleKey.legacy && !this._recoveryKey) {
+ const err = new Error('Your passphrase is needed once to update how this node keeps your identity');
+ err.reason = 'no_keys';
+ throw err;
+ }
+ return null;
+ }
+ const sealed = await K.resealLegacyIdentity(keys, this._bundleKey, this._recoveryKey, sealedFor);
+ this._newNodeBundle = sealed.bundleEnc;
+ this._newNodeBundleRecovery = sealed.bundleEncRecovery || null;
+ this._upgradedLegacy = true;
+ return keys;
+ }
+
+ /**
* This node's identity when the desktop application holds the keys.
*
* Kept by the application once it has it, so a bundle left on the node —
@@ -1316,8 +1361,16 @@ class MeshBayTransport {
throw _retiredBundleError();
}
try {
+ // An MBK2 bundle too (TRANSITIONAL): the application opens it with
+ // the legacy key it kept from the passphrase, and `settleNodeBundle`
+ // then replaces or withdraws it as browser access says.
pub = await P.openBundle(uid, pk, { bundleEnc: kpResp.bundle_enc });
} catch (e) {
+ if (String(e && e.message).includes('no_legacy_key')) {
+ const err = new Error('Your passphrase is needed once to update how this node keeps your identity');
+ err.reason = 'no_keys';
+ throw err;
+ }
// Sealed under a passphrase no longer in use: as in a browser, a
// passphrase change must report it, and a first join replaces it.
if (this._rewrapOnly) throw new Error('could not open the stored identity');
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/webrtc-test.html b/packages/meshbay-hub/src/meshbay_hub/static/webrtc-test.html
deleted file mode 100644
index 46003a7..0000000
--- a/packages/meshbay-hub/src/meshbay_hub/static/webrtc-test.html
+++ /dev/null
@@ -1,265 +0,0 @@
-<!DOCTYPE html>
-<html lang="en">
-<head>
- <meta charset="utf-8">
- <meta name="viewport" content="width=device-width, initial-scale=1">
- <title>MeshBay — WebRTC Spike Test</title>
- <style>
- *, *::before, *::after { box-sizing: border-box; }
- body { font-family: system-ui, sans-serif; margin: 0; background: #0f172a; color: #e2e8f0; }
- .container { max-width: 800px; margin: 32px auto; padding: 0 16px; }
- h1 { color: #38bdf8; font-size: 1.4em; }
- h2 { color: #94a3b8; font-size: 1.1em; margin-top: 2em; }
- .step { background: #1e293b; border: 1px solid #334155; border-radius: 8px;
- padding: 16px; margin: 12px 0; }
- .step.done { border-color: #22c55e; }
- .step.fail { border-color: #ef4444; }
- .step.active { border-color: #38bdf8; }
- input { padding: 8px 12px; border: 1px solid #475569; border-radius: 6px;
- background: #0f172a; color: #e2e8f0; font-size: 0.95em; margin: 4px; width: 240px; }
- button { padding: 8px 20px; background: #0ea5e9; color: #fff; border: none;
- border-radius: 6px; cursor: pointer; font-size: 0.95em; margin: 4px; }
- button:hover { background: #0284c7; }
- button:disabled { background: #475569; cursor: not-allowed; }
- #log { background: #020617; border: 1px solid #1e293b; border-radius: 8px;
- padding: 12px; font-family: monospace; font-size: 0.85em; line-height: 1.6;
- max-height: 400px; overflow-y: auto; white-space: pre-wrap; }
- .ok { color: #22c55e; }
- .err { color: #ef4444; }
- .info { color: #38bdf8; }
- .warn { color: #f59e0b; }
- .dim { color: #64748b; }
- .badge { display: inline-block; background: #22c55e; color: #0f172a; padding: 2px 8px;
- border-radius: 4px; font-size: 0.8em; font-weight: bold; margin-left: 8px; }
- .badge.fail { background: #ef4444; color: #fff; }
- </style>
-</head>
-<body>
-<div class="container">
- <h1>MeshBay — WebRTC DataChannel Spike Test</h1>
- <p class="dim">Phase 9.5 — E2E browser → NAT → node file transfer via WebRTC</p>
-
- <div class="step" id="step-login">
- <h2>1. Login to Hub</h2>
- <input id="username" placeholder="Username" value="bob">
- <input id="password" placeholder="Password" type="password" value="bob">
- <button id="btn-login" onclick="doLogin()">Login</button>
- <span id="login-status"></span>
- </div>
-
- <div class="step" id="step-connect">
- <h2>2. Connect to Node via WebRTC</h2>
- <input id="node-id" placeholder="Node ID">
- <input id="group-id" placeholder="Group ID (optional)">
- <button id="btn-connect" onclick="doConnect()" disabled>Connect</button>
- <span id="connect-status"></span>
- </div>
-
- <div class="step" id="step-transfer">
- <h2>3. File Transfer Test</h2>
- <button id="btn-index" onclick="doFetchIndex()" disabled>Fetch Index</button>
- <br>
- <input id="file-id" placeholder="File ID (blake3 hex, from node log)">
- <button id="btn-chunk" onclick="doFetchChunk()" disabled>Fetch Chunk</button>
- <span id="transfer-status"></span>
- </div>
-
- <h2>Log</h2>
- <div id="log"></div>
-</div>
-
-<script src="/transport.js?v=2"></script>
-<script>
-const HUB_URL = window.location.origin;
-const params = new URLSearchParams(window.location.search);
-let accessToken = null;
-let jwtToken = null;
-let transport = null;
-let fileIndex = null;
-let connecting = false;
-
-// Pre-fill from URL params
-if (params.get('user')) document.getElementById('username').value = params.get('user');
-if (params.get('pass')) document.getElementById('password').value = params.get('pass');
-if (params.get('node')) document.getElementById('node-id').value = params.get('node');
-if (params.get('group')) document.getElementById('group-id').value = params.get('group');
-if (params.get('file')) document.getElementById('file-id').value = params.get('file').replace(/\s+/g, '');
-
-// Auto-run if all params provided
-if (params.get('auto')) {
- setTimeout(async () => {
- await doLogin();
- if (accessToken) await doConnect();
- if (transport && transport.connected) {
- await doFetchIndex();
- if (document.getElementById('file-id').value) await doFetchChunk();
- }
- }, 500);
-}
-
-function logMsg(cls, text) {
- const el = document.getElementById('log');
- const line = document.createElement('span');
- line.className = cls;
- line.textContent = text + '\n';
- el.appendChild(line);
- el.scrollTop = el.scrollHeight;
-}
-
-function setStep(id, state) {
- const el = document.getElementById(id);
- el.className = 'step ' + state;
-}
-
-async function doLogin() {
- const user = document.getElementById('username').value;
- const pass = document.getElementById('password').value;
- logMsg('info', `Logging in as ${user}...`);
- setStep('step-login', 'active');
-
- try {
- const resp = await fetch(`${HUB_URL}/v1/users/login`, {
- method: 'POST',
- headers: { 'Content-Type': 'application/json' },
- body: JSON.stringify({ username: user, password: pass }),
- });
-
- if (!resp.ok) {
- const err = await resp.json();
- throw new Error(err.detail || resp.statusText);
- }
-
- const data = await resp.json();
- accessToken = data.access_token;
- jwtToken = data.access_token;
- logMsg('ok', `Login OK — token: ${accessToken.substring(0, 20)}...`);
- setStep('step-login', 'done');
- document.getElementById('login-status').innerHTML = '<span class="badge">OK</span>';
- document.getElementById('btn-connect').disabled = false;
- } catch (e) {
- logMsg('err', `Login FAILED: ${e.message}`);
- setStep('step-login', 'fail');
- document.getElementById('login-status').innerHTML = '<span class="badge fail">FAIL</span>';
- }
-}
-
-async function doConnect() {
- if (connecting) { logMsg('warn', 'Connect already in progress'); return; }
- const nodeId = document.getElementById('node-id').value;
- const groupId = document.getElementById('group-id').value;
- if (!nodeId) { logMsg('warn', 'Enter a node ID'); return; }
-
- connecting = true;
- if (transport) { transport.close(); transport = null; }
-
- logMsg('info', `Connecting to node ${nodeId.substring(0, 8)}... via WebRTC`);
- setStep('step-connect', 'active');
-
- try {
- transport = new MeshBayTransport(HUB_URL, accessToken);
-
- logMsg('dim', ' Creating RTCPeerConnection...');
- logMsg('dim', ' Creating DataChannel "mnp"...');
- logMsg('dim', ' Gathering ICE candidates...');
- logMsg('dim', ' Sending SDP offer to hub...');
-
- const t0 = performance.now();
- const ack = await transport.connect(nodeId, jwtToken, groupId);
- const elapsed = (performance.now() - t0).toFixed(0);
-
- logMsg('ok', `WebRTC connected in ${elapsed}ms`);
- logMsg('ok', ` MNP handshake_ack — node_pk: ${ack.node_pk?.substring(0, 16)}...`);
- logMsg('ok', ` DataChannel state: ${transport._channel?.readyState}`);
- setStep('step-connect', 'done');
- document.getElementById('connect-status').innerHTML = '<span class="badge">P2P OK</span>';
- document.getElementById('btn-index').disabled = false;
- document.getElementById('btn-chunk').disabled = false;
- } catch (e) {
- logMsg('err', `Connection FAILED: ${e.message}`);
- setStep('step-connect', 'fail');
- document.getElementById('connect-status').innerHTML = '<span class="badge fail">FAIL</span>';
- } finally {
- connecting = false;
- }
-}
-
-async function doFetchIndex() {
- logMsg('info', `Fetching Mesh Group Index... (channel: ${transport?._channel?.readyState})`);
- try {
- const t0 = performance.now();
- const indexBytes = await transport.fetchIndex();
- const elapsed = (performance.now() - t0).toFixed(0);
-
- logMsg('ok', `Index received: ${indexBytes.byteLength} bytes in ${elapsed}ms`);
-
- try {
- const envelope = msgpack_decode(indexBytes);
- logMsg('dim', ` type: ${envelope.type}, encrypted: ${envelope.encrypted}, version: ${envelope.version}`);
- logMsg('dim', ` group_id: ${envelope.group_id}`);
-
- if (envelope.encrypted) {
- logMsg('warn', ` Index is GEK-encrypted — browser decryption not implemented in spike`);
- logMsg('dim', ` ct_b64 length: ${envelope.ct_b64?.length || 0} chars`);
- logMsg('info', ` Spike workaround: enter a file_id manually or use Fetch First Chunk`);
- // Store envelope so chunk test can proceed with manual file_id
- fileIndex = { entries: [], envelope };
- } else {
- // Public group: decompress and parse
- logMsg('dim', ` Public index — data_b64 length: ${envelope.data_b64?.length || 0}`);
- fileIndex = { entries: [], envelope };
- }
- } catch (pe) {
- logMsg('warn', ` Could not parse index envelope: ${pe.message}`);
- }
- } catch (e) {
- logMsg('err', `Index fetch FAILED: ${e.message}`);
- }
-}
-
-async function doFetchChunk() {
- let fileId = document.getElementById('file-id').value.replace(/\s+/g, '');
-
- if (!fileId) {
- logMsg('warn', 'Enter a file_id (blake3 hex hash from node indexer log)');
- logMsg('dim', ' Look for "Initial scan complete" in the node terminal');
- logMsg('dim', ' Or run: python -c "import blake3; print(blake3.blake3(open(\'QE/demo-v3/shared_media/sample.txt\',\'rb\').read()).hexdigest())"');
- return;
- }
-
- logMsg('info', `Fetching chunk 0 of ${fileId.substring(0, 16)}... (channel: ${transport?._channel?.readyState})`);
-
- try {
- const t0 = performance.now();
- const chunkMsg = await transport.fetchChunk(fileId, 0);
- const elapsed = (performance.now() - t0).toFixed(0);
-
- if (chunkMsg.type === 'error') {
- logMsg('err', `Chunk fetch error: ${chunkMsg.detail}`);
- return;
- }
-
- logMsg('ok', `Chunk received in ${elapsed}ms:`);
- logMsg('ok', ` type: ${chunkMsg.type}`);
- logMsg('ok', ` chunk_index: ${chunkMsg.chunk_index}`);
- logMsg('ok', ` plaintext_size: ${chunkMsg.plaintext_size} bytes`);
- logMsg('ok', ` ct_b64 length: ${chunkMsg.ct_b64?.length || 0} chars`);
- logMsg('ok', ` nonce_b64: ${chunkMsg.nonce_b64?.substring(0, 16)}...`);
- logMsg('ok', ` sig_b64: ${chunkMsg.sig_b64?.substring(0, 16)}...`);
-
- logMsg('', '');
- logMsg('ok', '=== SPIKE TEST PASSED ===');
- logMsg('ok', 'Browser connected to node via WebRTC DataChannel.');
- logMsg('ok', 'MNP handshake, index sync, and file chunk transfer all work.');
- logMsg('ok', 'Data flowed P2P — hub was only used for signaling.');
-
- setStep('step-transfer', 'done');
- document.getElementById('transfer-status').innerHTML = '<span class="badge">E2E OK</span>';
- } catch (e) {
- logMsg('err', `Chunk fetch FAILED: ${e.message}`);
- setStep('step-transfer', 'fail');
- document.getElementById('transfer-status').innerHTML = '<span class="badge fail">FAIL</span>';
- }
-}
-</script>
-</body>
-</html>
diff --git a/packages/meshbay-hub/tests/test_bundle_key.py b/packages/meshbay-hub/tests/test_bundle_key.py
index 333f8f8..f6c99f8 100644
--- a/packages/meshbay-hub/tests/test_bundle_key.py
+++ b/packages/meshbay-hub/tests/test_bundle_key.py
@@ -156,7 +156,7 @@ def test_an_earlier_format_is_refused_by_name(tmp_path):
}
console.log(JSON.stringify(results));
""")
- assert out == [["retired", "bundle_format_retired"], ["retired", "bundle_format_retired"]]
+ assert out == [["legacy", "bundle_format_retired"], ["retired", "bundle_format_retired"]]
def test_two_devices_of_one_account_derive_the_same_playlist_key(tmp_path):
@@ -181,3 +181,49 @@ def test_the_playlist_key_is_no_node_key(tmp_path):
}));
""")
assert out["distinct"]
+
+
+def test_an_mbk2_bundle_is_opened_once_and_sealed_again_as_mbk3(tmp_path):
+ """
+ TRANSITIONAL. Nodes still hold bundles sealed under the passphrase's Argon2
+ key alone. The same Argon2 run that makes `M` makes that key, so the session
+ keeps it — decrypt only — and the identity is moved to MBK3 on the account's
+ next visit instead of the member being re-invited.
+ """
+ out = _run(tmp_path, """
+ argonCalls = 0;
+ const sk = await K().deriveBundleSessionKey('p', 'someone', 'uid-1', PEPPER, 1);
+ const calls = argonCalls;
+ // An MBK2 bundle as 0.16 wrote it: "MBK2" ‖ nonce ‖ AES-GCM(A), no AAD.
+ const a = await crypto.subtle.importKey('raw', await _bundleKeyBytes('p', 'someone'),
+ { name: 'AES-GCM' }, false, ['encrypt']);
+ const nonce = new Uint8Array(12).fill(3);
+ const plain = new TextEncoder().encode(JSON.stringify({ skEd: btoa('ED'), skX: btoa('XX') }));
+ const ct = new Uint8Array(await crypto.subtle.encrypt({ name: 'AES-GCM', iv: nonce }, a, plain));
+ const raw = new Uint8Array(4 + 12 + ct.length);
+ raw.set(new TextEncoder().encode('MBK2')); raw.set(nonce, 4); raw.set(ct, 16);
+ const mbk2 = btoa(String.fromCharCode(...raw));
+
+ const keys = await K().decryptLegacyBundle(mbk2, sk.legacy);
+ const resealed = await K().resealLegacyIdentity(keys, sk, null, { userId: 'uid-1', nodePk: 'NODE' });
+ const back = await K().decryptBundle(resealed.bundleEnc, await K().nodeBundleKey(sk, 'NODE'),
+ { userId: 'uid-1', nodePk: 'NODE' });
+ let otherPassphrase = 'opened';
+ const sk2 = await K().deriveBundleSessionKey('another', 'someone', 'uid-1', PEPPER, 1);
+ try { await K().decryptLegacyBundle(mbk2, sk2.legacy); } catch { otherPassphrase = 'refused'; }
+ let sealsUnderLegacy = 'yes';
+ try { await crypto.subtle.encrypt({ name: 'AES-GCM', iv: nonce }, sk.legacy, plain); }
+ catch { sealsUnderLegacy = 'no'; }
+ console.log(JSON.stringify({
+ calls, format: K().bundleFormat(mbk2), keys, newFormat: K().bundleFormat(resealed.bundleEnc),
+ back, otherPassphrase, sealsUnderLegacy, extractable: sk.legacy.extractable,
+ }));
+ """)
+ assert out["calls"] == 1, "keeping the legacy key must not cost a second Argon2 run"
+ assert out["format"] == "legacy"
+ assert out["keys"] == {"skEd": "RUQ=", "skX": "WFg="}
+ assert out["newFormat"] == "current"
+ assert out["back"] == out["keys"]
+ assert out["otherPassphrase"] == "refused"
+ assert out["sealsUnderLegacy"] == "no", "the legacy key opens; it never seals"
+ assert out["extractable"] is False
diff --git a/packages/meshbay-hub/tests/test_csv_cell.py b/packages/meshbay-hub/tests/test_csv_cell.py
new file mode 100644
index 0000000..ca38805
--- /dev/null
+++ b/packages/meshbay-hub/tests/test_csv_cell.py
@@ -0,0 +1,32 @@
+"""
+The audit export never hands a spreadsheet a formula.
+
+It carries text a member chose — a refused blob's kind, a file name. A cell that
+starts with `=`, `+`, `-` or `@` runs as a formula when the operator opens the
+file, so it is given a leading apostrophe, which spreadsheets read as text.
+"""
+
+import json
+import shutil
+import subprocess
+from pathlib import Path
+
+import pytest
+
+CSV = Path(__file__).resolve().parents[1] / "src" / "meshbay_hub" / "static" / "csv.js"
+
+pytestmark = pytest.mark.skipif(shutil.which("node") is None, reason="node unavailable")
+
+CASES = ['=HYPERLINK("http://x","y")', "+1+1", "-2+3", "@SUM(A1)", "\t=1", "plain",
+ 'with "quotes"', "a,b", "", None, 12]
+
+
+def test_a_cell_is_text_and_quoted_where_it_must_be(tmp_path):
+ script = tmp_path / "case.mjs"
+ script.write_text(
+ f"import {{ csvCell }} from '{CSV.as_uri()}';\n"
+ f"console.log(JSON.stringify({json.dumps(CASES)}.map(csvCell)));\n")
+ out = json.loads(subprocess.run(["node", str(script)], capture_output=True,
+ text=True, check=True).stdout)
+ assert out == ['"\'=HYPERLINK(""http://x"",""y"")"', "'+1+1", "'-2+3", "'@SUM(A1)",
+ "'\t=1", "plain", '"with ""quotes"""', '"a,b"', "", "", "12"]
diff --git a/packages/meshbay-hub/tests/test_desktop_keyring.py b/packages/meshbay-hub/tests/test_desktop_keyring.py
index 963ee55..e55e6d0 100644
--- a/packages/meshbay-hub/tests/test_desktop_keyring.py
+++ b/packages/meshbay-hub/tests/test_desktop_keyring.py
@@ -119,10 +119,33 @@ const v = JSON.parse(fs.readFileSync(input, 'utf8'));
await K.nodeBundleKey(sk, 'NODE-P'), { userId: v.userId, nodePk: 'NODE-P' });
out.sealed_here_opens_in_page = back.skX === pageId.skXB64;
+ // 3b. TRANSITIONAL: an MBK2 bundle, sealed under the Argon2 key alone as
+ // 0.16 wrote it, is opened with the legacy key kept beside M.
+ {
+ const nc = require('crypto');
+ const salt = nc.createHash('sha256').update(`meshbay:bundle:v2:${v.user}`).digest().subarray(0, 16);
+ const a = await argon2(v.password, salt,
+ { memory: 131072, passes: 3, parallelism: 1, tagLength: 32 });
+ const ed = nc.generateKeyPairSync('ed25519').privateKey.export({ format: 'der', type: 'pkcs8' });
+ const x = nc.generateKeyPairSync('x25519').privateKey.export({ format: 'der', type: 'pkcs8' });
+ const nonce = nc.randomBytes(12);
+ const c = nc.createCipheriv('aes-256-gcm', Buffer.from(a), nonce);
+ const body = Buffer.concat([c.update(JSON.stringify({ skEd: ed.toString('base64'),
+ skX: x.toString('base64') })), c.final()]);
+ const mbk2 = Buffer.concat([Buffer.from('MBK2'), nonce, body, c.getAuthTag()]).toString('base64');
+ out.legacy_open = ring.openBundle(v.userId, 'NODE-L', { bundleEnc: mbk2 });
+ out.legacy_kept = ring.identity(v.userId, 'NODE-L');
+ const keptLegacy = store.masters[v.userId].legacy;
+ delete store.masters[v.userId].legacy;
+ try { ring.openBundle(v.userId, 'NODE-M', { bundleEnc: mbk2 }); out.legacy_missing = 'opened'; }
+ catch (e) { out.legacy_missing = e.message; }
+ store.masters[v.userId].legacy = keptLegacy;
+ }
+
// 4. nothing but public keys come out of the keyring's answers.
out.identity_answer = ring.identity(v.userId, v.node);
out.retired = (() => { try { ring.openBundle(v.userId, 'NODE-R',
- { bundleEnc: Buffer.from('MBK2' + 'x'.repeat(40)).toString('base64') }); }
+ { bundleEnc: Buffer.from('y'.repeat(44)).toString('base64') }); }
catch (e) { return e.code; } })();
out.access_default = ring.browserAccess('someone-else');
ring.setBrowserAccess(v.userId, false);
@@ -284,3 +307,12 @@ def test_the_application_never_asks_its_own_crypto_for_argon2():
source = (KEYRING.parent / name).read_text(encoding="utf-8")
assert "crypto.argon2" not in source, name
assert "wasmArgon2(" in (KEYRING.parent / "main.js").read_text(encoding="utf-8")
+
+
+def test_an_mbk2_bundle_is_opened_with_the_kept_legacy_key(out):
+ """TRANSITIONAL. Kept unsealed, so the next settle replaces the node's copy
+ with MBK3 or withdraws it; without the legacy key the passphrase is asked
+ for, rather than the identity being given up."""
+ assert set(out["legacy_open"]) == {"pkEdB64", "pkXB64"}
+ assert out["legacy_kept"]["sealedWith"] is None
+ assert out["legacy_missing"] == "no_legacy_key"
diff --git a/packages/meshbay-hub/tests/test_desktop_shell.py b/packages/meshbay-hub/tests/test_desktop_shell.py
index 4426dd7..60aa32f 100644
--- a/packages/meshbay-hub/tests/test_desktop_shell.py
+++ b/packages/meshbay-hub/tests/test_desktop_shell.py
@@ -701,3 +701,26 @@ def test_an_account_made_in_the_application_starts_without_browser_access():
assert "platform.keys.createdHere(reg.userId)" in register
assert "userId" in (STATIC / "keyderive.js").read_text(encoding="utf-8").split(
"async function registerUser", 1)[1].split("\n}\n", 1)[0]
+
+
+def _handler(name: str) -> str:
+ source = _main()
+ start = source.index(f"handle('{name}'")
+ return source[start:source.index("\n });", start)]
+
+
+def test_a_finished_download_carries_the_mark_of_the_web():
+ """What a browser leaves on every download, so Windows applies SmartScreen
+ and Protected View. Only checkable here by reading: the stream exists on
+ NTFS alone, and this suite does not run on Windows."""
+ assert "markFromInternet(sink.path)" in _handler("save:end")
+ source = _main()
+ mark = source[source.index("function markFromInternet"):]
+ mark = mark[:mark.index("\n }\n")]
+ assert "Zone.Identifier" in mark and "ZoneId=3" in mark
+ assert "process.platform !== 'win32'" in mark
+
+
+def test_a_saved_name_cannot_hide_its_extension():
+ begin = _handler("save:begin")
+ assert "\\u202a-\\u202e" in begin and "\\u2066-\\u2069" in begin
diff --git a/packages/meshbay-hub/tests/test_group_name_checked.py b/packages/meshbay-hub/tests/test_group_name_checked.py
new file mode 100644
index 0000000..fbc8277
--- /dev/null
+++ b/packages/meshbay-hub/tests/test_group_name_checked.py
@@ -0,0 +1,33 @@
+"""
+A group name is checked where it is created.
+
+The column is 128 characters wide: a longer name was a database error on
+PostgreSQL and a silent truncation on SQLite. And the name is shown to other
+people — in their group list, in the invitation mail — so it carries no line
+breaks or other control characters, and no bidirectional override that makes it
+display as something other than what it is.
+"""
+
+import pytest
+from test_bundle_pepper import _login, _register
+
+
+@pytest.mark.asyncio
+@pytest.mark.parametrize("name,ok", [
+ ("Photos de famille", True),
+ ("x" * 128, True),
+ ("👨‍👩‍👧 Family", True), # a ZWJ sequence is a name, not a trick
+ ("x" * 129, False),
+ ("Films\nClick here", False),
+ ("tab\there", False),
+ ("evil‮gpj.exe", False),
+ (" ", False),
+])
+async def test_a_group_name(client, name, ok):
+ await _register(client, "group_namer")
+ token = (await _login(client, "group_namer"))["access_token"]
+ r = await client.post("/v1/groups", json={"name": name},
+ headers={"Authorization": f"Bearer {token}"})
+ assert (r.status_code < 300) is ok, (name, r.status_code, r.text)
+ if not ok:
+ assert r.status_code == 422
diff --git a/packages/meshbay-hub/tests/test_hub_work_is_bounded.py b/packages/meshbay-hub/tests/test_hub_work_is_bounded.py
new file mode 100644
index 0000000..40ea81d
--- /dev/null
+++ b/packages/meshbay-hub/tests/test_hub_work_is_bounded.py
@@ -0,0 +1,51 @@
+"""
+What an authenticated caller can make the hub do, bounded where it was not.
+
+An offer wrote an IP-log row — kept a year — before any check, for whatever
+string the caller named as a node, and carried an ICE list of any length to the
+node. A node could send `update_groups` as fast as it liked, each one a
+database read, with a list of any length.
+"""
+
+import pytest
+from meshbay_hub.db.models import IPLog
+from sqlalchemy import func, select
+from test_availability_between_members import _make_user
+
+
+def _offer(client, user, node_id, candidates):
+ return client.post(f"/v1/nodes/{node_id}/webrtc/offer",
+ json={"sdp": "v=0\r\n", "ice_candidates": candidates},
+ headers={"Authorization": f"Bearer {user['token']}"})
+
+
+@pytest.mark.asyncio
+async def test_an_offer_that_goes_nowhere_writes_no_log_row(client, db_session):
+ user = await _make_user(client, "offer_nowhere")
+ for i in range(5):
+ r = await _offer(client, user, f"not-a-node-{i}", [])
+ assert r.status_code == 404
+ rows = await db_session.scalar(
+ select(func.count()).select_from(IPLog).where(IPLog.event == "webrtc_offer"))
+ assert rows == 0
+
+
+@pytest.mark.asyncio
+async def test_the_ice_list_is_bounded(client):
+ from meshbay_hub.api.signaling import MAX_ICE_CANDIDATES
+ user = await _make_user(client, "offer_ice")
+ one = {"candidate": "candidate:1 1 udp 2122260223 192.0.2.1 50000 typ host",
+ "sdpMid": "0", "sdpMLineIndex": 0}
+ assert (await _offer(client, user, "nowhere", [one] * MAX_ICE_CANDIDATES)).status_code == 404
+ too_many = [one] * (MAX_ICE_CANDIDATES + 1)
+ assert (await _offer(client, user, "nowhere", too_many)).status_code == 422
+ big = {"candidate": "x" * 40_000}
+ assert (await _offer(client, user, "nowhere", [big])).status_code == 422
+
+
+def test_a_node_reloading_is_budgeted():
+ from meshbay_hub.api.revocation import UPDATE_GROUPS_BURST, _update_budget, _update_window
+ _update_window.clear()
+ assert all(_update_budget("node-a") for _ in range(UPDATE_GROUPS_BURST))
+ assert not _update_budget("node-a")
+ assert _update_budget("node-b"), "one node's budget is not another's"
diff --git a/packages/meshbay-hub/tests/test_known_browser.py b/packages/meshbay-hub/tests/test_known_browser.py
new file mode 100644
index 0000000..260f08e
--- /dev/null
+++ b/packages/meshbay-hub/tests/test_known_browser.py
@@ -0,0 +1,98 @@
+"""
+A stranger who knows your name locks only the browsers you never used.
+
+The sign-in lockout counts wrong passphrases per username: four an hour from
+anyone kept the owner out of every browser for as long as they cared to keep
+going. A browser that signed in to the account before presents a token and has
+a counter of its own, which nobody else can spend. The token is not a
+credential — the passphrase is still asked — and a passphrase re-checked inside
+an open session is not the sign-in counter's business at all.
+"""
+
+import pytest
+from meshbay_hub.db.models import KnownBrowser, User
+from sqlalchemy import func, select
+from test_bundle_pepper import KEY, _register
+
+WRONG = "w" * 44
+
+
+async def _sign_in(client, username, key=KEY, known=None):
+ body = {"username": username, "auth_key": key}
+ if known:
+ body["known_browser"] = known
+ return await client.post("/v1/users/login", json=body)
+
+
+async def _lock(client, username):
+ for _ in range(4):
+ await _sign_in(client, username, WRONG)
+ assert (await _sign_in(client, username)).status_code == 429
+
+
+@pytest.mark.asyncio
+async def test_a_known_browser_signs_in_through_a_strangers_lockout(client):
+ await _register(client, "known_owner")
+ token = (await _sign_in(client, "known_owner")).json()["known_browser"]
+
+ await _lock(client, "known_owner") # the stranger, without a token
+ r = await _sign_in(client, "known_owner", known=token)
+ assert r.status_code == 200, r.text
+ assert "known_browser" not in r.json(), "a known browser is not given a second token"
+
+
+@pytest.mark.asyncio
+async def test_another_accounts_token_is_no_way_round(client):
+ await _register(client, "known_alice")
+ await _register(client, "known_bobby")
+ alices = (await _sign_in(client, "known_alice")).json()["known_browser"]
+
+ await _lock(client, "known_bobby")
+ assert (await _sign_in(client, "known_bobby", known=alices)).status_code == 429
+
+
+@pytest.mark.asyncio
+async def test_a_known_browser_is_locked_by_its_own_failures(client):
+ """Whoever holds the token still guesses at the same rate."""
+ await _register(client, "known_guess")
+ token = (await _sign_in(client, "known_guess")).json()["known_browser"]
+ for _ in range(4):
+ assert (await _sign_in(client, "known_guess", WRONG, token)).status_code == 401
+ assert (await _sign_in(client, "known_guess", known=token)).status_code == 429
+ # ...and that spent nothing of a browser that has no token.
+ assert (await _sign_in(client, "known_guess")).status_code == 200
+
+
+@pytest.mark.asyncio
+async def test_a_locked_name_does_not_stop_its_owner_inside_a_session(client):
+ await _register(client, "known_inside")
+ session = (await _sign_in(client, "known_inside")).json()["access_token"]
+ await _lock(client, "known_inside")
+
+ r = await client.post("/v1/users/me/bundle-pepper", json={"auth_key": KEY},
+ headers={"Authorization": f"Bearer {session}"})
+ assert r.status_code == 200, r.text
+
+
+@pytest.mark.asyncio
+async def test_only_a_hash_is_kept_and_only_twenty(client, db_session):
+ await _register(client, "known_many")
+ tokens = [(await _sign_in(client, "known_many")).json()["known_browser"]
+ for _ in range(25)]
+ uid = (await db_session.execute(
+ select(User.id).where(User.username == "known_many"))).scalar_one()
+ rows = (await db_session.execute(
+ select(KnownBrowser).where(KnownBrowser.user_id == uid))).scalars().all()
+ assert len(rows) == 20
+ assert not any(t in {r.token_hash for r in rows} for t in tokens)
+
+
+@pytest.mark.asyncio
+async def test_erasing_the_account_forgets_its_browsers(client, db_session):
+ await _register(client, "known_gone")
+ login = (await _sign_in(client, "known_gone")).json()
+ r = await client.request("DELETE", "/v1/users/me", json={"auth_key": KEY},
+ headers={"Authorization": f"Bearer {login['access_token']}"})
+ assert r.status_code == 200, r.text
+ left = await db_session.scalar(select(func.count()).select_from(KnownBrowser))
+ assert left == 0
diff --git a/packages/meshbay-hub/tests/test_login_lockout.py b/packages/meshbay-hub/tests/test_login_lockout.py
index 601edbd..8f06d2d 100644
--- a/packages/meshbay-hub/tests/test_login_lockout.py
+++ b/packages/meshbay-hub/tests/test_login_lockout.py
@@ -131,8 +131,13 @@ async def test_a_burst_of_concurrent_guesses_gets_no_more_than_the_limit(client)
@pytest.mark.asyncio
-async def test_change_password_counts_on_the_same_row(client):
- """It checks the same passphrase, so it is the same oracle."""
+async def test_change_password_counts_on_the_sessions_own_row(client):
+ """
+ It checks the passphrase, so it is counted and locked like a sign-in — on a
+ row of the session's own. A stranger failing at sign-in must not stop the
+ owner changing their passphrase from a session they hold, and failures here
+ must not lock the owner's other browsers out of signing in.
+ """
await _register(client, "grace_test")
token = (await _login(client, "grace_test", RIGHT)).json()["access_token"]
auth = {"Authorization": f"Bearer {token}"}
@@ -145,7 +150,7 @@ async def test_change_password_counts_on_the_same_row(client):
r = await client.post("/v1/users/password", headers=auth, json={
"old_auth_key": RIGHT, "new_auth_key": "n" * 44})
assert r.status_code == 429, r.text
- assert (await _login(client, "grace_test", RIGHT)).status_code == 429
+ assert (await _login(client, "grace_test", RIGHT)).status_code == 200
@pytest.mark.asyncio
@@ -157,10 +162,17 @@ async def test_a_signed_in_session_is_told_its_own_lockout(client):
auth = {"Authorization": f"Bearer {token}"}
assert (await client.get("/v1/users/me", headers=auth)).json()["passphrase_locked_for"] == 0
- await _fail(client, "olivia_test", 4)
+ for _ in range(4):
+ await client.post("/v1/users/password", headers=auth, json={
+ "old_auth_key": WRONG, "new_auth_key": "n" * 44})
left = (await client.get("/v1/users/me", headers=auth)).json()["passphrase_locked_for"]
assert 3500 <= left <= 3600
+ # A stranger failing at sign-in is not this session's lockout.
+ await _fail(client, "olivia_test", 4)
+ other = (await _login(client, "olivia_test", RIGHT))
+ assert other.status_code == 429
+
@pytest.mark.asyncio
async def test_an_attempt_that_checks_no_passphrase_is_not_counted(client, db_session):
diff --git a/packages/meshbay-hub/tests/test_username_case.py b/packages/meshbay-hub/tests/test_username_case.py
new file mode 100644
index 0000000..42f4815
--- /dev/null
+++ b/packages/meshbay-hub/tests/test_username_case.py
@@ -0,0 +1,24 @@
+"""
+A username is unique whatever its case.
+
+Invitations and member management name people by username, so "Alice" beside
+"alice" is one person to whoever reads the list — and a second account under
+the other spelling is the way to be mistaken for them.
+"""
+
+import pytest
+from test_bundle_pepper import KEY
+
+
+async def _register(client, username, email):
+ return await client.post("/v1/users/register", json={
+ "username": username, "auth_key": KEY, "email": email})
+
+
+@pytest.mark.asyncio
+async def test_a_name_differing_only_by_case_is_taken(client):
+ assert (await _register(client, "alice_case", "a1@example.invalid")).status_code == 201
+ for other in ("Alice_case", "ALICE_CASE", "alice_CASE"):
+ r = await _register(client, other, "a2@example.invalid")
+ assert r.status_code == 409, (other, r.text)
+
diff --git a/packages/meshbay-node/src/meshbay_node/cli/groups.py b/packages/meshbay-node/src/meshbay_node/cli/groups.py
index 3fab4b9..fca4800 100644
--- a/packages/meshbay-node/src/meshbay_node/cli/groups.py
+++ b/packages/meshbay-node/src/meshbay_node/cli/groups.py
@@ -64,7 +64,7 @@ def group(args) -> None:
sys.exit(1)
if not args.target or not args.dir:
print("usage: meshbay-node group add <name> --dir <path> "
- "[--no-writable]")
+ "[--no-writable] [--open]")
print()
print("The group must already exist on the hub and be yours. This")
print("only tells the node to host it, and picks its first")
@@ -78,12 +78,19 @@ def group(args) -> None:
# is not a working group. Every root added *later* is read-only by
# default, which is the opposite rule and the right one there.
writable = args.writable is not False
+ # How people join is the operator's to say here, never read from the hub.
+ join_policy = "open" if getattr(args, "open", False) else "invite"
body = {"name": args.target, "shared_dir": args.dir,
- "writable": writable}
+ "writable": writable, "join_policy": join_policy}
out = _daemon_api(cfg, "/api/groups/attach", method="POST", body=body)
print(f"{out['name']} ({out['group_id'][:8]}) added to {out['config']}")
print(f" shared_dir {out['shared_dir']}"
f" ({'read-write' if writable else 'read-only'})")
+ print(f" join_policy {join_policy}")
+ if out.get("hub_join_policy") == "open" and join_policy != "open":
+ print()
+ print("The hub lists this group as open; this node admits by invitation")
+ print("only. To host it open, remove it and add it again with --open.")
print()
print("Tell the daemon to re-read its config, then give the group a key:")
print(" meshbay-node reload")
diff --git a/packages/meshbay-node/src/meshbay_node/cli/parser.py b/packages/meshbay-node/src/meshbay_node/cli/parser.py
index 29a9f63..cfc4704 100644
--- a/packages/meshbay-node/src/meshbay_node/cli/parser.py
+++ b/packages/meshbay-node/src/meshbay_node/cli/parser.py
@@ -87,6 +87,9 @@ def build_parser() -> argparse.ArgumentParser:
parser.add_argument("--no-removable", action="store_false",
dest="removable",
help="mark root as not removable (root set)")
+ parser.add_argument("--open", action="store_true",
+ help="group add: anyone the hub lists the group to may join "
+ "(default: by invitation only)")
parser.add_argument("--name", default=None,
help="root name (root add; defaults to directory basename)")
parser.add_argument("--log-level", default="INFO",
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py
index 45a6b9a..1777bc3 100644
--- a/packages/meshbay-node/src/meshbay_node/daemon.py
+++ b/packages/meshbay-node/src/meshbay_node/daemon.py
@@ -577,12 +577,12 @@ class NodeDaemon(EnrichmentMixin):
log.info("QUIC server disabled ([node] quic_enabled = false)")
# 8. Hub WebSocket (signaling + revocations + WebRTC offers)
- async def on_webrtc_offer(sdp, peer_id, ice_candidates):
+ async def on_webrtc_offer(sdp, peer_id, ice_candidates, user_id=""):
if not self._webrtc:
return None
try:
answer_sdp, answer_ice = await self._webrtc.handle_offer(
- sdp, peer_id)
+ sdp, peer_id, user_id)
log.info("WebRTC answer for peer=%s (%d peers)",
peer_id, self._webrtc.active_peers)
return (answer_sdp, answer_ice)
@@ -601,19 +601,8 @@ class NodeDaemon(EnrichmentMixin):
payload = _jwt.decode(
token, session.hub_pk_pem, algorithms=["EdDSA"],
options={"verify_exp": False})
- target = payload.get("target")
- tid = payload.get("target_id", "")
- if target == "user":
- denylist.deny_user(tid)
- elif target == "group":
- # H4: previously dropped on the floor, so "suspend a
- # group" was a hub-only gesture that no node enforced.
- denylist.deny_group(tid)
- self._drop_group_sessions(tid)
- elif target == "jti":
- denylist.deny_jti(tid)
- else:
- log.warning("Unknown revocation target: %r", target)
+ self._apply_revocation(denylist, payload.get("target"),
+ payload.get("target_id", ""))
except Exception as e:
log.warning("Invalid revocation token: %s", e)
@@ -1522,6 +1511,35 @@ class NodeDaemon(EnrichmentMixin):
except Exception:
pass
+ def _apply_revocation(self, denylist, target, target_id: str) -> None:
+ """
+ What a revocation the hub signed does on this node.
+
+ A revoked account or group is refused from now on, and its live sessions
+ are closed: a denylist entry alone stops the next connection and leaves
+ the current one streaming, downloading and chatting until it happens to
+ disconnect.
+ """
+ if target == "user":
+ denylist.deny_user(target_id)
+ self._drop_user_sessions(target_id)
+ elif target == "group":
+ denylist.deny_group(target_id)
+ self._drop_group_sessions(target_id)
+ elif target == "jti":
+ denylist.deny_jti(target_id)
+ else:
+ log.warning("Unknown revocation target: %r", target)
+
+ def _drop_user_sessions(self, user_id: str) -> None:
+ """Close every live session of a revoked account."""
+ if not self._webrtc or not user_id:
+ return
+ for session in list(self._webrtc._sessions.values()):
+ if getattr(session, "_user_id", None) == user_id:
+ spawn(session.close())
+ log.info("Dropped session for revoked account %s", user_id[:8])
+
def _drop_group_sessions(self, group_id: str) -> None:
"""Close live sessions for a revoked group (H4)."""
if not self._webrtc or not group_id:
diff --git a/packages/meshbay-node/src/meshbay_node/hub_client.py b/packages/meshbay-node/src/meshbay_node/hub_client.py
index 2ce8993..54f87f0 100644
--- a/packages/meshbay-node/src/meshbay_node/hub_client.py
+++ b/packages/meshbay-node/src/meshbay_node/hub_client.py
@@ -451,7 +451,8 @@ class HubClient:
"""Negotiate one WebRTC offer and return the answer, off the read loop."""
try:
answer = await on_webrtc_offer(
- msg["sdp"], msg["peer_id"], msg.get("ice_candidates", []))
+ msg["sdp"], msg["peer_id"], msg.get("ice_candidates", []),
+ str(msg.get("user_id") or ""))
except Exception as e:
log.warning("WebRTC offer from %s failed: %s",
str(msg.get("peer_id"))[:8], e)
diff --git a/packages/meshbay-node/src/meshbay_node/linkpreview.py b/packages/meshbay-node/src/meshbay_node/linkpreview.py
index 6e1e618..493ac65 100644
--- a/packages/meshbay-node/src/meshbay_node/linkpreview.py
+++ b/packages/meshbay-node/src/meshbay_node/linkpreview.py
@@ -13,16 +13,19 @@ hub:
apps (`_fetch_and_cache_poster`), over the same authorised path.
Because the node makes an outbound request to an address a *member* chose,
-this is an SSRF surface. `safe_url()` is the gate: http(s) only, no
+this is an SSRF surface. `check_url()` is the gate: http(s) only, no
credentials, the port restricted to the web set, and every resolved address
must be globally routable — no loopback, private, link-local, multicast or
-reserved range. Redirects are followed by hand so every hop is re-checked,
-and the address the connection actually landed on is re-checked against the
-same rule (`_reject_if_rebound`), so a name that resolves clean and then to
-something internal (rebinding) does not get its body read. A full pin —
-connect to the validated literal, verify the certificate for the name — is
-the remaining hardening. How many previews a member can trigger is
-rate-limited by the caller (`_do_link_preview_request`).
+reserved range. Redirects are followed by hand so every hop is re-checked.
+
+**The connection goes to the address that was checked** (`_PinnedBackend`):
+the name is resolved once, off the event loop, every answer is checked, and
+the socket is opened to that IP literal — TLS still verifies the certificate
+for the name. Resolving to check and letting the HTTP client resolve again to
+connect would let a name answer clean the first time and with a LAN address
+the second (DNS rebinding), and the request would be sent before anything
+could look. How many previews a member can trigger is rate-limited by the
+caller (`_do_link_preview_request`).
Nothing is stored durably: the caller keeps an in-memory TTL cache and the
OG image rides the existing `media_cache` thumb store (same as a poster).
@@ -38,6 +41,7 @@ from html.parser import HTMLParser
from io import BytesIO
from urllib.parse import urljoin, urlsplit
+import httpcore
import httpx
log = logging.getLogger(__name__)
@@ -77,7 +81,13 @@ def _addr_is_public(ip: str) -> bool:
def safe_url(url: str) -> str:
- """Return the URL unchanged if it is safe to fetch, else raise UnsafeURL."""
+ """
+ The URL unchanged if its shape is safe to fetch, else raise UnsafeURL.
+
+ Shape only — scheme, credentials, port, and the address when it is a
+ literal. A name is checked by `check_url` and, again, at connect time;
+ resolving here would block the event loop on a member's choice of name.
+ """
if not isinstance(url, str) or len(url) > 2048:
raise UnsafeURL("missing or oversized")
parts = urlsplit(url)
@@ -94,23 +104,102 @@ def safe_url(url: str) -> str:
raise UnsafeURL("bad port")
if port is not None and port not in _ALLOWED_PORTS:
raise UnsafeURL(f"port {port}")
- # An IP literal is checked directly; a name is resolved and every answer
- # must be public — a hostname with one public and one 127.0.0.1 record
- # would otherwise be a way in.
+ if _is_literal(host) and not _addr_is_public(host):
+ raise UnsafeURL(f"non-public address {host}")
+ return url
+
+
+def _is_literal(host: str) -> bool:
+ try:
+ ipaddress.ip_address(host)
+ return True
+ except ValueError:
+ return False
+
+
+_RESOLVE_TIMEOUT = 5.0
+
+
+async def resolve_public(host: str, port: int) -> str:
+ """
+ One public address for `host`, resolved off the event loop, or UnsafeURL.
+
+ Every answer must be public: a name with one public and one 127.0.0.1
+ record would otherwise be a way in.
+ """
+ if _is_literal(host):
+ if not _addr_is_public(host):
+ raise UnsafeURL(f"non-public address {host}")
+ return host
+ loop = asyncio.get_running_loop()
try:
- infos = socket.getaddrinfo(host, parts.port or (443 if parts.scheme == "https" else 80),
- proto=socket.IPPROTO_TCP)
- except socket.gaierror as e:
+ infos = await asyncio.wait_for(
+ loop.getaddrinfo(host, port, proto=socket.IPPROTO_TCP), _RESOLVE_TIMEOUT)
+ except (socket.gaierror, TimeoutError) as e:
raise UnsafeURL(f"cannot resolve: {e}")
- resolved = {info[4][0] for info in infos}
+ resolved = list(dict.fromkeys(info[4][0] for info in infos))
if not resolved:
raise UnsafeURL("resolves to nothing")
bad = [ip for ip in resolved if not _addr_is_public(ip)]
if bad:
raise UnsafeURL(f"non-public address {bad[0]}")
+ return resolved[0]
+
+
+async def check_url(url: str) -> str:
+ """`safe_url`, and the name's addresses checked too. The URL unchanged."""
+ safe_url(url)
+ parts = urlsplit(url)
+ await resolve_public(parts.hostname or "",
+ parts.port or (443 if parts.scheme == "https" else 80))
return url
+class _PinnedBackend(httpcore.AsyncNetworkBackend):
+ """
+ Opens every connection to an address `resolve_public` checked.
+
+ httpcore hands the backend the request's host; the TLS layer above still
+ uses that name for SNI and certificate verification, so pinning the socket
+ changes where it connects and nothing about whom it trusts.
+ """
+
+ def __init__(self) -> None:
+ self._inner = httpcore.AnyIOBackend()
+
+ async def connect_tcp(self, host, port, timeout=None, local_address=None,
+ socket_options=None):
+ ip = await resolve_public(host, port)
+ return await self._inner.connect_tcp(ip, port, timeout=timeout,
+ local_address=local_address,
+ socket_options=socket_options)
+
+ async def connect_unix_socket(self, path, timeout=None, socket_options=None):
+ raise UnsafeURL("no unix sockets")
+
+ async def sleep(self, seconds: float) -> None:
+ await self._inner.sleep(seconds)
+
+
+class _PinnedTransport(httpx.AsyncHTTPTransport):
+ """httpx's transport over a pool that connects through `_PinnedBackend`.
+
+ No proxy from the environment (`trust_env=False`): a proxy would resolve
+ the name itself, and the pin would bind nothing.
+ """
+
+ def __init__(self) -> None:
+ super().__init__(trust_env=False, retries=0)
+ self._pool = httpcore.AsyncConnectionPool(
+ ssl_context=httpx.create_ssl_context(trust_env=False),
+ network_backend=_PinnedBackend(), max_connections=10)
+
+
+def _new_client() -> httpx.AsyncClient:
+ return httpx.AsyncClient(transport=_PinnedTransport(), timeout=_TIMEOUT,
+ max_redirects=0, trust_env=False)
+
+
class _HeadParser(HTMLParser):
"""Collects <title> text and name/property→content from <meta> in <head>.
@@ -155,25 +244,6 @@ def _first(metas: dict[str, str], *keys: str) -> str | None:
return None
-def _reject_if_rebound(resp: httpx.Response) -> None:
- """
- `safe_url` validated the name's addresses; this checks the one the
- connection actually landed on, so a name that resolves clean and then to
- something internal (DNS rebinding) does not get its body read.
-
- Best-effort: the `network_stream` extension is not present on every
- transport (a MockTransport in tests has none), and its absence is not a
- failure — the pre-check and the per-hop redirect re-check still stand.
- """
- try:
- stream = resp.extensions.get("network_stream")
- addr = stream.get_extra_info("server_addr") if stream else None
- except Exception:
- return
- if addr and not _addr_is_public(str(addr[0])):
- raise UnsafeURL(f"connected to non-public address {addr[0]}")
-
-
async def _get(client: httpx.AsyncClient, url: str) -> httpx.Response:
"""
One GET with manual, re-validated redirects, **body not read**.
@@ -184,19 +254,14 @@ async def _get(client: httpx.AsyncClient, url: str) -> httpx.Response:
compressed stream of any size was held in memory first. The caps are the
only thing between a URL a member pasted and the node's memory.
"""
- current = safe_url(url)
+ current = await check_url(url)
for _ in range(_MAX_REDIRECTS + 1):
request = client.build_request("GET", current, headers={"User-Agent": _UA})
resp = await client.send(request, stream=True, follow_redirects=False)
- try:
- _reject_if_rebound(resp)
- except Exception:
- await resp.aclose()
- raise
if resp.is_redirect and "location" in resp.headers:
location = resp.headers["location"]
await resp.aclose()
- current = safe_url(urljoin(current, location))
+ current = await check_url(urljoin(current, location))
continue
return resp
raise UnsafeURL("too many redirects")
@@ -234,7 +299,7 @@ async def fetch_preview(url: str, *, client: httpx.AsyncClient | None = None) ->
"""
own = client is None
if own:
- client = httpx.AsyncClient(timeout=_TIMEOUT, max_redirects=0)
+ client = _new_client()
try:
return await asyncio.wait_for(_preview(client, url), _TOTAL_DEADLINE)
except (httpx.HTTPError, UnsafeURL, TimeoutError) as e:
@@ -246,7 +311,6 @@ async def fetch_preview(url: str, *, client: httpx.AsyncClient | None = None) ->
async def _preview(client: httpx.AsyncClient, url: str) -> dict | None:
- safe_url(url)
resp = await _get(client, url)
try:
ctype = resp.headers.get("content-type", "").split(";")[0].strip().lower()
@@ -272,7 +336,7 @@ async def _preview(client: httpx.AsyncClient, url: str) -> dict | None:
if image:
image = urljoin(final_url, image)
try:
- safe_url(image)
+ await check_url(image)
except UnsafeURL:
image = None
@@ -294,7 +358,7 @@ async def fetch_image(url: str, *, client: httpx.AsyncClient | None = None) -> b
"""Fetch and re-encode an OG image to a small JPEG. None on any failure."""
own = client is None
if own:
- client = httpx.AsyncClient(timeout=_TIMEOUT, max_redirects=0)
+ client = _new_client()
try:
raw = await asyncio.wait_for(_image_bytes(client, url), _TOTAL_DEADLINE)
except (httpx.HTTPError, UnsafeURL, TimeoutError) as e:
@@ -310,7 +374,6 @@ async def fetch_image(url: str, *, client: httpx.AsyncClient | None = None) -> b
async def _image_bytes(client: httpx.AsyncClient, url: str) -> bytes | None:
- safe_url(url)
resp = await _get(client, url)
try:
ctype = resp.headers.get("content-type", "").split(";")[0].strip().lower()
diff --git a/packages/meshbay-node/src/meshbay_node/media_probe.py b/packages/meshbay-node/src/meshbay_node/media_probe.py
index a6267d9..e9090ea 100644
--- a/packages/meshbay-node/src/meshbay_node/media_probe.py
+++ b/packages/meshbay-node/src/meshbay_node/media_probe.py
@@ -10,6 +10,12 @@ import asyncio
import json
from dataclasses import dataclass, field
+# How long ffprobe may take over one file's headers. The file is a member's
+# upload as often as the operator's own: one that keeps ffprobe busy must not
+# keep the stream request, the subtitle request or the enrichment slot that
+# asked for it — the same bound the index-time enrichment already put around it.
+FFPROBE_TIMEOUT_SECS = 30
+
_H264_PROFILES = {"Baseline": "42", "Main": "4d", "High": "64", "High 10": "6e"}
# Source video codecs whose MSE codec string is real but which no mainstream
@@ -158,7 +164,16 @@ async def probe_video(path: str) -> VideoProbe:
"-of", "json", path,
stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE,
)
- stdout, _ = await proc.communicate()
+ try:
+ stdout, _ = await asyncio.wait_for(proc.communicate(), FFPROBE_TIMEOUT_SECS)
+ except (TimeoutError, asyncio.CancelledError) as e:
+ # Killed, not abandoned: a cancelled wait leaves the process running,
+ # and a caller's own timeout (enrich.py) cancels exactly this wait.
+ proc.kill()
+ await proc.wait()
+ if isinstance(e, asyncio.CancelledError):
+ raise
+ raise RuntimeError(f"ffprobe timed out after {FFPROBE_TIMEOUT_SECS}s") from None
info = json.loads(stdout)
duration = float(info.get("format", {}).get("duration", 0))
diff --git a/packages/meshbay-node/src/meshbay_node/ops/groups.py b/packages/meshbay-node/src/meshbay_node/ops/groups.py
index 1d8003c..ac14c3a 100644
--- a/packages/meshbay-node/src/meshbay_node/ops/groups.py
+++ b/packages/meshbay-node/src/meshbay_node/ops/groups.py
@@ -9,7 +9,7 @@ from meshbay_common.crypto import generate_gek, wrap_gek_aes
from meshbay_node.config import DEFAULT_CONFIG_PATH
from meshbay_node.ops.core import OpError, _config, _group_ctx, _hub
-from meshbay_node.ops.node_toml import _find_group_range
+from meshbay_node.ops.node_toml import _find_group_range, toml_string
log = logging.getLogger("meshbay_node.ops")
@@ -142,17 +142,28 @@ async def list_groups(state: dict) -> dict:
return {"groups": out, "operator_paired": has_operator, "settings": settings}
+JOIN_POLICIES = ("invite", "open")
+
+
async def attach_group(state: dict, name: str, shared_dir: str,
- writable: bool = True) -> dict:
+ writable: bool = True, join_policy: str = "invite") -> dict:
"""
Write a new [[groups]] block into node.toml.
The name-to-id lookup happens here because this process is the one logged
into the hub. Nothing is created on the hub: the group already exists, this
only tells the node to host it.
+
+ `join_policy` is the operator's, given with this request, and `invite`
+ unless they say otherwise. The hub's own record of the group is not read
+ for it: a hub that could declare a group open would be handed its key by
+ anyone it sent. The hub's value is returned beside it, so a caller can say
+ when the two differ.
"""
if not name or not shared_dir:
raise OpError("name and shared_dir are required")
+ if join_policy not in JOIN_POLICIES:
+ raise OpError(f"join_policy must be one of {', '.join(JOIN_POLICIES)}")
config = _config(state)
hub = _hub(state)
try:
@@ -182,12 +193,12 @@ async def attach_group(state: dict, name: str, shared_dir: str,
raise OpError(f"Cannot create {path}: {e}") from e
conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH)
- join_policy = group.get("join_policy", "invite")
+ visibility = "public" if join_policy == "open" else "private"
block = (f'\n[[groups]]\n'
- f'id = "{group["id"]}"\n'
- f'name = "{group["name"]}"\n'
- f'visibility = "{group.get("visibility", "private")}"\n'
- f'join_policy = "{join_policy}"\n')
+ f'id = {toml_string(group["id"])}\n'
+ f'name = {toml_string(group["name"])}\n'
+ f'visibility = {toml_string(visibility)}\n'
+ f'join_policy = {toml_string(join_policy)}\n')
# No `upload_dir` here. `GroupConfig.__post_init__` still *reads* it, so an
# existing node.toml keeps working — but what it does on read is force every
# other root read-only and append that path as the one writable one, which
@@ -198,7 +209,7 @@ async def attach_group(state: dict, name: str, shared_dir: str,
block += (f'\n [[groups.roots]]\n'
# Forward slashes: a Windows path in a TOML basic string is a
# parse error (`\U`, `\a`, ... are escapes). pathlib reads `/`.
- f' path = "{path.as_posix()}"\n'
+ f' path = {toml_string(path.as_posix())}\n'
f' writable = {"true" if writable else "false"}\n')
try:
with conf_path.open("a", encoding="utf-8", newline="\n") as f:
@@ -208,7 +219,8 @@ async def attach_group(state: dict, name: str, shared_dir: str,
result = {"group_id": group["id"], "name": group["name"],
"shared_dir": str(path), "config": str(conf_path),
- "writable": writable,
+ "writable": writable, "join_policy": join_policy,
+ "hub_join_policy": group.get("join_policy", "invite"),
"note": "restart the node to pick it up"}
return result
diff --git a/packages/meshbay-node/src/meshbay_node/ops/node_toml.py b/packages/meshbay-node/src/meshbay_node/ops/node_toml.py
index f711f26..2407722 100644
--- a/packages/meshbay-node/src/meshbay_node/ops/node_toml.py
+++ b/packages/meshbay-node/src/meshbay_node/ops/node_toml.py
@@ -3,14 +3,50 @@
from __future__ import annotations
import re
+import tomllib
from pathlib import Path
from meshbay_node.ops.core import OpError
+def toml_string(value: str) -> str:
+ """A TOML basic string holding `value` exactly, quotes included.
+
+ Every string written into node.toml goes through here. A value with a quote
+ or a newline in it — a group name, a folder name, any of them chosen by
+ someone else — would otherwise end the string and write lines of its own.
+ """
+ out = ['"']
+ for ch in str(value):
+ if ch == '"':
+ out.append('\\"')
+ elif ch == "\\":
+ out.append("\\\\")
+ elif ord(ch) < 0x20 or ord(ch) == 0x7F:
+ out.append(f"\\u{ord(ch):04x}")
+ else:
+ out.append(ch)
+ out.append('"')
+ return "".join(out)
+
+
+def _string_value(line: str, key: str) -> str | None:
+ """The string `key` holds on this line, unescaped — or None.
+
+ Read as TOML, not by pattern: a value written by `toml_string` may carry an
+ escaped quote or backslash, which a `"([^"]*)"` pattern would cut short.
+ """
+ if not re.match(r"^\s*" + re.escape(key) + r"\s*=", line):
+ return None
+ try:
+ value = tomllib.loads(line.strip()).get(key)
+ except tomllib.TOMLDecodeError:
+ return None
+ return value if isinstance(value, str) else None
+
+
def _find_group_range(lines: list[str], group_id: str) -> tuple[int, int] | None:
"""Line range of a [[groups]] block by id: (start, end_exclusive)."""
- id_re = re.compile(r'^\s*id\s*=\s*"([^"]*)"')
block_starts: list[int] = []
for i, line in enumerate(lines):
if line.strip() == "[[groups]]":
@@ -24,8 +60,7 @@ def _find_group_range(lines: list[str], group_id: str) -> tuple[int, int] | None
boundary = k
break
for k in range(start + 1, boundary):
- m = id_re.match(lines[k])
- if m and m.group(1) == group_id:
+ if _string_value(lines[k], "id") == group_id:
return (start, boundary)
return None
@@ -61,8 +96,10 @@ def _update_node_toml(conf_path: Path, updates: dict) -> None:
if isinstance(value, bool):
return f"{key} = {'true' if value else 'false'}"
if isinstance(value, list):
- items = ", ".join(f'"{v}"' for v in value)
+ items = ", ".join(toml_string(v) for v in value)
return f"{key} = [{items}]"
+ if isinstance(value, str):
+ return f"{key} = {toml_string(value)}"
return f"{key} = {value}"
remaining = dict(updates)
@@ -115,7 +152,6 @@ def _remove_roots_block(conf_path: Path, group_id: str,
raise OpError(f"Group {group_id[:8]} not found in {conf_path}")
start, end = rng
- path_re = re.compile(r'^\s*path\s*=\s*"([^"]*)"')
roots_starts: list[int] = []
for i in range(start + 1, end):
if lines[i].strip() == "[[groups.roots]]":
@@ -124,10 +160,10 @@ def _remove_roots_block(conf_path: Path, group_id: str,
for j, rs in enumerate(roots_starts):
rs_end = roots_starts[j + 1] if j + 1 < len(roots_starts) else end
for k in range(rs, rs_end):
- m = path_re.match(lines[k])
- if m:
+ raw = _string_value(lines[k], "path")
+ if raw is not None:
try:
- p = str(Path(m.group(1)).expanduser().resolve())
+ p = str(Path(raw).expanduser().resolve())
except OSError:
continue
if p == resolved_path:
@@ -153,7 +189,6 @@ def _update_root_field(conf_path: Path, group_id: str,
raise OpError(f"Group {group_id[:8]} not found in {conf_path}")
start, end = rng
- path_re = re.compile(r'^\s*path\s*=\s*"([^"]*)"')
writable_re = re.compile(r'^\s*(writable|upload)\s*=')
removable_re = re.compile(r'^\s*removable\s*=')
roots_starts: list[int] = []
@@ -165,10 +200,10 @@ def _update_root_field(conf_path: Path, group_id: str,
rs_end = roots_starts[j + 1] if j + 1 < len(roots_starts) else end
found_path = False
for k in range(rs, rs_end):
- m = path_re.match(lines[k])
- if m:
+ raw = _string_value(lines[k], "path")
+ if raw is not None:
try:
- p = str(Path(m.group(1)).expanduser().resolve())
+ p = str(Path(raw).expanduser().resolve())
except OSError:
continue
if p == resolved_path:
diff --git a/packages/meshbay-node/src/meshbay_node/ops/roots.py b/packages/meshbay-node/src/meshbay_node/ops/roots.py
index e3e2781..480fe63 100644
--- a/packages/meshbay-node/src/meshbay_node/ops/roots.py
+++ b/packages/meshbay-node/src/meshbay_node/ops/roots.py
@@ -8,7 +8,12 @@ from pathlib import Path
from meshbay_node.config import DEFAULT_CONFIG_PATH
from meshbay_node.ops.core import OpError, _config, _group_ctx, _roster
-from meshbay_node.ops.node_toml import _insert_roots_block, _remove_roots_block, _update_root_field
+from meshbay_node.ops.node_toml import (
+ _insert_roots_block,
+ _remove_roots_block,
+ _update_root_field,
+ toml_string,
+)
from meshbay_node.roots import RootError, RootSet, off_disk
log = logging.getLogger("meshbay_node.ops")
@@ -46,11 +51,11 @@ async def add_root(state: dict, group_id: str, path: str, *,
raise OpError(f"Cannot create {added.path}: {e}") from e
conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH)
- root_block = f' [[groups.roots]]\n path = "{added.path.as_posix()}"'
+ root_block = f' [[groups.roots]]\n path = {toml_string(added.path.as_posix())}'
if name:
- root_block += f'\n name = "{added.name}"'
+ root_block += f'\n name = {toml_string(added.name)}'
if kind != "generic":
- root_block += f'\n kind = "{added.kind}"'
+ root_block += f'\n kind = {toml_string(added.kind)}'
if writable:
root_block += '\n writable = true'
if removable:
diff --git a/packages/meshbay-node/src/meshbay_node/roots.py b/packages/meshbay-node/src/meshbay_node/roots.py
index 89b0441..2de0708 100644
--- a/packages/meshbay-node/src/meshbay_node/roots.py
+++ b/packages/meshbay-node/src/meshbay_node/roots.py
@@ -30,6 +30,7 @@ from __future__ import annotations
import asyncio
import logging
+import os
import re
from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass, field
@@ -52,25 +53,77 @@ SAFE_UPLOAD_NAME = re.compile(
re.UNICODE)
-def _free_name(directory: Path, filename: str) -> str:
+# Files Windows Explorer acts on by itself when it shows a folder: a link's
+# icon, a folder's settings, a search connector. Placed by a member in a folder
+# the operator browses, any of them can make Explorer contact a server of the
+# member's choosing with the operator's Windows credentials — a known attack,
+# and why mail providers refuse the same types. Refused for every node: a
+# Linux node's folder may be shared to Windows machines.
+SHELL_ACTIVE_NAMES = frozenset({"desktop.ini"})
+SHELL_ACTIVE_SUFFIXES = (".lnk", ".url", ".scf", ".library-ms", ".searchconnector-ms")
+
+
+def shell_active(filename: str) -> bool:
+ name = filename.lower()
+ return name in SHELL_ACTIVE_NAMES or name.endswith(SHELL_ACTIVE_SUFFIXES)
+
+
+def _free_name(directory: Path, filename: str,
+ taken: frozenset[str] | set[str] = frozenset()) -> str:
"""
`filename`, or the first "name (n).ext" that is not taken.
- Never returns the name of a file that exists, so an upload cannot replace
- one — the property the per-user quarantine used to provide (C5a).
+ Never returns the name of a file that exists, nor one in `taken` — names
+ uploads in flight will publish under — so an upload cannot replace a file
+ or another upload (C5a).
"""
- if not (directory / filename).exists():
+ def free(name: str) -> bool:
+ return name not in taken and not (directory / name).exists()
+
+ if free(filename):
return filename
stem, dot, ext = filename.rpartition(".")
if not dot:
stem, ext = filename, ""
for n in range(2, 1000):
candidate = f"{stem} ({n}){dot}{ext}"
- if not (directory / candidate).exists():
+ if free(candidate):
return candidate
raise FileExistsError(filename)
+def publish_upload(part: Path, directory: Path, stored_name: str, filename: str,
+ taken: frozenset[str] | set[str] = frozenset()) -> str:
+ """
+ Move a finished `.part` to its name without ever replacing a file. The name
+ it was published under, which may not be `stored_name`.
+
+ A rename replaces whatever is at the target, and the target can appear
+ while the upload runs — the operator copying a file in, another group's
+ upload into a shared folder. A hard link refuses an existing target, so it
+ is the publication; where the filesystem has none (FAT, exFAT, some network
+ shares), the existence check and the rename are as close as it gets. A
+ taken name moves on to the next free one rather than failing the upload.
+ """
+ name = stored_name
+ for _ in range(8):
+ target = directory / name
+ try:
+ os.link(part, target)
+ except FileExistsError:
+ name = _free_name(directory, filename, taken)
+ continue
+ except OSError:
+ if target.exists():
+ name = _free_name(directory, filename, taken)
+ continue
+ part.rename(target)
+ return name
+ part.unlink()
+ return name
+ raise FileExistsError(stored_name)
+
+
def safe_subdir(roots: RootSet, rel: str) -> Path | None:
"""
Resolve a client-supplied directory inside one of the group's roots, or refuse.
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py
index 857db33..a2d9f47 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
@@ -120,7 +120,7 @@ class MusicMixin:
blob = await _transcode_audio_to_aac(file_path)
except Exception as e:
log.warning("Audio transcode failed for %s: %s", entry.id[:12], e)
- self._send({"type": "error", "detail": f"Transcode failed: {e}"})
+ self._send({"type": "error", "detail": "This track could not be converted"})
return
transcode_hash = blake3.blake3(blob).hexdigest()
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 4337e24..d6a5248 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
@@ -2,6 +2,7 @@
stream a session holds, and the ffmpeg pipeline behind it."""
import asyncio
+import contextlib
import logging
import time
@@ -34,6 +35,10 @@ log = logging.getLogger("meshbay_node.transport.webrtc_server")
# so the operator sets `max_concurrent_streams` under [node] in node.toml. This
# value applies when they have said nothing.
MAX_CONCURRENT_TRANSCODES = 8
+# Subtitle extractions one account may run at once. The player asks for one
+# track at a time; two covers a quick change of track. Each holds a transcode
+# slot for up to fifteen minutes on a long film.
+MAX_SUBTITLE_JOBS_PER_ACCOUNT = 2
STREAM_SEGMENT_SIZE = 256 * 1024
@@ -179,6 +184,36 @@ class StreamingMixin:
self._stream_task = asyncio.current_task()
await self._stream_video(msg)
+ @contextlib.contextmanager
+ def _account_share(self, kind: str, limit: int):
+ """
+ Hold one of this account's `limit` places for `kind`, or yield False.
+
+ Counted on the node, across every session of the account: a member's
+ devices and tabs share one allowance. The node's own account is not
+ counted — it is the operator's machine.
+ """
+ user = getattr(self, "_user_id", "") or ""
+ if not user or user == self._ctx.get("node_user_id"):
+ yield True
+ return
+ held = self._ctx.setdefault(f"_{kind}_by_account", {})
+ if held.get(user, 0) >= limit:
+ yield False
+ return
+ held[user] = held.get(user, 0) + 1
+ try:
+ yield True
+ finally:
+ held[user] -= 1
+ if held[user] <= 0:
+ held.pop(user, None)
+
+ def _streams_per_account(self) -> int:
+ """Half the node's viewers, rounded up: three screens in one home fit,
+ and no member alone takes every slot the operator set."""
+ return max(1, -(-self._stream_capacity() // 2))
+
def _transcode_semaphore(self) -> asyncio.Semaphore:
"""The node's stream budget, shared across every peer.
@@ -205,22 +240,28 @@ class StreamingMixin:
ctx = self._ctx
log.info("stream: waiting for a slot (%d of %d in use)",
ctx.get("_streams_in_flight", 0), self._stream_capacity())
- async with sem:
- # Counted here rather than read back out of the semaphore's private
- # `_value`: `set_capacity` needs to know how many slots are held in
- # order to resize without letting the pool overshoot, and a number
- # this code maintains itself is one that survives the semaphore
- # object being replaced underneath it.
- ctx["_streams_in_flight"] = ctx.get("_streams_in_flight", 0) + 1
- log.info("stream: slot acquired (%d of %d in use)",
- ctx["_streams_in_flight"], self._stream_capacity())
- try:
- await self._stream_video_inner(msg)
- finally:
- ctx["_streams_in_flight"] = max(
- 0, ctx.get("_streams_in_flight", 1) - 1)
- log.info("stream: slot released (%d of %d in use)",
+ with self._account_share("streams", self._streams_per_account()) as ok:
+ if not ok:
+ self._send({"type": "error",
+ "detail": "Too many videos playing from this account, "
+ "stop one and retry"})
+ return
+ async with sem:
+ # Counted here rather than read back out of the semaphore's private
+ # `_value`: `set_capacity` needs to know how many slots are held in
+ # order to resize without letting the pool overshoot, and a number
+ # this code maintains itself is one that survives the semaphore
+ # object being replaced underneath it.
+ ctx["_streams_in_flight"] = ctx.get("_streams_in_flight", 0) + 1
+ log.info("stream: slot acquired (%d of %d in use)",
ctx["_streams_in_flight"], self._stream_capacity())
+ try:
+ await self._stream_video_inner(msg)
+ finally:
+ ctx["_streams_in_flight"] = max(
+ 0, ctx.get("_streams_in_flight", 1) - 1)
+ log.info("stream: slot released (%d of %d in use)",
+ ctx["_streams_in_flight"], self._stream_capacity())
def _stream_capacity(self) -> int:
return self._ctx.get("max_concurrent_streams") or MAX_CONCURRENT_TRANSCODES
@@ -246,7 +287,10 @@ class StreamingMixin:
try:
probe = await _probe_video(str(file_path))
except Exception as e:
- self._send({"type": "error", "detail": f"Probe failed: {e}"})
+ # The cause to the operator's log; to the member, that it failed.
+ # ffmpeg's own words carry the operator's paths and versions.
+ log.warning("stream: probe failed for %s: %s", entry.id[:12], e)
+ self._send({"type": "error", "detail": "This video could not be read"})
return
codec_str = probe.codec
duration = probe.duration
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 70781eb..525f2a9 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
@@ -10,6 +10,7 @@ from meshbay_common.protocol import MNP
from meshbay_node.media_probe import probe_video as _probe_video
from meshbay_node.roots import off_disk
+from meshbay_node.transport.webrtc.apps.streaming import MAX_SUBTITLE_JOBS_PER_ACCOUNT
from meshbay_node.transport.webrtc.disk import _locate
from meshbay_node.transport.webrtc.media_tools import (
_extract_subtitle_to_webvtt,
@@ -135,10 +136,16 @@ class SubtitlesMixin:
return
budget = _subtitle_timeout_for(entry.size)
- async with sem:
- log.info("subtitle: extracting file=%s track=%d (slot taken, up to %.0fs)",
- file_id[:12], ordinal, budget)
- blob = await _extract_subtitle_to_webvtt(file_path, ordinal, budget)
+ with self._account_share("subtitles", MAX_SUBTITLE_JOBS_PER_ACCOUNT) as ok:
+ if not ok:
+ log.info("subtitle: refused, account at its extraction share")
+ self._send({"type": "error",
+ "detail": "Subtitles are already being prepared, retry shortly"})
+ return
+ async with sem:
+ log.info("subtitle: extracting file=%s track=%d (slot taken, up to %.0fs)",
+ file_id[:12], ordinal, budget)
+ blob = await _extract_subtitle_to_webvtt(file_path, ordinal, budget)
subtitle_hash = blake3.blake3(blob).hexdigest()
await media_cache.put_thumb(subtitle_hash, synthetic_id, blob)
@@ -158,8 +165,10 @@ class SubtitlesMixin:
except BaseException as e:
log.warning("subtitle: extract failed file=%s track=%d after %.1fs: %r",
file_id[:12], ordinal, time.monotonic() - t0, e)
+ # The cause is in the log line above; ffmpeg's own words carry the
+ # operator's paths and versions.
self._send({"type": "error",
- "detail": f"Subtitle extraction failed: {e}"})
+ "detail": "These subtitles could not be extracted"})
if isinstance(e, asyncio.CancelledError):
raise
finally:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py
index 107a43e..31cc83c 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py
@@ -7,7 +7,7 @@ import struct
import msgpack
from aiortc import RTCPeerConnection
-from meshbay_node.transport.webrtc.limits import MAX_MSG
+from meshbay_node.transport.webrtc.limits import MAX_MSG, UNPACK_LIMITS
def _extract_dtls_fingerprint(sdp: str) -> bytes:
@@ -77,7 +77,7 @@ class _DataChannelBuffer:
break
msg_bytes = bytes(self._buf[4:4 + length])
del self._buf[:4 + length]
- yield msgpack.unpackb(msg_bytes, raw=False)
+ yield msgpack.unpackb(msg_bytes, raw=False, **UNPACK_LIMITS)
def _get_remote_ip(pc: RTCPeerConnection) -> str:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py
index 26ec27c..493d7f0 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py
@@ -283,7 +283,18 @@ class ChatMixin:
# anyone on the node.
gctx = self._group_ctx()
chat_store = gctx.get("chat_store")
+ # The two fields that travel in clear beside the ciphertext (the sealed
+ # envelope carries its own). Stored and relayed to every member, so they
+ # are what they claim to be and no larger: a name as long as a username,
+ # a thread id as long as a message id. Anything else is dropped.
sender_name = msg.get("sender_name", "")
+ if not isinstance(sender_name, str) or len(sender_name) > 64:
+ sender_name = ""
+ thread_id = msg.get("thread_id")
+ id_like = (isinstance(thread_id, int) and not isinstance(thread_id, bool)
+ or isinstance(thread_id, str) and len(thread_id) <= 64)
+ if thread_id is not None and not id_like:
+ thread_id = None
# Two shapes, and keeping them apart is what makes this deployable.
#
@@ -329,7 +340,7 @@ class ChatMixin:
self._spawn(self._store_chat_message(
chat_store,
iteration=msg.get("iteration", 0), payload=raw,
- thread_id=msg.get("thread_id"), sender_name=sender_name,
+ thread_id=thread_id, sender_name=sender_name,
format=fmt, epoch=epoch, device=device, nonce=nonce, sig=sig,
))
@@ -340,7 +351,7 @@ class ChatMixin:
"sender_id": self._user_id,
"sender_name": sender_name,
"payload": payload,
- "thread_id": msg.get("thread_id"),
+ "thread_id": thread_id,
"timestamp": time.time(),
"format": fmt,
"epoch": epoch,
@@ -598,7 +609,7 @@ class ChatMixin:
the client asks, the node produces on demand, the asking device
caches — nothing durable here).
- `linkpreview.safe_url` is the SSRF gate: the URL a *member* chose
+ `linkpreview.check_url` is the SSRF gate: the URL a *member* chose
decides an outbound request from the operator's machine, so http(s)
only and the resolved address must be globally routable. Failure of
any kind — blocked, unreachable, not HTML, nothing worth showing —
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 882ddbc..ae3f0e1 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
@@ -139,7 +139,21 @@ class SessionCore:
log.info("WebRTC data received: %d bytes, msg #%d (peer=%s)",
len(message), self._msg_count, self._peer_id)
self._buffer.feed(message)
- for msg in self._buffer.messages():
+ decoded = self._buffer.messages()
+ while True:
+ try:
+ msg = next(decoded)
+ except StopIteration:
+ break
+ except ValueError as e:
+ # Over the size limit, or a container past its decode
+ # limit. The buffer still starts with that frame, so every
+ # later message would fail the same way: the session ends
+ # here. Only decoding is caught — a handler's own error is
+ # not a reason to drop the peer.
+ log.warning("Closing peer %s: %s", self._peer_id, e)
+ self._spawn(self.close())
+ break
self._handle_message(msg)
if _WEBRTC_TRACE:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py
index 7d458f4..86412e0 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py
@@ -2,7 +2,21 @@
CHUNK_SIZE = 1024 * 1024
-MAX_MSG = 64 * 1024 * 1024
+# The largest message a peer may send once it has proved the group key. The
+# largest a client really sends is a sealed playlist blob, 1 MiB (blobs.py);
+# chat is 64 KiB and an upload chunk 48 KiB. Eight times the largest, because a
+# message of many small objects decodes to several times its size in memory.
+MAX_MSG = 8 * 1024 * 1024
+
+# Per container, when a message is decoded: nothing a client sends comes near
+# them, and without them one message of tiny elements is one enormous list.
+UNPACK_LIMITS = {
+ "max_array_len": 100_000,
+ "max_map_len": 10_000,
+ "max_str_len": 1024 * 1024,
+ "max_bin_len": MAX_MSG,
+ "max_ext_len": 0,
+}
# What the `tr` on a chunk request turned out to be (see `_lease_of`).
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py
index 7ecb0a5..22c1690 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py
@@ -165,6 +165,7 @@ async def _seek_lands_at(file_path: Path, t: float, map_args: list[str]) -> floa
fd, tmp_name = tempfile.mkstemp(suffix=".mp4")
os.close(fd)
tmp_path = Path(tmp_name)
+ proc = probe = None
try:
proc = await asyncio.create_subprocess_exec(
platform.ffmpeg_cmd(), "-hide_banner", "-loglevel", "error", "-y",
@@ -189,6 +190,12 @@ async def _seek_lands_at(file_path: Path, t: float, map_args: list[str]) -> floa
log.warning("stream: seek probe failed at %.1fs: %r", t, e)
return None
finally:
+ # A timed-out wait leaves its process running; it is stopped here, not
+ # left to finish a seek nobody is waiting for.
+ for p in (proc, probe):
+ if p is not None and p.returncode is None:
+ p.kill()
+ await p.wait()
await _discard_scratch(tmp_path)
text = stdout.decode(errors="replace").strip().rstrip(",")
try:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
index 1c6d1ce..02a50af 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
@@ -8,7 +8,14 @@ from pathlib import Path
from meshbay_common.protocol import UPLOAD_PROBE_INDEX, file_upload_ack_wire, file_upload_payload
from meshbay_node import uploads as uploads_mod
-from meshbay_node.roots import SAFE_UPLOAD_NAME, RootSet, _free_name, off_disk
+from meshbay_node.roots import (
+ SAFE_UPLOAD_NAME,
+ RootSet,
+ _free_name,
+ off_disk,
+ publish_upload,
+ shell_active,
+)
from meshbay_node.transport.webrtc.disk import _append_chunk
from meshbay_node.transport.webrtc.limits import LEASE_NONE, LEASE_QUEUED
@@ -212,6 +219,9 @@ class UploadMixin:
if not SAFE_UPLOAD_NAME.match(filename):
_refuse("Invalid filename", "invalid_filename")
return
+ if shell_active(filename):
+ _refuse("This type of file is not accepted", "file_type_refused")
+ return
roots: RootSet | None = ctx.get("roots")
if not roots:
@@ -306,9 +316,13 @@ class UploadMixin:
# A shared directory means two people can send the same name. Refusing the
# second is safe but silly — everyone's camera produces IMG_1234.jpg — so
# a free name is found instead. Never a replacement.
+ # Names other uploads into this directory will publish under are taken
+ # too: none of them is on disk yet.
+ reserved = uploads.reserved_names(rel_dir)
stored_name = (state.stored_name if state
- else await off_disk(roots, _free_name, target_dir, filename))
- tmp_path = target_dir / f"{stored_name}{uploads_mod.PART_SUFFIX}"
+ else await off_disk(roots, _free_name, target_dir, filename, reserved))
+ tmp_path = (state.part_path if state and state.part_path
+ else target_dir / uploads_mod.part_name(stored_name))
final_path = target_dir / stored_name
if chunk_index == UPLOAD_PROBE_INDEX:
@@ -361,6 +375,22 @@ class UploadMixin:
await off_disk(roots, _append_chunk, tmp_path, chunk_bytes, chunk_index == 0)
uploads.advance(user_id, rel_dir, filename, chunk_index, len(chunk_bytes))
+ last = chunk_index + 1 >= total_chunks
+ if last:
+ # Published before the last ack, so the ack names the file as it is
+ # on disk: publication never replaces a file, and may have had to
+ # take another free name for this one.
+ uploads.drop(user_id, rel_dir, filename)
+ try:
+ stored_name = await off_disk(roots, publish_upload, tmp_path, target_dir,
+ stored_name, filename,
+ uploads.reserved_names(rel_dir))
+ except OSError as e:
+ log.warning("Upload %s could not be published: %s", stored_name, e)
+ _refuse("The file could not be stored", "store_failed")
+ return
+ final_path = target_dir / stored_name
+
self._send(file_upload_ack_wire(
gek, self._group_id or "",
upload_id=upload_id,
@@ -372,9 +402,7 @@ class UploadMixin:
dir=rel_dir,
))
- if chunk_index + 1 >= total_chunks:
- uploads.drop(user_id, rel_dir, filename)
- await off_disk(roots, tmp_path.rename, final_path)
+ if last:
log.info("Upload complete: %s (%d chunks, %d bytes)",
stored_name, total_chunks, state.bytes)
self._audit("file_upload", f"{rel_dir}/{stored_name}")
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
index 20636fa..0ca4b3d 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
@@ -70,7 +70,15 @@ log = logging.getLogger(__name__)
# so the cost to an operator grew with the number of people in their groups.
# Sized to be unreachable in ordinary use: a browser holds one connection per
# open group, and a handshake unfinished after a minute is not going to finish.
-MAX_PEER_SESSIONS = 64
+MAX_PEER_SESSIONS = 128
+# One account's share of them: half. A member of twenty groups hosted here,
+# with three devices and a spare tab, holds up to 52 (each device keeps up to
+# twelve connections for search and music, plus the open group page), so the
+# share never refuses real use — and no single member can hold more than half
+# of what the node will take. The node's own account is not counted against it:
+# it is the operator's machine. Measured, an idle connected session costs about
+# 0.15 MiB and one file descriptor.
+MAX_PEER_SESSIONS_PER_ACCOUNT = MAX_PEER_SESSIONS // 2
UNAUTHENTICATED_SESSION_TIMEOUT = 60 # seconds
@@ -238,8 +246,13 @@ class WebRTCTransport:
pass
return
+ def _sessions_of(self, user_id: str) -> int:
+ return sum(1 for s in self._sessions.values()
+ if (getattr(s, "_offer_user", "") or getattr(s, "_user_id", "") or "")
+ == user_id)
+
async def handle_offer(
- self, offer_sdp: str, peer_id: str,
+ self, offer_sdp: str, peer_id: str, user_id: str = "",
) -> tuple[str, list[dict]]:
"""
Process a WebRTC SDP offer from a browser client.
@@ -267,9 +280,18 @@ class WebRTCTransport:
log.warning("Refusing WebRTC offer: %d peer sessions already open",
len(self._sessions))
raise RuntimeError("Node is at its peer-connection limit")
+ # The account the hub authenticated for this offer. A hub that lied
+ # could only move the count between accounts; it can refuse offers
+ # outright already.
+ if (user_id and user_id != self._ctx.get("node_user_id")
+ and self._sessions_of(user_id) >= MAX_PEER_SESSIONS_PER_ACCOUNT):
+ log.warning("Refusing WebRTC offer: account %s already holds %d sessions",
+ user_id[:8], MAX_PEER_SESSIONS_PER_ACCOUNT)
+ raise RuntimeError("Account is at its peer-connection share")
pc = RTCPeerConnection(configuration=config)
session = WebRTCPeerSession(pc, self._ctx, peer_id=peer_id)
+ session._offer_user = user_id
self._sessions[peer_id] = session
self._reap_if_unauthenticated(peer_id)
diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py
index bbc4649..f810903 100644
--- a/packages/meshbay-node/src/meshbay_node/ui/app.py
+++ b/packages/meshbay-node/src/meshbay_node/ui/app.py
@@ -179,6 +179,7 @@ def create_ui_app(state: dict) -> FastAPI:
(payload.get("name") or "").strip(),
(payload.get("shared_dir") or "").strip(),
writable=bool(payload.get("writable", True)),
+ join_policy=str(payload.get("join_policy") or "invite"),
))
reload_fn = state.get("reload_fn")
if reload_fn:
diff --git a/packages/meshbay-node/src/meshbay_node/uploads.py b/packages/meshbay-node/src/meshbay_node/uploads.py
index dd31e1e..c485e39 100644
--- a/packages/meshbay-node/src/meshbay_node/uploads.py
+++ b/packages/meshbay-node/src/meshbay_node/uploads.py
@@ -24,6 +24,8 @@ build first.
from __future__ import annotations
+import re
+import secrets
import time
from collections.abc import Iterable
from dataclasses import dataclass, field
@@ -34,6 +36,26 @@ from pathlib import Path
# recognise one, and a second spelling of it would be a bug nobody could see.
PART_SUFFIX = ".part"
+
+def part_name(stored_name: str) -> str:
+ """The `.part` one upload writes: its final name, a tag of its own, `.part`.
+
+ Its own, because the final name alone is shared: two uploads that settled
+ on one name — two groups hosting one folder, each with its own lock —
+ would write one file, the second truncating the first.
+ """
+ return f"{stored_name}.{secrets.token_hex(4)}{PART_SUFFIX}"
+
+
+# What `part_name` writes, and the only thing the reaper deletes. A `.part`
+# without the node's tag is somebody else's — a browser's download in progress in
+# a shared folder, a copy the operator is making — and is never touched.
+_OWN_PART = re.compile(r"\.[0-9a-f]{8}" + re.escape(PART_SUFFIX) + r"$")
+
+
+def is_own_part(path: Path) -> bool:
+ return bool(_OWN_PART.search(path.name))
+
# How long a `.part` with no upload behind it is kept before it is deleted.
#
# Generous on purpose. The cost of waiting is disk; the cost of being wrong is
@@ -107,6 +129,15 @@ class PartialUploads:
def drop(self, user_id: str, rel_dir: str, filename: str) -> Partial | None:
return self._by_key.pop((user_id, rel_dir, filename), None)
+ def reserved_names(self, rel_dir: str) -> set[str]:
+ """The final names uploads in flight into `rel_dir` will take.
+
+ None of them exists on disk yet, so a name check that looked only at
+ the directory would hand the same name to a second upload.
+ """
+ return {state.stored_name for (_u, d, _f), state in self._by_key.items()
+ if d == rel_dir}
+
def __len__(self) -> int:
return len(self._by_key)
@@ -145,7 +176,7 @@ def orphaned_parts(candidates: Iterable[tuple[Path, float]],
"""
doomed: list[Path] = []
for path, mtime in candidates:
- if path.suffix != PART_SUFFIX:
+ if not is_own_part(path):
continue
if path in live:
continue
diff --git a/packages/meshbay-node/tests/golden/cli.json b/packages/meshbay-node/tests/golden/cli.json
index b397eeb..48d2f74 100644
--- a/packages/meshbay-node/tests/golden/cli.json
+++ b/packages/meshbay-node/tests/golden/cli.json
@@ -4,7 +4,7 @@
"asked": [],
"exit": 0,
"stderr": "",
- "stdout": "usage: meshbay-node [-h] [--hub-url HUB_URL] [--username USERNAME] [--dir DIR] [--yes]\n [--config CONFIG] [--group GROUP] [--link] [--writable] [--no-writable]\n [--removable] [--no-removable] [--name NAME]\n [--log-level {DEBUG,INFO,WARNING,ERROR}]\n [{init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}]\n [subcommand] [target] [value]\n\nMeshBay Node daemon\n\npositional arguments:\n {init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}\n init: provision config + keystore | reset: erase all node state | status:\n node state and keys | operator pair: pair a browser with this node |\n member list|invite|cancel|revoke|unpin | group list|add|remove | root\n list|add|remove|set|eject|plug | gek init|rotate | file list|rm | video\n rematch: re-resolve TMDB matches for a group's videos | chat\n status|rotate|encrypt-history|prune | denylist show|clear | stun\n list|add|remove|reset | transfers show|set|max-size|per-member: live\n transfer slots, the node-wide caps, the largest single upload, and how\n many one member may run at once in a group | reload: re-read node.toml\n (hot; systemd or the loopback API) | restart-daemon: restart the node\n (systemd unit, the Windows autostart launcher, or the service task,\n whichever applies) | autostart install|remove|start|stop|status (Windows:\n run meshbay-node at each sign-in, no admin) | service\n install|remove|start|stop|status (Windows: run at boot, before sign-in,\n needs admin once to install) | calibrate-argon2: benchmark\n subcommand 'pair' for operator; list|invite|revoke|unpin for member; list|add|remove\n for group; list|add|remove|set|eject|plug for root; init|rotate for gek;\n list|rm for file; rematch for video; show|clear for denylist;\n list|add|remove|reset for stun; show|set|max-size|per-member for\n transfers; install|remove|start|stop|status for autostart and for service\n target username for member invite|revoke|unpin (an optional e-mail label with\n --link, a link id for member cancel); group name for group add; file id\n for file rm; identifier for denylist clear; download cap for transfers\n set; size in GB for transfers max-size\n value the second value where a verb takes two: the upload cap for transfers set\n\noptions:\n -h, --help show this help message and exit\n --hub-url HUB_URL hub URL, for init (e.g. https://meshbay.org)\n --username USERNAME hub username, for init\n --dir DIR shared directory, for group add\n --yes skip the confirmation for destructive commands\n --config CONFIG Config file path\n --group GROUP group id (optional if only one is configured)\n --link member invite: an invitation link, for someone who may have no account yet\n (valid 7 days, single use)\n --writable root accepts member uploads (root add/set)\n --no-writable root is read-only (root add/set, group add)\n --removable mark root as removable (root set/add)\n --no-removable mark root as not removable (root set)\n --name NAME root name (root add; defaults to directory basename)\n --log-level {DEBUG,INFO,WARNING,ERROR}\n",
+ "stdout": "usage: meshbay-node [-h] [--hub-url HUB_URL] [--username USERNAME] [--dir DIR] [--yes]\n [--config CONFIG] [--group GROUP] [--link] [--writable] [--no-writable]\n [--removable] [--no-removable] [--open] [--name NAME]\n [--log-level {DEBUG,INFO,WARNING,ERROR}]\n [{init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}]\n [subcommand] [target] [value]\n\nMeshBay Node daemon\n\npositional arguments:\n {init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}\n init: provision config + keystore | reset: erase all node state | status:\n node state and keys | operator pair: pair a browser with this node |\n member list|invite|cancel|revoke|unpin | group list|add|remove | root\n list|add|remove|set|eject|plug | gek init|rotate | file list|rm | video\n rematch: re-resolve TMDB matches for a group's videos | chat\n status|rotate|encrypt-history|prune | denylist show|clear | stun\n list|add|remove|reset | transfers show|set|max-size|per-member: live\n transfer slots, the node-wide caps, the largest single upload, and how\n many one member may run at once in a group | reload: re-read node.toml\n (hot; systemd or the loopback API) | restart-daemon: restart the node\n (systemd unit, the Windows autostart launcher, or the service task,\n whichever applies) | autostart install|remove|start|stop|status (Windows:\n run meshbay-node at each sign-in, no admin) | service\n install|remove|start|stop|status (Windows: run at boot, before sign-in,\n needs admin once to install) | calibrate-argon2: benchmark\n subcommand 'pair' for operator; list|invite|revoke|unpin for member; list|add|remove\n for group; list|add|remove|set|eject|plug for root; init|rotate for gek;\n list|rm for file; rematch for video; show|clear for denylist;\n list|add|remove|reset for stun; show|set|max-size|per-member for\n transfers; install|remove|start|stop|status for autostart and for service\n target username for member invite|revoke|unpin (an optional e-mail label with\n --link, a link id for member cancel); group name for group add; file id\n for file rm; identifier for denylist clear; download cap for transfers\n set; size in GB for transfers max-size\n value the second value where a verb takes two: the upload cap for transfers set\n\noptions:\n -h, --help show this help message and exit\n --hub-url HUB_URL hub URL, for init (e.g. https://meshbay.org)\n --username USERNAME hub username, for init\n --dir DIR shared directory, for group add\n --yes skip the confirmation for destructive commands\n --config CONFIG Config file path\n --group GROUP group id (optional if only one is configured)\n --link member invite: an invitation link, for someone who may have no account yet\n (valid 7 days, single use)\n --writable root accepts member uploads (root add/set)\n --no-writable root is read-only (root add/set, group add)\n --removable mark root as removable (root set/add)\n --no-removable mark root as not removable (root set)\n --open group add: anyone the hub lists the group to may join (default: by\n invitation only)\n --name NAME root name (root add; defaults to directory basename)\n --log-level {DEBUG,INFO,WARNING,ERROR}\n",
"systemctl": []
},
"autostart no-such-sub": {
@@ -266,7 +266,7 @@
"asked": [],
"exit": 1,
"stderr": "",
- "stdout": "usage: meshbay-node group add <name> --dir <path> [--no-writable]\n\nThe group must already exist on the hub and be yours. This\nonly tells the node to host it, and picks its first\ndirectory, which accepts uploads unless --no-writable.\nAdd more with: meshbay-node root add <path> [--writable]\n",
+ "stdout": "usage: meshbay-node group add <name> --dir <path> [--no-writable] [--open]\n\nThe group must already exist on the hub and be yours. This\nonly tells the node to host it, and picks its first\ndirectory, which accepts uploads unless --no-writable.\nAdd more with: meshbay-node root add <path> [--writable]\n",
"systemctl": []
},
"group add g --dir /tmp/media --no-writable": {
@@ -275,6 +275,7 @@
"POST",
"/api/groups/attach",
{
+ "join_policy": "invite",
"name": "g",
"shared_dir": "/tmp/media",
"writable": false
@@ -284,7 +285,7 @@
"asked": [],
"exit": 0,
"stderr": "",
- "stdout": "g (g) added to <tmp>/node.toml\n shared_dir <tmp> (read-only)\n\nTell the daemon to re-read its config, then give the group a key:\n meshbay-node reload\n meshbay-node gek init --group g\n\nThe key is this group's own — members of your other groups cannot\nread it, and joining one says nothing about the other.\n",
+ "stdout": "g (g) added to <tmp>/node.toml\n shared_dir <tmp> (read-only)\n join_policy invite\n\nTell the daemon to re-read its config, then give the group a key:\n meshbay-node reload\n meshbay-node gek init --group g\n\nThe key is this group's own — members of your other groups cannot\nread it, and joining one says nothing about the other.\n",
"systemctl": []
},
"group list": {
@@ -423,7 +424,7 @@
"api": [],
"asked": [],
"exit": 2,
- "stderr": "usage: meshbay-node [-h] [--hub-url HUB_URL] [--username USERNAME] [--dir DIR] [--yes]\n [--config CONFIG] [--group GROUP] [--link] [--writable] [--no-writable]\n [--removable] [--no-removable] [--name NAME]\n [--log-level {DEBUG,INFO,WARNING,ERROR}]\n [{init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}]\n [subcommand] [target] [value]\nmeshbay-node: error: argument command: invalid choice: 'no-such-verb' (choose from init, reset, status, gek-init, gek, operator, member, group, root, file, video, chat, denylist, stun, transfers, reload, restart-daemon, autostart, service, calibrate-argon2)\n",
+ "stderr": "usage: meshbay-node [-h] [--hub-url HUB_URL] [--username USERNAME] [--dir DIR] [--yes]\n [--config CONFIG] [--group GROUP] [--link] [--writable] [--no-writable]\n [--removable] [--no-removable] [--open] [--name NAME]\n [--log-level {DEBUG,INFO,WARNING,ERROR}]\n [{init,reset,status,gek-init,gek,operator,member,group,root,file,video,chat,denylist,stun,transfers,reload,restart-daemon,autostart,service,calibrate-argon2}]\n [subcommand] [target] [value]\nmeshbay-node: error: argument command: invalid choice: 'no-such-verb' (choose from init, reset, status, gek-init, gek, operator, member, group, root, file, video, chat, denylist, stun, transfers, reload, restart-daemon, autostart, service, calibrate-argon2)\n",
"stdout": "",
"systemctl": []
},
diff --git a/packages/meshbay-node/tests/test_attach_from_the_hub.py b/packages/meshbay-node/tests/test_attach_from_the_hub.py
new file mode 100644
index 0000000..f4aaa25
--- /dev/null
+++ b/packages/meshbay-node/tests/test_attach_from_the_hub.py
@@ -0,0 +1,99 @@
+"""
+Hosting a group writes what the operator said, never what the hub says.
+
+`attach_group` looks the group up on the hub, because the node is the process
+signed in there. What it must not take from that answer is how people join: a
+hub able to declare a group open would be handed its key by anyone it sent
+(admission in `transport/webrtc/admission.py` admits a stranger to an open
+group). And every string it writes into node.toml is someone else's text — a
+group name chosen on the hub, a folder name — so none of it may end a TOML
+string and write lines of its own.
+"""
+
+import tomllib
+from pathlib import Path
+
+import pytest
+from meshbay_node import ops
+from meshbay_node.config import load_config
+from meshbay_node.ops.node_toml import toml_string
+
+GID = "0f8fad5b-d9cb-469f-a165-70867728950e"
+HOSTILE = 'Films"\n[node]\nui_port = 1\n# '
+
+
+class _Hub:
+ _session = object() # signed in
+
+ def __init__(self, group):
+ self._group = group
+
+ async def list_my_groups(self):
+ return [self._group]
+
+
+def _state(tmp_path: Path, group: dict) -> dict:
+ conf = tmp_path / "node.toml"
+ conf.write_text('[hub]\nurl = "https://hub.invalid"\nusername = "op"\n\n'
+ '[node]\nui_port = 18000\n', encoding="utf-8")
+ return {"config": load_config(conf), "config_path": str(conf), "hub": _Hub(group)}
+
+
+def _hosted(tmp_path: Path) -> dict:
+ parsed = tomllib.loads((tmp_path / "node.toml").read_text(encoding="utf-8"))
+ return parsed["groups"][0] | {"node": parsed["node"]}
+
+
+async def test_a_group_the_hub_calls_open_is_hosted_by_invitation(tmp_path):
+ state = _state(tmp_path, {"id": GID, "name": "Films", "visibility": "public",
+ "join_policy": "open"})
+ out = await ops.attach_group(state, "Films", str(tmp_path / "share"))
+
+ hosted = _hosted(tmp_path)
+ assert hosted["join_policy"] == "invite"
+ assert hosted["visibility"] == "private"
+ assert out["hub_join_policy"] == "open", "the caller is not told the two differ"
+
+
+async def test_the_operator_opens_it(tmp_path):
+ state = _state(tmp_path, {"id": GID, "name": "Films", "join_policy": "invite"})
+ await ops.attach_group(state, "Films", str(tmp_path / "share"), join_policy="open")
+
+ hosted = _hosted(tmp_path)
+ assert (hosted["join_policy"], hosted["visibility"]) == ("open", "public")
+
+
+async def test_an_unknown_policy_is_refused(tmp_path):
+ state = _state(tmp_path, {"id": GID, "name": "Films"})
+ with pytest.raises(ops.OpError):
+ await ops.attach_group(state, "Films", str(tmp_path / "share"),
+ join_policy="anyone")
+
+
+async def test_a_group_name_cannot_write_lines_into_node_toml(tmp_path):
+ state = _state(tmp_path, {"id": GID, "name": HOSTILE})
+ await ops.attach_group(state, HOSTILE, str(tmp_path / "share"))
+
+ hosted = _hosted(tmp_path)
+ assert hosted["name"] == HOSTILE
+ assert hosted["node"]["ui_port"] == 18000
+
+
+async def test_a_folder_name_cannot_either(tmp_path):
+ state = _state(tmp_path, {"id": GID, "name": "Films"})
+ await ops.attach_group(state, "Films", str(tmp_path / "share"))
+ state["config"] = load_config(tmp_path / "node.toml")
+ # Root names are refused with such characters already (roots.py); the
+ # path is not, and is the operator's own folder or a member's request.
+ weird = tmp_path / 'a "quoted"\\ folder\n[node]'
+ await ops.add_root(state, GID, str(weird), name="extra")
+
+ hosted = _hosted(tmp_path)
+ assert Path(hosted["roots"][-1]["path"]) == weird
+ assert hosted["node"]["ui_port"] == 18000
+
+
+@pytest.mark.parametrize("value", ["plain", 'q"uote', "back\\slash", "line\nbreak",
+ "tab\there", "del\x7f", "café 日本"])
+def test_every_string_reads_back_as_written(value):
+ assert tomllib.loads(f"v = {toml_string(value)}")["v"] == value
diff --git a/packages/meshbay-node/tests/test_chat_is_bounded.py b/packages/meshbay-node/tests/test_chat_is_bounded.py
index 3332601..c48f26d 100644
--- a/packages/meshbay-node/tests/test_chat_is_bounded.py
+++ b/packages/meshbay-node/tests/test_chat_is_bounded.py
@@ -204,3 +204,38 @@ async def test_one_member_at_their_limit_has_not_spent_anyone_elses(ctx, store):
await _flood(bob, ctx, 1)
assert not _errors(bob), "one member's flood silenced another"
assert _acks(bob)
+
+
+# ── the fields beside the ciphertext ─────────────────────────────────────────
+
+async def test_the_clear_fields_are_what_they_claim_and_no_larger(ctx, store):
+ """
+ `sender_name` and `thread_id` travel in clear beside the sealed envelope,
+ which carries its own. They are stored on the operator's disk and relayed
+ to every member, so a megabyte of name or a list for a thread id is dropped,
+ not kept.
+ """
+ alice = _session(ctx, store, user="alice", conn="c1")
+ bob = _session(ctx, store, user="bob", conn="c2")
+ msg = _message(alice, size=64)
+ msg["sender_name"] = "x" * (1024 * 1024)
+ msg["thread_id"] = list(range(10_000))
+ alice._do_chat_message(msg)
+ await _drain(ctx)
+
+ stored = (await store.get_recent(limit=1))[0]
+ assert stored.sender_name in ("", None)
+ assert stored.thread_id is None
+ relayed = [m for m in bob.sent if m.get("type") == "chat_msg"]
+ assert relayed and relayed[-1]["sender_name"] == "" and relayed[-1]["thread_id"] is None
+
+
+async def test_ordinary_clear_fields_pass_unchanged(ctx, store):
+ alice = _session(ctx, store, user="alice", conn="c1")
+ msg = _message(alice, size=64)
+ msg["sender_name"] = "Alice"
+ msg["thread_id"] = "42"
+ alice._do_chat_message(msg)
+ await _drain(ctx)
+ stored = (await store.get_recent(limit=1))[0]
+ assert (stored.sender_name, stored.thread_id) == ("Alice", "42")
diff --git a/packages/meshbay-node/tests/test_ffprobe_is_bounded.py b/packages/meshbay-node/tests/test_ffprobe_is_bounded.py
new file mode 100644
index 0000000..c96a19a
--- /dev/null
+++ b/packages/meshbay-node/tests/test_ffprobe_is_bounded.py
@@ -0,0 +1,61 @@
+"""
+ffprobe over a member's file is bounded, and a bounded wait stops the process.
+
+`probe_video` runs before every stream and subtitle request and in the
+enrichment pool. A file that keeps ffprobe busy must not hold any of them, and
+giving up on the wait is not enough: an asyncio subprocess whose wait was
+cancelled keeps running. These run a stand-in ffprobe that never answers and
+check both — the call returns, and the process is gone.
+"""
+
+import asyncio
+import os
+import sys
+import time
+
+import pytest
+from meshbay_node import media_probe, platform
+
+pytestmark = pytest.mark.skipif(sys.platform == "win32", reason="a POSIX shell stand-in")
+
+
+@pytest.fixture
+def hanging_ffprobe(tmp_path, monkeypatch):
+ pid_file = tmp_path / "pid"
+ tool = tmp_path / "ffprobe"
+ tool.write_text(f"#!/bin/sh\necho $$ > {pid_file}\nexec sleep 600\n")
+ tool.chmod(0o755)
+ monkeypatch.setattr(platform, "_ffprobe_path", str(tool))
+ return pid_file
+
+
+def _gone(pid: int) -> bool:
+ try:
+ os.kill(pid, 0)
+ except ProcessLookupError:
+ return True
+ # A zombie still answers kill(0); its state says it has exited.
+ try:
+ with open(f"/proc/{pid}/stat") as f:
+ return f.read().split()[2] == "Z"
+ except OSError:
+ return True
+
+
+async def test_a_probe_that_never_answers_times_out_and_is_stopped(hanging_ffprobe, monkeypatch):
+ monkeypatch.setattr(media_probe, "FFPROBE_TIMEOUT_SECS", 0.5)
+ started = time.monotonic()
+ with pytest.raises(RuntimeError, match="timed out"):
+ await media_probe.probe_video("/nonexistent/file.mkv")
+ assert time.monotonic() - started < 5
+ assert _gone(int(hanging_ffprobe.read_text()))
+
+
+async def test_a_caller_giving_up_stops_it_too(hanging_ffprobe):
+ with pytest.raises(TimeoutError):
+ await asyncio.wait_for(media_probe.probe_video("/nonexistent/file.mkv"), 0.5)
+ for _ in range(50):
+ if hanging_ffprobe.exists():
+ break
+ await asyncio.sleep(0.05)
+ assert _gone(int(hanging_ffprobe.read_text()))
diff --git a/packages/meshbay-node/tests/test_linkpreview.py b/packages/meshbay-node/tests/test_linkpreview.py
index 9fca186..3e6eaf7 100644
--- a/packages/meshbay-node/tests/test_linkpreview.py
+++ b/packages/meshbay-node/tests/test_linkpreview.py
@@ -6,12 +6,13 @@ decides an outbound request from the operator's machine. Anything that is not
a public http(s) address must be refused before a socket opens.
"""
+import asyncio
import socket
import httpx
import pytest
from meshbay_node import linkpreview
-from meshbay_node.linkpreview import UnsafeURL, safe_url
+from meshbay_node.linkpreview import UnsafeURL, check_url, safe_url
PUBLIC_IP = "93.184.216.34" # example.com, historically
@@ -43,9 +44,9 @@ def resolves_public(monkeypatch):
"javascript:alert(1)",
"not a url",
])
-def test_safe_url_refuses(url):
+async def test_check_url_refuses(url):
with pytest.raises(UnsafeURL):
- safe_url(url)
+ await check_url(url)
@pytest.mark.parametrize("url", [
@@ -72,11 +73,11 @@ def test_safe_url_allows_the_web_ports(url, resolves_public):
assert safe_url(url) == url
-def test_safe_url_accepts_a_public_host(resolves_public):
- assert safe_url("https://example.com/some/page") == "https://example.com/some/page"
+async def test_check_url_accepts_a_public_host(resolves_public):
+ assert await check_url("https://example.com/some/page") == "https://example.com/some/page"
-def test_safe_url_refuses_a_host_with_any_private_record(monkeypatch):
+async def test_check_url_refuses_a_host_with_any_private_record(monkeypatch):
def mixed(host, port, *a, **k):
return [
(socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, "", (PUBLIC_IP, port)),
@@ -84,7 +85,7 @@ def test_safe_url_refuses_a_host_with_any_private_record(monkeypatch):
]
monkeypatch.setattr(linkpreview.socket, "getaddrinfo", mixed)
with pytest.raises(UnsafeURL):
- safe_url("https://sneaky.example/x")
+ await check_url("https://sneaky.example/x")
# ── fetch_preview ──────────────────────────────────────────────────────────
@@ -273,3 +274,82 @@ async def test_a_declared_oversized_image_is_not_read(resolves_public):
async with _client(handler) as c:
assert await linkpreview.fetch_image("https://example.com/x.png", client=c) is None
assert counter["sent"] <= _CHUNK
+
+
+# ── The connection goes where the check said ───────────────────────────────
+
+class _Recorder:
+ """Stands in for the real socket layer under the pinned backend."""
+ def __init__(self):
+ self.hosts = []
+
+ async def connect_tcp(self, host, port, **kw):
+ self.hosts.append(host)
+ raise httpx.ConnectError("recorded, not connected")
+
+
+def _pinned_with(recorder):
+ backend = linkpreview._PinnedBackend()
+ backend._inner = recorder
+ return backend
+
+
+async def test_the_socket_is_opened_to_the_checked_address(resolves_public):
+ rec = _Recorder()
+ with pytest.raises(httpx.ConnectError):
+ await _pinned_with(rec).connect_tcp("example.com", 443)
+ assert rec.hosts == [PUBLIC_IP], "the name, not the checked address, was dialled"
+
+
+async def test_a_name_that_rebinds_never_reaches_the_lan(monkeypatch):
+ """
+ Answers clean when checked, then with a LAN address. Checked once and
+ dialled by name, the request would go to the LAN before anything looked;
+ resolved and checked by the backend that dials, it goes nowhere.
+ """
+ answers = iter([PUBLIC_IP, "192.168.1.1", "192.168.1.1"])
+
+ def rebinding(host, port, *a, **k):
+ return [(socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, "",
+ (next(answers), port))]
+ monkeypatch.setattr(linkpreview.socket, "getaddrinfo", rebinding)
+
+ await check_url("http://rebind.example/x") # the clean answer
+ rec = _Recorder()
+ with pytest.raises(UnsafeURL):
+ await _pinned_with(rec).connect_tcp("rebind.example", 80)
+ assert rec.hosts == []
+
+
+async def test_the_real_client_is_pinned():
+ """What `fetch_preview` uses when the caller gives no client."""
+ client = linkpreview._new_client()
+ try:
+ pool = client._transport._pool
+ assert isinstance(pool._network_backend, linkpreview._PinnedBackend)
+ assert client._trust_env is False, "a proxy from the environment would unpin it"
+ finally:
+ await client.aclose()
+
+
+async def test_resolving_does_not_hold_the_event_loop(monkeypatch):
+ import time as _time
+
+ def slow(host, port, *a, **k):
+ _time.sleep(0.4)
+ return [(socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, "", (PUBLIC_IP, port))]
+ monkeypatch.setattr(linkpreview.socket, "getaddrinfo", slow)
+
+ ticks = 0
+
+ async def ticker():
+ nonlocal ticks
+ while True:
+ await asyncio.sleep(0.02)
+ ticks += 1
+ t = asyncio.create_task(ticker())
+ try:
+ await check_url("https://slow.example/x")
+ finally:
+ t.cancel()
+ assert ticks >= 10, "the event loop stood still while a name resolved"
diff --git a/packages/meshbay-node/tests/test_member_capacity.py b/packages/meshbay-node/tests/test_member_capacity.py
new file mode 100644
index 0000000..16d112d
--- /dev/null
+++ b/packages/meshbay-node/tests/test_member_capacity.py
@@ -0,0 +1,151 @@
+"""
+What one member may hold of a node: a share, sized so that real use never
+meets it.
+
+The heaviest real member — twenty groups on this node, three devices and a
+spare tab — holds up to 52 peer sessions (twelve per device for search and
+music, plus the open group page). The node holds 128, and one account at most
+half. One account may play half the node's video slots and run two subtitle
+extractions. The node's own account is the operator's machine and is not
+counted. And a message is at most 8 MiB, decoded with a bound on every
+container, because one message of tiny elements decodes to many times its size.
+"""
+
+import struct
+from unittest.mock import MagicMock
+
+import msgpack
+import pytest
+from meshbay_node.transport.webrtc.channel import _DataChannelBuffer
+from meshbay_node.transport.webrtc.limits import MAX_MSG
+from meshbay_node.transport.webrtc_server import (
+ MAX_PEER_SESSIONS,
+ MAX_PEER_SESSIONS_PER_ACCOUNT,
+ WebRTCPeerSession,
+ WebRTCTransport,
+)
+
+
+class _Held:
+ def __init__(self, user):
+ self._offer_user = user
+ self._user_id = user
+
+ async def close(self):
+ pass
+
+
+def _transport(**ctx) -> WebRTCTransport:
+ tp = WebRTCTransport(sk_node=MagicMock(), hub_pk_pem=b"", gek=None,
+ roots=None, index=None)
+ tp._ctx.update(ctx)
+ return tp
+
+
+def test_the_shares_are_what_was_agreed():
+ assert MAX_PEER_SESSIONS == 128
+ assert MAX_PEER_SESSIONS_PER_ACCOUNT == 64
+ assert MAX_MSG == 8 * 1024 * 1024
+ # The heaviest real member fits with room to spare.
+ assert 4 * (12 + 1) < MAX_PEER_SESSIONS_PER_ACCOUNT
+
+
+@pytest.mark.asyncio
+async def test_one_account_cannot_hold_more_than_its_share():
+ tp = _transport()
+ for i in range(MAX_PEER_SESSIONS_PER_ACCOUNT):
+ tp._sessions[f"a-{i}"] = _Held("alice")
+
+ with pytest.raises(RuntimeError, match="share"):
+ await tp.handle_offer("v=0", "alice-one-more", "alice")
+ assert "alice-one-more" not in tp._sessions
+
+ # Somebody else still gets in: what stops this offer is not the share.
+ try:
+ await tp.handle_offer("v=0", "bob-first", "bob")
+ except RuntimeError as e:
+ assert "share" not in str(e) and "limit" not in str(e)
+ except Exception:
+ pass
+
+
+@pytest.mark.asyncio
+async def test_the_operators_own_account_is_not_counted():
+ tp = _transport(node_user_id="operator")
+ for i in range(MAX_PEER_SESSIONS_PER_ACCOUNT):
+ tp._sessions[f"o-{i}"] = _Held("operator")
+ try:
+ await tp.handle_offer("v=0", "operator-more", "operator")
+ except RuntimeError as e:
+ assert "share" not in str(e)
+ except Exception:
+ pass
+
+
+def _session(ctx, user):
+ s = WebRTCPeerSession.__new__(WebRTCPeerSession)
+ s._ctx = ctx
+ s._user_id = user
+ s.sent = []
+ s._send = s.sent.append
+ return s
+
+
+def test_an_accounts_devices_share_one_allowance():
+ ctx = {}
+ phone, laptop = _session(ctx, "alice"), _session(ctx, "alice")
+ with phone._account_share("streams", 2) as a, laptop._account_share("streams", 2) as b:
+ assert a and b
+ with phone._account_share("streams", 2) as c:
+ assert c is False
+ with _session(ctx, "bob")._account_share("streams", 2) as d:
+ assert d, "another member is not counted against alice"
+ with phone._account_share("streams", 2) as e:
+ assert e, "a place is given back when its stream ends"
+
+
+def test_the_operator_is_not_counted_for_streams_either():
+ ctx = {"node_user_id": "operator"}
+ s = _session(ctx, "operator")
+ with s._account_share("streams", 1) as a, s._account_share("streams", 1) as b:
+ assert a and b
+
+
+@pytest.mark.asyncio
+async def test_a_stream_past_the_accounts_share_is_refused_by_name():
+ ctx = {"max_concurrent_streams": 8, "_streams_by_account": {"alice": 4}}
+ s = _session(ctx, "alice")
+ assert s._streams_per_account() == 4
+ await s._stream_video({"file_id": "x"})
+ assert s.sent[-1]["type"] == "error"
+ assert "Too many videos" in s.sent[-1]["detail"]
+
+
+@pytest.mark.parametrize("cap,share", [(1, 1), (2, 1), (3, 2), (8, 4), (9, 5)])
+def test_the_stream_share_is_half_rounded_up(cap, share):
+ assert _session({"max_concurrent_streams": cap}, "a")._streams_per_account() == share
+
+
+def _frame(obj) -> bytes:
+ body = msgpack.packb(obj, use_bin_type=True)
+ return struct.pack(">I", len(body)) + body
+
+
+def test_the_largest_real_message_passes():
+ buf = _DataChannelBuffer()
+ buf.feed(_frame({"type": "user_blob_store", "blob": b"x" * (1024 * 1024 + 64)}))
+ assert len(list(buf.messages())) == 1
+
+
+def test_a_message_of_tiny_elements_is_refused():
+ buf = _DataChannelBuffer()
+ buf.feed(_frame({"type": "x", "items": [None] * 200_000}))
+ with pytest.raises(ValueError):
+ list(buf.messages())
+
+
+def test_a_message_over_the_limit_is_refused():
+ buf = _DataChannelBuffer()
+ buf.feed(struct.pack(">I", MAX_MSG + 1) + b"x")
+ with pytest.raises(ValueError):
+ list(buf.messages())
diff --git a/packages/meshbay-node/tests/test_member_errors_are_plain.py b/packages/meshbay-node/tests/test_member_errors_are_plain.py
new file mode 100644
index 0000000..e144e9b
--- /dev/null
+++ b/packages/meshbay-node/tests/test_member_errors_are_plain.py
@@ -0,0 +1,22 @@
+"""
+What a member is told when the media tools fail: that it failed.
+
+An exception's text from ffmpeg or ffprobe names the operator's paths, versions
+and the libraries the build has; the member who asked needs none of it, and the
+operator finds it in their log. Read from the source, because what matters is
+that no reply in these handlers carries an exception's text at all.
+"""
+
+import re
+from pathlib import Path
+
+APPS = Path(__file__).resolve().parents[1] / "src" / "meshbay_node" / "transport" / "webrtc"
+
+
+def test_no_reply_to_a_member_carries_an_exceptions_text():
+ offenders = []
+ for path in APPS.rglob("*.py"):
+ text = path.read_text(encoding="utf-8")
+ for m in re.finditer(r'"detail":\s*(f"[^"]*\{e\}[^"]*"|str\(e\))', text):
+ offenders.append(f"{path.name}: {m.group(0)}")
+ assert not offenders, offenders
diff --git a/packages/meshbay-node/tests/test_ops.py b/packages/meshbay-node/tests/test_ops.py
index 9eb8339..2c5a15a 100644
--- a/packages/meshbay-node/tests/test_ops.py
+++ b/packages/meshbay-node/tests/test_ops.py
@@ -405,7 +405,7 @@ def test_every_path_written_into_node_toml_goes_through_as_posix():
import re
source = ops_source()
# Every f-string interpolation that lands on the right of a TOML `path =`.
- writes = re.findall(r'path\s*=\s*\\?"\{([^}]+)\}', source)
+ writes = re.findall(r'path\s*=\s*\\?"?\{(?:toml_string\()?([^}]+)\}', source)
assert writes, "no TOML path writer found — did the config writer move?"
for expr in writes:
assert "as_posix()" in expr, (
diff --git a/packages/meshbay-node/tests/test_partial_uploads.py b/packages/meshbay-node/tests/test_partial_uploads.py
index 80ab454..52aa306 100644
--- a/packages/meshbay-node/tests/test_partial_uploads.py
+++ b/packages/meshbay-node/tests/test_partial_uploads.py
@@ -17,10 +17,12 @@ may inherit — or overwrite the position of — the other's.
"""
import os
+import re
import time
import types
from pathlib import Path
+import pytest
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
from meshbay_common.crypto import generate_gek
from meshbay_common.protocol import (
@@ -104,7 +106,7 @@ def _old(seconds: float) -> float:
NOW = 1_000_000.0
-FILM = Path("/roots/media/film.mkv.part")
+FILM = Path("/roots/media/film.mkv.0123abcd.part")
def test_a_part_nobody_is_writing_and_nobody_has_touched_is_deleted():
@@ -138,7 +140,7 @@ def test_the_same_name_in_another_directory_does_not_protect_it():
path rebuilt from a root and a relative directory would be a second
implementation that has to agree with the first for ever, and the state
records the path it is writing instead."""
- other = Path("/roots/archive/film.mkv.part")
+ other = Path("/roots/archive/film.mkv.0123abcd.part")
doomed = orphaned_parts([(other, _old(ORPHAN_AFTER_SECS + 1))],
live={FILM}, now=NOW)
assert doomed == [other]
@@ -157,15 +159,15 @@ def test_a_finished_file_is_not_a_candidate():
def test_a_file_from_the_future_is_left_alone():
"""A clock that went backwards is not evidence that a file is abandoned, and
deleting is not reversible."""
- doomed = orphaned_parts([(Path("/roots/media/a.part"), NOW + 10_000)],
+ doomed = orphaned_parts([(Path("/roots/media/a.0123abcd.part"), NOW + 10_000)],
live=set(), now=NOW)
assert doomed == []
def test_the_boundary_is_the_age_itself():
- at = [(Path("/roots/media/a.part"), _old(ORPHAN_AFTER_SECS))]
- just_under = [(Path("/roots/media/a.part"), _old(ORPHAN_AFTER_SECS - 1))]
- assert orphaned_parts(at, set(), NOW) == [Path("/roots/media/a.part")]
+ at = [(Path("/roots/media/a.0123abcd.part"), _old(ORPHAN_AFTER_SECS))]
+ just_under = [(Path("/roots/media/a.0123abcd.part"), _old(ORPHAN_AFTER_SECS - 1))]
+ assert orphaned_parts(at, set(), NOW) == [Path("/roots/media/a.0123abcd.part")]
assert orphaned_parts(just_under, set(), NOW) == []
@@ -243,9 +245,9 @@ def test_the_janitor_deletes_the_abandoned_and_keeps_the_rest(tmp_path):
one somebody is still writing stay, and a finished file is never a
candidate."""
root = _root(tmp_path, "media")
- old = _aged(root.path / "abandoned.mkv.part", ORPHAN_AFTER_SECS + 60)
- recent = _aged(root.path / "fresh.mkv.part", 30)
- live = _aged(root.path / "sending.mkv.part", ORPHAN_AFTER_SECS * 2)
+ old = _aged(root.path / "abandoned.mkv.0123abcd.part", ORPHAN_AFTER_SECS + 60)
+ recent = _aged(root.path / "fresh.mkv.0123abcd.part", 30)
+ live = _aged(root.path / "sending.mkv.0123abcd.part", ORPHAN_AFTER_SECS * 2)
finished = _aged(root.path / "done.mkv", ORPHAN_AFTER_SECS * 5)
uploads = PartialUploads()
@@ -262,7 +264,7 @@ def test_a_group_that_has_never_uploaded_anything_is_handled(tmp_path):
"""No `partial_uploads` in the context yet — it is created on first use, so
a node that has been up for five minutes has none."""
root = _root(tmp_path, "media")
- old = _aged(root.path / "left.mkv.part", ORPHAN_AFTER_SECS + 1)
+ old = _aged(root.path / "left.mkv.0123abcd.part", ORPHAN_AFTER_SECS + 1)
daemon = _daemon({"g1": {"roots": RootSet(roots=[root])}})
assert daemon._reap_once() == 1
assert not old.exists()
@@ -362,7 +364,7 @@ async def test_an_upload_in_flight_is_known_to_the_reaper(tmp_path):
total_chunks=2))
live = ctx["partial_uploads"].live_paths()
assert len(live) == 1
- assert next(iter(live)).name == "film.mkv.part"
+ assert re.fullmatch(r"film\.mkv\.[0-9a-f]{8}\.part", next(iter(live)).name)
assert next(iter(live)).exists()
@@ -489,3 +491,97 @@ async def test_an_upload_chunk_says_its_slot_is_in_use(tmp_path):
assert _errors(peer) == []
assert slots.leases["up-1"].used is True, (
"the node still believes nobody took this slot up, and will reclaim it")
+
+
+
+# ── never replacing a file ──────────────────────────────────────────────────
+
+async def test_two_members_sending_one_name_at_once_get_two_files(tmp_path):
+ """
+ Both chose a free name at chunk 0, and the free name was the same: the
+ final file did not exist yet, only the first one's `.part`. They then wrote
+ one `.part`, the second truncating the first, and published it twice.
+ """
+ ctx = _group_ctx(tmp_path)
+ alice, bob = _peer(ctx, "alice"), _peer(ctx, "bob")
+ for who, data in ((alice, b"hers-1"), (bob, b"his-1")):
+ await who._do_file_upload(sealed_upload(who, filename="IMG_1234.jpg", data=data,
+ chunk_index=0, total_chunks=2))
+ for who, data in ((alice, b"hers-2"), (bob, b"his-2")):
+ await who._do_file_upload(sealed_upload(who, filename="IMG_1234.jpg", data=data,
+ chunk_index=1, total_chunks=2))
+ assert _errors(alice) == [] and _errors(bob) == []
+
+ root = ctx["roots"].roots[0].path
+ assert (root / "IMG_1234.jpg").read_bytes() == b"hers-1hers-2"
+ assert (root / "IMG_1234 (2).jpg").read_bytes() == b"his-1his-2"
+ assert _acks(bob, ctx)[-1]["stored_as"] == "IMG_1234 (2).jpg"
+
+
+async def test_a_file_that_appears_during_an_upload_is_not_replaced(tmp_path):
+ """The operator copies a file in under the same name while a member's
+ upload is running. The rename at the end used to replace it."""
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ await peer._do_file_upload(sealed_upload(peer, filename="film.mkv", data=b"up-1",
+ chunk_index=0, total_chunks=2))
+ root = ctx["roots"].roots[0].path
+ (root / "film.mkv").write_bytes(b"the operator's")
+ await peer._do_file_upload(sealed_upload(peer, filename="film.mkv", data=b"up-2",
+ chunk_index=1, total_chunks=2))
+
+ assert _errors(peer) == []
+ assert (root / "film.mkv").read_bytes() == b"the operator's"
+ assert (root / "film (2).mkv").read_bytes() == b"up-1up-2"
+ assert _acks(peer, ctx)[-1]["stored_as"] == "film (2).mkv"
+ assert not list(root.glob("*.part")), "the part was left behind"
+
+
+def test_without_hard_links_a_file_is_still_not_replaced(tmp_path, monkeypatch):
+ """FAT, exFAT and some network shares have no hard links."""
+ from meshbay_node import roots as roots_mod
+
+ def no_links(*a, **k):
+ raise PermissionError("operation not permitted")
+ monkeypatch.setattr(roots_mod.os, "link", no_links)
+ (tmp_path / "a.txt").write_bytes(b"there first")
+ part = tmp_path / "a.txt.0123abcd.part"
+ part.write_bytes(b"uploaded")
+
+ name = roots_mod.publish_upload(part, tmp_path, "a.txt", "a.txt")
+ assert name == "a (2).txt"
+ assert (tmp_path / "a.txt").read_bytes() == b"there first"
+ assert (tmp_path / "a (2).txt").read_bytes() == b"uploaded"
+ assert not part.exists()
+
+
+# ── files Explorer acts on by itself ─────────────────────────────────────────
+
+
+@pytest.mark.parametrize("name", ["desktop.ini", "Desktop.INI", "photos.lnk", "site.url",
+ "x.scf", "Docs.library-ms", "s.searchConnector-ms"])
+async def test_a_file_explorer_acts_on_is_refused(tmp_path, name):
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ await peer._do_file_upload(sealed_upload(peer, filename=name, data=b"[x]"))
+ assert [m.get("code") for m in _errors(peer)] == ["file_type_refused"]
+ assert not any(ctx["roots"].roots[0].path.iterdir())
+
+
+async def test_an_ordinary_file_with_a_near_name_is_accepted(tmp_path):
+ ctx = _group_ctx(tmp_path)
+ peer = _peer(ctx)
+ await peer._do_file_upload(sealed_upload(peer, filename="url-notes.txt", data=b"x"))
+ assert _errors(peer) == []
+
+
+
+def test_a_part_the_node_did_not_write_is_never_deleted():
+ """A browser's download in progress, a copy the operator is making: a
+ `.part` without the node's tag is somebody else's, however old."""
+ doomed = orphaned_parts(
+ [(Path("/roots/media/film.mkv.part"), _old(ORPHAN_AFTER_SECS * 10)),
+ (Path("/roots/media/report.pdf.part"), _old(ORPHAN_AFTER_SECS * 10)),
+ (Path("/roots/media/film.mkv.0123abcd.part"), _old(ORPHAN_AFTER_SECS * 10))],
+ live=set(), now=NOW)
+ assert doomed == [Path("/roots/media/film.mkv.0123abcd.part")]
diff --git a/packages/meshbay-node/tests/test_revocation_closes_sessions.py b/packages/meshbay-node/tests/test_revocation_closes_sessions.py
new file mode 100644
index 0000000..7c8e79c
--- /dev/null
+++ b/packages/meshbay-node/tests/test_revocation_closes_sessions.py
@@ -0,0 +1,53 @@
+"""
+A revocation the hub signed closes what it revokes, not only what comes next.
+
+The denylist refuses the next connection. A revoked account's live sessions
+were left open — streaming, downloading, chatting — until they happened to end;
+only a group's revocation closed its sessions.
+"""
+
+import asyncio
+
+import pytest
+from meshbay_node.daemon import NodeDaemon
+from meshbay_node.transport.quic_server import Denylist
+
+
+class _Session:
+ def __init__(self, user_id, group_id):
+ self._user_id = user_id
+ self._group_id = group_id
+ self.closed = False
+
+ async def close(self):
+ self.closed = True
+
+
+def _daemon(sessions):
+ d = NodeDaemon.__new__(NodeDaemon)
+ d._webrtc = type("T", (), {"_sessions": sessions})()
+ return d
+
+
+@pytest.mark.asyncio
+async def test_a_revoked_account_is_disconnected(tmp_path):
+ mallory, alice = _Session("mallory", "g1"), _Session("alice", "g1")
+ phone = _Session("mallory", "g2")
+ d = _daemon({"a": mallory, "b": alice, "c": phone})
+ deny = Denylist(path=tmp_path / "deny.json")
+
+ d._apply_revocation(deny, "user", "mallory")
+ await asyncio.sleep(0.05)
+
+ assert mallory.closed and phone.closed, "every session of the account ends"
+ assert not alice.closed
+ assert deny.is_denied("mallory", "")
+
+
+@pytest.mark.asyncio
+async def test_a_revoked_group_is_still_disconnected(tmp_path):
+ a, b = _Session("alice", "g1"), _Session("alice", "g2")
+ d = _daemon({"a": a, "b": b})
+ d._apply_revocation(Denylist(path=tmp_path / "deny.json"), "group", "g1")
+ await asyncio.sleep(0.05)
+ assert a.closed and not b.closed