summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-08-14 19:35:37 +0200
committerChristophe Besson <cbesson@gmail.com>2026-08-14 19:35:37 +0200
commitc83a4f6ab0c8a83e8679e78427ae60dc29bb2c60 (patch)
treedea71c8e115742beaac5952c8c65481bbc130b07 /packages/meshbay-node
parentee6573c57f721db8550e34e1c1c79c5922c62a4b (diff)
parentd324792d68503109ab99616af6c85ee37045e169 (diff)
downloadmeshbay-c83a4f6ab0c8a83e8679e78427ae60dc29bb2c60.tar.gz
merge: Phase 11.5 security remediation, invite redesign, per-node identity
Brings in the security remediation branch. Three bodies of work, and what they changed about what this project may claim. Phase 11.5 closed the gap between the documents and the code: the unauthenticated node HTTP API and the TCP transport deleted, one handshake shared by the remaining two transports, mutual authentication, structured admin transcripts, upload confinement, group isolation, revocation that reaches nodes. Six critical and seven high findings closed, bounded, or deferred by decision. The invite redesign closed H3 and M3 — the last open High. The hub was the key directory: an inviter fetched the invitee's key from it and wrapped the group key for whatever came back, so a hub answering with its own key was handed the group key by an honest member following the protocol exactly. That lookup is gone. The node holds the group key and wraps it itself, for a key its recipient proves possession of, bound to an account by a one-time code the hub never sees. M3 fell out of the same work: node authority comes from a local roster, never from the hub. Per-node identity cut what remains of C4 down to one operator. A single keypair used to be copied to every node its owner joined; each node now gets its own, so cracking the bundle on one machine yields a key that is a stranger everywhere else — and on that machine, one that unlocks nothing its holder did not already serve. The bundle KDF moved to Argon2id 128 MB, and the hub stopped storing or publishing user keys at all. What this project may now say: the hub cannot read your content unless it ships you malicious client code. T3 remains, accepted (D1), and is what the native client removes. C4 is reduced, not closed, until 13.3. Chat is still plaintext at rest until Phase 15. Draft-v5 §2 states each claim against the adversary it holds against, which is the convention this branch exists to keep. Four defects were found by deploying it and using a browser, none by the test suite: a node going deaf on its hub socket, a token that predated group membership, a client reading values before they were assigned, and identity keys a browser held but never re-read. The lessons are recorded in CLAUDE.md. Tests: 343 across the three packages, plus QE/deploy/e2e.py — register, pair, invite, join, download, stream, second browser, revoke — run against the live deployment on a wiped hub and node.
Diffstat (limited to 'packages/meshbay-node')
-rw-r--r--packages/meshbay-node/src/meshbay_node/bundle_store.py15
-rw-r--r--packages/meshbay-node/src/meshbay_node/config.py59
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py502
-rw-r--r--packages/meshbay-node/src/meshbay_node/hub_client.py95
-rw-r--r--packages/meshbay-node/src/meshbay_node/keystore.py25
-rw-r--r--packages/meshbay-node/src/meshbay_node/roster.py400
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/__init__.py21
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/client.py148
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/http_server.py336
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/quic_client.py96
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/quic_server.py183
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/server.py286
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/tls_cert.py36
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py994
-rw-r--r--packages/meshbay-node/src/meshbay_node/ui/app.py419
-rw-r--r--packages/meshbay-node/tests/test_daemon.py77
-rw-r--r--packages/meshbay-node/tests/test_http_server.py227
-rw-r--r--packages/meshbay-node/tests/test_quic_transport.py13
-rw-r--r--packages/meshbay-node/tests/test_roster_pairing.py744
-rw-r--r--packages/meshbay-node/tests/test_security_regressions.py604
-rw-r--r--packages/meshbay-node/tests/test_transport.py222
-rw-r--r--packages/meshbay-node/tests/test_webrtc_transport.py368
22 files changed, 4076 insertions, 1794 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/bundle_store.py b/packages/meshbay-node/src/meshbay_node/bundle_store.py
index e7c6981..4cf3236 100644
--- a/packages/meshbay-node/src/meshbay_node/bundle_store.py
+++ b/packages/meshbay-node/src/meshbay_node/bundle_store.py
@@ -99,6 +99,21 @@ class BundleStore:
row = await cursor.fetchone()
return row[0] if row else None
+ async def delete_keypair(self, user_id: str) -> bool:
+ """
+ Drop someone's keypair bundle at their own request.
+
+ Backing keys up here is what lets a second browser recover them with the
+ password — and it is also what puts a PBKDF2-protected blob on every node
+ whose group they join (finding C4). Someone who does not need the first
+ should be able to withdraw the second, and not merely stop adding to it.
+ """
+ assert self._db
+ cur = await self._db.execute(
+ "DELETE FROM keypair_bundles WHERE user_id = ?", (user_id,))
+ await self._db.commit()
+ return cur.rowcount > 0
+
async def close(self) -> None:
if self._db:
await self._db.close()
diff --git a/packages/meshbay-node/src/meshbay_node/config.py b/packages/meshbay-node/src/meshbay_node/config.py
index a7a0785..f3752ea 100644
--- a/packages/meshbay-node/src/meshbay_node/config.py
+++ b/packages/meshbay-node/src/meshbay_node/config.py
@@ -26,28 +26,33 @@ url = "https://meshbay.org"
username = "myusername"
[node]
-port = 19000 # TCP+TLS (MNP v1)
-quic_port = 19010 # QUIC (MNP v2)
-http_port = 19001 # HTTP file API (public content)
-ui_port = 18000 # local web UI
+quic_port = 19010 # QUIC (MNP) — LAN, port-forwarded, hub-less direct access
+ui_port = 18000 # local admin UI (127.0.0.1 only)
-# Multiple groups — each with its own directory and ports
+# One-time codes. An invitation waits for someone to read their messages; an
+# operator pairing code is typed during the SSH session that printed it.
+invite_ttl_hours = 168 # 7 days
+pair_ttl_hours = 24
+
+# Browser and native clients reach this node over WebRTC DataChannel via hub
+# signaling — no inbound port to open. QUIC is the optional direct path.
+
+# Multiple groups — each with its own directory
[[groups]]
id = "" # set after joining
name = "My Media"
shared_dir = "/home/user/Media"
-port = 19000
quic_port = 19010
-http_port = 19001
[[groups]]
id = ""
name = "Public Archive"
shared_dir = "/home/user/Archive"
-port = 19002
quic_port = 19012
-http_port = 19003
-visibility = "public"
+visibility = "public" # discoverable on the hub
+# join_policy = "open" # anyone the hub says is a member gets the group key,
+ # with no pairing code. Only for groups where that is
+ # genuinely intended: it means the hub can join too.
[keystore]
# unlock_file = "~/.config/meshbay/unlock.key"
@@ -68,10 +73,13 @@ class HubConfig:
@dataclass
class NodeConfig:
- port: int = 19000
quic_port: int = 19010
- http_port: int = 19001
ui_port: int = 18000
+ # How long a one-time code stays usable. Invitations travel through a human
+ # conversation and are answered days later; operator pairing happens during
+ # the SSH session that printed it.
+ invite_ttl_hours: int = 168 # 7 days
+ pair_ttl_hours: int = 24
@dataclass
@@ -79,10 +87,16 @@ class GroupConfig:
id: str = ""
name: str = ""
shared_dir: str = ""
- visibility: str = "private" # public|private
- port: int = 19000 # TCP+TLS MNP port for this group
+ visibility: str = "private" # public|private — discoverability, not admission
+ # Admission. "invite" (default) means a newcomer needs a one-time pairing code
+ # before the node wraps the group key for them; "open" means the node pins
+ # whoever turns up first (TOFU) and serves them.
+ #
+ # Deliberately read from THIS file and never from the hub: a hub that could
+ # declare a group open would walk into any group it liked. Being findable
+ # (`visibility`) and being open (`join_policy`) are different questions.
+ join_policy: str = "invite" # invite|open
quic_port: int = 19010 # QUIC MNP port
- http_port: int = 19001 # HTTP file API port
@dataclass
@@ -121,10 +135,14 @@ def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config:
cfg.hub.username = hub.get("username", cfg.hub.username)
nd = raw.get("node", {})
- cfg.node.port = nd.get("port", cfg.node.port)
+ # `port` (TCP+TLS) and `http_port` no longer exist — both listeners were removed
+ # in Phase 11.5 (findings C1, C6). Regenerate node.toml with `meshbay-node init`.
cfg.node.quic_port = nd.get("quic_port", cfg.node.quic_port)
- cfg.node.http_port = nd.get("http_port", cfg.node.http_port)
cfg.node.ui_port = nd.get("ui_port", cfg.node.ui_port)
+ cfg.node.invite_ttl_hours = int(
+ nd.get("invite_ttl_hours", cfg.node.invite_ttl_hours))
+ cfg.node.pair_ttl_hours = int(
+ nd.get("pair_ttl_hours", cfg.node.pair_ttl_hours))
# Multi-group: [[groups]] array
if "groups" in raw:
@@ -134,9 +152,8 @@ def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config:
name=g.get("name", ""),
shared_dir=g.get("shared_dir", ""),
visibility=g.get("visibility", "private"),
- port=g.get("port", cfg.node.port),
+ join_policy=g.get("join_policy", "invite"),
quic_port=g.get("quic_port", cfg.node.quic_port),
- http_port=g.get("http_port", cfg.node.http_port),
))
# Back-compat: single [group] section
elif "group" in raw:
@@ -163,8 +180,8 @@ def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config:
cfg.hub.url = url
if user := os.environ.get("MESHBAY_USERNAME"):
cfg.hub.username = user
- if port := os.environ.get("MESHBAY_PORT"):
- cfg.node.port = int(port)
+ if port := os.environ.get("MESHBAY_QUIC_PORT"):
+ cfg.node.quic_port = int(port)
return cfg
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py
index fe12909..58fa99a 100644
--- a/packages/meshbay-node/src/meshbay_node/daemon.py
+++ b/packages/meshbay-node/src/meshbay_node/daemon.py
@@ -8,9 +8,9 @@ Startup sequence:
4. Fetch GEK bundle from hub (if group configured)
5. Start directory indexer (watchdog)
6. Create chat stores (one SQLite DB per group)
- 7. Create WebRTC transport (browser clients via DataChannel)
- 8. Start QUIC+TCP chunk servers (native clients)
- 9. Start HTTP file API (public content)
+ 7. Create WebRTC transport (browser + native clients via DataChannel)
+ 8. Start QUIC chunk server (LAN / port-forwarded / hub-less direct access)
+ 9. (Phase 11.5: the unauthenticated HTTP file API and the TCP+TLS server were removed)
10. Start hub WebSocket (signaling, revocations, WebRTC offers)
11. Start local web UI on node.ui_port (localhost only)
12. Run until SIGINT/SIGTERM
@@ -18,6 +18,9 @@ Startup sequence:
Usage:
meshbay-node # interactive password prompt
meshbay-node --config /path # custom config
+ meshbay-node status # node state + public key (works while stopped)
+ meshbay-node ui # print the local admin UI URL
+ meshbay-node gek-init # initialise the group key (no browser needed)
meshbay-node init # write example config + create keystore
meshbay-node --calibrate-argon2 # benchmark Argon2id, suggest parameters
"""
@@ -26,6 +29,7 @@ import asyncio
import base64
import json
import logging
+import os
import signal
import sys
from pathlib import Path
@@ -41,13 +45,12 @@ from meshbay_node.chat.store import ChatStore
from meshbay_node.config import Config, DEFAULT_CONFIG_PATH, load_config, write_example_config
from meshbay_node.hub_client import HubClient, HubConfig
from meshbay_node.indexer import DirectoryIndexer
-from meshbay_node.keystore import NodeKeys, load_or_create_keystore
+from meshbay_node.keystore import load_or_create_keystore
+from meshbay_node.roster import Roster
from meshbay_node.transport import (
- ChunkServer,
Denylist,
QUIC_AVAILABLE,
WEBRTC_AVAILABLE,
- create_http_app,
)
if QUIC_AVAILABLE:
@@ -103,22 +106,22 @@ class NodeDaemon:
"hub_url": config.hub.url,
"username": config.hub.username,
"groups": [g.name for g in config.groups],
- "node_port": config.node.port,
"quic_port": config.node.quic_port,
"endpoint_hint": None,
"indexes": {},
}
- self._tcp_server: ChunkServer | None = None
self._quic_server = None
self._webrtc = None
- self._denylist = Denylist() if Denylist else None
+ # Persisted so a restart does not silently un-revoke everyone (H4)
+ self._denylist = (
+ Denylist(path=config.data_dir / "denylist.json") if Denylist else None)
self._chat_stores: dict[str, ChatStore] = {}
self._audit_store: AuditStore | None = None
self._bundle_store: BundleStore | None = None
+ self._roster: Roster | None = None
self._indexers: list[DirectoryIndexer] = []
self._tasks: list[asyncio.Task] = []
self._hub: HubClient | None = None
- self._http_servers: list[uvicorn.Server] = []
async def run(self) -> None:
log.info("MeshBay Node starting up")
@@ -133,6 +136,17 @@ class NodeDaemon:
# 2. Start admin UI early (so operator can copy node key before hub login)
self._state["pk_node_ed25519"] = keys.pk_ed25519_b64
self._state["config"] = self._config
+ # Per-run token for the local admin UI (11.5.3). Not a password: it keeps
+ # other local processes and rebound browser pages out of an API that can
+ # re-initialise group keys.
+ ui_token = base64.urlsafe_b64encode(os.urandom(18)).decode().rstrip("=")
+ self._state["ui_token"] = ui_token
+ # Persisted so `meshbay-node ui` can open the browser. Nobody should ever
+ # have to copy a token out of a log or a terminal — that is not a workflow.
+ self._config.data_dir.mkdir(parents=True, exist_ok=True)
+ self._ui_token_file = self._config.data_dir / "ui-token"
+ self._ui_token_file.write_text(ui_token)
+ self._ui_token_file.chmod(0o600)
from meshbay_node.ui import create_ui_app
ui_app = create_ui_app(self._state)
ui_cfg = uvicorn.Config(
@@ -143,7 +157,7 @@ class NodeDaemon:
)
ui_server = uvicorn.Server(ui_cfg)
self._tasks.append(asyncio.create_task(ui_server.serve()))
- log.info("Admin UI at http://localhost:%d", self._config.node.ui_port)
+ log.info("Admin UI ready — open it with: meshbay-node ui")
# 3. Hub connection (Ed25519 auth — retries until node key is linked)
hub_cfg = HubConfig(
@@ -162,6 +176,14 @@ class NodeDaemon:
await self._bundle_store.open()
log.info("Bundle store opened: %s", data_dir / "bundles.db")
+ # 4b. Roster — who this node recognises and which keys are theirs.
+ # Node authority is established here, locally, and never learned from
+ # the hub: a hub that could name the operator's key could install
+ # itself as node administrator.
+ self._roster = Roster(db_path=data_dir / "roster.db")
+ await self._roster.open()
+ await self._roster.purge_expired()
+
# X25519 key material for GEK unwrapping
from cryptography.hazmat.primitives import serialization
sk_x_raw = keys.sk_x25519.private_bytes(
@@ -210,6 +232,10 @@ class NodeDaemon:
"gek": gek,
"shared_root": shared_root,
"index": indexer.index,
+ "visibility": group_cfg.visibility,
+ # Admission policy comes from node.toml, never from the hub:
+ # a hub that could declare a group open would be handed its key.
+ "join_policy": group_cfg.join_policy,
}
if not groups_ctx:
@@ -246,7 +272,11 @@ class NodeDaemon:
groups=groups_ctx,
denylist=denylist,
)
- self._webrtc._ctx["chat_store"] = first.get("chat_store")
+ # No global chat_store here: each group's store lives in
+ # groups_ctx[gid]["chat_store"] and is resolved per session via
+ # _group_ctx(). Assigning the first group's store transport-wide
+ # sent every group's chat to one database and served it back to
+ # members of every other group (finding H1).
self._webrtc._ctx["hub_ws"] = _WsSender(hub)
self._webrtc._ctx["node_user_id"] = session.user_id
self._webrtc._ctx["audit_store"] = self._audit_store
@@ -255,17 +285,27 @@ class NodeDaemon:
self._webrtc._ctx["pk_x25519_raw"] = pk_x_raw
self._webrtc._ctx["pk_x25519_b64"] = keys.pk_x25519_b64
- admin_pk = self._resolve_admin_pk(keys)
+ self._webrtc._ctx["roster"] = self._roster
+ self._webrtc._ctx["invite_ttl"] = (
+ self._config.node.invite_ttl_hours * 3600)
+ admin_pk = self._legacy_admin_pk()
+ paired = await self._roster.has_operator() if self._roster else False
if admin_pk:
self._webrtc._ctx["admin_pk_ed25519"] = admin_pk
- log.info("Admin Ed25519 key pinned for node sovereignty")
+ self._webrtc._ctx["has_admin_authority"] = paired
+ if paired or admin_pk:
+ sources = ([] if not paired else ["paired operator"]) + \
+ ([] if not admin_pk else ["node.toml admin_pk"])
+ log.info("Node authority: %s", " + ".join(sources))
else:
- log.warning("No admin_pk_ed25519 — admin operations disabled")
+ log.warning(
+ "No operator paired — invites and file deletion are "
+ "refused. Run: meshbay-node operator pair")
log.info("WebRTC transport ready")
else:
log.warning("WebRTC not available (aiortc not installed)")
- # 7. QUIC + TCP chunk servers
+ # 7. QUIC chunk server (LAN / port-forwarded / hub-less direct access)
if QUIC_AVAILABLE:
self._quic_server = QuicChunkServer(
sk_node=keys.sk_ed25519,
@@ -282,19 +322,6 @@ class NodeDaemon:
log.info("QUIC server on port %d (%d groups)",
self._config.node.quic_port, len(groups_ctx))
- self._tcp_server = ChunkServer(
- sk_node=keys.sk_ed25519,
- hub_pk_pem=session.hub_pk_pem,
- gek=first["gek"],
- shared_root=first["shared_root"],
- index=first["index"],
- host="0.0.0.0",
- port=self._config.node.port,
- groups=groups_ctx,
- )
- await self._tcp_server.start()
- log.info("TCP+TLS server on port %d", self._config.node.port)
-
# 8. Hub WebSocket (signaling + revocations + WebRTC offers)
async def on_webrtc_offer(sdp, peer_id, ice_candidates):
if not self._webrtc:
@@ -324,8 +351,15 @@ class NodeDaemon:
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)
except Exception as e:
log.warning("Invalid revocation token: %s", e)
@@ -338,37 +372,17 @@ class NodeDaemon:
self._tasks.append(ws_task)
log.info("Hub WS task started")
- # 9. HTTP file API (one per group)
- for gid, gctx in groups_ctx.items():
- group_cfg = next(
- (g for g in self._config.groups if g.id == gid), None)
- if not group_cfg:
- continue
- http_app = create_http_app(
- sk_node=keys.sk_ed25519,
- hub_pk_pem=session.hub_pk_pem,
- shared_root=gctx["shared_root"],
- index=gctx["index"],
- group_id=gid,
- group_name=group_cfg.name,
- gek=gctx.get("gek"),
- )
- http_cfg = uvicorn.Config(
- http_app,
- host="0.0.0.0",
- port=group_cfg.http_port,
- log_level="warning",
- )
- http_server = uvicorn.Server(http_cfg)
- self._http_servers.append(http_server)
- self._tasks.append(asyncio.create_task(http_server.serve()))
- log.info("HTTP API on port %d for group %s",
- group_cfg.http_port, group_cfg.name)
+ # 9. (removed in Phase 11.5) The per-group HTTP file API used to start here.
+ # It served the Mesh Group Index and raw plaintext files on 0.0.0.0 with no
+ # authentication, for private groups too — finding C1. Every client path now
+ # goes through the MNP handshake (JWT + group claim + GEK proof).
# 10. Update admin UI state (UI already running from step 2)
self._state["groups_ctx"] = groups_ctx
self._state["audit_store"] = self._audit_store
self._state["bundle_store"] = self._bundle_store
+ self._state["roster"] = self._roster
+ self._state["node_user_id"] = session.user_id
self._state["webrtc"] = self._webrtc
self._state["hub"] = hub
self._state["pk_x25519_raw"] = pk_x_raw
@@ -379,9 +393,15 @@ class NodeDaemon:
"yes" if self._webrtc else "no",
"yes" if self._quic_server else "no")
- # 11. Initial swarm registration
+ # 11. Initial swarm registration — PUBLIC groups only.
+ # Finding H7: registering every group's hashes hands the hub a content
+ # fingerprint of every private file on the node, which is exactly the
+ # metadata the "hub stores no content metadata" claim rules out. It also
+ # lets anyone confirm whether a known file exists in the network.
endpoint = f"webrtc:{self._config.node.quic_port}"
for gctx in groups_ctx.values():
+ if gctx.get("visibility") != "public":
+ continue
hashes = [e.id for e in gctx["index"].entries]
if hashes:
asyncio.ensure_future(self._register_swarm(hashes, endpoint))
@@ -403,14 +423,28 @@ class NodeDaemon:
return await hub.startup(endpoint_hint=None)
except _httpx.HTTPStatusError as e:
body = e.response.text if hasattr(e.response, 'text') else ''
- if e.response.status_code == 401 and "No node key" in body:
- self._state["status"] = "waiting_for_node_key"
- log.warning(
- "Node key not linked — open admin UI at "
- "http://localhost:%d, copy the key, and paste it in "
- "Settings > Link Node on the hub. Retrying in 30s...",
- self._config.node.ui_port,
- )
+ # Any 401 here needs a human at a browser, and the operator needs
+ # this daemon alive to read its public key out of the local admin
+ # UI. Exiting would take that UI down and strand them — which is
+ # exactly what happened when a node was started before its owner
+ # had registered.
+ if e.response.status_code == 401:
+ if "No node key" in body:
+ self._state["status"] = "waiting_for_node_key"
+ log.warning(
+ "Node key not linked. Open the admin UI, copy this "
+ "node's key, and paste it in Settings > Link Node on "
+ "%s. Retrying in 30s...",
+ self._config.hub.url,
+ )
+ else:
+ self._state["status"] = "waiting_for_account"
+ log.warning(
+ "Hub rejected the node credentials for user %r. "
+ "Register that account on %s first, then link this "
+ "node's key. Retrying in 30s...",
+ self._config.hub.username, self._config.hub.url,
+ )
await asyncio.sleep(30)
else:
raise
@@ -448,21 +482,25 @@ class NodeDaemon:
log.warning("No unwrappable GEK bundle found for group %s", group_id[:8])
return None
- def _resolve_admin_pk(self, keys: NodeKeys) -> Ed25519PublicKey | None:
- """Resolve the admin Ed25519 public key: config → auto-pin from node keystore."""
- if self._config.admin_pk_ed25519:
- try:
- raw = base64.b64decode(self._config.admin_pk_ed25519)
- return Ed25519PublicKey.from_public_bytes(raw)
- except Exception as e:
- log.error("Invalid admin_pk_ed25519 in config: %s", e)
- return None
+ def _legacy_admin_pk(self) -> Ed25519PublicKey | None:
+ """
+ The pre-roster way of naming the operator: `admin_pk_ed25519` in node.toml.
- pk = keys.sk_ed25519.public_key()
- from meshbay_common.crypto import pk_to_b64
- pk_b64 = pk_to_b64(pk)
- log.info("Auto-pinning admin key from node keystore: %s", pk_b64[:16])
- return pk
+ Still honoured so a deployment configured that way keeps working, but no
+ longer the only path — and the auto-pin that used to stand in for it is
+ gone. It pinned the node's *keystore* key while the browser signed with the
+ user's *identity* key, so admin operations failed closed with a signature
+ error that looked like a bug elsewhere (finding M3). An operator now pairs
+ a browser with `meshbay-node operator pair`.
+ """
+ if not self._config.admin_pk_ed25519:
+ return None
+ try:
+ raw = base64.b64decode(self._config.admin_pk_ed25519)
+ return Ed25519PublicKey.from_public_bytes(raw)
+ except Exception as e:
+ log.error("Invalid admin_pk_ed25519 in config: %s", e)
+ return None
async def _on_index_change(self, indexer: DirectoryIndexer) -> None:
"""Called when a DirectoryIndexer detects file changes."""
@@ -498,13 +536,25 @@ class NodeDaemon:
if pushed:
log.info("Index pushed to %d WebRTC peers", pushed)
- # 11.9 — Register file hashes with hub swarm table
- if self._hub and self._state.get("endpoint_hint"):
+ # 11.9 — Register file hashes with hub swarm table (public groups only, H7)
+ group_cfg = next(
+ (g for g in self._config.groups if g.id == group_id), None)
+ if (self._hub and self._state.get("endpoint_hint")
+ and group_cfg and group_cfg.visibility == "public"):
hashes = [e.id for e in idx.entries]
if hashes:
endpoint = f"webrtc:{self._config.node.quic_port}"
asyncio.ensure_future(self._register_swarm(hashes, endpoint))
+ 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:
+ return
+ for session in list(self._webrtc._sessions.values()):
+ if session._group_id == group_id:
+ asyncio.ensure_future(session.close())
+ log.info("Dropped session for revoked group %s", group_id[:8])
+
async def _register_swarm(self, hashes: list[str], endpoint: str) -> None:
try:
n = await self._hub.register_swarm(hashes, endpoint)
@@ -533,6 +583,9 @@ class NodeDaemon:
if self._bundle_store:
await self._bundle_store.close()
+ if self._roster:
+ await self._roster.close()
+
for store in self._chat_stores.values():
await store.close()
@@ -541,15 +594,68 @@ class NodeDaemon:
if self._quic_server:
await self._quic_server.stop()
- if self._tcp_server:
- await self._tcp_server.stop()
- for server in self._http_servers:
- server.should_exit = True
+ token_file = getattr(self, "_ui_token_file", None)
+ if token_file is not None:
+ token_file.unlink(missing_ok=True)
log.info("Node stopped")
+# ── CLI helpers ───────────────────────────────────────────────────────────────
+
+def _daemon_api(cfg: Config, path: str, method: str = "GET",
+ timeout: int = 30) -> dict:
+ """
+ Call the daemon's loopback API.
+
+ The daemon owns the roster, the hub session and the live group contexts, so
+ the CLI asks it to act rather than opening its databases behind its back. It
+ also means every operator action goes through the same authorization as the
+ admin UI (the per-run session token, 11.5.3).
+ """
+ import json as _json
+ import urllib.error
+ import urllib.parse
+ import urllib.request
+
+ token_file = cfg.data_dir / "ui-token"
+ if not token_file.exists():
+ print("Node is not running — start it with: meshbay-node")
+ sys.exit(1)
+
+ sep = "&" if "?" in path else "?"
+ url = (f"http://127.0.0.1:{cfg.node.ui_port}{path}"
+ f"{sep}t={token_file.read_text().strip()}")
+ try:
+ req = urllib.request.Request(url, method=method)
+ with urllib.request.urlopen(req, timeout=timeout) as r:
+ return _json.loads(r.read())
+ except urllib.error.HTTPError as e:
+ body = e.read().decode()[:300]
+ try:
+ detail = _json.loads(body).get("error", body)
+ except Exception:
+ detail = body
+ print(f"failed: {detail}")
+ sys.exit(1)
+ except Exception as e:
+ print(f"failed: {e}")
+ sys.exit(1)
+
+
+def _resolve_group(cfg: Config, group: str | None) -> str:
+ """The group argument, or the only configured one."""
+ if group:
+ return group
+ configured = [g.id for g in cfg.groups if g.id]
+ if len(configured) == 1:
+ return configured[0]
+ print("--group is required (several groups configured)"
+ if configured else "no group configured in node.toml")
+ sys.exit(1)
+
+
# ── Entry point ───────────────────────────────────────────────────────────────
def main() -> None:
@@ -557,16 +663,28 @@ def main() -> None:
parser = argparse.ArgumentParser(description="MeshBay Node daemon")
parser.add_argument("command", nargs="?",
- choices=["init", "calibrate-argon2"],
- help="init: write example config | calibrate-argon2: benchmark")
+ choices=["init", "status", "ui", "gek-init", "operator",
+ "member", "calibrate-argon2"],
+ help="init: write example config | status: node state and keys "
+ "| ui: print the admin UI URL | operator pair: pair a "
+ "browser with this node | member list|invite|revoke|unpin "
+ "| calibrate-argon2: benchmark")
+ parser.add_argument("subcommand", nargs="?",
+ help="'pair' for operator; list|invite|revoke|unpin for member")
+ parser.add_argument("target", nargs="?",
+ help="username, for member invite|revoke|unpin")
parser.add_argument("--config", type=Path, default=None,
help="Config file path")
+ parser.add_argument("--group", default=None,
+ help="group id (optional if only one is configured)")
parser.add_argument("--log-level", default="INFO",
choices=["DEBUG", "INFO", "WARNING", "ERROR"])
args = parser.parse_args()
+ # Query commands print a report; library logging would interleave with it.
+ quiet = args.command in ("status", "ui", "gek-init", "operator", "member")
logging.basicConfig(
- level=getattr(logging, args.log_level),
+ level=logging.ERROR if quiet else getattr(logging, args.log_level),
format="%(asctime)s %(levelname)-8s %(name)s: %(message)s",
)
@@ -579,6 +697,216 @@ def main() -> None:
calibrate_argon2()
return
+ if args.command == "status":
+ import json as _json
+ import urllib.request
+
+ cfg = load_config(args.config or DEFAULT_CONFIG_PATH)
+ print(f"hub {cfg.hub.url} (user {cfg.hub.username or '—'})")
+
+ # Read straight from the keystore: the operator needs this key to link the
+ # node, and that happens before the daemon can ever stay running.
+ try:
+ keys = load_or_create_keystore(
+ path=cfg.keystore.path, unlock_file=cfg.keystore.unlock_file)
+ print(f"node key {keys.pk_ed25519_b64}")
+ except Exception as e:
+ print(f"node key <keystore locked: {e}>")
+
+ token_file = cfg.data_dir / "ui-token"
+ live = None
+ if token_file.exists():
+ try:
+ url = (f"http://127.0.0.1:{cfg.node.ui_port}"
+ f"/api/status?t={token_file.read_text().strip()}")
+ with urllib.request.urlopen(url, timeout=3) as r:
+ live = _json.loads(r.read())
+ except Exception:
+ live = None
+
+ if live:
+ print(f"daemon running — {live.get('status')}")
+ print(f"node_id {live.get('endpoint_hint') or '—'}")
+ print(f"groups {live.get('group_count', 0)}"
+ f" files {live.get('total_files', 0)}"
+ f" peers {live.get('webrtc_peers', 0)}")
+ print(f"admin UI meshbay-node ui")
+ else:
+ print("daemon not running")
+
+ print(f"config {DEFAULT_CONFIG_PATH}")
+ if not cfg.groups:
+ print("groups none configured — create a group on the hub, then add")
+ print(" a [[groups]] entry with its id and shared_dir")
+ else:
+ for g in cfg.groups:
+ print(f" group {g.name} [{g.visibility}] {g.id or '<no id>'}")
+ print(f" {g.shared_dir or '<no shared_dir>'}")
+ # Node authority: the roster is the source of truth, node.toml the legacy
+ # form. Read the DB directly so this reports correctly while the daemon is
+ # stopped — the state an operator is most often in when checking.
+ import asyncio as _asyncio
+
+ from meshbay_node.roster import Roster as _Roster
+
+ async def _read_roster() -> tuple[list, int]:
+ r = _Roster(db_path=cfg.data_dir / "roster.db")
+ await r.open()
+ try:
+ return (await r.list_members()), len(await r.list_invites())
+ finally:
+ await r.close()
+
+ try:
+ members, pending = _asyncio.run(_read_roster())
+ except Exception as e:
+ members, pending = [], 0
+ print(f"roster <unreadable: {e}>")
+
+ operators = [m for m in members if m["role"] == "operator"
+ and m["status"] == "active"]
+ if operators:
+ for op in operators:
+ print(f"operator {op.get('username') or op['user_id'][:8]}"
+ f" key {(op.get('pk_ed25519') or '')[:16]}…"
+ f" paired {op.get('pinned_at', '?')}")
+ elif cfg.admin_pk_ed25519:
+ print("operator node.toml admin_pk_ed25519 (legacy)")
+ print(" run `meshbay-node operator pair` to replace it")
+ else:
+ print("operator NONE PAIRED — file deletion and member invites are")
+ print(" refused. Run: meshbay-node operator pair")
+ if pending:
+ print(f"invites {pending} pending code(s)")
+ return
+
+ if args.command == "member":
+ cfg = load_config(args.config or DEFAULT_CONFIG_PATH)
+ sub = args.subcommand or "list"
+
+ if sub == "list":
+ group = args.group or ""
+ out = _daemon_api(
+ cfg, f"/api/roster?group_id={group}" if group else "/api/roster")
+ identities = {i["user_id"]: i for i in out.get("identities", [])}
+
+ members = out.get("members", [])
+ if not members:
+ print("no members admitted yet")
+ print("invite someone: meshbay-node member invite <username>")
+ for m in members:
+ ident = identities.get(m["user_id"], {})
+ scope = m["group_id"][:8] if m["group_id"] else "node-wide"
+ print(f"{(ident.get('username') or m['user_id'])[:20]:20} "
+ f"{m['role']:9} {m['status']:8} {scope:10} "
+ f"pinned {ident.get('pinned_at', '?')} "
+ f"({ident.get('pinned_via', '?')})")
+
+ invites = out.get("invites", [])
+ if invites:
+ print()
+ for i in invites:
+ print(f"pending invite user {i['user_id'][:12]} "
+ f"group {(i['group_id'] or 'node-wide')[:8]} "
+ f"expires {i['expires_at']}")
+ return
+
+ if not args.target:
+ print(f"usage: meshbay-node member {sub} <username>")
+ sys.exit(1)
+
+ if sub == "invite":
+ group_id = _resolve_group(cfg, args.group)
+ out = _daemon_api(
+ cfg, f"/api/groups/{group_id}/invites?username={args.target}",
+ method="POST")
+ from meshbay_node.roster import write_code_file
+ path = write_code_file(cfg.data_dir, out["code"],
+ out.get("expires_at", ""), name="invite-code")
+ print(f"INVITATION CODE {out['code']}")
+ print(f"valid until {out.get('expires_at', '?')}")
+ print()
+ print(f"Send it to {args.target} however you normally talk. It works")
+ print("once, for that account only, and never passes through the hub.")
+ print("They enter it the first time they open the group — you do not")
+ print("need to be online then.")
+ print()
+ print(f"also written to {path}")
+ return
+
+ # revoke and unpin both name a person; the daemon resolves the account.
+ # It tries its own roster first and falls back to the hub, so a node that
+ # pinned someone before invitations carried a name is still manageable.
+ match = _daemon_api(cfg, f"/api/resolve?username={args.target}")
+
+ if sub == "revoke":
+ group_id = _resolve_group(cfg, args.group)
+ out = _daemon_api(
+ cfg, f"/api/members/{match['user_id']}/revoke?group_id={group_id}",
+ method="POST")
+ print(f"{args.target} revoked from {group_id[:8]}")
+ print("They stop receiving the group key on their next connection.")
+ print("They still hold the current one — rotate it:")
+ print(f" meshbay-node gek-init --group {group_id}")
+ return
+
+ if sub == "unpin":
+ _daemon_api(cfg, f"/api/members/{match['user_id']}/unpin", method="POST")
+ print(f"{args.target} unpinned — they can pair again with a new key")
+ print(f"issue a code: meshbay-node member invite {args.target}")
+ return
+
+ print("usage: meshbay-node member list|invite|revoke|unpin")
+ sys.exit(1)
+
+ if args.command == "gek-init":
+ cfg = load_config(args.config or DEFAULT_CONFIG_PATH)
+ group_id = _resolve_group(cfg, args.group)
+ out = _daemon_api(cfg, f"/api/groups/{group_id}/gek",
+ method="POST", timeout=60)
+
+ print(f"GEK ready for {group_id}")
+ print(f" {out.get('authorized_members', 0)} authorized member(s) — each "
+ f"receives the key on connect")
+ for err in out.get("errors") or []:
+ print(f" ! {err}")
+ return
+
+ if args.command == "operator":
+ if args.subcommand != "pair":
+ print("usage: meshbay-node operator pair")
+ sys.exit(1)
+
+ cfg = load_config(args.config or DEFAULT_CONFIG_PATH)
+ out = _daemon_api(cfg, "/api/operator/pair", method="POST")
+
+ from meshbay_node.roster import write_code_file
+ path = write_code_file(cfg.data_dir, out["code"], out.get("expires_at", ""))
+
+ print(f"PAIRING CODE {out['code']}")
+ print(f"valid until {out.get('expires_at', '?')}")
+ print()
+ print("Sign in to the web app as this node's operator, open one of your")
+ print("groups, go to the Members tab and enter the code there.")
+ print("It works once, for that account only, and authorizes invites and")
+ print("file deletion from that browser.")
+ print()
+ print(f"also written to {path}")
+ return
+
+ if args.command == "ui":
+ cfg = load_config(args.config or DEFAULT_CONFIG_PATH)
+ token_file = cfg.data_dir / "ui-token"
+ if not token_file.exists():
+ print("Node does not appear to be running — start it with: meshbay-node")
+ sys.exit(1)
+ print(f"http://127.0.0.1:{cfg.node.ui_port}"
+ f"/?t={token_file.read_text().strip()}")
+ print()
+ print("The UI listens on loopback only. From another machine:")
+ print(f" ssh -L {cfg.node.ui_port}:127.0.0.1:{cfg.node.ui_port} <this-host>")
+ return
+
cfg = load_config(args.config or DEFAULT_CONFIG_PATH)
if not cfg.hub.username:
print("Error: hub.username not set in config. Run: meshbay-node init")
diff --git a/packages/meshbay-node/src/meshbay_node/hub_client.py b/packages/meshbay-node/src/meshbay_node/hub_client.py
index 432af0a..d8417e3 100644
--- a/packages/meshbay-node/src/meshbay_node/hub_client.py
+++ b/packages/meshbay-node/src/meshbay_node/hub_client.py
@@ -13,6 +13,7 @@ No auth_key or password is ever stored on or transmitted from the node.
The hub issues a node-scoped JWT that cannot manage group membership.
"""
+import asyncio
import base64
import json
import logging
@@ -127,8 +128,8 @@ class HubClient:
access_token = data["access_token"]
decoded = jwt.decode(access_token, hub_pk_pem, algorithms=["EdDSA"])
- assert decoded["pk_user"] == self._keys.pk_ed25519_b64, \
- "Hub returned token for wrong public key"
+ # No pk_user claim to check any more: tokens carry no key. What binds this
+ # token to this node is the Ed25519 challenge it was issued against.
assert "jti" in decoded, "Hub token missing jti — hub is outdated"
assert decoded.get("scope") == "node", \
"Expected node-scoped token"
@@ -163,9 +164,19 @@ class HubClient:
raise RuntimeError("Not logged in")
await self.ensure_fresh_token()
+ # Proof of possession of the node key (M8) — same domain-separated shape
+ # as node_auth, so a signature for one can never satisfy the other.
+ timestamp = int(time.time())
+ message = (f"meshbay:node_announce:{self._session.user_id}:"
+ f"{self._keys.pk_ed25519_b64}:{timestamp}").encode()
+ signature = base64.b64encode(
+ self._keys.sk_ed25519.sign(message)).decode()
+
r = await self._http.post("/v1/nodes/announce", json={
"pk_node": self._keys.pk_ed25519_b64,
"endpoint_hint": endpoint_hint,
+ "timestamp": timestamp,
+ "signature": signature,
}, headers=self._session.auth_headers)
r.raise_for_status()
node_id = r.json()["node_id"]
@@ -219,9 +230,19 @@ class HubClient:
hub_url = self._session.hub_url.replace("https://", "wss://").replace("http://", "ws://")
ws_url = f"{hub_url}/v1/nodes/ws"
+ # Offers are handled off the read loop (see below), so keep a handle on
+ # the tasks to avoid them being garbage-collected mid-negotiation.
+ pending: set[asyncio.Task] = set()
+
while True:
try:
- async with websockets.connect(ws_url) as ws:
+ # Explicit keepalive: this connection is how a node stays visible
+ # to the hub, and a silently half-open socket looks exactly like a
+ # working one until someone notices the node has vanished.
+ async with websockets.connect(
+ ws_url, ping_interval=20, ping_timeout=20, close_timeout=5,
+ open_timeout=15,
+ ) as ws:
auth_msg = {
"type": "auth",
"token": self._session.access_token,
@@ -230,10 +251,20 @@ class HubClient:
if group_ids:
auth_msg["group_ids"] = group_ids
await ws.send(json.dumps(auth_msg))
- auth_resp = json.loads(await ws.recv())
+ # Bounded: a hub that accepts the socket and then says nothing
+ # — which is what it does for a few seconds while restarting —
+ # would otherwise park this task here forever, with the node
+ # running, silent, and invisible to everyone.
+ auth_resp = json.loads(
+ await asyncio.wait_for(ws.recv(), timeout=15))
if auth_resp.get("type") != "auth_ok":
- log.error("WS auth failed: %s", auth_resp)
- return
+ # Not fatal: the token may simply have expired while we
+ # were disconnected. Refresh on the next pass rather than
+ # ending the task, which used to strand the node for good.
+ log.warning("WS auth refused: %s — retrying in 5s", auth_resp)
+ await asyncio.sleep(5)
+ await self.ensure_fresh_token()
+ continue
self._ws = ws
log.info("Hub WS connected")
@@ -250,28 +281,60 @@ class HubClient:
on_revocation(msg.get("token", ""))
elif mtype == "webrtc_offer" and on_webrtc_offer:
- answer = await on_webrtc_offer(
- msg["sdp"], msg["peer_id"],
- msg.get("ice_candidates", []))
- if answer:
- await ws.send(json.dumps({
- "type": "webrtc_answer",
- "peer_id": msg["peer_id"],
- "sdp": answer[0],
- "ice_candidates": answer[1],
- }))
+ # Answered off the read loop on purpose. Awaiting the
+ # handler here meant one slow negotiation stopped the
+ # node reading this socket at all: no pings answered,
+ # no close frame noticed, no further offers served. A
+ # client that gave up mid-ICE left the node in
+ # CLOSE-WAIT, still running but invisible to the hub
+ # and unreachable by everyone, until it was restarted.
+ task = asyncio.create_task(
+ self._answer_offer(ws, on_webrtc_offer, msg))
+ pending.add(task)
+ task.add_done_callback(pending.discard)
elif mtype == "pong":
pass
except asyncio.CancelledError:
+ for task in pending:
+ task.cancel()
raise
except Exception as e:
log.warning("Hub WS disconnected: %s — reconnecting in 5s", e)
await asyncio.sleep(5)
+ else:
+ # A clean close ends the `async for` without raising. Say so, so a
+ # node that quietly stopped being reachable leaves a trace.
+ log.warning("Hub WS closed by the hub — reconnecting in 5s")
+ await asyncio.sleep(5)
finally:
self._ws = None
+ async def _answer_offer(self, ws, on_webrtc_offer, msg: dict) -> None:
+ """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", []))
+ except Exception as e:
+ log.warning("WebRTC offer from %s failed: %s",
+ str(msg.get("peer_id"))[:8], e)
+ return
+ if not answer:
+ return
+ try:
+ await ws.send(json.dumps({
+ "type": "webrtc_answer",
+ "peer_id": msg["peer_id"],
+ "sdp": answer[0],
+ "ice_candidates": answer[1],
+ }))
+ except Exception as e:
+ # The socket may have gone while we were negotiating; the client will
+ # retry, and the read loop is reconnecting.
+ log.warning("Could not deliver WebRTC answer to %s: %s",
+ str(msg.get("peer_id"))[:8], e)
+
# ── Swarm registration ─────────────────────────────────────────────────
async def register_swarm(self, content_hashes: list[str], endpoint: str) -> int:
diff --git a/packages/meshbay-node/src/meshbay_node/keystore.py b/packages/meshbay-node/src/meshbay_node/keystore.py
index 3777af0..59fc719 100644
--- a/packages/meshbay-node/src/meshbay_node/keystore.py
+++ b/packages/meshbay-node/src/meshbay_node/keystore.py
@@ -39,6 +39,12 @@ from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey
from meshbay_common.crypto import (
+ ARGON2_ITERATIONS,
+ ARGON2_LANES,
+ ARGON2_MEMORY_COST,
+ LEGACY_ARGON2_ITERATIONS,
+ LEGACY_ARGON2_LANES,
+ LEGACY_ARGON2_MEMORY_COST,
decrypt_keystore,
derive_keystore_key,
encrypt_keystore,
@@ -173,7 +179,18 @@ def load_keystore(
tag = base64.b64decode(envelope["tag_b64"])
ct = base64.b64decode(envelope["ciphertext_b64"])
- aes_key = derive_keystore_key(pwd, salt)
+ # Envelopes written before M2 carry no parameters and used the 64 MB profile.
+ params = envelope.get("argon2", {
+ "iterations": LEGACY_ARGON2_ITERATIONS,
+ "memory_cost": LEGACY_ARGON2_MEMORY_COST,
+ "lanes": LEGACY_ARGON2_LANES,
+ })
+ aes_key = derive_keystore_key(
+ pwd, salt,
+ iterations=params.get("iterations"),
+ memory_cost=params.get("memory_cost"),
+ lanes=params.get("lanes"),
+ )
try:
plaintext = decrypt_keystore(iv, ct, tag, aes_key)
except Exception:
@@ -205,6 +222,12 @@ def _write_keystore(path: Path, keys: NodeKeys, password: str) -> None:
envelope = {
"version": KEYSTORE_VERSION,
"argon2_salt_b64": base64.b64encode(salt).decode(),
+ # Recorded so parameters can be raised later without orphaning this file.
+ "argon2": {
+ "iterations": ARGON2_ITERATIONS,
+ "memory_cost": ARGON2_MEMORY_COST,
+ "lanes": ARGON2_LANES,
+ },
"iv_b64": base64.b64encode(iv).decode(),
"tag_b64": base64.b64encode(tag).decode(),
"ciphertext_b64": base64.b64encode(ct).decode(),
diff --git a/packages/meshbay-node/src/meshbay_node/roster.py b/packages/meshbay-node/src/meshbay_node/roster.py
new file mode 100644
index 0000000..6bda56b
--- /dev/null
+++ b/packages/meshbay-node/src/meshbay_node/roster.py
@@ -0,0 +1,400 @@
+"""
+Node roster — who this node recognises, and which keys are theirs.
+
+The node keeps its own answer to "may this person have the group key", derived from
+what the operator authorized locally. It is deliberately NOT derived from the hub:
+the hub decides group membership, and a hub that invents an account and mints a
+token for it would otherwise collect the GEK on connect. Hub membership is an input
+to the decision; it is not the decision.
+
+Three tables:
+
+ identities — one row per person, not per group. Someone paired for one group
+ needs no code for the next one on the same node.
+ members — role and status per (group, user).
+ invites — one-time pairing codes, stored as a hash. The code itself exists
+ only in the operator's hands and the invitee's.
+
+The code is what binds a public key to an account without asking the hub
+(finding H3). See `docs/invite-pairing-v1.md`.
+"""
+
+from __future__ import annotations
+
+import hashlib
+import logging
+import os
+import secrets
+from datetime import datetime, timedelta, timezone
+from pathlib import Path
+
+import aiosqlite
+
+log = logging.getLogger(__name__)
+
+# Crockford base32 without I, L, O and U: no character pair a human can confuse
+# when reading a code aloud or typing it from a phone screen.
+_ALPHABET = "0123456789ABCDEFGHJKMNPQRSTVWXYZ"
+CODE_LEN = 8 # 8 × 5 bits = 40 bits of entropy
+
+# Two different rhythms, so two different lifetimes.
+#
+# An invitation crosses a human conversation: it is sent by mail or message and
+# answered whenever the other person next looks. A day is not enough — the code
+# dies over a weekend and someone has to be at a browser, with the node online, to
+# issue another one.
+#
+# Operator pairing crosses an SSH session: the code is printed and typed minutes
+# later. There is no reason for it to outlive the sitting.
+#
+# The longer window costs little: a code is single use, bound to one account,
+# never seen by the hub, and 40 bits do not fall to guessing in a week against the
+# node-wide lockout.
+DEFAULT_INVITE_TTL = 7 * 24 * 3600 # seconds — member invitations
+DEFAULT_PAIR_TTL = 24 * 3600 # seconds — operator pairing
+
+_SCHEMA = """\
+CREATE TABLE IF NOT EXISTS identities (
+ user_id TEXT PRIMARY KEY,
+ username TEXT NOT NULL,
+ pk_ed25519 TEXT NOT NULL,
+ pk_x25519 TEXT NOT NULL,
+ pinned_at TEXT NOT NULL,
+ pinned_via TEXT NOT NULL
+);
+
+CREATE TABLE IF NOT EXISTS members (
+ group_id TEXT NOT NULL,
+ user_id TEXT NOT NULL,
+ role TEXT NOT NULL,
+ status TEXT NOT NULL,
+ approved_by TEXT NOT NULL,
+ approved_at TEXT NOT NULL,
+ PRIMARY KEY (group_id, user_id)
+);
+
+CREATE TABLE IF NOT EXISTS invites (
+ code_hash TEXT PRIMARY KEY,
+ group_id TEXT NOT NULL,
+ user_id TEXT NOT NULL,
+ username TEXT NOT NULL DEFAULT '',
+ role TEXT NOT NULL,
+ created_by TEXT NOT NULL,
+ created_at TEXT NOT NULL,
+ expires_at TEXT NOT NULL,
+ used_at TEXT
+);
+"""
+
+
+def generate_code() -> str:
+ """A fresh pairing code, formatted for a human to read out: XXXX-XXXX."""
+ raw = "".join(secrets.choice(_ALPHABET) for _ in range(CODE_LEN))
+ return f"{raw[:4]}-{raw[4:]}"
+
+
+def normalize_code(code: str) -> str:
+ """
+ Fold what a human typed onto what was generated.
+
+ Crockford's rules: case-insensitive, dashes and spaces are decoration, and the
+ excluded letters map onto the digits they resemble. Someone reading a code over
+ the phone should not be able to get it wrong in a way we could have absorbed.
+ """
+ out = []
+ for ch in code.upper():
+ if ch in "- \t":
+ continue
+ if ch in "IL":
+ out.append("1")
+ elif ch == "O":
+ out.append("0")
+ elif ch == "U":
+ out.append("V")
+ else:
+ out.append(ch)
+ return "".join(out)
+
+
+def hash_code(code: str) -> str:
+ """
+ Store codes hashed: a stolen roster DB must not yield usable invitations.
+
+ SHA-256 rather than a password KDF on purpose — the input is 40 bits of
+ uniformly random secret, not a human-chosen string, so there is nothing for a
+ slow hash to defend.
+ """
+ return hashlib.sha256(normalize_code(code).encode()).hexdigest()
+
+
+def _now() -> str:
+ return datetime.now(timezone.utc).isoformat(timespec="seconds")
+
+
+class Roster:
+ def __init__(self, db_path: Path):
+ self._db_path = db_path
+ self._db: aiosqlite.Connection | None = None
+
+ async def open(self) -> None:
+ self._db_path.parent.mkdir(parents=True, exist_ok=True)
+ self._db = await aiosqlite.connect(str(self._db_path))
+ self._db.row_factory = aiosqlite.Row
+ # WAL: the CLI writes invites (`operator pair`) while the daemon reads them.
+ await self._db.execute("PRAGMA journal_mode=WAL")
+ await self._db.executescript(_SCHEMA)
+ # invites.username was added after the first deployments: the name is what
+ # the operator types, and it cannot be recovered from the JWT because the
+ # hub does not put one there. CREATE TABLE IF NOT EXISTS will not add a
+ # column to a table that already exists.
+ async with self._db.execute("PRAGMA table_info(invites)") as cur:
+ columns = {r[1] for r in await cur.fetchall()}
+ if "username" not in columns:
+ await self._db.execute(
+ "ALTER TABLE invites ADD COLUMN username TEXT NOT NULL DEFAULT ''")
+ await self._db.commit()
+
+ async def close(self) -> None:
+ if self._db:
+ await self._db.close()
+ self._db = None
+
+ # ── Identities ───────────────────────────────────────────────────────────
+
+ async def pin_identity(
+ self,
+ user_id: str,
+ username: str,
+ pk_ed25519: str,
+ pk_x25519: str,
+ via: str,
+ ) -> None:
+ assert self._db
+ await self._db.execute(
+ "INSERT OR REPLACE INTO identities "
+ "(user_id, username, pk_ed25519, pk_x25519, pinned_at, pinned_via) "
+ "VALUES (?, ?, ?, ?, ?, ?)",
+ (user_id, username, pk_ed25519, pk_x25519, _now(), via),
+ )
+ await self._db.commit()
+
+ async def get_identity(self, user_id: str) -> dict | None:
+ assert self._db
+ async with self._db.execute(
+ "SELECT * FROM identities WHERE user_id = ?", (user_id,)
+ ) as cur:
+ row = await cur.fetchone()
+ return dict(row) if row else None
+
+ async def unpin(self, user_id: str) -> bool:
+ assert self._db
+ cur = await self._db.execute(
+ "DELETE FROM identities WHERE user_id = ?", (user_id,))
+ await self._db.commit()
+ return cur.rowcount > 0
+
+ async def list_identities(self) -> list[dict]:
+ assert self._db
+ async with self._db.execute(
+ "SELECT * FROM identities ORDER BY pinned_at"
+ ) as cur:
+ return [dict(r) for r in await cur.fetchall()]
+
+ # ── Authority ────────────────────────────────────────────────────────────
+
+ async def operator_pks(self) -> list[str]:
+ """
+ Base64 Ed25519 keys allowed to authorize admin operations on this node.
+
+ Read fresh on every check rather than cached: an unpin must take effect at
+ once, and this runs only on admin operations, which are rare.
+ """
+ assert self._db
+ async with self._db.execute(
+ "SELECT i.pk_ed25519 FROM identities i "
+ "JOIN members m ON m.user_id = i.user_id "
+ "WHERE m.role = 'operator' AND m.status = 'active'"
+ ) as cur:
+ return [r["pk_ed25519"] for r in await cur.fetchall()]
+
+ async def has_operator(self) -> bool:
+ return bool(await self.operator_pks())
+
+ async def is_authorized(self, group_id: str, user_id: str) -> bool:
+ """
+ May this person be handed the group key?
+
+ The node's own answer, not the hub's. Hub membership is what lets someone
+ reach the node; this is what decides whether the key is wrapped for them —
+ otherwise a hub that invents an account and mints a token for it would be
+ served the GEK on connect.
+
+ An operator is authorized for every group this node hosts: their authority
+ is node-wide and is recorded with an empty group_id.
+ """
+ assert self._db
+ async with self._db.execute(
+ "SELECT 1 FROM members WHERE user_id = ? AND status = 'active' "
+ "AND (group_id = ? OR (group_id = '' AND role = 'operator')) LIMIT 1",
+ (user_id, group_id),
+ ) as cur:
+ return await cur.fetchone() is not None
+
+ # ── Members ──────────────────────────────────────────────────────────────
+
+ async def set_member(
+ self,
+ group_id: str,
+ user_id: str,
+ role: str,
+ status: str,
+ approved_by: str,
+ ) -> None:
+ assert self._db
+ await self._db.execute(
+ "INSERT OR REPLACE INTO members "
+ "(group_id, user_id, role, status, approved_by, approved_at) "
+ "VALUES (?, ?, ?, ?, ?, ?)",
+ (group_id, user_id, role, status, approved_by, _now()),
+ )
+ await self._db.commit()
+
+ async def get_member(self, group_id: str, user_id: str) -> dict | None:
+ assert self._db
+ async with self._db.execute(
+ "SELECT * FROM members WHERE group_id = ? AND user_id = ?",
+ (group_id, user_id),
+ ) as cur:
+ row = await cur.fetchone()
+ return dict(row) if row else None
+
+ async def list_members(self, group_id: str | None = None) -> list[dict]:
+ assert self._db
+ sql = (
+ "SELECT m.*, i.username, i.pk_ed25519, i.pinned_at, i.pinned_via "
+ "FROM members m LEFT JOIN identities i ON i.user_id = m.user_id"
+ )
+ args: tuple = ()
+ if group_id is not None:
+ sql += " WHERE m.group_id = ?"
+ args = (group_id,)
+ async with self._db.execute(sql + " ORDER BY m.approved_at", args) as cur:
+ return [dict(r) for r in await cur.fetchall()]
+
+ async def set_status(self, group_id: str, user_id: str, status: str) -> bool:
+ assert self._db
+ cur = await self._db.execute(
+ "UPDATE members SET status = ? WHERE group_id = ? AND user_id = ?",
+ (status, group_id, user_id),
+ )
+ await self._db.commit()
+ return cur.rowcount > 0
+
+ # ── Invites ──────────────────────────────────────────────────────────────
+
+ async def create_invite(
+ self,
+ group_id: str,
+ user_id: str,
+ role: str,
+ created_by: str,
+ ttl: int = DEFAULT_INVITE_TTL,
+ username: str = "",
+ ) -> str:
+ """
+ Issue a one-time code. Returns it in the clear — this is the only moment it
+ exists outside the operator's hands; only its hash is kept.
+
+ Any earlier unused invite for the same person and group is dropped, so
+ re-inviting supersedes rather than accumulating valid codes.
+ """
+ assert self._db
+ await self._db.execute(
+ "DELETE FROM invites WHERE group_id = ? AND user_id = ? AND used_at IS NULL",
+ (group_id, user_id),
+ )
+ code = generate_code()
+ expires = datetime.now(timezone.utc) + timedelta(seconds=ttl)
+ await self._db.execute(
+ "INSERT INTO invites (code_hash, group_id, user_id, username, role, "
+ "created_by, created_at, expires_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
+ (hash_code(code), group_id, user_id, username, role, created_by, _now(),
+ expires.isoformat(timespec="seconds")),
+ )
+ await self._db.commit()
+ return code
+
+ async def consume_invite(self, code: str, user_id: str) -> dict | None:
+ """
+ Redeem a code for `user_id`, or return None.
+
+ Single use is enforced by the UPDATE's WHERE clause: two connections racing
+ the same code cannot both see `used_at IS NULL`, so exactly one wins.
+ """
+ assert self._db
+ code_hash = hash_code(code)
+ async with self._db.execute(
+ "SELECT * FROM invites WHERE code_hash = ?", (code_hash,)
+ ) as cur:
+ row = await cur.fetchone()
+ if not row:
+ return None
+
+ invite = dict(row)
+ if invite["used_at"] is not None:
+ return None
+ # A code is valid for exactly one account, so a leaked code cannot be
+ # redeemed by whoever finds it first.
+ if invite["user_id"] != user_id:
+ return None
+ if datetime.fromisoformat(invite["expires_at"]) < datetime.now(timezone.utc):
+ return None
+
+ cur = await self._db.execute(
+ "UPDATE invites SET used_at = ? WHERE code_hash = ? AND used_at IS NULL",
+ (_now(), code_hash),
+ )
+ await self._db.commit()
+ if cur.rowcount == 0:
+ return None
+ return invite
+
+ async def list_invites(self, include_used: bool = False) -> list[dict]:
+ assert self._db
+ sql = "SELECT * FROM invites"
+ if not include_used:
+ sql += " WHERE used_at IS NULL"
+ async with self._db.execute(sql + " ORDER BY created_at") as cur:
+ return [dict(r) for r in await cur.fetchall()]
+
+ async def purge_expired(self) -> int:
+ assert self._db
+ cur = await self._db.execute(
+ "DELETE FROM invites WHERE used_at IS NULL AND expires_at < ?",
+ (_now(),),
+ )
+ await self._db.commit()
+ return cur.rowcount
+
+
+async def open_roster(data_dir: Path) -> Roster:
+ roster = Roster(data_dir / "roster.db")
+ await roster.open()
+ return roster
+
+
+def write_code_file(data_dir: Path, code: str, expires_at: str,
+ name: str = "pair-code") -> Path:
+ """
+ Leave the code in a file as well as on stdout.
+
+ An operator working over SSH may not be able to copy out of their terminal,
+ and a code that can only be read off a scrolled-away screen is a dead end.
+ Pairing and invitation codes go to different files so one does not overwrite
+ the other.
+ """
+ path = data_dir / name
+ path.parent.mkdir(parents=True, exist_ok=True)
+ path.write_text(f"{code}\nexpires {expires_at}\n")
+ os.chmod(path, 0o600)
+ return path
diff --git a/packages/meshbay-node/src/meshbay_node/transport/__init__.py b/packages/meshbay-node/src/meshbay_node/transport/__init__.py
index df9c209..e423e35 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/__init__.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/__init__.py
@@ -1,7 +1,17 @@
-"""MeshBay Node transport layer — TCP+TLS (v1), QUIC (v2), WebRTC (browsers)."""
-from .server import ChunkServer
-from .client import ChunkClient
-from .http_server import create_http_app
+"""
+MeshBay Node transport layer — WebRTC DataChannel (primary), QUIC (direct/LAN).
+
+Transport decision (2026-08-13, second security review):
+ - WebRTC/ICE is the primary path for browser AND native clients. ICE/STUN is the
+ only NAT traversal validated on this project (2 ISPs, IPv4 STUN + IPv6, 4G CGNAT).
+ - QUIC is kept at parity for LAN, port-forwarded and hub-less `group://` access.
+ `punch_nat()` is a direct-connection helper, not a traversal stack.
+ - TCP+TLS (`server.py`/`client.py`) and the node HTTP file API (`http_server.py`)
+ were REMOVED in Phase 11.5. The HTTP API served private group indexes and
+ plaintext files with no authentication on 0.0.0.0 (finding C1); the TCP server
+ accepted a bare JWT with no GEK proof (finding C6). Neither is coming back —
+ every client path must go through the unified MNP handshake.
+"""
# QUIC transport (MNP v2) — requires aioquic>=1.0
try:
@@ -14,7 +24,7 @@ except ImportError:
Denylist = None # type: ignore[assignment,misc]
QUIC_AVAILABLE = False
-# WebRTC transport (browsers) — requires aiortc>=1.9
+# WebRTC transport (browsers + native clients) — requires aiortc>=1.9
try:
from .webrtc_server import WebRTCTransport, WebRTCPeerSession
WEBRTC_AVAILABLE = True
@@ -24,7 +34,6 @@ except ImportError:
WEBRTC_AVAILABLE = False
__all__ = [
- "ChunkServer", "ChunkClient", "create_http_app",
"QuicChunkServer", "QuicChunkClient", "Denylist", "QUIC_AVAILABLE",
"WebRTCTransport", "WebRTCPeerSession", "WEBRTC_AVAILABLE",
]
diff --git a/packages/meshbay-node/src/meshbay_node/transport/client.py b/packages/meshbay-node/src/meshbay_node/transport/client.py
deleted file mode 100644
index 63d50af..0000000
--- a/packages/meshbay-node/src/meshbay_node/transport/client.py
+++ /dev/null
@@ -1,148 +0,0 @@
-"""
-MeshBay — TCP+TLS chunk client (MNP v1).
-
-Used by the web client (or other nodes) to fetch files from a Mesh Node.
-Verifies Ed25519 chunk signatures using the node's public key from the hub.
-"""
-
-import asyncio
-import base64
-import logging
-import struct
-from pathlib import Path
-
-import blake3
-import msgpack
-from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey
-
-from meshbay_common import MNP_VERSION
-from meshbay_common.crypto import (
- chunk_key as derive_chunk_key,
- decrypt_chunk,
- verify_chunk_signature,
-)
-from meshbay_common.protocol import MNP
-from meshbay_node.transport.tls_cert import client_ssl_context
-
-log = logging.getLogger(__name__)
-
-MAX_MSG = 64 * 1024 * 1024
-
-
-async def _send(writer, obj):
- data = msgpack.packb(obj, use_bin_type=True)
- writer.write(struct.pack(">I", len(data)) + data)
- await writer.drain()
-
-async def _recv(reader):
- header = await reader.readexactly(4)
- length = struct.unpack(">I", header)[0]
- if length > MAX_MSG:
- raise ValueError(f"Message too large: {length}")
- return msgpack.unpackb(await reader.readexactly(length), raw=False)
-
-
-class ChunkClient:
- """
- Async client for fetching encrypted chunks from a ChunkServer.
-
- Usage:
- async with ChunkClient(host, port, jwt_token, gek, pk_node_b64) as client:
- data = await client.fetch_chunk(file_id, chunk_index=0)
- """
-
- def __init__(
- self,
- host: str,
- port: int,
- jwt_token: str,
- gek: bytes,
- pk_node_b64: str, # node's Ed25519 PK from hub — used for sig verification
- group_id: str = "",
- ):
- self._host = host
- self._port = port
- self._jwt_token = jwt_token
- self._gek = gek
- self._group_id = group_id
- self._pk_node = Ed25519PublicKey.from_public_bytes(
- base64.b64decode(pk_node_b64))
- self._reader: asyncio.StreamReader | None = None
- self._writer: asyncio.StreamWriter | None = None
-
- async def __aenter__(self):
- await self.connect()
- return self
-
- async def __aexit__(self, *_):
- await self.close()
-
- async def connect(self) -> None:
- ssl_ctx = client_ssl_context()
- self._reader, self._writer = await asyncio.open_connection(
- self._host, self._port, ssl=ssl_ctx)
-
- handshake_msg = {
- "type": MNP.HANDSHAKE,
- "v": MNP_VERSION,
- "token": self._jwt_token,
- }
- if self._group_id:
- handshake_msg["group_id"] = self._group_id
- await _send(self._writer, handshake_msg)
- ack = await _recv(self._reader)
- if ack.get("type") != MNP.HANDSHAKE_ACK:
- raise ConnectionError(f"Handshake rejected: {ack}")
- log.debug("Connected to node %s:%d", self._host, self._port)
-
- async def close(self) -> None:
- if self._writer:
- self._writer.close()
- await self._writer.wait_closed()
-
- async def fetch_index(self) -> bytes:
- """Request the Mesh Group Index. Returns raw wire bytes (encrypted)."""
- await _send(self._writer, {"type": MNP.INDEX_SYNC, "v": MNP_VERSION})
- msg = await _recv(self._reader)
- return base64.b64decode(msg["index_b64"])
-
- async def fetch_chunk(self, file_id: str, chunk_index: int) -> bytes:
- """
- Fetch, verify, and decrypt one chunk.
- Returns plaintext bytes.
- """
- await _send(self._writer, {
- "type": MNP.FILE_REQUEST,
- "v": MNP_VERSION,
- "file_id": file_id,
- "chunk_index": chunk_index,
- })
- msg = await _recv(self._reader)
-
- if msg.get("type") == "error":
- raise LookupError(msg.get("detail", "Unknown error"))
-
- ct = base64.b64decode(msg["ct_b64"])
- nonce = base64.b64decode(msg["nonce_b64"])
- ct_hash = base64.b64decode(msg["ct_hash_b64"])
- pt_hash = base64.b64decode(msg["pt_hash_b64"])
- sig = base64.b64decode(msg["sig_b64"])
- file_hash = base64.b64decode(msg["file_hash_b64"])
- ci = msg["chunk_index"]
-
- # 1. Verify Ed25519 signature
- verify_chunk_signature(self._pk_node, ci, nonce, ct_hash, sig)
-
- # 2. Verify ciphertext hash
- if blake3.blake3(ct).digest() != ct_hash:
- raise ValueError("Ciphertext hash mismatch")
-
- # 3. Decrypt
- ckey = derive_chunk_key(self._gek, file_hash, ci)
- plaintext = decrypt_chunk(ckey, nonce, ct)
-
- # 4. Verify plaintext hash
- if blake3.blake3(plaintext).digest() != pt_hash:
- raise ValueError("Plaintext hash mismatch after decryption")
-
- return plaintext
diff --git a/packages/meshbay-node/src/meshbay_node/transport/http_server.py b/packages/meshbay-node/src/meshbay_node/transport/http_server.py
deleted file mode 100644
index 151c2e8..0000000
--- a/packages/meshbay-node/src/meshbay_node/transport/http_server.py
+++ /dev/null
@@ -1,336 +0,0 @@
-"""
-MeshBay Node — HTTP file API (port 19001, public content).
-
-Serves public group content over standard HTTP so browsers can
-access files without any special protocol.
-
-Endpoints:
- GET / node info (JSON)
- GET /index public Mesh Group Index (JSON)
- GET /file/{file_id} full file download (streaming)
- GET /file/{file_id}/{chunk} single encrypted chunk (JSON)
- GET /hls/{file_id}/playlist.m3u8 HLS playlist
- GET /hls/{file_id}/{segment}.ts HLS segment (binary TS)
-
-Auth: Bearer JWT in Authorization header (or ?token= query param).
-For public groups: auth optional (anonymous browse allowed).
-For chunk download: auth required (JWT verified offline with hub PK).
-
-Note: this server handles PUBLIC content only (no GEK decryption).
-Private group content requires a client that can do ChaCha20 (Phase 5).
-"""
-
-import asyncio
-import base64
-import json
-import logging
-import os
-import struct
-import subprocess
-import tempfile
-from pathlib import Path
-
-import blake3
-import jwt
-from fastapi import FastAPI, Header, HTTPException, Query, Request
-from fastapi.responses import FileResponse, JSONResponse, StreamingResponse
-from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
-
-from meshbay_common import MNP_VERSION
-from meshbay_common.crypto import sign_chunk, pk_to_b64
-from meshbay_common.webcrypto import chunk_key_aes as derive_chunk_key, encrypt_chunk_aes as encrypt_chunk
-from meshbay_node import __version__
-from meshbay_node.indexer import GroupIndex
-from meshbay_node.indexer.group_index import GroupIndex
-
-log = logging.getLogger(__name__)
-
-CHUNK_SIZE = 1024 * 1024 # 1 MB
-HLS_SEGMENT_DURATION = 4 # seconds per HLS segment
-
-
-def create_http_app(
- sk_node: Ed25519PrivateKey,
- hub_pk_pem: bytes,
- shared_root: Path,
- index: GroupIndex,
- group_id: str,
- group_name: str,
- gek: bytes | None = None, # None for public groups
-) -> FastAPI:
- """
- Create the node's public HTTP API FastAPI app.
- Bind to 0.0.0.0:19001 (or configured port) for external access.
- """
- app = FastAPI(
- title="MeshBay Node HTTP API",
- version=__version__,
- docs_url=None,
- redoc_url=None,
- )
-
- # ── Auth helper ───────────────────────────────────────────────────────────
-
- def _verify_token_optional(
- authorization: str | None,
- token_param: str | None,
- ) -> dict | None:
- """Verify JWT if provided. Returns decoded payload or None."""
- raw = None
- if authorization and authorization.lower().startswith("bearer "):
- raw = authorization[7:]
- elif token_param:
- raw = token_param
- if not raw:
- return None
- try:
- return jwt.decode(raw, hub_pk_pem, algorithms=["EdDSA"])
- except Exception:
- return None
-
- def _require_token(
- authorization: str | None,
- token_param: str | None,
- ) -> dict:
- decoded = _verify_token_optional(authorization, token_param)
- if decoded is None:
- raise HTTPException(status_code=401, detail="Authentication required")
- return decoded
-
- # ── Node info ─────────────────────────────────────────────────────────────
-
- @app.get("/")
- async def node_info():
- return {
- "node_version": __version__,
- "mnp_version": MNP_VERSION,
- "group_id": group_id,
- "group_name": group_name,
- "file_count": index.count,
- "pk_node": pk_to_b64(sk_node.public_key()),
- }
-
- # ── Public index ──────────────────────────────────────────────────────────
-
- @app.get("/index")
- async def get_index(
- authorization: str | None = Header(default=None),
- token: str | None = Query(default=None),
- ):
- """Public Mesh Group Index as JSON. No auth required for public groups."""
- entries = [
- {
- "id": e.id,
- "name": e.name,
- "path": e.path,
- "size": e.size,
- "type": e.type,
- "duration": e.duration,
- }
- for e in index.entries
- ]
- return {
- "group_id": group_id,
- "group_name": group_name,
- "version": index.version,
- "entries": entries,
- }
-
- # ── Full file download (streaming) ────────────────────────────────────────
-
- @app.get("/file/{file_id}")
- async def download_file(
- file_id: str,
- authorization: str | None = Header(default=None),
- token: str | None = Query(default=None),
- ):
- """Stream an entire file. Public groups: no auth needed."""
- entry = index.get_entry(file_id)
- if not entry:
- raise HTTPException(status_code=404, detail="File not found in index")
-
- file_path = shared_root / entry.path / entry.name
- if not file_path.exists():
- raise HTTPException(status_code=404, detail="File not on disk")
-
- return FileResponse(
- path=str(file_path),
- filename=entry.name,
- media_type=_media_type(entry.name),
- )
-
- # ── Chunk endpoint (encrypted, for MNP-aware clients) ────────────────────
-
- @app.get("/file/{file_id}/{chunk_index}")
- async def get_chunk(
- file_id: str,
- chunk_index: int,
- authorization: str | None = Header(default=None),
- token: str | None = Query(default=None),
- ):
- """
- Serve one encrypted chunk (JSON). Auth required.
- Clients that understand MNP can decrypt with the GEK they got from the hub.
- """
- _require_token(authorization, token)
-
- entry = index.get_entry(file_id)
- if not entry:
- raise HTTPException(status_code=404, detail="File not found")
-
- file_path = shared_root / entry.path / entry.name
- if not file_path.exists():
- raise HTTPException(status_code=404, detail="File not on disk")
-
- # Read chunk
- with open(file_path, "rb") as f:
- f.seek(chunk_index * CHUNK_SIZE)
- plaintext = f.read(CHUNK_SIZE)
-
- if not plaintext:
- raise HTTPException(status_code=416, detail="Chunk out of range")
-
- file_hash = bytes.fromhex(entry.id)
- pt_hash = blake3.blake3(plaintext).digest()
-
- if gek:
- # Private group: encrypt chunk
- ckey = derive_chunk_key(gek, file_hash, chunk_index)
- nonce, ct = encrypt_chunk(ckey, plaintext)
- ct_hash = blake3.blake3(ct).digest()
- sig = sign_chunk(sk_node, chunk_index, nonce, ct_hash)
- return {
- "chunk_index": chunk_index,
- "plaintext_size": len(plaintext),
- "encrypted": True,
- "nonce_b64": base64.b64encode(nonce).decode(),
- "ct_b64": base64.b64encode(ct).decode(),
- "ct_hash_b64": base64.b64encode(ct_hash).decode(),
- "pt_hash_b64": base64.b64encode(pt_hash).decode(),
- "sig_b64": base64.b64encode(sig).decode(),
- "pk_node_b64": pk_to_b64(sk_node.public_key()),
- "file_hash_b64": base64.b64encode(file_hash).decode(),
- }
- else:
- # Public group: serve plaintext chunk (TLS provides transport encryption)
- pt_hash_b = blake3.blake3(plaintext).digest()
- sig_payload = chunk_index.to_bytes(4, "big") + bytes(12) + pt_hash_b
- sig = sk_node.sign(sig_payload)
- return {
- "chunk_index": chunk_index,
- "plaintext_size": len(plaintext),
- "encrypted": False,
- "data_b64": base64.b64encode(plaintext).decode(),
- "pt_hash_b64": base64.b64encode(pt_hash).decode(),
- "sig_b64": base64.b64encode(sig).decode(),
- "pk_node_b64": pk_to_b64(sk_node.public_key()),
- }
-
- # ── HLS streaming ─────────────────────────────────────────────────────────
-
- @app.get("/hls/{file_id}/playlist.m3u8")
- async def hls_playlist(
- file_id: str,
- authorization: str | None = Header(default=None),
- token: str | None = Query(default=None),
- ):
- """Generate HLS playlist for a video file."""
- entry = index.get_entry(file_id)
- if not entry or entry.type != "video":
- raise HTTPException(status_code=404, detail="Video file not found")
-
- file_path = shared_root / entry.path / entry.name
- if not file_path.exists():
- raise HTTPException(status_code=404, detail="File not on disk")
-
- duration = entry.duration or _probe_duration(file_path)
- if not duration:
- raise HTTPException(status_code=422, detail="Cannot determine video duration")
-
- n_segments = max(1, int(duration / HLS_SEGMENT_DURATION) + 1)
- token_param = f"?token={token}" if token else ""
-
- lines = [
- "#EXTM3U",
- "#EXT-X-VERSION:3",
- f"#EXT-X-TARGETDURATION:{HLS_SEGMENT_DURATION}",
- "#EXT-X-MEDIA-SEQUENCE:0",
- ]
- for i in range(n_segments):
- seg_dur = min(HLS_SEGMENT_DURATION, duration - i * HLS_SEGMENT_DURATION)
- if seg_dur <= 0:
- break
- lines.append(f"#EXTINF:{seg_dur:.3f},")
- lines.append(f"/hls/{file_id}/{i}.ts{token_param}")
- lines.append("#EXT-X-ENDLIST")
-
- return StreamingResponse(
- iter(["\n".join(lines)]),
- media_type="application/vnd.apple.mpegurl",
- )
-
- @app.get("/hls/{file_id}/{segment_index}.ts")
- async def hls_segment(
- file_id: str,
- segment_index: int,
- authorization: str | None = Header(default=None),
- token: str | None = Query(default=None),
- ):
- """Serve one HLS segment as MPEG-TS via ffmpeg transcoding."""
- entry = index.get_entry(file_id)
- if not entry or entry.type != "video":
- raise HTTPException(status_code=404, detail="Video not found")
-
- file_path = shared_root / entry.path / entry.name
- if not file_path.exists():
- raise HTTPException(status_code=404, detail="File not on disk")
-
- start_time = segment_index * HLS_SEGMENT_DURATION
-
- async def generate():
- proc = await asyncio.create_subprocess_exec(
- "ffmpeg", "-hide_banner", "-loglevel", "error",
- "-ss", str(start_time),
- "-i", str(file_path),
- "-t", str(HLS_SEGMENT_DURATION),
- "-c:v", "copy", "-c:a", "copy",
- "-f", "mpegts", "pipe:1",
- stdout=asyncio.subprocess.PIPE,
- stderr=asyncio.subprocess.DEVNULL,
- )
- assert proc.stdout
- while chunk := await proc.stdout.read(65536):
- yield chunk
- await proc.wait()
-
- return StreamingResponse(generate(), media_type="video/mp2t")
-
- return app
-
-
-# ── Helpers ───────────────────────────────────────────────────────────────────
-
-def _media_type(filename: str) -> str:
- ext = Path(filename).suffix.lower()
- return {
- ".mp4": "video/mp4", ".mkv": "video/x-matroska",
- ".webm": "video/webm", ".avi": "video/x-msvideo",
- ".mp3": "audio/mpeg", ".flac": "audio/flac",
- ".ogg": "audio/ogg", ".opus": "audio/opus",
- ".jpg": "image/jpeg", ".png": "image/png",
- ".pdf": "application/pdf",
- }.get(ext, "application/octet-stream")
-
-
-def _probe_duration(path: Path) -> float | None:
- """Use ffprobe to get video duration in seconds."""
- try:
- result = subprocess.run(
- ["ffprobe", "-v", "quiet", "-print_format", "json",
- "-show_format", str(path)],
- capture_output=True, text=True, timeout=10,
- )
- data = json.loads(result.stdout)
- return float(data["format"]["duration"])
- except Exception:
- return None
diff --git a/packages/meshbay-node/src/meshbay_node/transport/quic_client.py b/packages/meshbay-node/src/meshbay_node/transport/quic_client.py
index 288465f..9102085 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/quic_client.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/quic_client.py
@@ -11,6 +11,7 @@ TLS cert is self-signed; we use CERT_NONE equivalent in QUIC config.
import asyncio
import base64
import logging
+import os
import struct
from pathlib import Path
@@ -26,6 +27,32 @@ from meshbay_common import MNP_VERSION
from meshbay_common.crypto import verify_chunk_signature
from meshbay_common.webcrypto import chunk_key_aes as derive_chunk_key, decrypt_chunk_aes as decrypt_chunk
from meshbay_common.protocol import MNP
+from meshbay_common.handshake import (
+ NONCE_LEN,
+ ROLE_CLIENT,
+ ROLE_NODE,
+ handshake_transcript,
+ make_proof,
+ quic_binding,
+ verify_proof,
+)
+
+
+def _peer_cert_der(proto) -> bytes | None:
+ """
+ The server certificate as seen by the client — the channel-binding anchor.
+
+ Spike 11.5.6: aioquic 1.3.0 exposes no RFC 5705 exporter, and the peer
+ certificate only through a private attribute. Returns None when it is absent;
+ callers decide, because absence is not always an error — see below.
+ """
+ from cryptography.hazmat.primitives import serialization
+
+ tls = getattr(getattr(proto, "_quic", None), "tls", None)
+ cert = getattr(tls, "_peer_certificate", None) if tls is not None else None
+ if cert is None:
+ return None
+ return cert.public_bytes(serialization.Encoding.DER)
log = logging.getLogger(__name__)
@@ -99,6 +126,7 @@ class QuicChunkClient:
pk_node_b64: str,
local_port: int = 0, # 0 = OS picks; set for hole punching (Port-Restricted)
group_id: str = "",
+ peer_cert_der: bytes | None = None,
session_ticket: object | None = None,
):
self._host = host
@@ -112,6 +140,9 @@ class QuicChunkClient:
self._proto: _MNPClientProtocol | None = None
self._cm = None
self._ctrl_stream = 0
+ # 11.5.6 binding anchor. Travels with the session ticket: on a resumed
+ # TLS session the server does not re-send its certificate.
+ self._peer_cert_der: bytes | None = peer_cert_der
self._session_ticket = session_ticket
async def __aenter__(self):
@@ -146,19 +177,70 @@ class QuicChunkClient:
)
self._proto = await self._cm.__aenter__()
- handshake_msg = {
- "type": MNP.HANDSHAKE,
+ nonce_c = os.urandom(NONCE_LEN)
+ self._proto._send(self._ctrl_stream, {
+ "type": MNP.HANDSHAKE,
+ "v": MNP_VERSION,
+ "token": self._jwt_token,
+ "group_id": self._group_id,
+ "nonce": base64.b64encode(nonce_c).decode(),
+ })
+
+ reply = await self._proto._recv(self._ctrl_stream)
+ if reply.get("type") != MNP.HANDSHAKE_CHALLENGE:
+ raise ConnectionError(f"QUIC handshake rejected: {reply}")
+
+ nonce_s = base64.b64decode(reply["nonce"])
+
+ # On a RESUMED TLS session the server does not re-send its certificate, so
+ # there is nothing live to bind to. The session ticket is cryptographically
+ # derived from the original handshake, so binding to the certificate seen
+ # then is sound — but only if we actually saw one. We never fall back to an
+ # unbound proof: that would silently drop MitM detection (L4).
+ cert_der = _peer_cert_der(self._proto)
+ if cert_der is not None:
+ self._peer_cert_der = cert_der
+ elif getattr(self, "_peer_cert_der", None) is None:
+ raise ConnectionError(
+ "QUIC peer certificate unavailable and none cached from a prior "
+ "session — refusing to handshake without channel binding")
+ binding = quic_binding(self._peer_cert_der)
+
+ self._proto._send(self._ctrl_stream, {
+ "type": MNP.HANDSHAKE_RESPONSE,
"v": MNP_VERSION,
- "token": self._jwt_token,
- }
- if self._group_id:
- handshake_msg["group_id"] = self._group_id
- self._proto._send(self._ctrl_stream, handshake_msg)
+ "proof": base64.b64encode(make_proof(
+ self._gek, ROLE_CLIENT, self._group_id,
+ nonce_c, nonce_s, binding)).decode(),
+ })
+
ack = await self._proto._recv(self._ctrl_stream)
if ack.get("type") != MNP.HANDSHAKE_ACK:
raise ConnectionError(f"QUIC handshake rejected: {ack}")
+
+ # Authenticate the node before trusting anything it serves (C3).
+ if not verify_proof(
+ self._gek, base64.b64decode(ack.get("proof", "")), ROLE_NODE,
+ self._group_id, nonce_c, nonce_s, binding,
+ ):
+ raise ConnectionError("Node failed to prove GEK possession")
+
+ transcript = handshake_transcript(
+ ROLE_NODE, self._group_id, nonce_c, nonce_s, binding)
+ try:
+ Ed25519PublicKey.from_public_bytes(
+ base64.b64decode(ack["node_pk"])
+ ).verify(base64.b64decode(ack["sig"]), transcript)
+ except Exception as exc:
+ raise ConnectionError(f"Node signature invalid: {exc}") from exc
+
log.debug("QUIC connected to %s:%d", self._host, self._port)
+ @property
+ def peer_cert_der(self) -> bytes | None:
+ """Binding anchor to carry alongside a saved session ticket (11.5.6)."""
+ return self._peer_cert_der
+
async def close(self) -> None:
if self._cm:
await self._cm.__aexit__(None, None, None)
diff --git a/packages/meshbay-node/src/meshbay_node/transport/quic_server.py b/packages/meshbay-node/src/meshbay_node/transport/quic_server.py
index f439e62..ed3925d 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/quic_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/quic_server.py
@@ -20,6 +20,7 @@ The transport is the only change — all crypto, auth, and message types stay th
import asyncio
import base64
import logging
+import os
import struct
import subprocess
from pathlib import Path
@@ -34,6 +35,17 @@ from aioquic.quic.events import QuicEvent, StreamDataReceived, StreamReset
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
from meshbay_common import MNP_VERSION
+from meshbay_common.handshake import (
+ NONCE_LEN,
+ ROLE_CLIENT,
+ ROLE_NODE,
+ HandshakeError,
+ authorize_token,
+ handshake_transcript,
+ make_proof,
+ quic_binding,
+ verify_proof,
+)
from meshbay_common.crypto import (
sign_chunk,
pk_to_b64,
@@ -41,7 +53,6 @@ from meshbay_common.crypto import (
from meshbay_common.webcrypto import chunk_key_aes as derive_chunk_key, encrypt_chunk_aes as encrypt_chunk
from meshbay_common.protocol import MNP
from meshbay_node.indexer import GroupIndex
-from meshbay_node.transport.tls_cert import server_ssl_context
log = logging.getLogger(__name__)
@@ -51,22 +62,69 @@ ALPN = ["meshbay-mnp"]
class Denylist:
- """Shared denylist for revoked users and invalidated JWTs."""
+ """
+ Denylist for revoked users, groups and invalidated JWTs.
- def __init__(self):
+ Finding H4: revocations used to live only in memory, so a node restart silently
+ un-revoked everyone, and group revocations were dropped entirely — the hub
+ signed and broadcast them but the node's handler only understood "user" and
+ "jti". Now persisted to disk and group targets are honoured.
+ """
+
+ def __init__(self, path: Path | None = None):
self.user_ids: set[str] = set()
+ self.group_ids: set[str] = set()
self.jtis: set[str] = set()
+ self._path = path
+ self._load()
- def is_denied(self, user_id: str, jti: str) -> bool:
- return user_id in self.user_ids or jti in self.jtis
+ def is_denied(self, user_id: str, jti: str, group_id: str = "") -> bool:
+ return (user_id in self.user_ids
+ or jti in self.jtis
+ or (bool(group_id) and group_id in self.group_ids))
def deny_user(self, user_id: str) -> None:
self.user_ids.add(user_id)
log.info("Denied user: %s", user_id[:8])
+ self._save()
+
+ def deny_group(self, group_id: str) -> None:
+ self.group_ids.add(group_id)
+ log.info("Denied group: %s", group_id[:8])
+ self._save()
def deny_jti(self, jti: str) -> None:
self.jtis.add(jti)
log.info("Denied jti: %s", jti[:8])
+ self._save()
+
+ def _load(self) -> None:
+ if not self._path or not self._path.exists():
+ return
+ try:
+ import json
+ data = json.loads(self._path.read_text())
+ self.user_ids = set(data.get("users", []))
+ self.group_ids = set(data.get("groups", []))
+ self.jtis = set(data.get("jtis", []))
+ log.info("Denylist loaded: %d users, %d groups, %d jtis",
+ len(self.user_ids), len(self.group_ids), len(self.jtis))
+ except Exception as e:
+ log.warning("Could not load denylist from %s: %s", self._path, e)
+
+ def _save(self) -> None:
+ if not self._path:
+ return
+ try:
+ import json
+ self._path.parent.mkdir(parents=True, exist_ok=True)
+ self._path.write_text(json.dumps({
+ "users": sorted(self.user_ids),
+ "groups": sorted(self.group_ids),
+ "jtis": sorted(self.jtis),
+ }))
+ except Exception as e:
+ log.warning("Could not persist denylist to %s: %s", self._path, e)
# ── Wire helpers ──────────────────────────────────────────────────────────────
@@ -111,6 +169,9 @@ class _MNPServerProtocol(QuicConnectionProtocol):
self._user_id: str | None = None
self._group_id: str | None = None
self._buffers: dict[int, _StreamBuffer] = {}
+ self._nonce_client: bytes = b""
+ self._gek_challenge: bytes | None = None
+ self._pending = None
def quic_event_received(self, event: QuicEvent) -> None:
if isinstance(event, StreamDataReceived):
@@ -130,6 +191,8 @@ class _MNPServerProtocol(QuicConnectionProtocol):
try:
if mtype == MNP.HANDSHAKE:
self._do_handshake_sync(stream_id, msg)
+ elif mtype == MNP.HANDSHAKE_RESPONSE:
+ self._do_handshake_response_sync(stream_id, msg)
elif self._user_id is None:
self._send(stream_id, {"type": "error", "detail": "Handshake required"})
elif mtype == MNP.INDEX_SYNC:
@@ -147,44 +210,117 @@ class _MNPServerProtocol(QuicConnectionProtocol):
self._send(stream_id, {"type": "error", "detail": str(e)})
def _do_handshake_sync(self, stream_id: int, msg: dict) -> None:
- token = msg.get("token", "")
- group_id = msg.get("group_id", "")
+ """
+ Authorization half of the unified handshake (11.5.4).
+
+ This used to be a second, weaker copy of the WebRTC logic: group_id was
+ optional (so omitting it skipped the membership check entirely — M1),
+ node-scoped daemon tokens were accepted as client tokens (M9), and the
+ checks could drift from the WebRTC path independently. All of that now
+ comes from meshbay_common.handshake, shared with WebRTC.
+
+ NOT YET DONE — finding C6 remains open on this transport: there is still no
+ GEK proof here, so a forged or stolen token reaches the node and can inject
+ chat without holding the group key. The challenge/response and mutual node
+ proof (quic_binding() is written and unit-tested for exactly this) are the
+ remaining work in 11.5.4/5/6.
+ """
try:
- decoded = jwt.decode(token, self._ctx["hub_pk_pem"], algorithms=["EdDSA"])
- except Exception as e:
- self._send(stream_id, {"type": "error", "detail": f"Invalid JWT: {e}"})
+ peer = authorize_token(
+ msg.get("token", ""),
+ self._ctx["hub_pk_pem"],
+ group_id=msg.get("group_id", ""),
+ hosted_groups=self._ctx.get("groups"),
+ denylist=self._ctx.get("denylist"),
+ )
+ except HandshakeError as refusal:
+ self._send(stream_id, {"type": "error", "detail": str(refusal)})
+ self._quic.close()
+ return
+
+ try:
+ self._nonce_client = base64.b64decode(msg.get("nonce", ""))
+ except Exception:
+ self._nonce_client = b""
+ if len(self._nonce_client) < NONCE_LEN:
+ self._send(stream_id, {"type": "error", "detail": "Client nonce required"})
self._quic.close()
return
- denylist = self._ctx.get("denylist")
- if denylist and denylist.is_denied(decoded.get("sub", ""), decoded.get("jti", "")):
- self._send(stream_id, {"type": "error", "detail": "Token revoked"})
+ gctx = self._ctx["groups"][peer.group_id] if "groups" in self._ctx else self._ctx
+ if not gctx.get("gek"):
+ self._send(stream_id, {
+ "type": "error",
+ "detail": "Group encryption not initialized — contact node operator",
+ })
+ self._quic.close()
+ return
+
+ # Decoded but NOT authenticated: authentication is the GEK proof below.
+ self._pending = peer
+ self._gek_challenge = os.urandom(NONCE_LEN)
+ self._send(stream_id, {
+ "type": MNP.HANDSHAKE_CHALLENGE,
+ "v": MNP_VERSION,
+ "nonce": base64.b64encode(self._gek_challenge).decode(),
+ })
+
+ def _do_handshake_response_sync(self, stream_id: int, msg: dict) -> None:
+ """Verify the client's GEK proof, then prove the node in return (C6, C3)."""
+ if not self._gek_challenge or self._pending is None:
+ self._send(stream_id, {"type": "error", "detail": "No pending handshake challenge"})
+ return
+
+ peer = self._pending
+ gctx = self._ctx["groups"][peer.group_id] if "groups" in self._ctx else self._ctx
+ gek = gctx.get("gek")
+ if not gek:
+ self._send(stream_id, {"type": "error", "detail": "Group encryption not initialized"})
self._quic.close()
return
- if group_id and group_id not in decoded.get("groups", []):
- self._send(stream_id, {"type": "error", "detail": "Not a member of this group"})
+ binding = self._ctx.get("server_cert_der")
+ if not binding:
+ # Refuse rather than fall back to an unbound proof (L4).
+ self._send(stream_id, {"type": "error", "detail": "Channel binding unavailable"})
self._quic.close()
return
+ binding = quic_binding(binding)
- if group_id and "groups" in self._ctx and group_id not in self._ctx["groups"]:
- self._send(stream_id, {"type": "error", "detail": "Group not hosted on this node"})
+ try:
+ proof = base64.b64decode(msg.get("proof", ""))
+ except Exception:
+ self._send(stream_id, {"type": "error", "detail": "Invalid proof encoding"})
+ return
+
+ if not verify_proof(gek, proof, ROLE_CLIENT, peer.group_id,
+ self._nonce_client, self._gek_challenge, binding):
+ self._send(stream_id, {"type": "error", "detail": "GEK proof failed"})
self._quic.close()
return
- self._user_id = decoded["sub"]
- self._group_id = group_id
+ self._user_id = peer.user_id
+ self._group_id = peer.group_id
peers = self._ctx.get("_peers")
if peers is not None:
peers[self._user_id] = self
- log.info("QUIC handshake OK — user=%s group=%s", self._user_id[:8], group_id[:8] if group_id else "none")
+ transcript = handshake_transcript(
+ ROLE_NODE, peer.group_id, self._nonce_client, self._gek_challenge, binding)
+ node_proof = make_proof(
+ gek, ROLE_NODE, peer.group_id, self._nonce_client, self._gek_challenge, binding)
+
+ log.info("QUIC handshake OK — user=%s group=%s",
+ self._user_id[:8], self._group_id[:8])
self._send(stream_id, {
"type": MNP.HANDSHAKE_ACK,
"v": MNP_VERSION,
"node_pk": pk_to_b64(self._ctx["sk_node"].public_key()),
+ "proof": base64.b64encode(node_proof).decode(),
+ "sig": base64.b64encode(self._ctx["sk_node"].sign(transcript)).decode(),
})
+ self._gek_challenge = None
def _group_ctx(self) -> dict:
"""Resolve the active group context (multi-group or legacy single-group)."""
@@ -408,6 +544,13 @@ class QuicChunkServer:
generate_self_signed_cert(self._cert_path, self._key_path)
config = QuicConfiguration(is_client=False, alpn_protocols=ALPN)
config.load_cert_chain(str(self._cert_path), str(self._key_path))
+
+ # Channel-binding anchor for the handshake proof (11.5.6). Read from our own
+ # cert file — no aioquic internals needed on this side.
+ from cryptography import x509
+ from cryptography.hazmat.primitives import serialization as _ser
+ self._ctx["server_cert_der"] = x509.load_pem_x509_certificate(
+ self._cert_path.read_bytes()).public_bytes(_ser.Encoding.DER)
return config
def _store_ticket(self, ticket: Any) -> None:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/server.py b/packages/meshbay-node/src/meshbay_node/transport/server.py
deleted file mode 100644
index b77f1f2..0000000
--- a/packages/meshbay-node/src/meshbay_node/transport/server.py
+++ /dev/null
@@ -1,286 +0,0 @@
-"""
-MeshBay Node — TCP+TLS chunk server (MNP v1).
-
-Serves encrypted file chunks to authenticated clients over TLS.
-Each connection:
- 1. Client sends MNP handshake with JWT bearer token
- 2. Server verifies JWT offline (hub PK cached)
- 3. Client sends chunk requests
- 4. Server reads from disk, encrypts on-the-fly, signs, sends
-
-Wire protocol: length-prefixed msgpack (4-byte big-endian length header).
-All messages carry {"type": ..., "v": MNP_VERSION}.
-"""
-
-import asyncio
-import base64
-import logging
-import struct
-import time
-from pathlib import Path
-
-import blake3
-import jwt
-import msgpack
-from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
-
-from meshbay_common import MNP_VERSION
-from meshbay_common.crypto import (
- chunk_key as derive_chunk_key,
- encrypt_chunk,
- sign_chunk,
- pk_to_b64,
-)
-from meshbay_common.protocol import MNP
-from meshbay_node.indexer import GroupIndex
-from meshbay_node.transport.tls_cert import server_ssl_context
-
-log = logging.getLogger(__name__)
-
-CHUNK_SIZE = 1024 * 1024 # 1 MB
-MAX_MSG = 64 * 1024 * 1024 # 64 MB max message size (safety)
-
-
-# ── Wire helpers ──────────────────────────────────────────────────────────────
-
-async def _send(writer: asyncio.StreamWriter, obj: dict) -> None:
- data = msgpack.packb(obj, use_bin_type=True)
- writer.write(struct.pack(">I", len(data)) + data)
- await writer.drain()
-
-async def _recv(reader: asyncio.StreamReader) -> dict:
- header = await reader.readexactly(4)
- length = struct.unpack(">I", header)[0]
- if length > MAX_MSG:
- raise ValueError(f"Message too large: {length}")
- data = await reader.readexactly(length)
- return msgpack.unpackb(data, raw=False)
-
-
-# ── Chunk serving ─────────────────────────────────────────────────────────────
-
-def _serve_chunk(
- sk_node: Ed25519PrivateKey,
- gek: bytes,
- file_path: Path,
- file_hash: bytes,
- chunk_index: int,
-) -> dict:
- """Read, encrypt, sign one chunk. Blocking — run in executor."""
- with open(file_path, "rb") as f:
- f.seek(chunk_index * CHUNK_SIZE)
- plaintext = f.read(CHUNK_SIZE)
-
- pt_hash = blake3.blake3(plaintext).digest()
- ckey = derive_chunk_key(gek, file_hash, chunk_index)
- nonce, ct = encrypt_chunk(ckey, plaintext)
- ct_hash = blake3.blake3(ct).digest()
- sig = sign_chunk(sk_node, chunk_index, nonce, ct_hash)
-
- return {
- "type": MNP.FILE_CHUNK,
- "v": MNP_VERSION,
- "chunk_index": chunk_index,
- "plaintext_size": len(plaintext),
- "nonce_b64": base64.b64encode(nonce).decode(),
- "ct_b64": base64.b64encode(ct).decode(),
- "ct_hash_b64": base64.b64encode(ct_hash).decode(),
- "pt_hash_b64": base64.b64encode(pt_hash).decode(),
- "sig_b64": base64.b64encode(sig).decode(),
- "pk_node_b64": pk_to_b64(sk_node.public_key()),
- "file_hash_b64": base64.b64encode(file_hash).decode(),
- }
-
-
-# ── Connection handler ────────────────────────────────────────────────────────
-
-class _ConnectionHandler:
- def __init__(
- self,
- reader: asyncio.StreamReader,
- writer: asyncio.StreamWriter,
- sk_node: Ed25519PrivateKey,
- hub_pk_pem: bytes,
- gek: bytes,
- shared_root: Path,
- index: GroupIndex,
- groups: dict[str, dict] | None = None,
- ):
- self._reader = reader
- self._writer = writer
- self._sk_node = sk_node
- self._hub_pk_pem = hub_pk_pem
- self._gek = gek
- self._shared_root = shared_root
- self._index = index
- self._groups = groups
- self._peer = writer.get_extra_info("peername")
- self._user_id: str | None = None
- self._group_id: str | None = None
-
- async def handle(self) -> None:
- try:
- await self._handshake()
- await self._serve_loop()
- except asyncio.IncompleteReadError:
- log.debug("[%s] Client disconnected", self._peer)
- except Exception as e:
- log.warning("[%s] Error: %s", self._peer, e)
- await _send(self._writer, {"type": "error", "detail": str(e)})
- finally:
- self._writer.close()
-
- async def _handshake(self) -> None:
- msg = await _recv(self._reader)
- if msg.get("type") != MNP.HANDSHAKE:
- raise ValueError(f"Expected handshake, got {msg.get('type')!r}")
-
- token = msg.get("token", "")
- group_id = msg.get("group_id", "")
- try:
- decoded = jwt.decode(token, self._hub_pk_pem, algorithms=["EdDSA"])
- except Exception as e:
- raise PermissionError(f"Invalid JWT: {e}") from e
-
- if group_id and group_id not in decoded.get("groups", []):
- raise PermissionError("Not a member of this group")
-
- if group_id and self._groups and group_id not in self._groups:
- raise PermissionError("Group not hosted on this node")
-
- self._user_id = decoded["sub"]
- self._group_id = group_id
-
- if group_id and self._groups and group_id in self._groups:
- ctx = self._groups[group_id]
- self._gek = ctx["gek"]
- self._shared_root = ctx["shared_root"]
- self._index = ctx["index"]
-
- log.info("[%s] Handshake OK — user=%s group=%s", self._peer, self._user_id[:8], group_id[:8] if group_id else "none")
-
- await _send(self._writer, {
- "type": MNP.HANDSHAKE_ACK,
- "v": MNP_VERSION,
- "node_pk": pk_to_b64(self._sk_node.public_key()),
- })
-
- async def _serve_loop(self) -> None:
- loop = asyncio.get_event_loop()
- while True:
- msg = await _recv(self._reader)
- mtype = msg.get("type")
-
- if mtype == MNP.INDEX_SYNC:
- wire = self._index.serialize()
- await _send(self._writer, {
- "type": MNP.INDEX_SYNC,
- "v": MNP_VERSION,
- "index_b64": base64.b64encode(wire).decode(),
- })
-
- elif mtype == MNP.FILE_REQUEST:
- file_id = msg["file_id"]
- chunk_index = msg["chunk_index"]
-
- entry = self._index.get_entry(file_id)
- if entry is None:
- await _send(self._writer, {
- "type": "error",
- "detail": f"File not found: {file_id[:8]}",
- })
- continue
-
- file_path = self._shared_root / entry.path / entry.name
- if not file_path.exists():
- await _send(self._writer, {
- "type": "error", "detail": "File not on disk"})
- continue
-
- file_hash = bytes.fromhex(entry.id)
- chunk = await loop.run_in_executor(
- None, _serve_chunk,
- self._sk_node, self._gek, file_path, file_hash, chunk_index)
- await _send(self._writer, chunk)
-
- else:
- log.warning("[%s] Unknown message type: %s", self._peer, mtype)
-
-
-# ── Server ────────────────────────────────────────────────────────────────────
-
-class ChunkServer:
- """
- Async TCP+TLS server that serves encrypted file chunks.
-
- Usage:
- server = ChunkServer(
- host="0.0.0.0", port=19000,
- sk_node=sk, hub_pk_pem=pk_pem,
- gek=gek, shared_root=Path("/data"),
- index=group_index,
- )
- await server.start()
- # ... when shutting down:
- await server.stop()
- """
-
- def __init__(
- self,
- sk_node: Ed25519PrivateKey,
- hub_pk_pem: bytes,
- gek: bytes,
- shared_root: Path,
- index: GroupIndex,
- host: str = "0.0.0.0",
- port: int = 19000,
- cert_path: Path | None = None,
- key_path: Path | None = None,
- groups: dict[str, dict] | None = None,
- ):
- self._sk_node = sk_node
- self._hub_pk_pem = hub_pk_pem
- self._gek = gek
- self._shared_root = shared_root
- self._index = index
- self._host = host
- self._port = port
- self._cert_path = cert_path
- self._key_path = key_path
- self._groups = groups
- self._server: asyncio.Server | None = None
-
- @property
- def port(self) -> int:
- return self._port
-
- async def start(self) -> None:
- ssl_ctx = server_ssl_context(
- cert_path=self._cert_path or Path.home() / ".config/meshbay/node_tls.crt",
- key_path=self._key_path or Path.home() / ".config/meshbay/node_tls.key",
- )
- self._server = await asyncio.start_server(
- self._handle_connection,
- host=self._host,
- port=self._port,
- ssl=ssl_ctx,
- )
- log.info("ChunkServer listening on %s:%d (TLS)", self._host, self._port)
-
- async def stop(self) -> None:
- if self._server:
- self._server.close()
- await self._server.wait_closed()
- self._server = None
- log.info("ChunkServer stopped")
-
- async def _handle_connection(
- self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter
- ) -> None:
- handler = _ConnectionHandler(
- reader, writer,
- self._sk_node, self._hub_pk_pem,
- self._gek, self._shared_root, self._index,
- groups=self._groups,
- )
- await handler.handle()
diff --git a/packages/meshbay-node/src/meshbay_node/transport/tls_cert.py b/packages/meshbay-node/src/meshbay_node/transport/tls_cert.py
index 1354ac9..374fd08 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/tls_cert.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/tls_cert.py
@@ -1,17 +1,18 @@
"""
-Self-signed TLS certificate generation for the node.
+Self-signed TLS certificate generation for the node's QUIC listener.
The cert is used for transport confidentiality only.
Node identity is verified via Ed25519 PK (from hub), not TLS cert chain.
-Clients connect with ssl.CERT_NONE + verify Ed25519 at the MNP handshake layer.
-Certificate is generated once and cached at ~/.config/meshbay/node_tls.pem/.key.
+Phase 11.5 note: the certificate hash is also the intended channel-binding anchor for
+the QUIC handshake proof (11.5.6), since QUIC has no DTLS fingerprint to bind to.
+
+Certificate is generated once and cached at ~/.config/meshbay/node_tls.crt/.key.
"""
import logging
import os
from pathlib import Path
-import ssl
import datetime
import ipaddress
@@ -69,27 +70,6 @@ def generate_self_signed_cert(
return cert_path, key_path
-def server_ssl_context(
- cert_path: Path = DEFAULT_CERT,
- key_path: Path = DEFAULT_KEY,
-) -> ssl.SSLContext:
- """SSL context for the node's TCP server."""
- if not cert_path.exists() or not key_path.exists():
- generate_self_signed_cert(cert_path, key_path)
-
- ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
- ctx.load_cert_chain(certfile=cert_path, keyfile=key_path)
- ctx.minimum_version = ssl.TLSVersion.TLSv1_3
- return ctx
-
-
-def client_ssl_context() -> ssl.SSLContext:
- """
- SSL context for clients connecting to a node.
- CERT_NONE because we verify node identity via Ed25519 PK at the MNP layer.
- """
- ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_CLIENT)
- ctx.check_hostname = False
- ctx.verify_mode = ssl.CERT_NONE
- ctx.minimum_version = ssl.TLSVersion.TLSv1_3
- return ctx
+# `server_ssl_context()` / `client_ssl_context()` were removed in Phase 11.5 along with
+# the TCP+TLS transport they served. QUIC builds its own QuicConfiguration and calls
+# generate_self_signed_cert() directly.
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 13e90c8..fe4e3c2 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
@@ -28,7 +28,9 @@ import hashlib
import hmac
import logging
import os
+import re
import struct
+import time
from pathlib import Path
from typing import Any
@@ -41,16 +43,68 @@ from cryptography.hazmat.primitives.asymmetric.ed25519 import (
)
from meshbay_common import MNP_VERSION
-from meshbay_common.crypto import pk_to_b64
+from meshbay_common.handshake import (
+ NONCE_LEN,
+ ROLE_CLIENT,
+ ROLE_NODE,
+ HandshakeError,
+ authorize_token,
+ handshake_transcript,
+ make_proof,
+ verify_proof,
+ webrtc_binding,
+)
+from meshbay_common.adminop import (
+ ADMIN_CHALLENGE_TTL,
+ OP_FILE_DELETE,
+ OP_INVITE_CREATE,
+ admin_transcript,
+)
+from meshbay_common.crypto import pk_to_b64, wrap_gek_aes
+from meshbay_common.join import (
+ JOIN_TTL,
+ ROLE_MEMBER,
+ ROLE_OPERATOR,
+ join_transcript,
+)
from meshbay_common.webcrypto import chunk_key_aes, encrypt_chunk_aes
from meshbay_common.protocol import MNP
from meshbay_node.indexer import GroupIndex
+from meshbay_node.roster import DEFAULT_INVITE_TTL
log = logging.getLogger(__name__)
CHUNK_SIZE = 1024 * 1024
MAX_MSG = 64 * 1024 * 1024
+# Upload limits (finding C5a). Uploads used to land directly in the shared root under
+# a name the client chose, overwriting whatever was already there — which both violated
+# node sovereignty and defeated the delete authorization (overwrite a file, become its
+# recorded uploader, then delete it legitimately).
+MAX_UPLOAD_BYTES = 4 * 1024 * 1024 * 1024 # 4 GB per file
+
+# Budget for an unauthenticated peer: enough for a handshake and a bundle fetch,
+# nowhere near enough to be a memory-exhaustion primitive (H6).
+PRE_HANDSHAKE_MAX_MSG = 64 * 1024
+# ffmpeg is spawned per stream request; without a cap any member can fork-bomb
+# the node by requesting many streams at once (H6).
+MAX_CONCURRENT_TRANSCODES = 2
+# Bundle fetches are served in the pre-proof window (C4). Bounded and audited
+# until the native client removes remote keypair bundles entirely.
+MAX_PRE_PROOF_FETCHES = 4
+# Pairing codes carry 40 bits and are single-use, but a connection must not be
+# allowed to sit there guessing. Failures are audited, so a grind is visible.
+MAX_JOIN_ATTEMPTS = 5
+# Per-connection limits alone would not bind an attacker who can open connections
+# at will — and the adversary who can mint tokens for any account is the hub. So
+# failed pairings are also counted node-wide over a window.
+MAX_JOIN_FAILURES_WINDOW = 20
+JOIN_FAILURE_WINDOW = 600 # seconds
+UPLOAD_DIR_NAME = ".uploads"
+# Conservative allowlist: also what keeps markup out of filenames, which the node admin
+# UI used to render unescaped (finding H2).
+SAFE_UPLOAD_NAME = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._ -]{0,127}$")
+
def _extract_dtls_fingerprint(sdp: str) -> bytes:
"""Extract the DTLS SHA-256 fingerprint from SDP as raw 32 bytes."""
@@ -119,10 +173,18 @@ def _pack(obj: dict) -> bytes:
class _DataChannelBuffer:
- """Accumulate DataChannel messages and extract length-prefixed msgpack."""
+ """
+ Accumulate DataChannel messages and extract length-prefixed msgpack.
+
+ Finding H6: the limit was a flat 64 MB applied even before the handshake, so an
+ unauthenticated peer could announce a 64 MB frame and dribble bytes into it,
+ holding that much memory per connection. Until a peer has proved GEK
+ possession it gets a small budget; the large one is for file uploads.
+ """
- def __init__(self):
+ def __init__(self, max_message: int = MAX_MSG):
self._buf = bytearray()
+ self.max_message = max_message
def feed(self, data: bytes):
self._buf.extend(data)
@@ -130,7 +192,7 @@ class _DataChannelBuffer:
def messages(self):
while len(self._buf) >= 4:
length = struct.unpack(">I", self._buf[:4])[0]
- if length > MAX_MSG:
+ if length > self.max_message:
raise ValueError(f"Message too large: {length}")
if len(self._buf) < 4 + length:
break
@@ -162,15 +224,25 @@ class WebRTCPeerSession:
self._pc = pc
self._ctx = node_ctx
self._channel: RTCDataChannel | None = None
- self._buffer = _DataChannelBuffer()
+ self._buffer = _DataChannelBuffer(max_message=PRE_HANDSHAKE_MAX_MSG)
+ self._pre_proof_fetches = 0
self._user_id: str | None = None
self._group_id: str | None = None
self._peer_id: str = peer_id
self._remote_ip: str = ""
self._username: str = ""
- self._pk_user: str = ""
+ # Set from the roster: the key this node pinned for this account. Never
+ # from the JWT — the hub picks what goes in there.
+ self._pinned_pk: str = ""
self._gek_challenge: bytes | None = None
- self._admin_challenges: dict[str, bytes] = {}
+ # Same value as the GEK challenge, but kept for the life of the connection:
+ # a join_request is signed over it, and it must stay verifiable after the
+ # handshake clears the challenge (an operator pairs while already connected).
+ self._nonce_node: bytes = b""
+ self._join_attempts = 0
+ self._nonce_client: bytes = b""
+ self._admin_ops: dict[str, dict] = {} # op_id → pending admin operation
+ self._uploads: dict[str, dict] = {} # filename → {next_index, bytes}
def _setup_channel(self, channel: RTCDataChannel) -> None:
self._channel = channel
@@ -191,10 +263,30 @@ class WebRTCPeerSession:
self._do_handshake(msg)
elif mtype == MNP.HANDSHAKE_RESPONSE:
self._do_handshake_response(msg)
- elif mtype == MNP.GEK_BUNDLE_FETCH and self._gek_challenge is not None:
- asyncio.ensure_future(self._do_gek_bundle_fetch())
- elif mtype == MNP.KEYPAIR_BUNDLE_FETCH and self._gek_challenge is not None:
- asyncio.ensure_future(self._do_keypair_bundle_fetch())
+ elif mtype in (MNP.GEK_BUNDLE_FETCH, MNP.KEYPAIR_BUNDLE_FETCH) \
+ and self._gek_challenge is not None:
+ # Served before the GEK proof by necessity: the client needs its
+ # wrapped bundle in order to compute the proof. That window is a
+ # disclosure surface (C4) — a hub that forges a JWT reaches it — so
+ # it is bounded and audited here, and closed properly when clients
+ # stop storing keypair bundles on other people's nodes.
+ self._pre_proof_fetches += 1
+ if self._pre_proof_fetches > MAX_PRE_PROOF_FETCHES:
+ self._audit_auth_failed(
+ getattr(self, "_pending_group", ""), "pre-proof fetch flood")
+ self._send({"type": "error", "detail": "Too many requests"})
+ return
+ self._audit_pre_proof_fetch(mtype)
+ if mtype == MNP.GEK_BUNDLE_FETCH:
+ asyncio.ensure_future(self._do_gek_bundle_fetch())
+ else:
+ asyncio.ensure_future(self._do_keypair_bundle_fetch())
+ elif mtype == MNP.JOIN_REQUEST and self._nonce_node:
+ # Valid both before the GEK proof (a new member has no GEK to prove
+ # with) and after it (an operator pairing a browser is already
+ # connected). Authority comes from the pairing code and the
+ # signature, never from the session state.
+ asyncio.ensure_future(self._do_join_request(msg))
elif self._user_id is None:
self._send({"type": "error", "detail": "Handshake required"})
elif mtype == MNP.INDEX_SYNC:
@@ -213,17 +305,21 @@ class WebRTCPeerSession:
self._do_file_delete(msg)
elif mtype == MNP.ADMIN_RESPONSE:
self._do_admin_response(msg)
- elif mtype == MNP.GEK_BUNDLE_STORE:
- asyncio.ensure_future(self._do_gek_bundle_store(msg))
+ elif mtype == MNP.INVITE_CREATE:
+ self._do_invite_create(msg)
elif mtype == MNP.KEYPAIR_BUNDLE_STORE:
asyncio.ensure_future(self._do_keypair_bundle_store(msg))
+ elif mtype == MNP.KEYPAIR_BUNDLE_DELETE:
+ asyncio.ensure_future(self._do_keypair_bundle_delete())
elif mtype == MNP.STREAM_REQUEST:
asyncio.ensure_future(self._stream_video(msg))
else:
log.warning("Unknown MNP message type on DataChannel: %s", mtype)
except Exception as e:
- log.error("Error handling %s on DataChannel: %s", mtype, e)
- self._send({"type": "error", "detail": str(e)})
+ # Log the detail locally; send the peer a generic message. Exception
+ # text here carries filesystem paths and internal state (finding L3).
+ log.error("Error handling %s on DataChannel: %s", mtype, e, exc_info=True)
+ self._send({"type": "error", "detail": "Request failed"})
def _audit(self, event: str, detail: str = "") -> None:
audit = self._ctx.get("audit_store")
@@ -239,57 +335,73 @@ class WebRTCPeerSession:
detail=detail,
))
+ def _channel_binding(self) -> bytes:
+ """Both DTLS fingerprints, so a proof is valid on this connection only."""
+ offer_fp = b""
+ answer_fp = b""
+ if self._pc.remoteDescription:
+ offer_fp = _extract_dtls_fingerprint(self._pc.remoteDescription.sdp)
+ if self._pc.localDescription:
+ answer_fp = _extract_dtls_fingerprint(self._pc.localDescription.sdp)
+ if not offer_fp or not answer_fp:
+ return b""
+ return webrtc_binding(offer_fp, answer_fp)
+
def _do_handshake(self, msg: dict) -> None:
- token = msg.get("token", "")
group_id = msg.get("group_id", "")
try:
- decoded = jwt.decode(token, self._ctx["hub_pk_pem"], algorithms=["EdDSA"])
- except Exception as e:
- self._send({"type": "error", "detail": f"Invalid JWT: {e}"})
- self._audit_auth_failed(group_id, str(e))
- return
-
- denylist = self._ctx.get("denylist")
- if denylist and denylist.is_denied(decoded.get("sub", ""), decoded.get("jti", "")):
- self._send({"type": "error", "detail": "Token revoked"})
- return
-
- if group_id and group_id not in decoded.get("groups", []):
- self._send({"type": "error", "detail": "Not a member of this group"})
+ peer = authorize_token(
+ msg.get("token", ""),
+ self._ctx["hub_pk_pem"],
+ group_id=group_id,
+ hosted_groups=self._ctx.get("groups"),
+ denylist=self._ctx.get("denylist"),
+ )
+ except HandshakeError as refusal:
+ # HandshakeError messages are authored to be peer-safe, unlike arbitrary
+ # exception text (L3) — the client needs to know *why* it was refused.
+ self._send({"type": "error", "detail": str(refusal),
+ "code": getattr(refusal, "code", "")})
+ self._audit_auth_failed(group_id, str(refusal))
return
- if group_id and "groups" in self._ctx and group_id not in self._ctx["groups"]:
- self._send({"type": "error", "detail": "Group not hosted on this node"})
+ try:
+ self._nonce_client = base64.b64decode(msg.get("nonce", ""))
+ except Exception:
+ self._nonce_client = b""
+ if len(self._nonce_client) < NONCE_LEN:
+ # The client nonce is what makes the NODE's proof fresh (C3). Without
+ # it a recorded ack could be replayed by an impersonating peer.
+ self._send({"type": "error", "detail": "Client nonce required"})
return
- # Store decoded JWT data but DO NOT set self._user_id yet —
- # the user is not authenticated until they prove GEK possession.
- self._pending_sub = decoded["sub"]
- self._pending_group = group_id
- self._pending_username = decoded.get("username", "")
- self._pending_pk_user = decoded.get("pk_user", "")
+ # Decoded, but NOT authenticated: that happens on the GEK proof.
+ self._pending_sub = peer.user_id
+ self._pending_group = peer.group_id
+ self._pending_username = peer.username
- ctx = self._ctx
- if "groups" in ctx and group_id:
- gctx = ctx["groups"].get(group_id, ctx)
- else:
- gctx = ctx
- gek = gctx.get("gek")
-
- nonce = os.urandom(32)
- self._gek_challenge = nonce
- challenge = {
- "type": MNP.HANDSHAKE_CHALLENGE,
- "v": MNP_VERSION,
- "nonce": base64.b64encode(nonce).decode(),
- }
- if not gek:
+ gctx = self._ctx["groups"][peer.group_id] if "groups" in self._ctx else self._ctx
+ if not gctx.get("gek"):
self._send({
"type": "error",
"detail": "Group encryption not initialized — contact node operator",
})
return
- self._send(challenge)
+
+ self._gek_challenge = os.urandom(NONCE_LEN)
+ self._nonce_node = self._gek_challenge
+ self._send({
+ "type": MNP.HANDSHAKE_CHALLENGE,
+ "v": MNP_VERSION,
+ "nonce": base64.b64encode(self._gek_challenge).decode(),
+ # Announced here because a first-time joiner needs it *before* the
+ # ack: join_request signs a transcript naming this node, and someone
+ # who has never held the GEK cannot complete the handshake to learn
+ # it. Unverified at this point — the ack proves it, the client checks
+ # the two match, and a wrong value only makes our own verification
+ # fail. It is never a substitute for the ack's proof and signature.
+ "node_pk": self._node_pk_b64(),
+ })
def _do_handshake_response(self, msg: dict) -> None:
if not self._gek_challenge or not hasattr(self, "_pending_sub"):
@@ -297,61 +409,71 @@ class WebRTCPeerSession:
return
group_id = self._pending_group
- ctx = self._ctx
- if "groups" in ctx and group_id:
- gctx = ctx["groups"].get(group_id, ctx)
- else:
- gctx = ctx
+ gctx = self._ctx["groups"][group_id] if "groups" in self._ctx else self._ctx
gek = gctx.get("gek")
-
if not gek:
self._send({"type": "error", "detail": "Group encryption not initialized"})
self._gek_challenge = None
return
- proof = msg.get("proof", "")
try:
- proof_bytes = base64.b64decode(proof)
+ proof_bytes = base64.b64decode(msg.get("proof", ""))
except Exception:
self._send({"type": "error", "detail": "Invalid proof encoding"})
return
- offer_fp = b""
- answer_fp = b""
- if self._pc.remoteDescription:
- offer_fp = _extract_dtls_fingerprint(self._pc.remoteDescription.sdp)
- if self._pc.localDescription:
- answer_fp = _extract_dtls_fingerprint(self._pc.localDescription.sdp)
+ binding = self._channel_binding()
+ if not binding:
+ # Refuse rather than fall back to an unbound proof (L4).
+ self._send({"type": "error", "detail": "Channel binding unavailable"})
+ self._gek_challenge = None
+ self._audit_auth_failed(group_id, "no channel binding")
+ return
- data = self._gek_challenge + offer_fp + answer_fp
- expected = hmac.new(gek, data, hashlib.sha256).digest()
- if not hmac.compare_digest(proof_bytes, expected):
+ if not verify_proof(gek, proof_bytes, ROLE_CLIENT, group_id,
+ self._nonce_client, self._gek_challenge, binding):
self._send({"type": "error", "detail": "GEK proof failed"})
self._gek_challenge = None
self._audit_auth_failed(group_id, "GEK HMAC mismatch")
return
+ self._complete_handshake(gek, binding)
self._gek_challenge = None
- self._complete_handshake()
- def _complete_handshake(self) -> None:
+ def _complete_handshake(self, gek: bytes, binding: bytes) -> None:
+ # Authenticated peers may send large frames (file uploads); unauthenticated
+ # ones may not (H6).
+ self._buffer.max_message = MAX_MSG
self._user_id = self._pending_sub
self._group_id = self._pending_group
self._username = self._pending_username
- self._pk_user = self._pending_pk_user
+ asyncio.ensure_future(self._load_pinned_pk())
- peers = self._ctx.get("_peers")
- if peers is not None:
- peers[self._user_id] = self
+ self._peer_registry()[self._user_id] = self
node_user_id = self._ctx.get("node_user_id")
log.info("WebRTC handshake OK — user=%s group=%s",
self._user_id[:8],
self._group_id[:8] if self._group_id else "none")
+ # The node proves itself too (C3): possession of the GEK over the client's
+ # nonce, plus a signature over the same transcript with its long-term key.
+ # Previously the client received an unverifiable node_pk and trusted
+ # is_node_admin from whoever answered — so a peer that had hijacked
+ # signaling could serve a forged index, chat history and permissions.
+ node_transcript = handshake_transcript(
+ ROLE_NODE, self._group_id or "", self._nonce_client,
+ self._gek_challenge or b"", binding)
+ node_proof = make_proof(
+ gek, ROLE_NODE, self._group_id or "", self._nonce_client,
+ self._gek_challenge or b"", binding)
+
ack = {
"type": MNP.HANDSHAKE_ACK,
"v": MNP_VERSION,
"node_pk": pk_to_b64(self._ctx["sk_node"].public_key()),
+ "proof": base64.b64encode(node_proof).decode(),
+ "sig": base64.b64encode(
+ self._ctx["sk_node"].sign(node_transcript)).decode(),
"is_node_admin": bool(node_user_id and self._user_id == node_user_id),
}
if node_user_id:
@@ -388,65 +510,41 @@ class WebRTCPeerSession:
else:
self._send({"type": MNP.GEK_BUNDLE_RESP, "v": MNP_VERSION, "found": False})
- async def _do_gek_bundle_store(self, msg: dict) -> None:
- """Store a wrapped GEK bundle for a target user (admin operation)."""
- bundle_store = self._ctx.get("bundle_store")
- if not bundle_store:
- self._send({"type": "error", "detail": "Bundle store not available"})
- return
-
- target_user_id = msg.get("user_id", "")
- group_id = msg.get("group_id") or self._group_id
- pk_eph = msg.get("pk_eph_b64", "")
- nonce = msg.get("nonce_b64", "")
- wrapped = msg.get("wrapped_b64", "")
+ def _do_invite_create(self, msg: dict) -> None:
+ """
+ Issue a one-time pairing code for someone the operator wants to admit.
- if not target_user_id or not pk_eph or not nonce or not wrapped or not group_id:
- self._send({"type": "error", "detail": "Missing bundle fields"})
+ Replaces the old invite path, where the inviter fetched the invitee's
+ public key from the hub and wrapped the group key for whatever came back
+ (H3). The node now needs nothing but a name: it will wrap the key itself,
+ later, for a key the invitee proves they hold.
+ """
+ roster = self._ctx.get("roster")
+ if roster is None:
+ self._send({"type": "error", "detail": "Roster not available"})
return
- await bundle_store.store(group_id, target_user_id, pk_eph, nonce, wrapped)
- log.info("GEK bundle stored: group=%s user=%s", group_id[:8], target_user_id[:8])
- self._audit("gek_bundle_store", f"target={target_user_id[:8]}")
-
- self._send({
- "type": "ack", "v": MNP_VERSION,
- "detail": "gek_bundle_stored",
- "user_id": target_user_id,
- })
-
- # Auto-activate GEK if the bundle is for the node operator
- node_user_id = self._ctx.get("node_user_id")
- if node_user_id and target_user_id == node_user_id and group_id:
- await self._try_activate_gek(group_id, target_user_id)
-
- async def _try_activate_gek(self, group_id: str, user_id: str) -> None:
- """Unwrap and activate GEK for the node when the operator's bundle arrives."""
- from meshbay_common.crypto import unwrap_gek_aes
-
- bundle_store = self._ctx.get("bundle_store")
- sk_x_raw = self._ctx.get("sk_x25519_raw")
- pk_x_raw = self._ctx.get("pk_x25519_raw")
- if not bundle_store or not sk_x_raw or not pk_x_raw:
+ invitee_id = msg.get("user_id", "")
+ group_id = msg.get("group_id") or self._group_id
+ if not invitee_id or not group_id:
+ self._send({"type": "error", "detail": "Missing user_id or group_id"})
return
-
- bundle = await bundle_store.fetch(group_id, user_id)
- if not bundle:
+ if group_id != self._group_id:
+ self._send({"type": "error", "detail": "Wrong group for this session"})
return
- try:
- gek = unwrap_gek_aes(bundle, sk_x_raw, pk_x_raw)
- except Exception as e:
- log.warning("Failed to unwrap GEK for auto-activation: %s", e)
+ if not self._has_admin_authority():
+ self._send({
+ "type": "error",
+ "detail": "No operator paired — run `meshbay-node operator pair`",
+ })
return
- groups = self._ctx.get("groups")
- if groups and group_id in groups:
- groups[group_id]["gek"] = gek
- log.info("GEK auto-activated for group %s", group_id[:8])
- elif "gek" in self._ctx:
- self._ctx["gek"] = gek
- log.info("GEK auto-activated (single-group mode)")
+ self._issue_admin_challenge(OP_INVITE_CREATE, invitee_id, {
+ "group_id": group_id,
+ "user_id": invitee_id,
+ "username": str(msg.get("username", ""))[:64],
+ })
async def _do_keypair_bundle_fetch(self) -> None:
"""Serve the caller's encrypted keypair bundle during the handshake window."""
@@ -491,6 +589,287 @@ class WebRTCPeerSession:
"detail": "keypair_bundle_stored",
})
+ # ── Pairing and join (H3, M3) ────────────────────────────────────────────
+
+ def _join_refuse(self, reason: str, audit_detail: str = "") -> None:
+ self._join_attempts += 1
+ # Node-wide window, shared across connections: reconnecting must not reset
+ # the budget.
+ now = time.time()
+ failures = [t for t in self._ctx.get("join_failures", [])
+ if now - t < JOIN_FAILURE_WINDOW]
+ failures.append(now)
+ self._ctx["join_failures"] = failures
+ self._audit_join("join_refused", audit_detail or reason)
+ self._send({
+ "type": MNP.JOIN_RESULT,
+ "v": MNP_VERSION,
+ "ok": False,
+ "reason": reason,
+ })
+
+ def _audit_join(self, event: str, detail: str) -> None:
+ audit = self._ctx.get("audit_store")
+ if not audit:
+ return
+ self._remote_ip = self._remote_ip or _get_remote_ip(self._pc)
+ asyncio.ensure_future(audit.log_event(
+ user_id=self._user_id or getattr(self, "_pending_sub", "unknown"),
+ event=event,
+ ip=self._remote_ip,
+ username=self._username or getattr(self, "_pending_username", ""),
+ group_id=self._group_id or getattr(self, "_pending_group", "") or "",
+ detail=detail,
+ ))
+
+ async def _do_join_request(self, msg: dict) -> None:
+ """
+ Pin an identity, or recognise one already pinned.
+
+ The client signs its own Ed25519 and X25519 keys together with the node's
+ nonce, so the identity key vouches for the encryption key — that is what
+ will make it safe for the node to wrap the GEK for a key that arrived over
+ the wire instead of one fetched from the hub's directory (H3).
+
+ A first pairing needs a one-time code, which the hub never sees. Afterwards
+ the pin is the credential and a changed key is refused outright, the same
+ rule the client applies to `pk_node` (11.5.8).
+ """
+ roster = self._ctx.get("roster")
+ if roster is None:
+ self._send({"type": "error", "detail": "Roster not available"})
+ return
+
+ if self._join_attempts >= MAX_JOIN_ATTEMPTS:
+ self._send({"type": "error", "detail": "Too many attempts"})
+ return
+
+ now = time.time()
+ recent = [t for t in self._ctx.get("join_failures", [])
+ if now - t < JOIN_FAILURE_WINDOW]
+ if len(recent) >= MAX_JOIN_FAILURES_WINDOW:
+ self._audit_join("join_throttled", f"{len(recent)} failures in window")
+ self._send({"type": "error", "detail": "Pairing temporarily locked"})
+ return
+
+ user_id = self._user_id or getattr(self, "_pending_sub", "")
+ username = self._username or getattr(self, "_pending_username", "")
+ if not user_id:
+ self._send({"type": "error", "detail": "Handshake required"})
+ return
+
+ pk_ed_b64 = msg.get("pk_ed25519", "")
+ pk_x_b64 = msg.get("pk_x25519", "")
+ code = msg.get("code", "")
+ ts = msg.get("ts", 0)
+
+ try:
+ pk_ed_raw = base64.b64decode(pk_ed_b64)
+ pk_x_raw = base64.b64decode(pk_x_b64)
+ if len(pk_ed_raw) != 32 or len(pk_x_raw) != 32:
+ raise ValueError
+ pk_ed = Ed25519PublicKey.from_public_bytes(pk_ed_raw)
+ except Exception:
+ self._join_refuse("invalid_keys")
+ return
+
+ if not isinstance(ts, int) or abs(time.time() - ts) > JOIN_TTL:
+ self._join_refuse("stale_request")
+ return
+
+ # An empty group_id means operator pairing, which is node-wide. Anything
+ # else must be the group this connection authenticated to — a signature
+ # obtained for one group must not name another.
+ group_id = msg.get("group_id", "") or ""
+ session_group = self._group_id or getattr(self, "_pending_group", "") or ""
+ if group_id and group_id != session_group:
+ self._join_refuse("group_mismatch")
+ return
+
+ transcript = join_transcript(
+ node_pk_b64=self._node_pk_b64(),
+ group_id=group_id,
+ user_id=user_id,
+ pk_ed25519_b64=pk_ed_b64,
+ pk_x25519_b64=pk_x_b64,
+ nonce_node=self._nonce_node,
+ ts=ts,
+ )
+ try:
+ sig = base64.b64decode(msg.get("sig", ""))
+ except Exception:
+ self._join_refuse("invalid_signature_encoding")
+ return
+ if not self._verify_sig(pk_ed, transcript, sig):
+ self._join_refuse("signature_invalid")
+ return
+
+ known = await roster.get_identity(user_id)
+ if known:
+ if known["pk_ed25519"] != pk_ed_b64 or known["pk_x25519"] != pk_x_b64:
+ # The blocking warning, raised where it matters: whoever this is
+ # holds a different key than the person the operator paired.
+ self._join_refuse(
+ "key_changed",
+ f"pinned={known['pk_ed25519'][:16]} presented={pk_ed_b64[:16]}")
+ return
+ # An operator's row is node-wide (empty group), so a lookup for the
+ # group they happen to be opening finds nothing. Fall back to it, or
+ # the client is told it has no role on a node it administers.
+ member = (await roster.get_member(group_id, user_id)
+ or await roster.get_member("", user_id))
+ await self._join_ok(
+ user_id, pk_x_raw, session_group,
+ role=member["role"] if member else "",
+ recognised=True,
+ )
+ return
+
+ if not code:
+ if self._group_join_policy(session_group) == "open":
+ # An open-join group admits anyone the hub calls a member, so a
+ # code would protect nothing — the hub can walk in through the
+ # front door. Pin what turns up and say so in the audit log.
+ await self._pin_and_admit(
+ roster, user_id, username, pk_ed_b64, pk_x_b64,
+ group_id=session_group, role=ROLE_MEMBER,
+ approved_by="open-join", via="tofu")
+ await self._join_ok(user_id, pk_x_raw, session_group,
+ role=ROLE_MEMBER, recognised=False)
+ return
+ self._join_refuse("code_required")
+ return
+
+ invite = await roster.consume_invite(code, user_id)
+ if not invite:
+ self._join_refuse("code_invalid")
+ return
+
+ await self._pin_and_admit(
+ # The name comes from the invitation, not from the token: the hub does
+ # not put a username claim in a JWT, so pinning from the session alone
+ # left the roster nameless and `member revoke <name>` unable to match.
+ roster, user_id, invite["username"] or username, pk_ed_b64, pk_x_b64,
+ group_id=invite["group_id"], role=invite["role"],
+ approved_by=invite["created_by"], via="code")
+ # The roster row comes from the invitation; the key comes from the
+ # connection. An operator pairing is node-wide (empty group), but they
+ # redeemed the code while opening a group and expect to read it — and
+ # is_authorized() already grants an operator every group on this node.
+ await self._join_ok(user_id, pk_x_raw, session_group or invite["group_id"],
+ role=invite["role"], recognised=False)
+
+ def _group_join_policy(self, group_id: str) -> str:
+ """
+ Admission policy for a group, read from the node's own configuration.
+
+ Never from the hub: a hub that could declare a group open would be handed
+ the key to it (§3.4 of docs/invite-pairing-v1.md).
+ """
+ gctx = (self._ctx.get("groups") or {}).get(group_id) or {}
+ return gctx.get("join_policy", "invite")
+
+ async def _pin_and_admit(
+ self, roster, user_id: str, username: str, pk_ed_b64: str, pk_x_b64: str,
+ *, group_id: str, role: str, approved_by: str, via: str,
+ ) -> None:
+ await roster.pin_identity(
+ user_id=user_id, username=username,
+ pk_ed25519=pk_ed_b64, pk_x25519=pk_x_b64, via=via,
+ )
+ await roster.set_member(
+ group_id=group_id, user_id=user_id, role=role,
+ status="active", approved_by=approved_by,
+ )
+ if role == ROLE_OPERATOR:
+ self._ctx["has_admin_authority"] = True
+
+ log.info("Identity pinned (%s): user=%s role=%s", via, user_id[:8], role)
+ self._audit_join("join_pinned", f"role={role} via={via}")
+
+ async def _join_ok(
+ self, user_id: str, pk_x_raw: bytes, group_id: str,
+ *, role: str, recognised: bool,
+ ) -> None:
+ """
+ Answer a join, wrapping the group key for the key the caller just proved.
+
+ This is the H3 fix. The inviter used to fetch the invitee's public key from
+ the hub and wrap the GEK for whatever came back, so a hub that answered
+ with its own key was handed the group key by an honest member following the
+ protocol exactly. The node now wraps for a key that arrived from its owner
+ over an authenticated channel, bound to a pinned identity.
+ """
+ reply = {
+ "type": MNP.JOIN_RESULT,
+ "v": MNP_VERSION,
+ "ok": True,
+ "recognised": recognised,
+ "role": role,
+ }
+
+ roster = self._ctx["roster"]
+ if group_id and not await roster.is_authorized(group_id, user_id):
+ # Pinned on this node, but not admitted to this group. Hub membership
+ # alone must not produce a key.
+ reply["gek"] = False
+ reply["reason"] = "not_authorized_for_group"
+ self._send(reply)
+ self._audit_join("join_no_gek", f"group={group_id[:8]} not authorized")
+ return
+
+ gctx = (self._ctx.get("groups") or {}).get(group_id) or {}
+ gek = gctx.get("gek")
+ if not gek:
+ reply["gek"] = False
+ reply["reason"] = "no_gek"
+ self._send(reply)
+ return
+
+ bundle = wrap_gek_aes(gek, pk_x_raw)
+ reply["gek"] = True
+ reply["pk_eph_b64"] = bundle["pk_eph_b64"]
+ reply["nonce_b64"] = bundle["nonce_b64"]
+ reply["wrapped_b64"] = bundle["wrapped_b64"]
+ self._send(reply)
+ self._audit_join("gek_wrapped", f"group={group_id[:8]}")
+
+ async def _do_keypair_bundle_delete(self) -> None:
+ """
+ Withdraw our own key backup from this node.
+
+ Only ever our own: the user_id comes from the authenticated session, never
+ from the message. Someone who does not want a second browser should not be
+ leaving a PBKDF2-protected blob on every node they have ever joined (C4),
+ and turning the setting off has to remove what is already there — not just
+ stop adding to it.
+ """
+ bundle_store = self._ctx.get("bundle_store")
+ if not bundle_store:
+ self._send({"type": "error", "detail": "Bundle store not available"})
+ return
+
+ removed = await bundle_store.delete_keypair(self._user_id)
+ if removed:
+ log.info("Keypair bundle withdrawn by user=%s", self._user_id[:8])
+ self._audit("keypair_bundle_delete")
+ self._send({"type": "ack", "v": MNP_VERSION,
+ "detail": "keypair_bundle_deleted", "removed": removed})
+
+ def _audit_pre_proof_fetch(self, mtype: str) -> None:
+ """Record bundle access made before the GEK proof (C4)."""
+ audit = self._ctx.get("audit_store")
+ if not audit:
+ return
+ self._remote_ip = self._remote_ip or _get_remote_ip(self._pc)
+ asyncio.ensure_future(audit.log_event(
+ user_id=getattr(self, "_pending_sub", "unknown"),
+ event="pre_proof_fetch",
+ ip=self._remote_ip,
+ group_id=getattr(self, "_pending_group", "") or "",
+ detail=mtype,
+ ))
+
def _audit_auth_failed(self, group_id: str, reason: str) -> None:
audit = self._ctx.get("audit_store")
if audit:
@@ -508,6 +887,20 @@ class WebRTCPeerSession:
return self._ctx["groups"][self._group_id]
return self._ctx
+ def _peer_registry(self) -> dict:
+ """
+ Connected peers for THIS group only.
+
+ Finding H1: this used to live on the shared transport context, so a chat
+ message was broadcast to every peer on the node regardless of which group
+ they had authenticated to.
+ """
+ return self._group_ctx().setdefault("_peers", {})
+
+ def _user_names(self) -> dict:
+ """Display-name cache, per group — same leak as _peer_registry (H1)."""
+ return self._group_ctx().setdefault("_user_names", {})
+
def _do_index_sync(self) -> None:
ctx = self._group_ctx()
idx = ctx["index"]
@@ -555,6 +948,17 @@ class WebRTCPeerSession:
self._audit("file_download", entry.name)
def _do_stream_segment(self, msg: dict) -> None:
+ asyncio.ensure_future(self._do_stream_segment_async(msg))
+
+ async def _do_stream_segment_async(self, msg: dict) -> None:
+ """
+ Legacy HLS segment extraction (superseded by stream_req/MSE).
+
+ Finding H6: this ran subprocess.run(..., timeout=30) directly inside the
+ event loop, so a single request stalled the whole daemon — every peer,
+ every group — for up to thirty seconds. Now async and under the same
+ transcode semaphore as _stream_video.
+ """
ctx = self._group_ctx()
file_id = msg["file_id"]
segment_index = msg["segment_index"]
@@ -570,21 +974,34 @@ class WebRTCPeerSession:
self._send({"type": "error", "detail": "File not on disk"})
return
- import subprocess
+ sem = self._ctx.get("_transcode_sem")
+ if sem is None:
+ sem = asyncio.Semaphore(MAX_CONCURRENT_TRANSCODES)
+ self._ctx["_transcode_sem"] = sem
+
try:
- result = subprocess.run(
- ["ffmpeg", "-hide_banner", "-loglevel", "error",
- "-ss", str(segment_index * segment_duration),
- "-i", str(file_path),
- "-t", str(segment_duration),
- "-c:v", "copy", "-c:a", "copy",
- "-f", "mpegts", "pipe:1"],
- capture_output=True, timeout=30,
- )
- if result.returncode != 0 or not result.stdout:
+ async with sem:
+ proc = await asyncio.create_subprocess_exec(
+ "ffmpeg", "-hide_banner", "-loglevel", "error",
+ "-ss", str(segment_index * segment_duration),
+ "-i", str(file_path),
+ "-t", str(segment_duration),
+ "-c:v", "copy", "-c:a", "copy",
+ "-f", "mpegts", "pipe:1",
+ stdout=asyncio.subprocess.PIPE,
+ stderr=asyncio.subprocess.DEVNULL,
+ )
+ try:
+ stdout, _ = await asyncio.wait_for(proc.communicate(), timeout=30)
+ except asyncio.TimeoutError:
+ proc.kill()
+ await proc.wait()
+ self._send({"type": "error", "detail": "Segment extraction timed out"})
+ return
+ if proc.returncode != 0 or not stdout:
self._send({"type": "error", "detail": "Segment extraction failed"})
return
- segment_data = result.stdout
+ segment_data = stdout
except Exception:
self._send({"type": "error", "detail": "Segment extraction failed"})
return
@@ -599,11 +1016,14 @@ class WebRTCPeerSession:
})
def _do_chat_message(self, msg: dict) -> None:
- chat_store = self._ctx.get("chat_store")
+ # Per-group store — see _peer_registry() and finding H1. Reading chat_store
+ # off the shared transport context sent every group's messages to the first
+ # group's database, and served them back to anyone on the node.
+ chat_store = self._group_ctx().get("chat_store")
payload = msg.get("payload", "")
sender_name = msg.get("sender_name", "")
if sender_name:
- self._ctx.setdefault("_user_names", {})[self._user_id] = sender_name
+ self._user_names()[self._user_id] = sender_name
if chat_store:
raw = payload.encode() if isinstance(payload, str) else payload
asyncio.ensure_future(chat_store.save_message(
@@ -614,7 +1034,7 @@ class WebRTCPeerSession:
sender_name=sender_name,
))
- peers = self._ctx.get("_peers", {})
+ peers = self._peer_registry()
broadcast = {
"type": MNP.CHAT_MESSAGE,
"v": MNP_VERSION,
@@ -647,7 +1067,7 @@ class WebRTCPeerSession:
self._audit("chat_message")
def _do_chat_history(self, msg: dict) -> None:
- chat_store = self._ctx.get("chat_store")
+ chat_store = self._group_ctx().get("chat_store")
if not chat_store:
self._send({
"type": MNP.CHAT_HISTORY_RESPONSE,
@@ -662,7 +1082,7 @@ class WebRTCPeerSession:
async def _send_chat_history(self, chat_store, since: float, limit: int) -> None:
msgs = await chat_store.get_messages(since=since, limit=limit)
- names = self._ctx.get("_user_names", {})
+ names = self._user_names()
self._send({
"type": MNP.CHAT_HISTORY_RESPONSE,
"v": MNP_VERSION,
@@ -691,24 +1111,55 @@ class WebRTCPeerSession:
self._send({"type": "error", "detail": "Missing filename or data"})
return
+ if not SAFE_UPLOAD_NAME.match(filename):
+ self._send({"type": "error", "detail": "Invalid filename"})
+ return
+
shared_root = ctx.get("shared_root")
if not shared_root:
self._send({"type": "error", "detail": "No shared directory"})
return
- upload_dir = shared_root / ".uploads"
- upload_dir.mkdir(exist_ok=True)
- safe_name = filename.replace("/", "_").replace("\\", "_").replace("..", "_")
- tmp_path = upload_dir / f"{safe_name}.part"
+ # Per-user quarantine: a member can only ever write inside their own directory,
+ # so they cannot overwrite the operator's files or another member's (C5a).
+ rel_dir = f"{UPLOAD_DIR_NAME}/{self._user_id}"
+ user_dir = shared_root / UPLOAD_DIR_NAME / self._user_id
+ user_dir.mkdir(parents=True, exist_ok=True)
+ tmp_path = user_dir / f"{filename}.part"
+ final_path = user_dir / filename
+
+ state = self._uploads.get(filename)
+ if chunk_index == 0:
+ if final_path.exists():
+ self._send({"type": "error", "detail": "File already exists"})
+ return
+ state = {"next_index": 0, "bytes": 0}
+ self._uploads[filename] = state
+ elif state is None:
+ self._send({"type": "error", "detail": "Upload not started"})
+ return
+
+ # Reject out-of-order or replayed chunks — otherwise chunk_index>0 appends
+ # blindly to whatever .part file is already on disk.
+ if chunk_index != state["next_index"]:
+ self._send({"type": "error", "detail": "Unexpected chunk index"})
+ return
if isinstance(data, str):
chunk_bytes = base64.b64decode(data)
else:
chunk_bytes = bytes(data)
- mode = "ab" if chunk_index > 0 else "wb"
- with open(tmp_path, mode) as f:
+ if state["bytes"] + len(chunk_bytes) > MAX_UPLOAD_BYTES:
+ self._uploads.pop(filename, None)
+ tmp_path.unlink(missing_ok=True)
+ self._send({"type": "error", "detail": "Upload exceeds size limit"})
+ return
+
+ with open(tmp_path, "wb" if chunk_index == 0 else "ab") as f:
f.write(chunk_bytes)
+ state["next_index"] = chunk_index + 1
+ state["bytes"] += len(chunk_bytes)
self._send({
"type": MNP.FILE_UPLOAD_ACK,
@@ -718,21 +1169,30 @@ class WebRTCPeerSession:
})
if chunk_index + 1 >= total_chunks:
- final_path = shared_root / safe_name
+ self._uploads.pop(filename, None)
tmp_path.rename(final_path)
- log.info("Upload complete: %s (%d chunks)", safe_name, total_chunks)
- self._audit("file_upload", safe_name)
- self._register_uploader(ctx, safe_name)
+ log.info("Upload complete: %s (%d chunks, %d bytes)",
+ filename, total_chunks, state["bytes"])
+ self._audit("file_upload", f"{rel_dir}/{filename}")
+ self._register_uploader(ctx, rel_dir, filename)
+
+ def _register_uploader(self, ctx: dict, rel_dir: str, filename: str) -> None:
+ """
+ Tag the index entry with the uploader's identity after upload completes.
- def _register_uploader(self, ctx: dict, filename: str) -> None:
- """Tag the index entry with the uploader's user_id after upload completes."""
+ The key recorded here is the one this node pinned, not the one the token
+ carried. `pk_user` was a hub-chosen claim, and it decided who could later
+ delete the file: a hub issuing a token naming its own key could delete
+ anyone's uploads on any node. Deletion is supposed to be authorized by the
+ node, and this closes the last place where it was not.
+ """
idx = ctx.get("index")
if not idx:
return
for entry in idx.entries:
- if entry.name == filename and entry.path == "":
+ if entry.name == filename and entry.path == rel_dir:
entry.uploader_id = self._user_id
- entry.uploader_pk = self._pk_user
+ entry.uploader_pk = self._pinned_pk
return
def _do_file_delete(self, msg: dict) -> None:
@@ -747,28 +1207,115 @@ class WebRTCPeerSession:
self._send({"type": "error", "detail": "File not found"})
return
- admin_pk = self._ctx.get("admin_pk_ed25519")
has_uploader_pk = bool(entry.uploader_pk)
- if not admin_pk and not has_uploader_pk:
+ if not self._has_admin_authority() and not has_uploader_pk:
self._send({"type": "error", "detail": "No authorized key for deletion"})
return
- challenge = os.urandom(32)
- self._admin_challenges[file_id] = challenge
+ self._issue_admin_challenge(OP_FILE_DELETE, file_id)
+
+ # ── Admin operation challenge/response (finding H5) ──────────────────────
+
+ def _node_pk_b64(self) -> str:
+ return pk_to_b64(self._ctx["sk_node"].public_key())
+
+ def _issue_admin_challenge(
+ self, op: str, subject: str, payload: dict | None = None,
+ ) -> None:
+ """
+ Ask the client to authorize `op` on `subject` with its Ed25519 identity key.
+
+ The client is sent the transcript *fields*, not opaque bytes, so it can
+ rebuild and inspect what it signs. The node keeps the authoritative copy and
+ rebuilds the transcript itself at verification time — nothing signed is ever
+ taken from the response message.
+ """
+ nonce = os.urandom(32)
+ ts = int(time.time())
+ op_id = base64.b64encode(os.urandom(16)).decode()
+ self._admin_ops[op_id] = {
+ "op": op, "subject": subject, "nonce": nonce, "ts": ts,
+ "payload": payload or {},
+ }
self._send({
"type": MNP.ADMIN_CHALLENGE,
"v": MNP_VERSION,
- "challenge": base64.b64encode(challenge).decode(),
- "file_id": file_id,
+ "op_id": op_id,
+ "op": op,
+ "subject": subject,
+ "nonce": base64.b64encode(nonce).decode(),
+ "ts": ts,
+ "node_pk": self._node_pk_b64(),
+ "group_id": self._group_id or "",
})
+ @staticmethod
+ def _verify_sig(pk: Ed25519PublicKey | None, transcript: bytes, sig: bytes) -> bool:
+ if pk is None:
+ return False
+ try:
+ pk.verify(sig, transcript)
+ return True
+ except Exception:
+ return False
+
+ async def _load_pinned_pk(self) -> None:
+ """Remember which key this node pinned for the peer we just authenticated."""
+ roster = self._ctx.get("roster")
+ if roster is None or not self._user_id:
+ return
+ ident = await roster.get_identity(self._user_id)
+ if ident:
+ self._pinned_pk = ident["pk_ed25519"]
+
+ def _has_admin_authority(self) -> bool:
+ """
+ Cheap synchronous pre-check: is there anyone who could authorize this?
+
+ Only decides whether to issue a challenge at all — the gate is
+ `_verify_admin_sig`. The flag is set at startup and refreshed in-process
+ when an operator pairs.
+ """
+ return bool(self._ctx.get("admin_pk_ed25519")
+ or self._ctx.get("has_admin_authority"))
+
+ async def _verify_admin_sig(self, transcript: bytes, sig: bytes) -> bool:
+ """
+ Check a signature against every key holding node-operator authority.
+
+ Read from the roster on each call rather than cached: revoking a paired
+ browser must take effect immediately, and admin operations are rare enough
+ that a SQLite read costs nothing. `admin_pk_ed25519` in node.toml is still
+ honoured so an existing deployment keeps working until its operator pairs
+ (M3) — it is the legacy form of the same statement.
+ """
+ legacy = self._ctx.get("admin_pk_ed25519")
+ if self._verify_sig(legacy, transcript, sig):
+ return True
+
+ roster = self._ctx.get("roster")
+ if roster is None:
+ return False
+ for pk_b64 in await roster.operator_pks():
+ try:
+ pk = Ed25519PublicKey.from_public_bytes(base64.b64decode(pk_b64))
+ except Exception:
+ continue
+ if self._verify_sig(pk, transcript, sig):
+ return True
+ return False
+
def _do_admin_response(self, msg: dict) -> None:
- file_id = msg.get("file_id", "")
+ op_id = msg.get("op_id", "")
sig_b64 = msg.get("signature", "")
- challenge = self._admin_challenges.pop(file_id, None)
- if not challenge:
- self._send({"type": "error", "detail": "No pending admin challenge"})
+ pending = self._admin_ops.pop(op_id, None)
+ if not pending:
+ self._send({"type": "error", "detail": "No pending admin operation"})
+ return
+
+ if time.time() - pending["ts"] > ADMIN_CHALLENGE_TTL:
+ self._send({"type": "error", "detail": "Admin challenge expired"})
return
try:
@@ -777,40 +1324,96 @@ class WebRTCPeerSession:
self._send({"type": "error", "detail": "Invalid signature encoding"})
return
+ transcript = admin_transcript(
+ op=pending["op"],
+ node_pk_b64=self._node_pk_b64(),
+ group_id=self._group_id or "",
+ subject=pending["subject"],
+ nonce=pending["nonce"],
+ ts=pending["ts"],
+ )
+
+ if pending["op"] == OP_FILE_DELETE:
+ asyncio.ensure_future(
+ self._admin_exec_file_delete(pending, transcript, sig_bytes))
+ elif pending["op"] == OP_INVITE_CREATE:
+ asyncio.ensure_future(
+ self._admin_exec_invite_create(pending, transcript, sig_bytes))
+ else:
+ self._send({"type": "error", "detail": "Unknown admin operation"})
+
+ async def _admin_exec_file_delete(
+ self, pending: dict, transcript: bytes, sig: bytes,
+ ) -> None:
+ file_id = pending["subject"]
ctx = self._group_ctx()
entry = ctx["index"].get_entry(file_id)
if not entry:
self._send({"type": "error", "detail": "File not found"})
return
- verified = False
-
- # Try admin key (locally pinned)
- admin_pk = self._ctx.get("admin_pk_ed25519")
- if admin_pk:
- try:
- admin_pk.verify(sig_bytes, challenge)
- verified = True
- except Exception:
- pass
-
- # Try uploader key (stored at upload time)
- if not verified and entry.uploader_pk:
+ uploader_pk = None
+ if entry.uploader_pk:
try:
- uploader_key = Ed25519PublicKey.from_public_bytes(
+ uploader_pk = Ed25519PublicKey.from_public_bytes(
base64.b64decode(entry.uploader_pk))
- uploader_key.verify(sig_bytes, challenge)
- verified = True
except Exception:
- pass
+ uploader_pk = None
- if not verified:
+ # Node operator, or the user who uploaded this file — verified by the key
+ # recorded at upload time, never by a JWT claim (the hub controls those).
+ if not (await self._verify_admin_sig(transcript, sig)
+ or self._verify_sig(uploader_pk, transcript, sig)):
self._send({"type": "error", "detail": "Signature verification failed"})
self._audit("admin_auth_failed", f"file_delete:{file_id[:16]}")
return
self._exec_file_delete(ctx, file_id, entry)
+ async def _admin_exec_invite_create(
+ self, pending: dict, transcript: bytes, sig: bytes,
+ ) -> None:
+ # Node operator only. A group admin who does not run the node has no
+ # authority over who this node admits (deny by default). Delegation is
+ # designed but deferred — see §6.2 of docs/invite-pairing-v1.md.
+ if not await self._verify_admin_sig(transcript, sig):
+ self._send({"type": "error", "detail": "Signature verification failed"})
+ self._audit("admin_auth_failed", f"invite_create:{pending['subject'][:16]}")
+ return
+
+ roster = self._ctx.get("roster")
+ if roster is None:
+ self._send({"type": "error", "detail": "Roster not available"})
+ return
+
+ payload = pending["payload"]
+ code = await roster.create_invite(
+ group_id=payload["group_id"],
+ user_id=payload["user_id"],
+ role=ROLE_MEMBER,
+ created_by=self._user_id or "",
+ ttl=self._ctx.get("invite_ttl", DEFAULT_INVITE_TTL),
+ username=payload.get("username", ""),
+ )
+ invites = await roster.list_invites()
+ expires = next(
+ (i["expires_at"] for i in invites
+ if i["user_id"] == payload["user_id"]
+ and i["group_id"] == payload["group_id"]), "")
+
+ log.info("Invite created: group=%s user=%s",
+ payload["group_id"][:8], payload["user_id"][:8])
+ self._audit("invite_create", f"target={payload['user_id'][:8]}")
+ # The code exists in the clear exactly here and in the operator's hands.
+ self._send({
+ "type": MNP.INVITE_RESULT,
+ "v": MNP_VERSION,
+ "code": code,
+ "expires_at": expires,
+ "user_id": payload["user_id"],
+ "username": payload.get("username", ""),
+ })
+
def _exec_file_delete(self, ctx: dict, file_id: str, entry) -> None:
file_path = ctx["shared_root"] / entry.path / entry.name
if file_path.exists():
@@ -827,6 +1430,20 @@ class WebRTCPeerSession:
async def _stream_video(self, msg: dict) -> None:
"""Stream a video file as fMP4 segments via MSE-compatible output."""
+ # One ffmpeg per request with no cap lets any member exhaust the node's
+ # CPU and process table (H6). The semaphore lives on the transport context
+ # so it is shared across all peers, not per-session.
+ sem = self._ctx.get("_transcode_sem")
+ if sem is None:
+ sem = asyncio.Semaphore(MAX_CONCURRENT_TRANSCODES)
+ self._ctx["_transcode_sem"] = sem
+ if sem.locked() and sem._value <= 0:
+ self._send({"type": "error", "detail": "Server busy, retry shortly"})
+ return
+ async with sem:
+ await self._stream_video_inner(msg)
+
+ async def _stream_video_inner(self, msg: dict) -> None:
ctx = self._group_ctx()
file_id = msg.get("file_id", "")
entry = ctx["index"].get_entry(file_id)
@@ -914,9 +1531,8 @@ class WebRTCPeerSession:
async def close(self) -> None:
self._audit("disconnect")
- peers = self._ctx.get("_peers")
- if peers and self._user_id:
- peers.pop(self._user_id, None)
+ if self._user_id:
+ self._peer_registry().pop(self._user_id, None)
await self._pc.close()
diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py
index b4885af..28654df 100644
--- a/packages/meshbay-node/src/meshbay_node/ui/app.py
+++ b/packages/meshbay-node/src/meshbay_node/ui/app.py
@@ -9,13 +9,14 @@ FastAPI app providing:
- API endpoints for all data (JSON)
Served only on 127.0.0.1 — not exposed to the network.
-No authentication required (localhost only).
+Gated by a per-run session token (11.5.3) — printed at daemon startup.
"""
import base64
import json
import logging
import time
+from html import escape
from pathlib import Path
from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Query
@@ -23,6 +24,7 @@ from fastapi.responses import HTMLResponse, JSONResponse
from meshbay_node import __version__
from meshbay_common.crypto import generate_gek, wrap_gek_aes
+from meshbay_common.join import ROLE_MEMBER, ROLE_OPERATOR
log = logging.getLogger(__name__)
@@ -35,6 +37,55 @@ def create_ui_app(state: dict) -> FastAPI:
redoc_url=None,
)
+ @app.middleware("http")
+ async def _require_session_token(request, call_next):
+ """
+ Gate the admin UI behind a per-run token (11.5.3).
+
+ "localhost only" is weaker than it sounds: any process on the machine can
+ reach it, and a page in the operator's browser can reach it too via DNS
+ rebinding. Since this API can re-initialise a group's GEK and read the
+ audit log, an unauthenticated loopback service is a privilege boundary
+ waiting to be crossed. The token is printed at startup and accepted as
+ ?t= or the X-MeshBay-Token header.
+ """
+ from fastapi.responses import PlainTextResponse
+
+ token = state.get("ui_token")
+ if token:
+ supplied = (request.query_params.get("t")
+ or request.headers.get("X-MeshBay-Token"))
+ if supplied != token:
+ return PlainTextResponse("Forbidden", status_code=403)
+ return await call_next(request)
+
+ @app.middleware("http")
+ async def _security_headers(request, call_next):
+ """
+ Defence in depth behind the escaping fixes for H2. This UI is unauthenticated
+ on loopback, so script execution here equals full control of the node admin API.
+
+ Note what this does and does not do: the page relies on inline <script>, so
+ script-src must allow 'unsafe-inline' and CSP therefore does NOT prevent an
+ injected script from running. Escaping is the actual fix. What CSP buys is
+ containment — connect-src/img-src/form-action 'self'|'none' stop an injected
+ script from exfiltrating the audit log or config to an external host.
+ """
+ response = await call_next(request)
+ response.headers["Content-Security-Policy"] = (
+ "default-src 'none'; "
+ "style-src 'unsafe-inline'; "
+ "script-src 'unsafe-inline'; "
+ "connect-src 'self'; "
+ "img-src 'self' data:; "
+ "form-action 'none'; "
+ "frame-ancestors 'none'; "
+ "base-uri 'none'"
+ )
+ response.headers["X-Content-Type-Options"] = "nosniff"
+ response.headers["Referrer-Policy"] = "no-referrer"
+ return response
+
# ── JSON API ─────────────────────────────────────────────────────────────
@app.get("/api/status")
@@ -48,7 +99,6 @@ def create_ui_app(state: dict) -> FastAPI:
"status": state.get("status", "starting"),
"hub_url": state.get("hub_url", ""),
"username": state.get("username", ""),
- "node_port": state.get("node_port", 0),
"quic_port": state.get("quic_port", 0),
"endpoint_hint": state.get("endpoint_hint"),
"group_count": len(groups_ctx),
@@ -154,9 +204,7 @@ def create_ui_app(state: dict) -> FastAPI:
return {
"hub_url": config.hub.url,
"username": config.hub.username,
- "node_port": config.node.port,
"quic_port": config.node.quic_port,
- "http_port": config.node.http_port,
"ui_port": config.node.ui_port,
"data_dir": str(config.data_dir),
"groups": [
@@ -170,11 +218,159 @@ def create_ui_app(state: dict) -> FastAPI:
],
}
+ # ── Operator pairing (localhost only) ──────────────────────────────────
+
+ @app.post("/api/operator/pair")
+ async def operator_pair():
+ """
+ Issue a one-time code that pairs a browser as this node's operator.
+
+ The code is the whole point: it binds the operator's browser identity key
+ to their account without asking the hub, which is what stops a hub from
+ naming itself node administrator (M3, and the same substitution as H3).
+ It is returned once and stored only as a hash.
+ """
+ roster = state.get("roster")
+ user_id = state.get("node_user_id")
+ if not roster or not user_id:
+ return JSONResponse({"error": "Node not connected to hub yet"}, 503)
+
+ config = state.get("config")
+ ttl = (config.node.pair_ttl_hours if config else 24) * 3600
+ code = await roster.create_invite(
+ group_id="", # operator authority is node-wide
+ user_id=user_id,
+ role=ROLE_OPERATOR,
+ created_by="local-cli",
+ ttl=ttl,
+ username=(config.hub.username if config else ""),
+ )
+ invites = await roster.list_invites()
+ expires = next((i["expires_at"] for i in invites
+ if i["user_id"] == user_id and i["role"] == ROLE_OPERATOR), "")
+ return {"code": code, "expires_at": expires, "user_id": user_id}
+
+ @app.get("/api/roster")
+ async def api_roster(group_id: str = ""):
+ roster = state.get("roster")
+ if not roster:
+ return {"identities": [], "members": [], "invites": []}
+ return {
+ "identities": await roster.list_identities(),
+ "members": await roster.list_members(group_id or None),
+ "invites": await roster.list_invites(),
+ }
+
+ @app.post("/api/groups/{group_id}/invites")
+ async def create_invite(group_id: str, username: str):
+ """
+ Issue an invitation code from the CLI, without a browser.
+
+ The hub is asked for the account id and nothing else — never for a key.
+ A hub that answered with the wrong account would produce an invite whose
+ code it never learns, since the code goes to a human out of band.
+ """
+ roster = state.get("roster")
+ groups_ctx = state.get("groups_ctx", {})
+ if not roster:
+ return JSONResponse({"error": "Roster not available"}, 503)
+ if group_id not in groups_ctx:
+ return JSONResponse({"error": "Group not hosted on this node"}, 404)
+
+ hub = state.get("hub")
+ if not hub or not hub._session:
+ return JSONResponse({"error": "Hub not connected"}, 503)
+ try:
+ account = await hub.get_user_pubkeys(username)
+ except Exception as e:
+ return JSONResponse({"error": f"Unknown user {username!r}: {e}"}, 404)
+
+ config = state.get("config")
+ ttl = (config.node.invite_ttl_hours if config else 168) * 3600
+ code = await roster.create_invite(
+ group_id=group_id,
+ user_id=account["user_id"],
+ role=ROLE_MEMBER,
+ created_by="local-cli",
+ ttl=ttl,
+ username=username,
+ )
+ invites = await roster.list_invites()
+ expires = next((i["expires_at"] for i in invites
+ if i["user_id"] == account["user_id"]
+ and i["group_id"] == group_id), "")
+ return {"code": code, "expires_at": expires,
+ "username": username, "user_id": account["user_id"]}
+
+ @app.get("/api/resolve")
+ async def resolve_user(username: str):
+ """
+ Map a username to an account id for the CLI.
+
+ The roster answers first — it is the node's own record. The hub is the
+ fallback for identities pinned before invitations carried a name, and for
+ people admitted through an open-join group. Only an account id comes back;
+ no key is ever taken from here.
+ """
+ roster = state.get("roster")
+ if roster:
+ for ident in await roster.list_identities():
+ if ident["username"] == username:
+ return {"user_id": ident["user_id"], "source": "roster"}
+ hub = state.get("hub")
+ if hub and hub._session:
+ try:
+ account = await hub.get_user_pubkeys(username)
+ return {"user_id": account["user_id"], "source": "hub"}
+ except Exception:
+ pass
+ return JSONResponse({"error": f"Unknown user {username!r}"}, 404)
+
+ @app.post("/api/members/{user_id}/revoke")
+ async def revoke_member(user_id: str, group_id: str):
+ """
+ Stop serving the group key to someone.
+
+ Takes effect on their next connection: the key is wrapped on demand, so
+ there is no stored bundle left behind that would outlive this. Rotating
+ the group key is still required — they hold the current one.
+ """
+ roster = state.get("roster")
+ if not roster:
+ return JSONResponse({"error": "Roster not available"}, 503)
+ if not await roster.set_status(group_id, user_id, "revoked"):
+ return JSONResponse({"error": "No such member in that group"}, 404)
+ log.info("Member revoked: user=%s group=%s", user_id[:8], group_id[:8])
+ return {"status": "revoked", "user_id": user_id, "group_id": group_id,
+ "reminder": "rotate the group key: meshbay-node gek-init"}
+
+ @app.post("/api/members/{user_id}/unpin")
+ async def unpin_member(user_id: str):
+ """Forget a pinned identity, so the person can pair again with a new key."""
+ roster = state.get("roster")
+ if not roster:
+ return JSONResponse({"error": "Roster not available"}, 503)
+ if not await roster.unpin(user_id):
+ return JSONResponse({"error": "No such pinned identity"}, 404)
+ log.info("Identity unpinned: user=%s", user_id[:8])
+ return {"status": "unpinned", "user_id": user_id}
+
# ── GEK initialization (operator only, localhost) ──────────────────────
@app.post("/api/groups/{group_id}/gek")
async def init_gek(group_id: str):
- """Generate GEK, wrap for all group members, store, and activate."""
+ """
+ Generate the group key and activate it.
+
+ It used to be wrapped here for every member, using public keys fetched from
+ the hub — which is H3 with the node as the victim instead of the inviter: a
+ hub answering with its own key was handed the group key by the node itself.
+
+ Nothing is pre-wrapped for members now. Each member's copy is produced when
+ they connect, for a key they proved they hold (`join_request`). Only the
+ node's own copy is stored, so the daemon can reload the key across restarts
+ without the operator's browser.
+ """
groups_ctx = state.get("groups_ctx", {})
if group_id not in groups_ctx:
return JSONResponse({"error": "Group not hosted on this node"}, 404)
@@ -187,48 +383,15 @@ def create_ui_app(state: dict) -> FastAPI:
if not bundle_store:
return JSONResponse({"error": "Bundle store not available"}, 503)
- await hub.ensure_fresh_token()
- session = hub._session
- members_resp = await hub._http.get(
- f"/v1/groups/{group_id}/members",
- headers=session.auth_headers,
- )
- if not members_resp.is_success:
- return JSONResponse(
- {"error": f"Failed to fetch members: {members_resp.status_code}"}, 502)
- members = members_resp.json().get("members", [])
- if not members:
- return JSONResponse({"error": "No members in group"}, 400)
-
existing_gek = groups_ctx[group_id].get("gek")
gek = existing_gek or generate_gek()
+ errors: list[str] = []
- wrapped_count = 0
- errors = []
- for member in members:
- username = member["username"]
- user_id = member["user_id"]
- try:
- pk_data = await hub.get_user_pubkeys(username)
- pk_x_raw = base64.b64decode(pk_data["pk_x25519"])
- bundle = wrap_gek_aes(gek, pk_x_raw)
- await bundle_store.store(
- group_id, user_id,
- bundle["pk_eph_b64"], bundle["nonce_b64"], bundle["wrapped_b64"],
- )
- wrapped_count += 1
- log.info("GEK wrapped for %s (%s)", username, user_id[:8])
- except Exception as e:
- errors.append(f"{username}: {e}")
- log.warning("Failed to wrap GEK for %s: %s", username, e)
+ roster = state.get("roster")
+ authorized = len(await roster.list_members(group_id)) if roster else 0
- if wrapped_count == 0:
- return JSONResponse(
- {"error": "Failed to wrap GEK for any member", "details": errors}, 500)
-
- # Also store a copy wrapped for the node keystore X25519 key
- # so the daemon can reload GEK on restart without the operator's browser keys
- config = state.get("config")
+ # Store a copy wrapped for the node keystore X25519 key so the daemon can
+ # reload the GEK on restart without the operator's browser keys.
node_user_id = hub._session.user_id if hub._session else None
pk_x_node_raw = state.get("pk_x25519_raw")
if pk_x_node_raw and node_user_id:
@@ -239,13 +402,14 @@ def create_ui_app(state: dict) -> FastAPI:
node_bundle["pk_eph_b64"], node_bundle["nonce_b64"],
node_bundle["wrapped_b64"],
)
- log.info("GEK also wrapped for node keystore (daemon reload)")
+ log.info("GEK wrapped for node keystore (daemon reload)")
except Exception as e:
+ errors.append(f"node keystore: {e}")
log.warning("Failed to wrap GEK for node keystore: %s", e)
groups_ctx[group_id]["gek"] = gek
- log.info("GEK initialized for group %s — wrapped for %d/%d members",
- group_id[:8], wrapped_count, len(members))
+ log.info("GEK initialized for group %s — %d authorized member(s) will "
+ "receive it on connect", group_id[:8], authorized)
webrtc = state.get("webrtc")
if webrtc and "groups" in webrtc._ctx and group_id in webrtc._ctx["groups"]:
@@ -254,8 +418,7 @@ def create_ui_app(state: dict) -> FastAPI:
return {
"status": "ok",
"group_id": group_id,
- "wrapped_count": wrapped_count,
- "total_members": len(members),
+ "authorized_members": authorized,
"errors": errors,
}
@@ -311,11 +474,21 @@ def create_ui_app(state: dict) -> FastAPI:
@app.get("/", response_class=HTMLResponse)
async def root():
- return _render_page(state)
+ # Roster reads are async and the page renderer is not, so gather here.
+ roster = state.get("roster")
+ roster_view = None
+ if roster:
+ identities = {i["user_id"]: i for i in await roster.list_identities()}
+ roster_view = {
+ "identities": identities,
+ "members": await roster.list_members(),
+ "invites": await roster.list_invites(),
+ }
+ return _render_page(state, roster_view)
@app.get("/audit", response_class=HTMLResponse)
async def audit_page():
- return _render_audit_page()
+ return _render_audit_page(state.get("ui_token", ""))
return app
@@ -330,7 +503,71 @@ def _fmt_size(n: int) -> str:
return f"{n / (1024 * 1024 * 1024):.2f} GB"
-def _render_page(state: dict) -> str:
+def _render_roster(roster_view: dict | None) -> str:
+ """
+ Who this node recognises, and which keys are theirs.
+
+ Every value here is escaped: usernames come from the hub and pass through the
+ roster, so they are attacker-influenced text on the operator's own admin page
+ (the H2 rule applies to them exactly as it does to filenames).
+ """
+ if roster_view is None:
+ return '<p class="muted">Roster unavailable</p>'
+
+ identities = roster_view["identities"]
+ rows = ""
+ for m in roster_view["members"]:
+ ident = identities.get(m["user_id"], {})
+ scope = escape(m["group_id"][:8]) if m["group_id"] else "node-wide"
+ status_color = "#22c55e" if m["status"] == "active" else "#ef4444"
+ rows += (
+ f"<tr><td>{escape(str(ident.get('username') or m['user_id']))}</td>"
+ f"<td>{escape(str(m['role']))}</td>"
+ f"<td><span class='badge' style='background:{status_color}'>"
+ f"{escape(str(m['status']))}</span></td>"
+ f"<td>{scope}</td>"
+ f"<td><code>{escape(str(ident.get('pk_ed25519', ''))[:16])}…</code></td>"
+ f"<td>{escape(str(ident.get('pinned_at', '?')))} "
+ f"({escape(str(ident.get('pinned_via', '?')))})</td></tr>"
+ )
+ if not rows:
+ rows = ('<tr><td colspan="6" class="muted">Nobody admitted yet — '
+ 'run <code>meshbay-node member invite &lt;username&gt;</code></td></tr>')
+
+ invite_rows = ""
+ for i in roster_view["invites"]:
+ invite_rows += (
+ f"<tr><td><code>{escape(str(i['user_id'])[:16])}</code></td>"
+ f"<td>{escape(str(i['group_id'][:8] or 'node-wide'))}</td>"
+ f"<td>{escape(str(i['role']))}</td>"
+ f"<td>{escape(str(i['expires_at']))}</td></tr>"
+ )
+ invites_html = ""
+ if invite_rows:
+ invites_html = f"""
+ <details style="margin-top:10px"><summary>Pending invitations</summary>
+ <table>
+ <thead><tr><th>Account</th><th>Group</th><th>Role</th><th>Expires</th></tr></thead>
+ <tbody>{invite_rows}</tbody>
+ </table>
+ </details>"""
+
+ return f"""
+ <table>
+ <thead><tr><th>User</th><th>Role</th><th>Status</th><th>Scope</th>
+ <th>Identity key</th><th>Pinned</th></tr></thead>
+ <tbody>{rows}</tbody>
+ </table>
+ {invites_html}
+ <p class="muted" style="margin-top:8px">
+ Codes are issued from the CLI: <code>meshbay-node operator pair</code>,
+ <code>meshbay-node member invite &lt;username&gt;</code>. They never pass
+ through the hub.
+ </p>"""
+
+
+def _render_page(state: dict, roster_view: dict | None = None) -> str:
+ token_js = json.dumps(state.get("ui_token", ""))
status = state.get("status", "starting")
indexes = state.get("indexes", {})
groups_ctx = state.get("groups_ctx", {})
@@ -356,12 +593,16 @@ def _render_page(state: dict) -> str:
fcount = idx.count if idx else 0
total_size = sum(e.size for e in idx.entries) if idx else 0
+ # Everything interpolated below is attacker-controlled: filenames come from
+ # uploads by any group member. Rendering them raw was a stored XSS into the
+ # unauthenticated localhost admin UI, i.e. full control of the node admin API
+ # from the operator's browser (finding H2).
file_rows = ""
if idx:
for e in sorted(idx.entries, key=lambda x: x.name):
file_rows += (
- f"<tr><td>{e.name}</td><td>{e.type}</td>"
- f"<td>{_fmt_size(e.size)}</td><td>{e.path or '/'}</td></tr>"
+ f"<tr><td>{escape(e.name)}</td><td>{escape(e.type)}</td>"
+ f"<td>{_fmt_size(e.size)}</td><td>{escape(e.path or '/')}</td></tr>"
)
has_gek = bool(ctx.get("gek"))
@@ -385,14 +626,14 @@ def _render_page(state: dict) -> str:
groups_html += f"""
<div class="card">
- <h3>{name}
- <span class="badge" style="background:#6366f1">{vis}</span>
+ <h3>{escape(str(name))}
+ <span class="badge" style="background:#6366f1">{escape(str(vis))}</span>
{gek_badge}
</h3>
- <p><b>Directory:</b> <code>{shared}</code></p>
+ <p><b>Directory:</b> <code>{escape(str(shared))}</code></p>
<p><b>Files:</b> {fcount} &mdash; <b>Total:</b> {_fmt_size(total_size)}</p>
{gek_action}
- <p class="muted">ID: {gid}</p>
+ <p class="muted">ID: {escape(gid)}</p>
<details><summary>File list</summary>
<table>
<thead><tr><th>Name</th><th>Type</th><th>Size</th><th>Path</th></tr></thead>
@@ -408,10 +649,10 @@ def _render_page(state: dict) -> str:
from meshbay_node.transport.webrtc_server import _get_remote_ip
ip = session._remote_ip or _get_remote_ip(session._pc)
peers_html += (
- f"<tr><td>{session._username or session._user_id or '—'}</td>"
- f"<td>{ip or '—'}</td>"
- f"<td>{session._group_id[:8] if session._group_id else '—'}</td>"
- f"<td>{session._pc.connectionState}</td></tr>"
+ f"<tr><td>{escape(session._username or session._user_id or '—')}</td>"
+ f"<td>{escape(ip or '—')}</td>"
+ f"<td>{escape(session._group_id[:8] if session._group_id else '—')}</td>"
+ f"<td>{escape(session._pc.connectionState)}</td></tr>"
)
if not peers_html:
peers_html = '<tr><td colspan="4" class="muted">No connected peers</td></tr>'
@@ -487,14 +728,16 @@ def _render_page(state: dict) -> str:
<tbody>{peers_html}</tbody>
</table>
+ <h2>Roster</h2>
+ {_render_roster(roster_view)}
+
<h2>Groups</h2>
{groups_html or '<p class="muted">No groups configured</p>'}
<h2>Node Configuration</h2>
<div class="card">
<p><b>Hub:</b> {state.get("hub_url", "—")}</p>
- <p><b>QUIC port:</b> {state.get("quic_port", "—")} &mdash;
- <b>TCP port:</b> {state.get("node_port", "—")}</p>
+ <p><b>QUIC port:</b> {state.get("quic_port", "—")}</p>
<p><b>Node ID:</b> <code>{state.get("endpoint_hint") or "—"}</code></p>
</div>
@@ -524,17 +767,18 @@ def _render_page(state: dict) -> str:
</div>
</div>
<script>
+const TOKEN = {token_js};
async function initGEK(groupId) {{
const btn = document.getElementById('gek-btn-' + groupId.slice(0,8));
const status = document.getElementById('gek-status-' + groupId.slice(0,8));
if (btn) btn.disabled = true;
if (status) status.textContent = 'Initializing...';
try {{
- const resp = await fetch('/api/groups/' + groupId + '/gek', {{ method: 'POST' }});
+ const resp = await fetch('/api/groups/' + groupId + '/gek?t=' + TOKEN, {{ method: 'POST' }});
const data = await resp.json();
if (resp.ok) {{
- if (status) status.textContent = 'GEK initialized — wrapped for '
- + data.wrapped_count + '/' + data.total_members + ' members';
+ if (status) status.textContent = 'GEK initialized — '
+ + data.authorized_members + ' authorized member(s) get it on connect';
if (status) status.style.color = '#22c55e';
setTimeout(() => location.reload(), 2000);
}} else {{
@@ -554,8 +798,11 @@ setTimeout(()=>location.reload(), 10000);
</html>"""
-def _render_audit_page() -> str:
- return """<!DOCTYPE html>
+def _render_audit_page(token: str = "") -> str:
+ return _AUDIT_HTML.replace("__TOKEN__", json.dumps(token))
+
+
+_AUDIT_HTML = """<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="utf-8">
@@ -591,7 +838,7 @@ def _render_audit_page() -> str:
<body>
<div class="container">
<h1>Audit Log</h1>
- <nav><a href="/">Dashboard</a><a href="/audit">Audit Log</a></nav>
+ <nav><a id="navHome" href="/">Dashboard</a><a id="navAudit" href="/audit">Audit Log</a></nav>
<div class="filters">
<select id="eventFilter">
@@ -620,23 +867,39 @@ def _render_audit_page() -> str:
</table>
</div>
<script>
+const TOKEN = __TOKEN__;
async function load() {
const ev = document.getElementById('eventFilter').value;
const limit = document.getElementById('limitSelect').value;
- let url = '/api/audit?limit=' + limit;
+ let url = '/api/audit?limit=' + limit + (TOKEN ? '&t=' + TOKEN : '');
if (ev) url += '&event=' + ev;
const r = await fetch(url);
const data = await r.json();
const tbody = document.getElementById('tbody');
document.getElementById('count').textContent = data.entries.length + ' entries';
- tbody.innerHTML = data.entries.map(e => {
- const t = new Date(e.timestamp * 1000).toLocaleString();
- return '<tr><td>' + t + '</td><td>' + e.event + '</td><td>'
- + (e.username || e.user_id.slice(0,8)) + '</td><td>'
- + (e.ip || '—') + '</td><td>'
- + (e.group_id ? e.group_id.slice(0,8) : '—') + '</td><td>'
- + (e.detail || '') + '</td></tr>';
- }).join('');
+ // textContent, not innerHTML: e.detail carries filenames chosen by group members
+ // (finding H2). Building this row with string concatenation was a stored XSS.
+ tbody.replaceChildren(...data.entries.map(e => {
+ const tr = document.createElement('tr');
+ const cells = [
+ new Date(e.timestamp * 1000).toLocaleString(),
+ e.event,
+ e.username || (e.user_id || '').slice(0, 8),
+ e.ip || '—',
+ e.group_id ? e.group_id.slice(0, 8) : '—',
+ e.detail || '',
+ ];
+ for (const value of cells) {
+ const td = document.createElement('td');
+ td.textContent = value;
+ tr.appendChild(td);
+ }
+ return tr;
+ }));
+}
+for (const [id, href] of [['navHome','/'],['navAudit','/audit']]) {
+ const el = document.getElementById(id);
+ if (el && TOKEN) el.href = href + '?t=' + TOKEN;
}
document.getElementById('eventFilter').onchange = load;
document.getElementById('limitSelect').onchange = load;
diff --git a/packages/meshbay-node/tests/test_daemon.py b/packages/meshbay-node/tests/test_daemon.py
index 1c5a07e..60ef11f 100644
--- a/packages/meshbay-node/tests/test_daemon.py
+++ b/packages/meshbay-node/tests/test_daemon.py
@@ -21,7 +21,6 @@ from meshbay_node.config import Config, HubConfig, NodeConfig, GroupConfig, Keys
from meshbay_node.daemon import NodeDaemon
from meshbay_node.indexer import DirectoryIndexer
-
def _mock_keystore_keys(sk_ed):
"""Create a mock keystore with real Ed25519 + X25519 key material."""
sk_x = X25519PrivateKey.generate()
@@ -35,23 +34,33 @@ def _mock_keystore_keys(sk_ed):
mock_keys.pk_x25519_b64 = base64.b64encode(pk_x_raw).decode()
return mock_keys
+def _free_port() -> int:
+ """
+ A port nobody else in the session is on.
+
+ These tests start the real admin UI server. Hardcoding 28000 made them fail
+ with EADDRINUSE whenever another test file had a node running — which is why
+ the full suite failed while each file passed on its own.
+ """
+ import socket
+ with socket.socket() as s:
+ s.bind(("127.0.0.1", 0))
+ return s.getsockname()[1]
+
@pytest.fixture
def sk_hub():
return Ed25519PrivateKey.generate()
-
@pytest.fixture
def hub_pk_pem(sk_hub):
return sk_hub.public_key().public_bytes(
serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo)
-
@pytest.fixture
def gek():
return generate_gek()
-
@pytest.fixture
def shared_dir(tmp_path):
d = tmp_path / "shared"
@@ -60,26 +69,22 @@ def shared_dir(tmp_path):
(d / "hello.txt").write_bytes(b"hello daemon test " * 50)
return d
-
@pytest.fixture
def node_config(tmp_path, shared_dir):
return Config(
hub=HubConfig(url="http://localhost:9999", username="testuser"),
- node=NodeConfig(port=29000, quic_port=29010, http_port=29001, ui_port=28000),
+ node=NodeConfig(quic_port=_free_port(), ui_port=_free_port()),
groups=[GroupConfig(
id="g" * 32,
name="test-group",
shared_dir=str(shared_dir),
visibility="private",
- port=29000,
quic_port=29010,
- http_port=29001,
)],
keystore=KeystoreConfig(path=tmp_path / "keystore.enc"),
data_dir=tmp_path / "data",
)
-
@pytest.mark.asyncio
async def test_daemon_creates_chat_store(tmp_path, node_config, gek, hub_pk_pem):
"""Daemon creates ChatStore for each group and shuts down cleanly."""
@@ -130,7 +135,14 @@ async def test_daemon_creates_chat_store(tmp_path, node_config, gek, hub_pk_pem)
assert daemon._chat_stores[group_id]._db is not None
if daemon._webrtc:
- assert "chat_store" in daemon._webrtc._ctx
+ # Finding H1: chat_store must live in the per-group context, never on
+ # the shared transport context. Hoisting the first group's store
+ # transport-wide sent every group's chat to one database and served it
+ # back to members of every other group.
+ assert "chat_store" not in daemon._webrtc._ctx
+ groups_ctx = daemon._webrtc._ctx["groups"]
+ assert groups_ctx[group_id]["chat_store"] is daemon._chat_stores[group_id]
+
assert "hub_ws" in daemon._webrtc._ctx
assert "node_user_id" in daemon._webrtc._ctx
assert daemon._webrtc._ctx["node_user_id"] == "user123"
@@ -146,7 +158,6 @@ async def test_daemon_creates_chat_store(tmp_path, node_config, gek, hub_pk_pem)
for store in daemon._chat_stores.values():
assert store._db is None
-
@pytest.mark.asyncio
async def test_daemon_no_groups_exits(tmp_path):
"""Daemon with no valid groups exits cleanly."""
@@ -188,19 +199,18 @@ async def test_daemon_no_groups_exits(tmp_path):
assert len(daemon._chat_stores) == 0
-
@pytest.mark.asyncio
async def test_daemon_index_change_pushes_to_peers(tmp_path, shared_dir, gek, hub_pk_pem):
"""Index change callback pushes updated index to WebRTC peers."""
config = Config(
hub=HubConfig(url="http://localhost:9999", username="testuser"),
- node=NodeConfig(port=29000, quic_port=29010, http_port=29001, ui_port=28000),
+ node=NodeConfig(quic_port=_free_port(), ui_port=_free_port()),
groups=[GroupConfig(
id="a" * 32,
name="test-group",
shared_dir=str(shared_dir),
visibility="private",
- port=29000, quic_port=29010, http_port=29001,
+ quic_port=29010,
)],
keystore=KeystoreConfig(path=tmp_path / "keystore.enc"),
data_dir=tmp_path / "data",
@@ -232,10 +242,45 @@ async def test_daemon_index_change_pushes_to_peers(tmp_path, shared_dir, gek, hu
assert msg["group_id"] == "a" * 32
assert len(msg["entries"]) == indexer.index.count
+ # Finding H7: this group is private, so its content hashes must NOT be
+ # registered with the hub. The test previously asserted the opposite —
+ # publishing a fingerprint of every private file was treated as expected
+ # behaviour. Index push to members is unaffected (asserted above).
+ await asyncio.sleep(0.1)
+ daemon._hub.register_swarm.assert_not_called()
+
+@pytest.mark.asyncio
+async def test_daemon_index_change_registers_swarm_for_public_group(
+ tmp_path, shared_dir, gek, hub_pk_pem):
+ """Public groups still register content hashes with the hub swarm (H7)."""
+ config = Config(
+ hub=HubConfig(url="http://localhost:9999", username="testuser"),
+ node=NodeConfig(quic_port=_free_port(), ui_port=_free_port()),
+ groups=[GroupConfig(
+ id="a" * 32,
+ name="public-group",
+ shared_dir=str(shared_dir),
+ visibility="public",
+ quic_port=29010,
+ )],
+ keystore=KeystoreConfig(path=tmp_path / "keystore.enc"),
+ data_dir=tmp_path / "data",
+ )
+ daemon = NodeDaemon(config)
+ daemon._hub = AsyncMock()
+ daemon._hub.register_swarm = AsyncMock(return_value=2)
+ daemon._state["endpoint_hint"] = "node123"
+
+ indexer = DirectoryIndexer(
+ root=shared_dir, group_id="a" * 32,
+ sk_node=Ed25519PrivateKey.generate(), gek=gek)
+ await indexer.initial_scan()
+
+ await daemon._on_index_change(indexer)
+
await asyncio.sleep(0.1)
daemon._hub.register_swarm.assert_called_once()
- call_args = daemon._hub.register_swarm.call_args
- assert len(call_args[0][0]) == indexer.index.count
+ assert len(daemon._hub.register_swarm.call_args[0][0]) == indexer.index.count
@pytest.mark.asyncio
diff --git a/packages/meshbay-node/tests/test_http_server.py b/packages/meshbay-node/tests/test_http_server.py
deleted file mode 100644
index d4ccc32..0000000
--- a/packages/meshbay-node/tests/test_http_server.py
+++ /dev/null
@@ -1,227 +0,0 @@
-"""Tests for the node HTTP file API."""
-
-import asyncio
-import base64
-import json
-import os
-import time
-import pytest
-import jwt
-import httpx
-from pathlib import Path
-from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
-from cryptography.hazmat.primitives import serialization
-
-from meshbay_common.crypto import generate_gek, pk_to_b64
-from meshbay_node.indexer import DirectoryIndexer
-from meshbay_node.transport.http_server import create_http_app
-
-
-@pytest.fixture
-def sk_node():
- return Ed25519PrivateKey.generate()
-
-@pytest.fixture
-def sk_hub():
- return Ed25519PrivateKey.generate()
-
-@pytest.fixture
-def hub_pk_pem(sk_hub):
- return sk_hub.public_key().public_bytes(
- serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo)
-
-@pytest.fixture
-def gek():
- return generate_gek()
-
-@pytest.fixture
-def shared_dir(tmp_path):
- d = tmp_path / "shared"
- d.mkdir()
- (d / "video.mp4").write_bytes(os.urandom(3 * 1024 * 1024)) # 3MB
- (d / "doc.pdf").write_bytes(os.urandom(512 * 1024))
- (d / "song.mp3").write_bytes(os.urandom(256 * 1024))
- return d
-
-def make_token(sk_hub, pk_node_b64, ttl=3600):
- sk_pem = sk_hub.private_bytes(
- serialization.Encoding.PEM, serialization.PrivateFormat.PKCS8,
- serialization.NoEncryption())
- now = int(time.time())
- return jwt.encode({
- "iss": "test-hub", "sub": "user-001",
- "pk_user": pk_node_b64, "hub_id": "test-hub",
- "jti": "test-jti", "iat": now, "exp": now + ttl,
- }, sk_pem, algorithm="EdDSA")
-
-
-@pytest.mark.asyncio
-async def test_node_info(sk_node, sk_hub, hub_pk_pem, gek, shared_dir):
- indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek)
- await indexer.initial_scan()
-
- app = create_http_app(
- sk_node=sk_node, hub_pk_pem=hub_pk_pem,
- shared_root=shared_dir, index=indexer.index,
- group_id="test-group", group_name="Test Group",
- )
- async with httpx.AsyncClient(
- transport=httpx.ASGITransport(app=app), base_url="http://test"
- ) as c:
- r = await c.get("/")
- assert r.status_code == 200
- data = r.json()
- assert data["group_id"] == "test-group"
- assert data["file_count"] == 3
- assert "pk_node" in data
-
-
-@pytest.mark.asyncio
-async def test_public_index(sk_node, sk_hub, hub_pk_pem, shared_dir):
- """Public group: index accessible without auth."""
- indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None)
- await indexer.initial_scan()
-
- app = create_http_app(
- sk_node=sk_node, hub_pk_pem=hub_pk_pem,
- shared_root=shared_dir, index=indexer.index,
- group_id="pub-group", group_name="Public Group",
- gek=None,
- )
- async with httpx.AsyncClient(
- transport=httpx.ASGITransport(app=app), base_url="http://test"
- ) as c:
- r = await c.get("/index")
- assert r.status_code == 200
- data = r.json()
- assert len(data["entries"]) == 3
- names = {e["name"] for e in data["entries"]}
- assert "video.mp4" in names
- assert "doc.pdf" in names
-
-
-@pytest.mark.asyncio
-async def test_file_download(sk_node, sk_hub, hub_pk_pem, shared_dir):
- """Full file download via HTTP."""
- indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None)
- await indexer.initial_scan()
-
- app = create_http_app(
- sk_node=sk_node, hub_pk_pem=hub_pk_pem,
- shared_root=shared_dir, index=indexer.index,
- group_id="g", group_name="G",
- )
- entry = next(e for e in indexer.index.entries if e.name == "doc.pdf")
- original = (shared_dir / "doc.pdf").read_bytes()
-
- async with httpx.AsyncClient(
- transport=httpx.ASGITransport(app=app), base_url="http://test"
- ) as c:
- r = await c.get(f"/file/{entry.id}")
- assert r.status_code == 200
- assert r.content == original
-
-
-@pytest.mark.asyncio
-async def test_chunk_public_group(sk_node, sk_hub, hub_pk_pem, shared_dir):
- """Public group chunk: plaintext, signed, auth required."""
- indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None)
- await indexer.initial_scan()
-
- app = create_http_app(
- sk_node=sk_node, hub_pk_pem=hub_pk_pem,
- shared_root=shared_dir, index=indexer.index,
- group_id="g", group_name="G", gek=None,
- )
- entry = next(e for e in indexer.index.entries if e.name == "video.mp4")
- token = make_token(sk_hub, pk_to_b64(sk_node.public_key()))
-
- async with httpx.AsyncClient(
- transport=httpx.ASGITransport(app=app), base_url="http://test"
- ) as c:
- r = await c.get(f"/file/{entry.id}/0",
- headers={"Authorization": f"Bearer {token}"})
- assert r.status_code == 200
- chunk = r.json()
- assert chunk["encrypted"] is False
- assert chunk["chunk_index"] == 0
- assert "data_b64" in chunk
-
- # Verify the chunk data matches original
- original = (shared_dir / "video.mp4").read_bytes()
- data = base64.b64decode(chunk["data_b64"])
- assert data == original[:len(data)]
-
-
-@pytest.mark.asyncio
-async def test_chunk_private_group(sk_node, sk_hub, hub_pk_pem, gek, shared_dir):
- """Private group chunk: encrypted with GEK."""
- from meshbay_common.webcrypto import chunk_key_aes as derive_chunk_key, decrypt_chunk_aes as decrypt_chunk
- import blake3
-
- indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek)
- await indexer.initial_scan()
-
- app = create_http_app(
- sk_node=sk_node, hub_pk_pem=hub_pk_pem,
- shared_root=shared_dir, index=indexer.index,
- group_id="g", group_name="G", gek=gek,
- )
- entry = next(e for e in indexer.index.entries if e.name == "doc.pdf")
- token = make_token(sk_hub, pk_to_b64(sk_node.public_key()))
-
- async with httpx.AsyncClient(
- transport=httpx.ASGITransport(app=app), base_url="http://test"
- ) as c:
- r = await c.get(f"/file/{entry.id}/0",
- headers={"Authorization": f"Bearer {token}"})
- assert r.status_code == 200
- chunk = r.json()
- assert chunk["encrypted"] is True
-
- # Decrypt and verify
- file_hash = base64.b64decode(chunk["file_hash_b64"])
- nonce = base64.b64decode(chunk["nonce_b64"])
- ct = base64.b64decode(chunk["ct_b64"])
- ckey = derive_chunk_key(gek, file_hash, 0)
- plaintext = decrypt_chunk(ckey, nonce, ct)
- original = (shared_dir / "doc.pdf").read_bytes()
- assert plaintext == original[:len(plaintext)]
-
-
-@pytest.mark.asyncio
-async def test_chunk_requires_auth(sk_node, hub_pk_pem, shared_dir):
- """Chunk endpoint rejects unauthenticated requests."""
- indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None)
- await indexer.initial_scan()
- app = create_http_app(
- sk_node=sk_node, hub_pk_pem=hub_pk_pem,
- shared_root=shared_dir, index=indexer.index,
- group_id="g", group_name="G",
- )
- entry = indexer.index.entries[0]
-
- async with httpx.AsyncClient(
- transport=httpx.ASGITransport(app=app), base_url="http://test"
- ) as c:
- r = await c.get(f"/file/{entry.id}/0") # no token
- assert r.status_code == 401
-
-
-@pytest.mark.asyncio
-async def test_unknown_file_404(sk_node, hub_pk_pem, sk_hub, shared_dir):
- indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None)
- await indexer.initial_scan()
- app = create_http_app(
- sk_node=sk_node, hub_pk_pem=hub_pk_pem,
- shared_root=shared_dir, index=indexer.index,
- group_id="g", group_name="G",
- )
- token = make_token(sk_hub, pk_to_b64(sk_node.public_key()))
-
- async with httpx.AsyncClient(
- transport=httpx.ASGITransport(app=app), base_url="http://test"
- ) as c:
- r = await c.get("/file/nonexistent-hash/0",
- headers={"Authorization": f"Bearer {token}"})
- assert r.status_code == 404
diff --git a/packages/meshbay-node/tests/test_quic_transport.py b/packages/meshbay-node/tests/test_quic_transport.py
index 0c1a1cd..93ab1b0 100644
--- a/packages/meshbay-node/tests/test_quic_transport.py
+++ b/packages/meshbay-node/tests/test_quic_transport.py
@@ -49,7 +49,8 @@ def make_jwt(sk_hub, pk_node_b64, ttl=3600, groups=None):
"iss": "test-hub", "sub": "user-001",
"pk_user": pk_node_b64, "hub_id": "test-hub",
"jti": "test-jti", "iat": now, "exp": now + ttl,
- "groups": groups or [],
+ # group_id is mandatory (M1), so default tokens are members of "g".
+ "groups": groups if groups is not None else ["g"],
}, sk_pem, algorithm="EdDSA")
@@ -80,6 +81,7 @@ async def test_quic_chunk_roundtrip(sk_node, sk_hub, gek, shared_dir, tmp_path):
host="127.0.0.1", port=19100,
jwt_token=token, gek=gek,
pk_node_b64=pk_to_b64(sk_node.public_key()),
+ group_id="g",
) as client:
chunk0 = await client.fetch_chunk(entry.id, chunk_index=0)
chunk1 = await client.fetch_chunk(entry.id, chunk_index=1)
@@ -116,6 +118,7 @@ async def test_quic_fetch_index(sk_node, sk_hub, gek, shared_dir, tmp_path):
host="127.0.0.1", port=19101,
jwt_token=token, gek=gek,
pk_node_b64=pk_to_b64(sk_node.public_key()),
+ group_id="g",
) as client:
wire = await client.fetch_index()
recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek)
@@ -220,10 +223,12 @@ async def test_quic_session_resumption(sk_node, sk_hub, gek, shared_dir, tmp_pat
async with QuicChunkClient(
host="127.0.0.1", port=19104,
jwt_token=token, gek=gek, pk_node_b64=pk_b64,
+ group_id="g",
) as client:
wire = await client.fetch_index()
assert GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek).count == 2
saved_ticket = client.session_ticket
+ saved_cert = client.peer_cert_der
# Allow server to process the close
await asyncio.sleep(0.1)
@@ -233,6 +238,10 @@ async def test_quic_session_resumption(sk_node, sk_hub, gek, shared_dir, tmp_pat
host="127.0.0.1", port=19104,
jwt_token=token, gek=gek, pk_node_b64=pk_b64,
session_ticket=saved_ticket,
+ # A resumed session carries no certificate, so the binding anchor from the
+ # original handshake travels with the ticket (11.5.6).
+ peer_cert_der=saved_cert,
+ group_id="g",
) as client:
wire = await client.fetch_index()
assert GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek).count == 2
@@ -269,6 +278,7 @@ async def test_quic_denylist_blocks_user(sk_node, sk_hub, gek, shared_dir, tmp_p
async with QuicChunkClient(
host="127.0.0.1", port=19105,
jwt_token=token, gek=gek, pk_node_b64=pk_b64,
+ group_id="g",
) as client:
wire = await client.fetch_index()
assert GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek).count == 2
@@ -281,6 +291,7 @@ async def test_quic_denylist_blocks_user(sk_node, sk_hub, gek, shared_dir, tmp_p
async with QuicChunkClient(
host="127.0.0.1", port=19105,
jwt_token=token, gek=gek, pk_node_b64=pk_b64,
+ group_id="g",
) as client:
await client.fetch_index()
diff --git a/packages/meshbay-node/tests/test_roster_pairing.py b/packages/meshbay-node/tests/test_roster_pairing.py
new file mode 100644
index 0000000..665c060
--- /dev/null
+++ b/packages/meshbay-node/tests/test_roster_pairing.py
@@ -0,0 +1,744 @@
+"""
+Roster and operator pairing (M3, and the mechanism that will close H3).
+
+Negative assertions, per the posture set in Phase 11.5: each test states an attack
+or a mistake that must not work. The one to keep an eye on is
+`test_daemon_does_not_auto_pin_keystore_key` — the auto-pin is what made node
+sovereignty inert as shipped, and it fails closed, so nothing else in the suite
+notices if it comes back.
+
+See `docs/invite-pairing-v1.md`.
+"""
+
+import base64
+import time
+from pathlib import Path
+
+import pytest
+from cryptography.hazmat.primitives import serialization
+from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
+from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey
+
+from meshbay_common.crypto import generate_gek, pk_to_b64, unwrap_gek_aes
+from meshbay_common.join import ROLE_MEMBER, ROLE_OPERATOR, join_transcript
+from meshbay_node.indexer.group_index import GroupIndex
+from meshbay_node.roster import Roster, hash_code, normalize_code
+from meshbay_node.transport.webrtc_server import WebRTCPeerSession
+
+
+# ── Fixtures ──────────────────────────────────────────────────────────────────
+
+@pytest.fixture
+async def roster(tmp_path):
+ r = Roster(db_path=tmp_path / "roster.db")
+ await r.open()
+ yield r
+ await r.close()
+
+
+def _keypair_full():
+ """(sk_ed, pk_ed_b64, pk_x_b64, sk_x) — the X25519 secret is needed to unwrap."""
+ sk_ed = Ed25519PrivateKey.generate()
+ sk_x = X25519PrivateKey.generate()
+ pk_ed_b64 = pk_to_b64(sk_ed.public_key())
+ pk_x_b64 = base64.b64encode(
+ sk_x.public_key().public_bytes(
+ encoding=serialization.Encoding.Raw,
+ format=serialization.PublicFormat.Raw,
+ )
+ ).decode()
+ return sk_ed, pk_ed_b64, pk_x_b64, sk_x
+
+
+def _keypair():
+ sk_ed, pk_ed_b64, pk_x_b64, _ = _keypair_full()
+ return sk_ed, pk_ed_b64, pk_x_b64
+
+
+def _session(tmp_path: Path, roster, user_id: str = "grenet",
+ group_id: str | None = None, gek: bytes | None = None,
+ join_policy: str = "invite") -> WebRTCPeerSession:
+ """A peer session with the join path wired and sending stubbed out."""
+ shared_root = tmp_path / "shared"
+ shared_root.mkdir(exist_ok=True)
+ index = GroupIndex(group_id="g" * 32, sk_node=Ed25519PrivateKey.generate())
+
+ session = WebRTCPeerSession.__new__(WebRTCPeerSession)
+ session._ctx = {
+ "shared_root": shared_root,
+ "index": index,
+ "sk_node": index.sk_node,
+ "roster": roster,
+ }
+ if group_id:
+ session._ctx["groups"] = {
+ group_id: {
+ "gek": gek,
+ "shared_root": shared_root,
+ "index": index,
+ "join_policy": join_policy,
+ },
+ }
+ session._group_id = group_id
+ session._user_id = user_id
+ session._username = user_id
+ session._pk_user = ""
+ session._uploads = {}
+ session._join_attempts = 0
+ session._nonce_node = b"\x11" * 32
+ session._remote_ip = ""
+ session.sent = []
+ session._send = session.sent.append
+ session._audit = lambda *a, **k: None
+ return session
+
+
+def _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64, code="", user_id="grenet",
+ group_id="", nonce=None, ts=None):
+ ts = int(time.time()) if ts is None else ts
+ transcript = join_transcript(
+ node_pk_b64=session._node_pk_b64(),
+ group_id=group_id,
+ user_id=user_id,
+ pk_ed25519_b64=pk_ed_b64,
+ pk_x25519_b64=pk_x_b64,
+ nonce_node=nonce if nonce is not None else session._nonce_node,
+ ts=ts,
+ )
+ return {
+ "type": "join_request",
+ "group_id": group_id,
+ "pk_ed25519": pk_ed_b64,
+ "pk_x25519": pk_x_b64,
+ "code": code,
+ "ts": ts,
+ "sig": base64.b64encode(sk_ed.sign(transcript)).decode(),
+ }
+
+
+def _last(session):
+ return session.sent[-1] if session.sent else {}
+
+
+# ── Roster ────────────────────────────────────────────────────────────────────
+
+async def test_invite_is_single_use(roster):
+ code = await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli")
+ assert await roster.consume_invite(code, "grenet") is not None
+ assert await roster.consume_invite(code, "grenet") is None, (
+ "a pairing code must not be redeemable twice")
+
+
+async def test_invite_is_bound_to_one_account(roster):
+ """A leaked code must be useless to whoever finds it."""
+ code = await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli")
+ assert await roster.consume_invite(code, "eve") is None
+ assert await roster.consume_invite(code, "grenet") is not None
+
+
+async def test_expired_invite_is_refused(roster):
+ code = await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli", ttl=-1)
+ assert await roster.consume_invite(code, "grenet") is None
+
+
+async def test_reinvite_supersedes_the_previous_code(roster):
+ first = await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli")
+ second = await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli")
+ assert await roster.consume_invite(first, "grenet") is None
+ assert await roster.consume_invite(second, "grenet") is not None
+
+
+async def test_codes_are_not_stored_in_the_clear(roster, tmp_path):
+ code = await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli")
+ rows = await roster.list_invites()
+ assert rows and rows[0]["code_hash"] != normalize_code(code)
+ assert rows[0]["code_hash"] == hash_code(code)
+
+
+def test_code_normalization_absorbs_human_error():
+ """Someone reading a code aloud must not be able to get it wrong."""
+ assert normalize_code("k7m2-qx4p") == normalize_code("K7M2QX4P")
+ assert normalize_code("O1IL") == "0111"
+ assert normalize_code(" k7m2 qx4p ") == "K7M2QX4P"
+
+
+async def test_operator_pks_reflect_unpinning(roster):
+ _, pk_ed_b64, pk_x_b64 = _keypair()
+ await roster.pin_identity("grenet", "grenet", pk_ed_b64, pk_x_b64, "code")
+ await roster.set_member("", "grenet", ROLE_OPERATOR, "active", "local-cli")
+ assert await roster.operator_pks() == [pk_ed_b64]
+
+ await roster.unpin("grenet")
+ assert await roster.operator_pks() == [], (
+ "authority must disappear with the pin, without a daemon restart")
+
+
+# ── Join / pairing over MNP ───────────────────────────────────────────────────
+
+async def test_pairing_with_a_valid_code_pins_the_identity(tmp_path, roster):
+ session = _session(tmp_path, roster)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+ code = await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli")
+
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64, code=code))
+
+ assert _last(session).get("ok") is True
+ pinned = await roster.get_identity("grenet")
+ assert pinned["pk_ed25519"] == pk_ed_b64
+ assert await roster.operator_pks() == [pk_ed_b64]
+
+
+async def test_pairing_without_a_code_is_refused(tmp_path, roster):
+ """Fails closed: an unknown identity gets nothing until someone authorizes it."""
+ session = _session(tmp_path, roster)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+
+ await session._do_join_request(_join_msg(session, sk_ed, pk_ed_b64, pk_x_b64))
+
+ assert _last(session).get("ok") is False
+ assert _last(session).get("reason") == "code_required"
+ assert await roster.get_identity("grenet") is None
+
+
+async def test_wrong_code_pins_nothing(tmp_path, roster):
+ session = _session(tmp_path, roster)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+ await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli")
+
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64, code="ZZZZ-ZZZZ"))
+
+ assert _last(session).get("reason") == "code_invalid"
+ assert await roster.get_identity("grenet") is None
+
+
+async def test_signature_must_cover_the_presented_keys(tmp_path, roster):
+ """
+ The heart of it: the X25519 key is only trustworthy because the Ed25519
+ identity signed it. Swapping in another encryption key after signing must fail.
+ """
+ session = _session(tmp_path, roster)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+ code = await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli")
+
+ msg = _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64, code=code)
+ _, _, attacker_pk_x = _keypair()
+ msg["pk_x25519"] = attacker_pk_x
+
+ await session._do_join_request(msg)
+
+ assert _last(session).get("reason") == "signature_invalid"
+ assert await roster.get_identity("grenet") is None
+
+
+async def test_join_cannot_be_replayed_onto_another_connection(tmp_path, roster):
+ session = _session(tmp_path, roster)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+ code = await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli")
+
+ # Signed against a nonce this connection never issued.
+ msg = _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64, code=code,
+ nonce=b"\x99" * 32)
+ await session._do_join_request(msg)
+
+ assert _last(session).get("reason") == "signature_invalid"
+ assert await roster.get_identity("grenet") is None
+
+
+async def test_pinned_identity_presenting_a_new_key_is_refused(tmp_path, roster):
+ """
+ 11.5.8's rule, applied to people: a changed key is refused outright rather
+ than warned about, and clearing it is a deliberate operator action.
+ """
+ session = _session(tmp_path, roster)
+ _, old_pk_ed, old_pk_x = _keypair()
+ await roster.pin_identity("grenet", "grenet", old_pk_ed, old_pk_x, "code")
+
+ sk_ed2, new_pk_ed, new_pk_x = _keypair()
+ await session._do_join_request(
+ _join_msg(session, sk_ed2, new_pk_ed, new_pk_x, code="ANY-CODE"))
+
+ assert _last(session).get("reason") == "key_changed"
+ assert (await roster.get_identity("grenet"))["pk_ed25519"] == old_pk_ed
+
+
+async def test_attempts_are_bounded(tmp_path, roster):
+ session = _session(tmp_path, roster)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+ await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli")
+
+ for _ in range(6):
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64, code="AAAA-AAAA"))
+
+ assert any(m.get("detail") == "Too many attempts" for m in session.sent), (
+ "a connection must not be able to sit there guessing codes")
+
+
+async def test_failures_are_counted_across_connections(tmp_path, roster):
+ """
+ The adversary who can mint a token for any account is the hub, and it can
+ reconnect at will — so a per-connection budget alone would bound nothing.
+ """
+ shared_ctx = None
+ for _ in range(6):
+ session = _session(tmp_path, roster)
+ if shared_ctx is None:
+ shared_ctx = session._ctx
+ else:
+ session._ctx = shared_ctx # same node, new connection
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+ for _ in range(4):
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64, code="AAAA-AAAA"))
+
+ assert any(m.get("detail") == "Pairing temporarily locked"
+ for m in session.sent), (
+ "reconnecting must not reset the pairing budget")
+
+
+async def test_group_id_cannot_name_another_group(tmp_path, roster):
+ session = _session(tmp_path, roster)
+ session._group_id = "a" * 32
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64, group_id="b" * 32))
+
+ assert _last(session).get("reason") == "group_mismatch"
+
+
+# ── H3: the node wraps the group key, and only for people it admitted ─────────
+
+GROUP = "g" * 32
+
+
+async def test_node_wraps_the_gek_for_the_key_the_member_proved(tmp_path, roster):
+ """
+ The H3 fix. Nobody fetches a public key from the hub: the node encrypts the
+ group key for the X25519 key the joiner signed with their pinned identity, so
+ a hub substituting a key of its own has nothing to substitute into.
+ """
+ gek = generate_gek()
+ session = _session(tmp_path, roster, user_id="bob", group_id=GROUP, gek=gek)
+ sk_ed, pk_ed_b64, pk_x_b64, sk_x = _keypair_full()
+
+ code = await roster.create_invite(GROUP, "bob", ROLE_MEMBER, "grenet")
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64, code=code,
+ user_id="bob", group_id=GROUP))
+
+ reply = _last(session)
+ assert reply["ok"] is True and reply["gek"] is True
+
+ pk_x_raw = base64.b64decode(pk_x_b64)
+ sk_x_raw = sk_x.private_bytes(
+ encoding=serialization.Encoding.Raw,
+ format=serialization.PrivateFormat.Raw,
+ encryption_algorithm=serialization.NoEncryption(),
+ )
+ assert unwrap_gek_aes(reply, sk_x_raw, pk_x_raw) == gek
+
+
+async def test_hub_membership_alone_yields_no_key(tmp_path, roster):
+ """
+ A hub can invent an account, add it to a group and mint it a token. What it
+ cannot do is put it on the node's roster — so the key never leaves.
+ """
+ gek = generate_gek()
+ session = _session(tmp_path, roster, user_id="eve", group_id=GROUP, gek=gek)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+
+ # Pinned on this node (say, for another group) but never admitted to this one.
+ await roster.pin_identity("eve", "eve", pk_ed_b64, pk_x_b64, "code")
+
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64,
+ user_id="eve", group_id=GROUP))
+
+ reply = _last(session)
+ assert reply.get("gek") is False
+ assert reply.get("reason") == "not_authorized_for_group"
+ assert "wrapped_b64" not in reply
+
+
+async def test_open_join_group_admits_without_a_code(tmp_path, roster):
+ """§3.4: where anyone may join, a code protects nothing and is not required."""
+ gek = generate_gek()
+ session = _session(tmp_path, roster, user_id="newcomer", group_id=GROUP,
+ gek=gek, join_policy="open")
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64,
+ user_id="newcomer", group_id=GROUP))
+
+ reply = _last(session)
+ assert reply["ok"] is True and reply["gek"] is True
+ pinned = await roster.get_identity("newcomer")
+ assert pinned["pinned_via"] == "tofu"
+
+
+async def test_invite_only_group_still_demands_a_code(tmp_path, roster):
+ """Being public (discoverable) is not being open (admitting anyone)."""
+ gek = generate_gek()
+ session = _session(tmp_path, roster, user_id="newcomer", group_id=GROUP,
+ gek=gek, join_policy="invite")
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64,
+ user_id="newcomer", group_id=GROUP))
+
+ assert _last(session).get("reason") == "code_required"
+ assert await roster.get_identity("newcomer") is None
+
+
+async def test_unknown_group_is_invite_only(tmp_path, roster):
+ """
+ Fail closed: a group whose policy the node cannot read is treated as
+ invite-only, never as open.
+ """
+ session = _session(tmp_path, roster, user_id="newcomer")
+ session._group_id = "unconfigured-group"
+ assert session._group_join_policy("unconfigured-group") == "invite"
+ assert session._group_join_policy("") == "invite"
+
+
+def test_join_policy_is_carried_from_node_config():
+ """
+ The policy reaches the transport from node.toml. If it ever came from the hub
+ instead, a hub could declare any group open and be handed its key.
+ """
+ daemon_src = (Path(__file__).parent.parent
+ / "src" / "meshbay_node" / "daemon.py").read_text()
+ assert '"join_policy": group_cfg.join_policy' in daemon_src
+
+ config_src = (Path(__file__).parent.parent
+ / "src" / "meshbay_node" / "config.py").read_text()
+ assert "join_policy" in config_src, "GroupConfig must carry the admission policy"
+
+
+async def test_revoked_member_stops_receiving_the_key(tmp_path, roster):
+ """
+ Wrapping on demand is what makes revocation work. A stored bundle survived
+ revocation; this does not. (Rotating the GEK is still required — the
+ ex-member has the old one.)
+ """
+ gek = generate_gek()
+ session = _session(tmp_path, roster, user_id="bob", group_id=GROUP, gek=gek)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+ await roster.pin_identity("bob", "bob", pk_ed_b64, pk_x_b64, "code")
+ await roster.set_member(GROUP, "bob", ROLE_MEMBER, "active", "grenet")
+
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64,
+ user_id="bob", group_id=GROUP))
+ assert _last(session)["gek"] is True
+
+ await roster.set_status(GROUP, "bob", "revoked")
+ session.sent.clear()
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64,
+ user_id="bob", group_id=GROUP))
+ assert _last(session).get("gek") is False
+
+
+# ── What a first-time joiner can know ─────────────────────────────────────────
+
+def test_challenge_carries_node_pk_in_source():
+ """
+ Belt and braces for the above: the field must be in the message the node
+ builds, whatever the surrounding handshake does.
+ """
+ source = (Path(__file__).parent.parent
+ / "src" / "meshbay_node" / "transport" / "webrtc_server.py").read_text()
+ challenge = source[source.find("MNP.HANDSHAKE_CHALLENGE,"):]
+ challenge = challenge[:challenge.find("})")]
+ assert "node_pk" in challenge, (
+ "the challenge must announce the node key — a first-time joiner cannot "
+ "learn it any other way, and join_request signs it")
+
+
+async def test_a_key_pinned_by_one_node_is_worthless_at_another(tmp_path, roster):
+ """
+ The whole point of per-node identity: node A's operator who cracks the bundle
+ on their own disk holds a key node B has never seen. Presenting it there is a
+ first contact like any other — it needs a code from B's operator.
+ """
+ gek = generate_gek()
+ node_b = _session(tmp_path, roster, user_id="bob", group_id=GROUP, gek=gek)
+
+ # The key bob uses at node A. Node B's roster knows nothing about it.
+ sk_ed_a, pk_ed_a, pk_x_a = _keypair()
+
+ await node_b._do_join_request(
+ _join_msg(node_b, sk_ed_a, pk_ed_a, pk_x_a,
+ user_id="bob", group_id=GROUP))
+
+ assert _last(node_b).get("reason") == "code_required"
+ assert await roster.get_identity("bob") is None
+
+
+async def test_the_stolen_key_cannot_be_forced_in_with_someone_elses_code(
+ tmp_path, roster):
+ """And a code issued for another account does not help either."""
+ gek = generate_gek()
+ session = _session(tmp_path, roster, user_id="eve", group_id=GROUP, gek=gek)
+ sk_ed, pk_ed, pk_x = _keypair()
+ code = await roster.create_invite(GROUP, "bob", ROLE_MEMBER, "grenet")
+
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed, pk_x, code=code,
+ user_id="eve", group_id=GROUP))
+
+ assert _last(session).get("reason") == "code_invalid"
+ assert await roster.get_identity("eve") is None
+
+
+# ── Code lifetimes ────────────────────────────────────────────────────────────
+
+async def test_invitations_outlive_pairing_codes(roster):
+ """
+ An invitation crosses a human conversation; a pairing code crosses an SSH
+ session. A day was long enough for the second and not for the first — a code
+ that dies over a weekend means someone has to be at a browser to reissue it.
+ """
+ from meshbay_node.roster import DEFAULT_INVITE_TTL, DEFAULT_PAIR_TTL
+
+ assert DEFAULT_INVITE_TTL == 7 * 24 * 3600
+ assert DEFAULT_PAIR_TTL == 24 * 3600
+ assert DEFAULT_INVITE_TTL > DEFAULT_PAIR_TTL
+
+
+def test_code_lifetimes_are_configurable(tmp_path):
+ """The operator decides, not the default."""
+ from meshbay_node.config import load_config
+
+ path = tmp_path / "node.toml"
+ path.write_text(
+ '[hub]\nurl = "https://example.org"\nusername = "grenet"\n'
+ "[node]\ninvite_ttl_hours = 72\npair_ttl_hours = 2\n"
+ )
+ cfg = load_config(path)
+ assert cfg.node.invite_ttl_hours == 72
+ assert cfg.node.pair_ttl_hours == 2
+
+ default = load_config(tmp_path / "missing.toml")
+ assert default.node.invite_ttl_hours == 168
+ assert default.node.pair_ttl_hours == 24
+
+
+async def test_expiry_is_enforced_at_redemption(tmp_path, roster):
+ """Purging is housekeeping; the check that matters happens on use."""
+ session = _session(tmp_path, roster)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+ code = await roster.create_invite("", "grenet", ROLE_OPERATOR, "local-cli", ttl=-1)
+
+ await session._do_join_request(
+ _join_msg(session, sk_ed, pk_ed_b64, pk_x_b64, code=code))
+
+ assert _last(session).get("reason") == "code_invalid"
+ assert await roster.get_identity("grenet") is None
+
+
+# ── M3: where node authority comes from ───────────────────────────────────────
+
+async def test_admin_signature_verified_against_the_paired_key(tmp_path, roster):
+ session = _session(tmp_path, roster)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+ await roster.pin_identity("grenet", "grenet", pk_ed_b64, pk_x_b64, "code")
+ await roster.set_member("", "grenet", ROLE_OPERATOR, "active", "local-cli")
+
+ transcript = b"meshbay:admin:v1 whatever"
+ assert await session._verify_admin_sig(transcript, sk_ed.sign(transcript))
+
+ stranger = Ed25519PrivateKey.generate()
+ assert not await session._verify_admin_sig(
+ transcript, stranger.sign(transcript))
+
+
+async def test_unpinned_operator_loses_authority_immediately(tmp_path, roster):
+ """No caching: revoking a paired browser must not need a daemon restart."""
+ session = _session(tmp_path, roster)
+ sk_ed, pk_ed_b64, pk_x_b64 = _keypair()
+ await roster.pin_identity("grenet", "grenet", pk_ed_b64, pk_x_b64, "code")
+ await roster.set_member("", "grenet", ROLE_OPERATOR, "active", "local-cli")
+
+ transcript = b"meshbay:admin:v1 whatever"
+ assert await session._verify_admin_sig(transcript, sk_ed.sign(transcript))
+
+ await roster.unpin("grenet")
+ assert not await session._verify_admin_sig(transcript, sk_ed.sign(transcript))
+
+
+# ── Operator surface (slice 3) ────────────────────────────────────────────────
+
+def _ui_client(tmp_path, roster, **extra):
+ from fastapi.testclient import TestClient
+
+ from meshbay_node.config import Config
+ from meshbay_node.ui.app import create_ui_app
+
+ state = {
+ "status": "running", "groups_ctx": {GROUP: {"gek": b"k" * 32}},
+ "indexes": {}, "ui_token": "tok", "roster": roster,
+ "node_user_id": "grenet", "config": Config(),
+ }
+ state.update(extra)
+ return TestClient(create_ui_app(state)), state
+
+
+async def test_revoke_endpoint_stops_authorization(tmp_path, roster):
+ client, _ = _ui_client(tmp_path, roster)
+ _, pk_ed_b64, pk_x_b64 = _keypair()
+ await roster.pin_identity("bob", "bob", pk_ed_b64, pk_x_b64, "code")
+ await roster.set_member(GROUP, "bob", ROLE_MEMBER, "active", "grenet")
+ assert await roster.is_authorized(GROUP, "bob")
+
+ resp = client.post(f"/api/members/bob/revoke?group_id={GROUP}&t=tok")
+ assert resp.status_code == 200
+ assert "gek-init" in resp.json()["reminder"], (
+ "revocation must remind the operator to rotate the key they still hold")
+ assert not await roster.is_authorized(GROUP, "bob")
+
+
+async def test_unpin_endpoint_allows_repairing(tmp_path, roster):
+ client, _ = _ui_client(tmp_path, roster)
+ _, pk_ed_b64, pk_x_b64 = _keypair()
+ await roster.pin_identity("bob", "bob", pk_ed_b64, pk_x_b64, "code")
+
+ assert client.post("/api/members/bob/unpin?t=tok").status_code == 200
+ assert await roster.get_identity("bob") is None
+ assert client.post("/api/members/bob/unpin?t=tok").status_code == 404
+
+
+async def test_operator_surface_needs_the_session_token(tmp_path, roster):
+ """11.5.3 applies to every one of these: they change who may hold the key."""
+ client, _ = _ui_client(tmp_path, roster)
+ for path in ("/api/roster",
+ "/api/operator/pair",
+ f"/api/members/bob/revoke?group_id={GROUP}",
+ "/api/members/bob/unpin",
+ f"/api/groups/{GROUP}/invites?username=bob"):
+ method = client.get if path == "/api/roster" else client.post
+ assert method(path).status_code == 403, f"{path} reachable without a token"
+
+
+async def test_cli_invite_asks_the_hub_for_an_account_never_a_key(tmp_path, roster):
+ """
+ The CLI resolves a username to an account id through the hub, and stops there.
+ A key fetched from the hub is what H3 was; an account id is not a secret and
+ a wrong one produces an invite whose code the hub never learns.
+ """
+ class _Hub:
+ _session = object()
+
+ async def get_user_pubkeys(self, username):
+ return {"user_id": f"id-of-{username}",
+ "pk_x25519": "SHOULD-NOT-BE-USED",
+ "pk_ed25519": "SHOULD-NOT-BE-USED"}
+
+ client, _ = _ui_client(tmp_path, roster, hub=_Hub())
+ resp = client.post(f"/api/groups/{GROUP}/invites?username=bob&t=tok")
+ assert resp.status_code == 200
+ body = resp.json()
+ assert body["user_id"] == "id-of-bob"
+
+ invites = await roster.list_invites()
+ assert [i["user_id"] for i in invites] == ["id-of-bob"]
+ # Whatever the hub said about keys was never stored anywhere.
+ assert "SHOULD-NOT-BE-USED" not in str(invites)
+ assert await roster.get_identity("id-of-bob") is None
+
+
+def _run_cli(monkeypatch, tmp_path, argv, responses):
+ """Drive the real CLI with the daemon API stubbed, capturing the calls."""
+ import sys as _sys
+
+ from meshbay_node import daemon as _daemon
+
+ calls = []
+
+ def fake_api(cfg, path, method="GET", timeout=30):
+ calls.append((method, path))
+ for key, value in responses.items():
+ if key in path:
+ return value
+ return {}
+
+ monkeypatch.setattr(_daemon, "_daemon_api", fake_api)
+
+ conf = tmp_path / "node.toml"
+ conf.write_text(
+ f'data_dir = "{tmp_path}"\n'
+ '[hub]\nurl = "https://example.org"\nusername = "grenet"\n'
+ f'[[groups]]\nid = "{GROUP}"\nname = "demo"\n'
+ f'shared_dir = "{tmp_path}"\n'
+ )
+ monkeypatch.setattr(_sys, "argv",
+ ["meshbay-node", *argv, "--config", str(conf)])
+ try:
+ _daemon.main()
+ except SystemExit as e:
+ calls.append(("exit", e.code))
+ return calls
+
+
+def test_cli_member_commands_reach_the_right_endpoints(monkeypatch, tmp_path, capsys):
+ resolved = {"user_id": "u-bob", "source": "roster"}
+
+ calls = _run_cli(monkeypatch, tmp_path, ["member", "revoke", "bob"],
+ {"/api/resolve": resolved,
+ "revoke": {"status": "revoked", "reminder": "gek-init"}})
+ assert ("POST", f"/api/members/u-bob/revoke?group_id={GROUP}") in calls
+ # The operator is told the revocation does not take back the key they hold.
+ assert "rotate" in capsys.readouterr().out.lower()
+
+ calls = _run_cli(monkeypatch, tmp_path, ["member", "unpin", "bob"],
+ {"/api/resolve": resolved, "unpin": {"status": "unpinned"}})
+ assert ("POST", "/api/members/u-bob/unpin") in calls
+
+
+def test_cli_resolves_a_name_before_acting(monkeypatch, tmp_path):
+ """
+ The name has to be turned into an account first, and the node's own roster is
+ asked before the hub. A JWT carries no username, so an identity pinned without
+ an invitation has none — the hub fallback is what keeps it manageable.
+ """
+ calls = _run_cli(monkeypatch, tmp_path, ["member", "revoke", "bob"],
+ {"/api/resolve": {"user_id": "u-bob", "source": "hub"},
+ "revoke": {"status": "revoked", "reminder": "gek-init"}})
+
+ assert ("GET", "/api/resolve?username=bob") == calls[0], (
+ "the CLI must resolve the name before acting on anyone")
+ assert ("POST", f"/api/members/u-bob/revoke?group_id={GROUP}") in calls
+
+def test_daemon_does_not_auto_pin_keystore_key():
+ """
+ M3: the daemon used to auto-pin its own keystore key as the admin key, while
+ the browser signs with the user's identity key. Different keys, so every
+ privileged operation failed closed with a signature error that looked like a
+ bug elsewhere — and the demo only worked because a deploy script overwrote it.
+
+ Authority now comes from the roster, or from an explicit node.toml value.
+ """
+ source = (Path(__file__).parent.parent
+ / "src" / "meshbay_node" / "daemon.py").read_text()
+ assert "Auto-pinning admin key" not in source
+ assert "_resolve_admin_pk" not in source, (
+ "the auto-pin resolver is back — node authority must be established "
+ "locally by pairing, never inferred from the node's own keystore (M3)")
+
+
+def test_admin_authority_is_never_fetched_from_the_hub():
+ """
+ The fix M3 invites: ask the hub which key belongs to the operator. That would
+ hand a malicious hub the node — the same substitution as H3, one level deeper.
+ """
+ source = (Path(__file__).parent.parent
+ / "src" / "meshbay_node" / "daemon.py").read_text()
+ admin_region = source[source.find("_legacy_admin_pk"):]
+ assert "pubkeys" not in admin_region.split("def ")[1], (
+ "node authority must never be resolved through a hub lookup")
diff --git a/packages/meshbay-node/tests/test_security_regressions.py b/packages/meshbay-node/tests/test_security_regressions.py
new file mode 100644
index 0000000..dcd9cf6
--- /dev/null
+++ b/packages/meshbay-node/tests/test_security_regressions.py
@@ -0,0 +1,604 @@
+"""
+Phase 11.5 security regression tests.
+
+Each test here encodes a finding from `second-review.md`. They are negative tests:
+they assert that an attack does NOT work. The pre-11.5 code passed 209 feature
+tests while every one of these attacks succeeded — the suite only ever exercised
+happy paths, never an authorization boundary.
+
+If one of these starts failing, a fix has been reverted. Do not "fix" the test.
+"""
+
+import base64
+import struct
+from pathlib import Path
+
+import pytest
+from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
+
+from meshbay_common.protocol import IndexEntry
+from meshbay_node.indexer.group_index import GroupIndex
+from meshbay_node.transport.webrtc_server import WebRTCPeerSession
+
+
+def _safe_name_re():
+ """
+ Imported lazily so that a missing allowlist fails the two tests that need it,
+ rather than aborting collection of the whole module and hiding every other
+ finding's result.
+ """
+ from meshbay_node.transport.webrtc_server import SAFE_UPLOAD_NAME
+ return SAFE_UPLOAD_NAME
+
+
+# ── C1: the unauthenticated HTTP file API must stay deleted ───────────────────
+
+def test_http_file_api_is_gone():
+ """
+ C1: transport/http_server.py served GET /index and GET /file/{id} on 0.0.0.0
+ with no authentication, for private groups too. It was deleted rather than
+ patched. Re-adding any module that serves file bytes outside the MNP handshake
+ reintroduces a full confidentiality bypass.
+ """
+ with pytest.raises(ImportError):
+ import meshbay_node.transport.http_server # noqa: F401
+
+ import meshbay_node.transport as transport
+ assert not hasattr(transport, "create_http_app")
+
+
+def test_tcp_transport_is_gone():
+ """C6: the TCP+TLS server accepted a bare JWT with no GEK proof."""
+ with pytest.raises(ImportError):
+ import meshbay_node.transport.server # noqa: F401
+
+ import meshbay_node.transport as transport
+ assert not hasattr(transport, "ChunkServer")
+
+
+def test_daemon_exposes_no_plaintext_listener():
+ """
+ C1: the daemon must not bind anything that serves content without a handshake.
+ NodeConfig no longer carries an HTTP port at all.
+ """
+ from meshbay_node.config import NodeConfig, GroupConfig
+
+ assert "http_port" not in NodeConfig.__dataclass_fields__
+ assert "http_port" not in GroupConfig.__dataclass_fields__
+ assert "port" not in NodeConfig.__dataclass_fields__
+
+
+# ── C5a: upload filename allowlist ───────────────────────────────────────────
+
+@pytest.mark.parametrize("name", [
+ "../../etc/passwd",
+ "..\\windows\\system32",
+ "/absolute/path",
+ "<img src=x onerror=alert(1)>", # the H2 stored-XSS vector
+ 'name";DROP TABLE x;--',
+ ".hidden",
+ "",
+ "a" * 200,
+ "file\x00.mp4",
+ "sub/dir/file.mp4",
+])
+def test_upload_rejects_unsafe_filenames(name):
+ """C5a/H2: only a conservative allowlist may reach the filesystem."""
+ assert not _safe_name_re().match(name), f"should be rejected: {name!r}"
+
+
+@pytest.mark.parametrize("name", [
+ "movie.mp4",
+ "My Holiday Video.mkv",
+ "report-2026.pdf",
+ "track_01.flac",
+])
+def test_upload_accepts_ordinary_filenames(name):
+ """The allowlist must not break normal use."""
+ assert _safe_name_re().match(name), f"should be accepted: {name!r}"
+
+
+def _session(tmp_path: Path, user_id: str) -> WebRTCPeerSession:
+ """A peer session wired to a real shared root, with sending stubbed out."""
+ shared_root = tmp_path / "shared"
+ shared_root.mkdir(exist_ok=True)
+ index = GroupIndex(group_id="g" * 32, sk_node=Ed25519PrivateKey.generate())
+ ctx = {"shared_root": shared_root, "index": index, "sk_node": index.sk_node}
+
+ session = WebRTCPeerSession.__new__(WebRTCPeerSession)
+ session._ctx = ctx
+ session._group_id = None
+ session._user_id = user_id
+ session._pk_user = ""
+ session._uploads = {}
+ session.sent = []
+ session._send = session.sent.append
+ session._audit = lambda *a, **k: None
+ return session
+
+
+def test_upload_cannot_overwrite_another_members_file(tmp_path):
+ """
+ C5a: uploads used to land in the shared root under a client-chosen name and
+ overwrite whatever was there. That let any member destroy the operator's files,
+ and — by becoming the recorded uploader of the replaced file — delete them
+ through the uploader path, bypassing the Ed25519 admin challenge entirely.
+ """
+ victim = _session(tmp_path, "victim-user")
+ shared_root = victim._ctx["shared_root"]
+
+ original = shared_root / "important.mp4"
+ original.write_bytes(b"operator's original content")
+
+ attacker = _session(tmp_path, "attacker-user")
+ attacker._do_file_upload({
+ "filename": "important.mp4",
+ "chunk_index": 0,
+ "total_chunks": 1,
+ "data": base64.b64encode(b"attacker content").decode(),
+ })
+
+ assert original.read_bytes() == b"operator's original content"
+ uploaded = shared_root / ".uploads" / "attacker-user" / "important.mp4"
+ assert uploaded.exists(), "upload should be quarantined, not dropped"
+ assert uploaded.read_bytes() == b"attacker content"
+
+
+def test_upload_rejects_out_of_order_chunks(tmp_path):
+ """C5a: chunk_index > 0 used to append blindly to any .part file on disk."""
+ session = _session(tmp_path, "user-1")
+ session._do_file_upload({
+ "filename": "movie.mp4", "chunk_index": 3, "total_chunks": 5,
+ "data": base64.b64encode(b"spliced").decode(),
+ })
+ assert any(m.get("type") == "error" for m in session.sent)
+
+
+def test_upload_second_attempt_cannot_replace_own_completed_file(tmp_path):
+ """C5a: even the original uploader goes through a fresh name, not an overwrite."""
+ session = _session(tmp_path, "user-1")
+ payload = {"filename": "movie.mp4", "chunk_index": 0, "total_chunks": 1,
+ "data": base64.b64encode(b"first").decode()}
+ session._do_file_upload(dict(payload))
+ session.sent.clear()
+
+ session._do_file_upload(dict(payload))
+ assert any(m.get("type") == "error" for m in session.sent)
+ stored = session._ctx["shared_root"] / ".uploads" / "user-1" / "movie.mp4"
+ assert stored.read_bytes() == b"first"
+
+
+# ── H1: group isolation ──────────────────────────────────────────────────────
+
+def test_chat_store_and_peers_are_per_group(tmp_path):
+ """
+ H1: chat_store and the peer registry were read from the shared transport
+ context, so on a multi-group node every group's messages went to the first
+ group's database and were served back to members of every other group.
+ """
+ index_a = GroupIndex(group_id="a" * 32, sk_node=Ed25519PrivateKey.generate())
+ index_b = GroupIndex(group_id="b" * 32, sk_node=Ed25519PrivateKey.generate())
+ groups = {
+ "a" * 32: {"chat_store": "STORE_A", "index": index_a, "shared_root": tmp_path},
+ "b" * 32: {"chat_store": "STORE_B", "index": index_b, "shared_root": tmp_path},
+ }
+ ctx = {"groups": groups}
+
+ sess_a = WebRTCPeerSession.__new__(WebRTCPeerSession)
+ sess_a._ctx, sess_a._group_id, sess_a._user_id = ctx, "a" * 32, "alice"
+
+ sess_b = WebRTCPeerSession.__new__(WebRTCPeerSession)
+ sess_b._ctx, sess_b._group_id, sess_b._user_id = ctx, "b" * 32, "bob"
+
+ assert sess_a._group_ctx()["chat_store"] == "STORE_A"
+ assert sess_b._group_ctx()["chat_store"] == "STORE_B"
+
+ sess_a._peer_registry()["alice"] = sess_a
+ sess_b._peer_registry()["bob"] = sess_b
+
+ # Alice's broadcast target set must not contain Bob, who is in another group.
+ assert "bob" not in sess_a._peer_registry()
+ assert "alice" not in sess_b._peer_registry()
+
+ sess_a._user_names()["alice"] = "Alice"
+ assert "alice" not in sess_b._user_names()
+
+
+def test_daemon_sets_no_global_chat_store(tmp_path):
+ """H1: the daemon must not hoist one group's chat store onto the transport."""
+ source = (Path(__file__).parent.parent
+ / "src" / "meshbay_node" / "daemon.py").read_text()
+ assert '_ctx["chat_store"]' not in source, (
+ "daemon must not assign a transport-wide chat_store — it leaks chat "
+ "across groups (H1)"
+ )
+
+
+# ── H2: node admin UI escaping ───────────────────────────────────────────────
+
+def test_no_member_can_hand_the_node_key_material(tmp_path):
+ """
+ C5b, strengthened by the invite redesign (docs/invite-pairing-v1.md).
+
+ This test used to assert that `gek_bundle_store` answered with an admin
+ challenge and stored nothing without an operator signature. The message is now
+ gone entirely: the node holds the GEK and wraps it itself, so no member ever
+ submits key material, authorized or not. Deleting the path is a stronger
+ guarantee than gating it, which is why the assertion changed rather than the
+ behaviour regressing.
+ """
+ from meshbay_common.protocol import MNP as _MNP
+
+ assert not hasattr(_MNP, "GEK_BUNDLE_STORE"), (
+ "the member-supplied bundle message is back — the node must never accept "
+ "key material over MNP (C5b)"
+ )
+
+ source = (Path(__file__).parent.parent
+ / "src" / "meshbay_node" / "transport" / "webrtc_server.py").read_text()
+ assert "_do_gek_bundle_store" not in source
+ assert "_admin_exec_bundle_store" not in source
+
+
+def test_unknown_message_stores_nothing(tmp_path):
+ """A peer sending the retired message must not reach any storage path."""
+ session = _session(tmp_path, "ordinary-member")
+ session._group_id = None
+ session._admin_ops = {}
+
+ stored = []
+
+ class _Store:
+ async def store(self, *args):
+ stored.append(args)
+
+ session._ctx["bundle_store"] = _Store()
+ session._handle_message({
+ "type": "gek_bundle_store",
+ "user_id": "victim", "group_id": "g" * 32,
+ "pk_eph_b64": "AA==", "nonce_b64": "AA==", "wrapped_b64": "AA==",
+ })
+
+ assert stored == [], "a retired message type still reached the bundle store"
+
+
+def test_gek_auto_activation_is_gone():
+ """
+ C5b: the node used to unwrap and adopt any bundle addressed to the operator.
+ Since the operator's X25519 public key is public, any member could hand the
+ node a GEK of their choosing. Nothing arriving over MNP may set a live GEK.
+ """
+ source = (Path(__file__).parent.parent / "src" / "meshbay_node"
+ / "transport" / "webrtc_server.py").read_text()
+ assert "_try_activate_gek" not in source
+ assert 'unwrap_gek_aes' not in source, (
+ "the MNP path must not unwrap a GEK — activation is local-admin only"
+ )
+
+
+# ── H5: admin challenge is bound, not a blind signing oracle ─────────────────
+
+def _transcript(**kw):
+ from meshbay_common.adminop import admin_transcript
+ base = dict(op="file_delete", node_pk_b64="NODEPK", group_id="g" * 32,
+ subject="file-1", nonce=b"\x01" * 32, ts=1_700_000_000)
+ base.update(kw)
+ return admin_transcript(**base)
+
+
+def test_admin_transcript_is_domain_separated():
+ """H5: signatures here can never be valid in another MeshBay protocol."""
+ assert _transcript().startswith(b"meshbay:admin:v1")
+
+
+@pytest.mark.parametrize("field,value", [
+ ("op", "invite_create"),
+ ("subject", "file-2"),
+ ("node_pk_b64", "OTHERNODE"),
+ ("group_id", "h" * 32),
+ ("nonce", b"\x02" * 32),
+ ("ts", 1_700_000_001),
+])
+def test_admin_transcript_binds_every_field(field, value):
+ """
+ H5: a signature must not carry over to another operation, subject, node,
+ group, challenge or moment in time.
+ """
+ assert _transcript() != _transcript(**{field: value}), (
+ f"transcript ignores {field} — signature would be reusable"
+ )
+
+
+def test_admin_transcript_is_unambiguous():
+ """
+ H5/L4: fields are length-prefixed. With plain concatenation a crafted subject
+ could impersonate the following field and two different operations would
+ produce identical signed bytes.
+ """
+ a = _transcript(subject="file-1", group_id="g")
+ b = _transcript(subject="1", group_id="gfile-")
+ assert a != b, "concatenation is ambiguous — length prefixes missing"
+
+
+def test_admin_signature_does_not_transfer_between_operations(tmp_path):
+ """
+ H5: the concrete attack. A signature collected to delete a file must not
+ authorize storing a GEK bundle.
+ """
+ from meshbay_common.adminop import OP_FILE_DELETE, OP_INVITE_CREATE
+
+ sk_admin = Ed25519PrivateKey.generate()
+ delete_transcript = _transcript(op=OP_FILE_DELETE)
+ signature = sk_admin.sign(delete_transcript)
+
+ invite_transcript = _transcript(op=OP_INVITE_CREATE)
+ with pytest.raises(Exception):
+ sk_admin.public_key().verify(signature, invite_transcript)
+
+
+def test_admin_challenge_expires(tmp_path):
+ """H5: a stale challenge must not be usable."""
+ import time as _time
+ from meshbay_common.adminop import ADMIN_CHALLENGE_TTL, OP_FILE_DELETE
+
+ session = _session(tmp_path, "operator")
+ session._group_id = None
+ session._admin_ops = {
+ "op-1": {
+ "op": OP_FILE_DELETE, "subject": "file-1", "nonce": b"\x00" * 32,
+ "ts": int(_time.time()) - ADMIN_CHALLENGE_TTL - 5, "payload": {},
+ }
+ }
+ session._do_admin_response({"op_id": "op-1", "signature": ""})
+ assert any(m.get("type") == "error" and "expired" in m.get("detail", "").lower()
+ for m in session.sent)
+
+
+def test_denylist_persists_and_honours_groups(tmp_path):
+ """
+ H4: revocations lived only in memory, so a node restart silently un-revoked
+ everyone, and 'group' targets were dropped entirely — the hub signed and
+ broadcast them, the node's handler understood only 'user' and 'jti'.
+ """
+ from meshbay_node.transport import Denylist
+
+ path = tmp_path / "denylist.json"
+ first = Denylist(path=path)
+ first.deny_group("g-revoked")
+ first.deny_user("u-revoked")
+ first.deny_jti("j-revoked")
+
+ # A fresh instance stands in for a daemon restart.
+ reloaded = Denylist(path=path)
+ assert reloaded.is_denied("", "", "g-revoked"), "group revocation not honoured"
+ assert reloaded.is_denied("u-revoked", "")
+ assert reloaded.is_denied("", "j-revoked")
+ assert not reloaded.is_denied("someone", "other", "g-allowed")
+
+
+def test_swarm_registration_skips_private_groups():
+ """
+ H7: the daemon registered content hashes for every group, private included,
+ handing the hub a fingerprint of every private file. The bug was masked by a
+ mis-mounted route, so fixing the route without this filter would have turned a
+ dormant leak into a live one.
+ """
+ source = (Path(__file__).parent.parent / "src" / "meshbay_node"
+ / "daemon.py").read_text()
+ assert 'visibility' in source and '_register_swarm' in source
+ # Both registration sites must gate on public visibility.
+ for marker in ['gctx.get("visibility") != "public"',
+ 'group_cfg.visibility == "public"']:
+ assert marker in source, f"swarm registration not gated: {marker}"
+
+
+def test_keystore_argon2_is_production_strength():
+ """M2: the keystore protects the node's private keys and sat at 64 MB."""
+ from meshbay_common.crypto import ARGON2_MEMORY_COST
+ assert ARGON2_MEMORY_COST >= 262144
+
+
+def test_keystore_records_argon2_params_for_migration(tmp_path):
+ """
+ M2: raising the parameters must not orphan existing keystores, so each
+ envelope records the parameters it was written with.
+ """
+ import json
+ from meshbay_node.keystore import create_keystore, load_keystore
+
+ path = tmp_path / "keystore.enc"
+ created = create_keystore(path=path, password="correct horse battery")
+ envelope = json.loads(path.read_text())
+ assert envelope["argon2"]["memory_cost"] >= 262144
+
+ reopened = load_keystore(path=path, password="correct horse battery")
+ assert reopened.pk_ed25519_b64 == created.pk_ed25519_b64
+
+
+def test_legacy_keystore_still_opens(tmp_path):
+ """M2: a keystore written under the 64 MB profile must still unlock."""
+ import base64 as _b64
+ import json
+ import msgpack
+ from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey
+ from meshbay_common.crypto import (
+ LEGACY_ARGON2_ITERATIONS, LEGACY_ARGON2_LANES, LEGACY_ARGON2_MEMORY_COST,
+ derive_keystore_key, encrypt_keystore, pk_to_b64, sk_to_b64,
+ )
+ from meshbay_node.keystore import load_keystore
+
+ sk_ed, sk_x = Ed25519PrivateKey.generate(), X25519PrivateKey.generate()
+ payload = msgpack.packb({
+ "sk_ed25519_b64": sk_to_b64(sk_ed),
+ "sk_x25519_b64": sk_to_b64(sk_x),
+ }, use_bin_type=True)
+
+ salt = b"\x01" * 16
+ key = derive_keystore_key(
+ "legacy-pass", salt,
+ iterations=LEGACY_ARGON2_ITERATIONS,
+ memory_cost=LEGACY_ARGON2_MEMORY_COST,
+ lanes=LEGACY_ARGON2_LANES,
+ )
+ iv, ct, tag = encrypt_keystore(payload, key)
+
+ path = tmp_path / "legacy.enc"
+ # No "argon2" key — exactly how pre-M2 envelopes look.
+ path.write_text(json.dumps({
+ "version": 1,
+ "argon2_salt_b64": _b64.b64encode(salt).decode(),
+ "iv_b64": _b64.b64encode(iv).decode(),
+ "tag_b64": _b64.b64encode(tag).decode(),
+ "ciphertext_b64": _b64.b64encode(ct).decode(),
+ }))
+
+ keys = load_keystore(path=path, password="legacy-pass")
+ assert keys.pk_ed25519_b64 == pk_to_b64(sk_ed.public_key())
+ assert keys.pk_x25519_b64 == pk_to_b64(sk_x.public_key())
+
+
+def test_dead_gek_protocol_constants_removed():
+ """L1: the node never serves a GEK; the message types should not suggest it."""
+ from meshbay_common.protocol import MNP
+ assert not hasattr(MNP, "GEK_REQUEST")
+ assert not hasattr(MNP, "GEK_RESPONSE")
+
+
+def test_peer_errors_do_not_leak_internals():
+ """
+ L3: arbitrary exception text carries filesystem paths and internal state, so
+ the catch-all handler must not relay it.
+
+ Deliberately narrow: HandshakeError messages ARE sent to the peer, because a
+ client needs to know why it was refused, and those strings are authored for
+ that purpose. The check targets the generic `except Exception as e` path.
+ """
+ source = (Path(__file__).parent.parent / "src" / "meshbay_node"
+ / "transport" / "webrtc_server.py").read_text()
+ assert '"detail": str(e)' not in source, (
+ "generic exception text relayed to peer — use a fixed message"
+ )
+ # And the catch-all must still exist, sending something opaque.
+ assert '"detail": "Request failed"' in source
+
+
+def test_pre_handshake_message_budget_is_small():
+ """
+ H6: the frame limit was a flat 64 MB applied before authentication, so an
+ unauthenticated peer could announce a huge frame and dribble bytes into it.
+ """
+ from meshbay_node.transport.webrtc_server import (
+ MAX_MSG, PRE_HANDSHAKE_MAX_MSG, _DataChannelBuffer,
+ )
+ assert PRE_HANDSHAKE_MAX_MSG <= 1024 * 1024
+ assert PRE_HANDSHAKE_MAX_MSG < MAX_MSG
+
+ buf = _DataChannelBuffer(max_message=PRE_HANDSHAKE_MAX_MSG)
+ buf.feed(struct.pack(">I", PRE_HANDSHAKE_MAX_MSG + 1) + b"x")
+ with pytest.raises(ValueError):
+ list(buf.messages())
+
+
+def test_stream_segment_is_not_synchronous():
+ """
+ H6: _do_stream_segment ran subprocess.run(timeout=30) inside the event loop,
+ stalling every peer on the node for up to thirty seconds per request.
+
+ Asserts the property (the worker is a coroutine, ffmpeg is spawned through
+ asyncio) rather than grepping for "subprocess.run" — which also matches the
+ comment that documents the old behaviour.
+ """
+ import ast
+ import inspect
+ from meshbay_node.transport.webrtc_server import WebRTCPeerSession
+
+ assert inspect.iscoroutinefunction(WebRTCPeerSession._do_stream_segment_async)
+
+ source = (Path(__file__).parent.parent / "src" / "meshbay_node"
+ / "transport" / "webrtc_server.py").read_text()
+ tree = ast.parse(source)
+ blocking = [
+ node for node in ast.walk(tree)
+ if isinstance(node, ast.Call)
+ and isinstance(node.func, ast.Attribute)
+ and node.func.attr == "run"
+ and isinstance(node.func.value, ast.Name)
+ and node.func.value.id == "subprocess"
+ ]
+ assert not blocking, "blocking subprocess.run() in the event loop"
+ assert "_transcode_sem" in source, "ffmpeg spawns must be capped"
+
+
+def test_pre_proof_fetches_are_bounded():
+ """C4: the pre-proof bundle window is a disclosure surface; bound it."""
+ from meshbay_node.transport.webrtc_server import MAX_PRE_PROOF_FETCHES
+ assert 0 < MAX_PRE_PROOF_FETCHES <= 10
+
+
+def test_node_admin_ui_requires_token():
+ """
+ 11.5.3: "localhost only" is not authentication. Any local process — or a
+ rebound browser page — could re-initialise a group's GEK and read the audit log.
+ """
+ from fastapi.testclient import TestClient
+ from meshbay_node.ui.app import create_ui_app
+
+ app = create_ui_app({"status": "running", "groups_ctx": {},
+ "indexes": {}, "ui_token": "secret-token"})
+ client = TestClient(app)
+
+ assert client.get("/api/status").status_code == 403
+ assert client.get("/api/status?t=wrong").status_code == 403
+ assert client.get("/api/config?t=wrong").status_code == 403
+ assert client.get("/api/status?t=secret-token").status_code == 200
+ assert client.get(
+ "/api/status", headers={"X-MeshBay-Token": "secret-token"}
+ ).status_code == 200
+
+
+def test_admin_ui_escapes_filenames(tmp_path):
+ """
+ H2: filenames are chosen by any group member and were rendered into the
+ localhost admin UI unescaped, giving script execution against an
+ unauthenticated admin API.
+ """
+ from meshbay_node.ui.app import _render_page
+
+ payload = '<img src=x onerror="fetch(1)">'
+ index = GroupIndex(group_id="g" * 32, sk_node=Ed25519PrivateKey.generate())
+ index.add_entry(IndexEntry(
+ id="0" * 64, name=payload, path="", size=1, type="video", added_at=0,
+ ))
+
+ html = _render_page({
+ "status": "running",
+ "groups_ctx": {"g" * 32: {"index": index, "shared_root": tmp_path}},
+ "indexes": {"g" * 32: index},
+ })
+
+ assert payload not in html, "filename rendered unescaped — stored XSS (H2)"
+ assert "&lt;img" in html, "filename should appear escaped"
+
+
+def test_admin_ui_escapes_roster_usernames(tmp_path):
+ """
+ H2 again, for the roster: usernames originate at the hub and land on the
+ operator's own admin page, which can re-key groups and read the audit log.
+ """
+ from meshbay_node.ui.app import _render_page
+
+ payload = '<img src=x onerror="fetch(1)">'
+ html = _render_page(
+ {"status": "running", "groups_ctx": {}, "indexes": {}},
+ {
+ "identities": {"u1": {"user_id": "u1", "username": payload,
+ "pk_ed25519": "AAA", "pinned_at": "now",
+ "pinned_via": "code"}},
+ "members": [{"group_id": "", "user_id": "u1", "role": "operator",
+ "status": "active"}],
+ "invites": [],
+ },
+ )
+
+ assert payload not in html, "username rendered unescaped — stored XSS (H2)"
+ assert "&lt;img" in html
diff --git a/packages/meshbay-node/tests/test_transport.py b/packages/meshbay-node/tests/test_transport.py
deleted file mode 100644
index 0e70d72..0000000
--- a/packages/meshbay-node/tests/test_transport.py
+++ /dev/null
@@ -1,222 +0,0 @@
-"""
-Integration test: ChunkServer ↔ ChunkClient over TLS.
-
-Starts a real TLS server on localhost, connects a client,
-fetches index and a chunk, verifies signature+hash+decryption.
-"""
-
-import asyncio
-import base64
-import os
-import time
-import jwt
-import pytest
-from pathlib import Path
-from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
-from cryptography.hazmat.primitives import serialization
-
-from meshbay_common.crypto import generate_gek, pk_to_b64
-from meshbay_node.indexer import DirectoryIndexer, GroupIndex
-from meshbay_node.transport.server import ChunkServer
-from meshbay_node.transport.client import ChunkClient
-
-
-@pytest.fixture
-def sk_node():
- return Ed25519PrivateKey.generate()
-
-@pytest.fixture
-def sk_hub():
- return Ed25519PrivateKey.generate()
-
-@pytest.fixture
-def gek():
- return generate_gek()
-
-@pytest.fixture
-def shared_dir(tmp_path):
- d = tmp_path / "shared"
- d.mkdir()
- (d / "test.mp4").write_bytes(os.urandom(2 * 1024 * 1024)) # 2 MB
- (d / "small.txt").write_bytes(b"hello meshbay " * 100)
- return d
-
-def make_jwt(sk_hub, pk_node_b64, user_id="user-001", ttl=3600, groups=None):
- sk_pem = sk_hub.private_bytes(
- serialization.Encoding.PEM,
- serialization.PrivateFormat.PKCS8,
- serialization.NoEncryption(),
- )
- now = int(time.time())
- return jwt.encode({
- "iss": "test-hub", "sub": user_id,
- "pk_user": pk_node_b64, "hub_id": "test-hub",
- "jti": "test-jti",
- "iat": now, "exp": now + ttl,
- "groups": groups or [],
- }, sk_pem, algorithm="EdDSA")
-
-
-@pytest.mark.asyncio
-async def test_chunk_server_client_roundtrip(
- sk_node, sk_hub, gek, shared_dir, tmp_path):
- """Full integration: server serves a chunk, client verifies and decrypts."""
-
- # Build index
- indexer = DirectoryIndexer(
- root=shared_dir, group_id="g", sk_node=sk_node, gek=gek)
- await indexer.initial_scan()
- assert indexer.index.count == 2
-
- # Hub PK for JWT verification
- hub_pk_pem = sk_hub.public_key().public_bytes(
- serialization.Encoding.PEM,
- serialization.PublicFormat.SubjectPublicKeyInfo)
-
- # TLS cert in tmp dir
- cert_path = tmp_path / "node.crt"
- key_path = tmp_path / "node.key"
-
- server = ChunkServer(
- sk_node=sk_node,
- hub_pk_pem=hub_pk_pem,
- gek=gek,
- shared_root=shared_dir,
- index=indexer.index,
- host="127.0.0.1",
- port=0, # OS picks a free port
- cert_path=cert_path,
- key_path=key_path,
- )
- await server.start()
- port = server._server.sockets[0].getsockname()[1]
-
- token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key()))
-
- # Find the large test file in the index
- entry = next(e for e in indexer.index.entries if e.name == "test.mp4")
-
- async with ChunkClient(
- host="127.0.0.1",
- port=port,
- jwt_token=token,
- gek=gek,
- pk_node_b64=pk_to_b64(sk_node.public_key()),
- ) as client:
- # Fetch first chunk
- chunk0 = await client.fetch_chunk(entry.id, chunk_index=0)
- assert len(chunk0) == 1024 * 1024 # first 1MB of 2MB file
-
- # Fetch second chunk
- chunk1 = await client.fetch_chunk(entry.id, chunk_index=1)
- assert len(chunk1) == 1024 * 1024 # second 1MB
-
- # Reassembled file matches original
- original = (shared_dir / "test.mp4").read_bytes()
- assert chunk0 + chunk1 == original
-
- await server.stop()
-
-
-@pytest.mark.asyncio
-async def test_invalid_jwt_rejected(sk_node, sk_hub, gek, shared_dir, tmp_path):
- hub_pk_pem = sk_hub.public_key().public_bytes(
- serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo)
-
- indexer = DirectoryIndexer(root=shared_dir, group_id="g",
- sk_node=sk_node, gek=gek)
- await indexer.initial_scan()
-
- cert_path = tmp_path / "node.crt"
- key_path = tmp_path / "node.key"
-
- server = ChunkServer(
- sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek,
- shared_root=shared_dir, index=indexer.index,
- host="127.0.0.1", port=0,
- cert_path=cert_path, key_path=key_path,
- )
- await server.start()
- port = server._server.sockets[0].getsockname()[1]
-
- # Use a different hub key to sign the token
- sk_other_hub = Ed25519PrivateKey.generate()
- bad_token = make_jwt(sk_other_hub, pk_to_b64(sk_node.public_key()))
-
- with pytest.raises(Exception):
- async with ChunkClient(
- host="127.0.0.1", port=port,
- jwt_token=bad_token, gek=gek,
- pk_node_b64=pk_to_b64(sk_node.public_key()),
- ) as client:
- pass
-
- await server.stop()
-
-
-@pytest.mark.asyncio
-async def test_wrong_group_rejected(sk_node, sk_hub, gek, shared_dir, tmp_path):
- """TCP+TLS server rejects a client whose JWT groups don't include the requested group_id."""
- hub_pk_pem = sk_hub.public_key().public_bytes(
- serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo)
-
- indexer = DirectoryIndexer(root=shared_dir, group_id="g",
- sk_node=sk_node, gek=gek)
- await indexer.initial_scan()
-
- cert_path = tmp_path / "node.crt"
- key_path = tmp_path / "node.key"
-
- server = ChunkServer(
- sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek,
- shared_root=shared_dir, index=indexer.index,
- host="127.0.0.1", port=0,
- cert_path=cert_path, key_path=key_path,
- )
- await server.start()
- port = server._server.sockets[0].getsockname()[1]
-
- token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key()), groups=["group-a"])
-
- with pytest.raises(ConnectionError, match="rejected"):
- async with ChunkClient(
- host="127.0.0.1", port=port,
- jwt_token=token, gek=gek,
- pk_node_b64=pk_to_b64(sk_node.public_key()),
- group_id="group-b",
- ) as client:
- pass
-
- await server.stop()
-
-
-@pytest.mark.asyncio
-async def test_fetch_index(sk_node, sk_hub, gek, shared_dir, tmp_path):
- hub_pk_pem = sk_hub.public_key().public_bytes(
- serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo)
- indexer = DirectoryIndexer(root=shared_dir, group_id="g",
- sk_node=sk_node, gek=gek)
- await indexer.initial_scan()
-
- cert_path = tmp_path / "node.crt"
- key_path = tmp_path / "node.key"
-
- server = ChunkServer(
- sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek,
- shared_root=shared_dir, index=indexer.index,
- host="127.0.0.1", port=0,
- cert_path=cert_path, key_path=key_path,
- )
- await server.start()
- port = server._server.sockets[0].getsockname()[1]
- token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key()))
-
- async with ChunkClient(
- host="127.0.0.1", port=port, jwt_token=token,
- gek=gek, pk_node_b64=pk_to_b64(sk_node.public_key()),
- ) as client:
- wire = await client.fetch_index()
- recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek)
- assert recovered.count == 2
-
- await server.stop()
diff --git a/packages/meshbay-node/tests/test_webrtc_transport.py b/packages/meshbay-node/tests/test_webrtc_transport.py
index 693a68b..93cd3fd 100644
--- a/packages/meshbay-node/tests/test_webrtc_transport.py
+++ b/packages/meshbay-node/tests/test_webrtc_transport.py
@@ -19,7 +19,9 @@ import jwt
import msgpack
import pytest
from cryptography.hazmat.primitives import serialization
-from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
+from cryptography.hazmat.primitives.asymmetric.ed25519 import (
+ Ed25519PrivateKey, Ed25519PublicKey,
+)
from aiortc import RTCPeerConnection, RTCSessionDescription
from meshbay_common import MNP_VERSION
@@ -29,10 +31,24 @@ from meshbay_common.crypto import (
wrap_gek,
wrap_gek_aes,
unwrap_gek,
+ unwrap_gek_aes,
)
from meshbay_common.webcrypto import chunk_key_aes, decrypt_chunk_aes
from meshbay_common.protocol import MNP
+TEST_GROUP = "g"
+
+from meshbay_common.handshake import (
+ NONCE_LEN, ROLE_CLIENT, ROLE_NODE, handshake_transcript,
+ make_proof, verify_proof, webrtc_binding,
+)
+from meshbay_common.adminop import (
+ OP_FILE_DELETE,
+ OP_INVITE_CREATE,
+ admin_transcript,
+)
+from meshbay_common.join import ROLE_MEMBER, ROLE_OPERATOR, join_transcript
from meshbay_node.bundle_store import BundleStore
+from meshbay_node.roster import Roster
from meshbay_node.indexer import DirectoryIndexer
from meshbay_node.transport.webrtc_server import WebRTCTransport
@@ -72,6 +88,9 @@ def _hub_pk_pem(sk_hub):
def _make_jwt(sk_hub, groups=None, pk_user="test"):
+ # group_id is mandatory now (M1), so the default token must be a member
+ # of the group the tests connect to. Tests that exercise refusal pass
+ # groups=[...] explicitly.
sk_pem = sk_hub.private_bytes(
serialization.Encoding.PEM,
serialization.PrivateFormat.PKCS8,
@@ -82,10 +101,25 @@ def _make_jwt(sk_hub, groups=None, pk_user="test"):
"iss": "test-hub", "sub": "user-001",
"pk_user": pk_user, "hub_id": "test-hub",
"jti": "test-jti-webrtc", "iat": now, "exp": now + 3600,
- "groups": groups or [],
+ "groups": groups if groups is not None else [TEST_GROUP],
}, sk_pem, algorithm="EdDSA")
+def _transcript_from(challenge_msg: dict) -> bytes:
+ """
+ Rebuild the signed transcript from an admin_challenge, the way a real client
+ does — from the announced fields, never from opaque bytes on the wire (H5).
+ """
+ return admin_transcript(
+ op=challenge_msg["op"],
+ node_pk_b64=challenge_msg["node_pk"],
+ group_id=challenge_msg["group_id"],
+ subject=challenge_msg["subject"],
+ nonce=base64.b64decode(challenge_msg["nonce"]),
+ ts=challenge_msg["ts"],
+ )
+
+
def _pack(obj: dict) -> bytes:
data = msgpack.packb(obj, use_bin_type=True)
return struct.pack(">I", len(data)) + data
@@ -103,36 +137,80 @@ def _extract_dtls_fp(sdp: str) -> bytes:
return b""
-async def _handshake_with_gek_proof(channel, received, sk_hub, gek, groups=None,
- browser_pc=None):
- """Send handshake, handle GEK challenge, return handshake_ack."""
- token = _make_jwt(sk_hub, groups=groups)
+async def _do_mnp_handshake(channel, received, token, gek, pc, group_id):
+ """
+ Client half of the unified handshake (11.5.4): client nonce, length-prefixed
+ role-bound transcript, and verification of the node's own proof + signature.
+ """
+ nonce_c = os.urandom(NONCE_LEN)
channel.send(_pack({
- "type": MNP.HANDSHAKE,
- "v": MNP_VERSION,
- "token": token,
+ "type": MNP.HANDSHAKE, "v": MNP_VERSION,
+ "token": token, "group_id": group_id,
+ "nonce": base64.b64encode(nonce_c).decode(),
}))
msg = await asyncio.wait_for(received.get(), timeout=5.0)
- if msg["type"] == MNP.HANDSHAKE_CHALLENGE:
- nonce = base64.b64decode(msg["nonce"])
- offer_fp = b""
- answer_fp = b""
- if browser_pc:
- offer_fp = _extract_dtls_fp(browser_pc.localDescription.sdp)
- answer_fp = _extract_dtls_fp(browser_pc.remoteDescription.sdp)
- proof = hmac.new(gek, nonce + offer_fp + answer_fp, hashlib.sha256).digest()
- channel.send(_pack({
- "type": MNP.HANDSHAKE_RESPONSE,
- "v": MNP_VERSION,
- "proof": base64.b64encode(proof).decode(),
- }))
- msg = await asyncio.wait_for(received.get(), timeout=5.0)
+ if msg["type"] != MNP.HANDSHAKE_CHALLENGE:
+ return msg
+
+ nonce_s = base64.b64decode(msg["nonce"])
+ binding = webrtc_binding(
+ _extract_dtls_fp(pc.localDescription.sdp),
+ _extract_dtls_fp(pc.remoteDescription.sdp),
+ )
+ proof = make_proof(gek, ROLE_CLIENT, group_id, nonce_c, nonce_s, binding)
+ channel.send(_pack({
+ "type": MNP.HANDSHAKE_RESPONSE, "v": MNP_VERSION,
+ "proof": base64.b64encode(proof).decode(),
+ }))
+ ack = await asyncio.wait_for(received.get(), timeout=5.0)
+
+ if ack.get("type") == MNP.HANDSHAKE_ACK:
+ # The client must authenticate the node too (C3).
+ assert verify_proof(
+ gek, base64.b64decode(ack["proof"]), ROLE_NODE,
+ group_id, nonce_c, nonce_s, binding), "node proof invalid"
+ transcript = handshake_transcript(
+ ROLE_NODE, group_id, nonce_c, nonce_s, binding)
+ Ed25519PublicKey.from_public_bytes(
+ base64.b64decode(ack["node_pk"])
+ ).verify(base64.b64decode(ack["sig"]), transcript)
+ return ack
+
+
+async def _handshake_with_gek_proof(channel, received, sk_hub, gek, groups=None,
+ browser_pc=None, group_id=TEST_GROUP):
+ """Send handshake, handle GEK challenge, return handshake_ack."""
+ token = _make_jwt(sk_hub, groups=groups or [group_id])
+ msg = await _do_mnp_handshake(
+ channel, received, token, gek, browser_pc, group_id)
assert msg["type"] == MNP.HANDSHAKE_ACK
return msg
-async def _setup_peer(transport, sk_hub, gek, peer_id, jwt_sub="user-001", sk_user=None):
- """Create a peer connection, perform handshake with GEK proof, return (pc, channel, queue)."""
+def _token(sk_hub, jwt_sub, peer_id, group_id, pk_user="test"):
+ """A hub-issued user token, as the browser would present it."""
+ sk_h_pem = sk_hub.private_bytes(
+ serialization.Encoding.PEM,
+ serialization.PrivateFormat.PKCS8,
+ serialization.NoEncryption(),
+ )
+ now = int(time.time())
+ return jwt.encode({
+ "iss": "test-hub", "sub": jwt_sub,
+ "pk_user": pk_user, "hub_id": "test-hub",
+ "jti": f"jti-{peer_id}", "iat": now, "exp": now + 3600,
+ "groups": [group_id], "scope": "user",
+ }, sk_h_pem, algorithm="EdDSA")
+
+
+async def _open_channel(transport, peer_id):
+ """
+ Signaling only: a live DataChannel with no MNP handshake performed.
+
+ Separate from `_setup_peer` because someone joining a group for the first time
+ cannot complete the handshake — they have no GEK to prove — and the join has to
+ happen in that window.
+ """
pc = RTCPeerConnection()
q = asyncio.Queue()
buf = bytearray()
@@ -161,6 +239,13 @@ async def _setup_peer(transport, sk_hub, gek, peer_id, jwt_sub="user-001", sk_us
answer_sdp, _ = await transport.handle_offer(pc.localDescription.sdp, peer_id)
await pc.setRemoteDescription(RTCSessionDescription(sdp=answer_sdp, type="answer"))
await asyncio.wait_for(ready.wait(), timeout=5.0)
+ return pc, ch, q
+
+
+async def _setup_peer(transport, sk_hub, gek, peer_id, jwt_sub="user-001", sk_user=None,
+ group_id=TEST_GROUP):
+ """Create a peer connection, perform handshake with GEK proof, return (pc, channel, queue)."""
+ pc, ch, q = await _open_channel(transport, peer_id)
pk_user = "test"
if sk_user:
@@ -169,31 +254,9 @@ async def _setup_peer(transport, sk_hub, gek, peer_id, jwt_sub="user-001", sk_us
serialization.Encoding.Raw, serialization.PublicFormat.Raw)
).decode()
- sk_h_pem = sk_hub.private_bytes(
- serialization.Encoding.PEM,
- serialization.PrivateFormat.PKCS8,
- serialization.NoEncryption(),
- )
- now = int(time.time())
- token = jwt.encode({
- "iss": "test-hub", "sub": jwt_sub,
- "pk_user": pk_user, "hub_id": "test-hub",
- "jti": f"jti-{peer_id}", "iat": now, "exp": now + 3600,
- "groups": [],
- }, sk_h_pem, algorithm="EdDSA")
+ token = _token(sk_hub, jwt_sub, peer_id, group_id, pk_user)
- ch.send(_pack({"type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token}))
- msg = await asyncio.wait_for(q.get(), timeout=5.0)
- if msg["type"] == MNP.HANDSHAKE_CHALLENGE:
- nonce = base64.b64decode(msg["nonce"])
- offer_fp = _extract_dtls_fp(pc.localDescription.sdp)
- answer_fp = _extract_dtls_fp(pc.remoteDescription.sdp)
- proof = hmac.new(gek, nonce + offer_fp + answer_fp, hashlib.sha256).digest()
- ch.send(_pack({
- "type": MNP.HANDSHAKE_RESPONSE, "v": MNP_VERSION,
- "proof": base64.b64encode(proof).decode(),
- }))
- msg = await asyncio.wait_for(q.get(), timeout=5.0)
+ msg = await _do_mnp_handshake(ch, q, token, gek, pc, group_id)
assert msg["type"] == MNP.HANDSHAKE_ACK
return pc, ch, q
@@ -676,6 +739,8 @@ async def test_webrtc_wrong_gek_proof_rejected(sk_node, sk_hub, gek, shared_dir)
token = _make_jwt(sk_hub)
channel.send(_pack({
"type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token,
+ "group_id": TEST_GROUP,
+ "nonce": base64.b64encode(os.urandom(NONCE_LEN)).decode(),
}))
challenge = await asyncio.wait_for(received.get(), timeout=5.0)
@@ -734,6 +799,8 @@ async def test_webrtc_dtls_channel_binding_detects_mitm(sk_node, sk_hub, gek, sh
token = _make_jwt(sk_hub)
channel.send(_pack({
"type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token,
+ "group_id": TEST_GROUP,
+ "nonce": base64.b64encode(os.urandom(NONCE_LEN)).decode(),
}))
challenge = await asyncio.wait_for(received.get(), timeout=5.0)
@@ -783,13 +850,13 @@ async def test_webrtc_admin_challenge_response(sk_node, sk_hub, gek, shared_dir)
challenge_msg = await asyncio.wait_for(received.get(), timeout=5.0)
assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE
- assert challenge_msg["file_id"] == entry.id
+ assert challenge_msg["op"] == OP_FILE_DELETE
+ assert challenge_msg["subject"] == entry.id
- challenge = base64.b64decode(challenge_msg["challenge"])
- signature = sk_admin.sign(challenge)
+ signature = sk_admin.sign(_transcript_from(challenge_msg))
channel.send(_pack({
"type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION,
- "file_id": entry.id,
+ "op_id": challenge_msg["op_id"],
"signature": base64.b64encode(signature).decode(),
}))
@@ -832,11 +899,10 @@ async def test_webrtc_admin_bad_signature_rejected(sk_node, sk_hub, gek, shared_
challenge_msg = await asyncio.wait_for(received.get(), timeout=5.0)
assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE
- challenge = base64.b64decode(challenge_msg["challenge"])
- bad_sig = sk_attacker.sign(challenge)
+ bad_sig = sk_attacker.sign(_transcript_from(challenge_msg))
channel.send(_pack({
"type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION,
- "file_id": entry.id,
+ "op_id": challenge_msg["op_id"],
"signature": base64.b64encode(bad_sig).decode(),
}))
@@ -914,14 +980,14 @@ async def test_webrtc_uploader_delete_requires_challenge(sk_node, sk_hub, gek, s
challenge_msg = await asyncio.wait_for(received.get(), timeout=5.0)
assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE
- assert challenge_msg["file_id"] == entry.id
+ assert challenge_msg["op"] == OP_FILE_DELETE
+ assert challenge_msg["subject"] == entry.id
# Sign with uploader's Ed25519 key
- challenge = base64.b64decode(challenge_msg["challenge"])
- signature = sk_uploader.sign(challenge)
+ signature = sk_uploader.sign(_transcript_from(challenge_msg))
channel.send(_pack({
"type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION,
- "file_id": entry.id,
+ "op_id": challenge_msg["op_id"],
"signature": base64.b64encode(signature).decode(),
}))
@@ -979,11 +1045,10 @@ async def test_webrtc_uploader_impersonation_blocked(sk_node, sk_hub, gek, share
assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE
# Sign with user B's key (wrong key)
- challenge = base64.b64decode(challenge_msg["challenge"])
- bad_sig = sk_user_b.sign(challenge)
+ bad_sig = sk_user_b.sign(_transcript_from(challenge_msg))
channel.send(_pack({
"type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION,
- "file_id": entry.id,
+ "op_id": challenge_msg["op_id"],
"signature": base64.b64encode(bad_sig).decode(),
}))
@@ -1014,54 +1079,122 @@ def x25519_keypair():
@pytest.mark.asyncio
-async def test_gek_bundle_store_and_fetch(sk_node, sk_hub, gek, shared_dir,
- tmp_path, x25519_keypair):
- """GEK bundle stored on node via DataChannel, then fetched during handshake."""
+async def test_invite_then_join_delivers_the_gek(sk_node, sk_hub, gek, shared_dir,
+ tmp_path, x25519_keypair):
+ """
+ The whole invite flow over a real DataChannel, end to end.
+
+ The operator asks for a code; the invitee — who has never held the group key
+ and therefore cannot complete the GEK proof — redeems it in the pre-proof
+ window and the node wraps the key for the X25519 key they just proved they
+ hold. At no point is a public key fetched from the hub, which is the point:
+ that lookup was H3.
+ """
hub_pk_pem = _hub_pk_pem(sk_hub)
indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek)
await indexer.initial_scan()
- bundle_store = BundleStore(db_path=tmp_path / "bundles.db")
- await bundle_store.open()
+ roster = Roster(db_path=tmp_path / "roster.db")
+ await roster.open()
transport = WebRTCTransport(
sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek,
shared_root=shared_dir, index=indexer.index,
stun_servers=[],
)
- transport._ctx["bundle_store"] = bundle_store
+ transport._ctx["roster"] = roster
+ transport._ctx["has_admin_authority"] = True
+ transport._ctx["groups"] = {
+ TEST_GROUP: {"gek": gek, "shared_root": shared_dir, "index": indexer.index},
+ }
+
+ # A paired operator, as `meshbay-node operator pair` would have left it.
+ sk_admin = Ed25519PrivateKey.generate()
+ admin_pk_b64 = pk_to_b64(sk_admin.public_key())
+ await roster.pin_identity("user-001", "grenet", admin_pk_b64, "AA==", "code")
+ await roster.set_member("", "user-001", ROLE_OPERATOR, "active", "local-cli")
- # Connect as admin and store a GEK bundle for user-002
pc_admin, ch_admin, q_admin = await _setup_peer(
transport, sk_hub, gek, "peer-admin")
- sk_x_raw, pk_x_raw = x25519_keypair
- bundle = wrap_gek(gek, pk_x_raw)
+ # 1. The operator asks the node for an invitation code.
+ ch_admin.send(_pack({
+ "type": MNP.INVITE_CREATE, "v": MNP_VERSION,
+ "user_id": "user-002", "group_id": TEST_GROUP, "username": "bob",
+ }))
+ challenge_msg = await asyncio.wait_for(q_admin.get(), timeout=5.0)
+ assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE
+ assert challenge_msg["op"] == OP_INVITE_CREATE
+ assert challenge_msg["subject"] == "user-002"
ch_admin.send(_pack({
- "type": MNP.GEK_BUNDLE_STORE,
- "v": MNP_VERSION,
- "user_id": "user-002",
- "group_id": "g",
- "pk_eph_b64": bundle["pk_eph_b64"],
- "nonce_b64": bundle["nonce_b64"],
- "wrapped_b64": bundle["wrapped_b64"],
+ "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION,
+ "op_id": challenge_msg["op_id"],
+ "signature": base64.b64encode(
+ sk_admin.sign(_transcript_from(challenge_msg))).decode(),
}))
- ack = await asyncio.wait_for(q_admin.get(), timeout=5.0)
- assert ack["type"] == "ack"
- assert ack["detail"] == "gek_bundle_stored"
+ invite = await asyncio.wait_for(q_admin.get(), timeout=5.0)
+ assert invite["type"] == MNP.INVITE_RESULT
+ code = invite["code"]
+ assert code and len(code) == 9 # XXXX-XXXX
- # Verify bundle was persisted
- stored = await bundle_store.fetch("g", "user-002")
- assert stored is not None
- assert stored["pk_eph_b64"] == bundle["pk_eph_b64"]
+ # 2. Bob connects. He cannot prove GEK possession — he has never had it — so
+ # he redeems the code in the pre-proof window instead.
+ sk_x_raw, pk_x_raw = x25519_keypair
+ sk_bob_ed = Ed25519PrivateKey.generate()
+ pc_bob, ch_bob, q_bob = await _open_channel(transport, "peer-bob")
- # Unwrap to verify it's correct
- recovered = unwrap_gek(stored, sk_x_raw, pk_x_raw)
- assert recovered == gek
+ nonce_c = os.urandom(NONCE_LEN)
+ ch_bob.send(_pack({
+ "type": MNP.HANDSHAKE, "v": MNP_VERSION,
+ "token": _token(sk_hub, "user-002", "peer-bob", TEST_GROUP),
+ "group_id": TEST_GROUP,
+ "nonce": base64.b64encode(nonce_c).decode(),
+ }))
+ challenge = await asyncio.wait_for(q_bob.get(), timeout=5.0)
+ assert challenge["type"] == MNP.HANDSHAKE_CHALLENGE
+ nonce_s = base64.b64decode(challenge["nonce"])
- await bundle_store.close()
+ # Bob signs a transcript naming the node, and he cannot complete the handshake
+ # that would prove its key — he has no GEK yet. So he has to be able to learn
+ # it from the challenge; taking it from the test's own knowledge of sk_node
+ # would hide the fact that a real client cannot.
+ assert challenge["node_pk"] == pk_to_b64(sk_node.public_key()), (
+ "the challenge must announce the node key to a first-time joiner")
+ node_pk_b64 = challenge["node_pk"]
+
+ pk_ed_b64 = pk_to_b64(sk_bob_ed.public_key())
+ pk_x_b64 = base64.b64encode(pk_x_raw).decode()
+ ts = int(time.time())
+ transcript = join_transcript(
+ node_pk_b64=node_pk_b64,
+ group_id=TEST_GROUP, user_id="user-002",
+ pk_ed25519_b64=pk_ed_b64, pk_x25519_b64=pk_x_b64,
+ nonce_node=nonce_s, ts=ts,
+ )
+ ch_bob.send(_pack({
+ "type": MNP.JOIN_REQUEST, "v": MNP_VERSION,
+ "group_id": TEST_GROUP,
+ "pk_ed25519": pk_ed_b64, "pk_x25519": pk_x_b64,
+ "code": code, "ts": ts,
+ "sig": base64.b64encode(sk_bob_ed.sign(transcript)).decode(),
+ }))
+
+ result = await asyncio.wait_for(q_bob.get(), timeout=5.0)
+ assert result["type"] == MNP.JOIN_RESULT
+ assert result["ok"] is True
+ assert result["gek"] is True
+ assert result["role"] == ROLE_MEMBER
+
+ # 3. The key really is the group key, and only Bob's secret opens it.
+ assert unwrap_gek_aes(result, sk_x_raw, pk_x_raw) == gek
+
+ # 4. The code is spent.
+ assert await roster.consume_invite(code, "user-002") is None
+
+ await roster.close()
await pc_admin.close()
+ await pc_bob.close()
await transport.close_all()
@@ -1124,8 +1257,10 @@ async def test_gek_bundle_fetch_during_handshake(sk_node, sk_hub, gek, shared_di
# Step 1: Send handshake with group_id so _pending_group is set
token = _make_jwt(sk_hub, groups=["g"])
+ nonce_c = os.urandom(NONCE_LEN)
channel.send(_pack({
"type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g",
+ "nonce": base64.b64encode(nonce_c).decode(),
}))
msg = await asyncio.wait_for(received.get(), timeout=5.0)
assert msg["type"] == MNP.HANDSHAKE_CHALLENGE
@@ -1140,11 +1275,12 @@ async def test_gek_bundle_fetch_during_handshake(sk_node, sk_hub, gek, shared_di
recovered_gek = unwrap_gek(bundle_resp, sk_x_raw, pk_x_raw)
assert recovered_gek == gek
- nonce = base64.b64decode(msg["nonce"])
- offer_fp = _extract_dtls_fp(browser_pc.localDescription.sdp)
- answer_fp = _extract_dtls_fp(browser_pc.remoteDescription.sdp)
- proof = hmac.new(recovered_gek, nonce + offer_fp + answer_fp,
- hashlib.sha256).digest()
+ nonce_s = base64.b64decode(msg["nonce"])
+ binding = webrtc_binding(
+ _extract_dtls_fp(browser_pc.localDescription.sdp),
+ _extract_dtls_fp(browser_pc.remoteDescription.sdp),
+ )
+ proof = make_proof(recovered_gek, ROLE_CLIENT, "g", nonce_c, nonce_s, binding)
# Step 4: Complete handshake
channel.send(_pack({
@@ -1229,8 +1365,10 @@ async def test_keypair_bundle_store_and_fetch(sk_node, sk_hub, gek, shared_dir,
await asyncio.wait_for(ready.wait(), timeout=5.0)
token = _make_jwt(sk_hub, groups=["g"])
+ nonce_c = os.urandom(NONCE_LEN)
channel.send(_pack({
"type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g",
+ "nonce": base64.b64encode(nonce_c).decode(),
}))
msg = await asyncio.wait_for(received.get(), timeout=5.0)
assert msg["type"] == MNP.HANDSHAKE_CHALLENGE
@@ -1296,8 +1434,10 @@ async def test_keypair_bundle_fetch_not_found(sk_node, sk_hub, gek, shared_dir,
await asyncio.wait_for(ready.wait(), timeout=5.0)
token = _make_jwt(sk_hub, groups=["g"])
+ nonce_c = os.urandom(NONCE_LEN)
channel.send(_pack({
"type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g",
+ "nonce": base64.b64encode(nonce_c).decode(),
}))
msg = await asyncio.wait_for(received.get(), timeout=5.0)
assert msg["type"] == MNP.HANDSHAKE_CHALLENGE
@@ -1313,9 +1453,18 @@ async def test_keypair_bundle_fetch_not_found(sk_node, sk_hub, gek, shared_dir,
@pytest.mark.asyncio
-async def test_gek_auto_activate_on_node_bundle_store(sk_node, sk_hub, gek, shared_dir,
+async def test_gek_not_auto_activated_on_bundle_store(sk_node, sk_hub, gek, shared_dir,
tmp_path, x25519_keypair):
- """Storing the node operator's GEK bundle auto-activates GEK (AES variant)."""
+ """
+ A GEK bundle arriving over MNP must NOT become the node's live key (C5b).
+
+ This test previously asserted the opposite: storing a bundle addressed to the
+ node operator auto-activated it, with no signature required. Because the
+ operator's X25519 public key is public — the node publishes it in handshake_ack
+ — any group member could wrap a key of their own choosing for it and take over
+ the group, locking every legitimate member out. GEK activation now happens only
+ through the node's local admin UI or CLI.
+ """
hub_pk_pem = _hub_pk_pem(sk_hub)
indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek)
await indexer.initial_scan()
@@ -1324,7 +1473,8 @@ async def test_gek_auto_activate_on_node_bundle_store(sk_node, sk_hub, gek, shar
bundle_store = BundleStore(db_path=tmp_path / "bundles.db")
await bundle_store.open()
- new_gek = generate_gek()
+ attacker_gek = generate_gek()
+ assert attacker_gek != gek
transport = WebRTCTransport(
sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek,
@@ -1336,14 +1486,18 @@ async def test_gek_auto_activate_on_node_bundle_store(sk_node, sk_hub, gek, shar
transport._ctx["sk_x25519_raw"] = sk_x_raw
transport._ctx["pk_x25519_raw"] = pk_x_raw
transport._ctx["pk_x25519_b64"] = base64.b64encode(pk_x_raw).decode()
+ transport._ctx["admin_pk_ed25519"] = Ed25519PrivateKey.generate().public_key()
pc_admin, ch_admin, q_admin = await _setup_peer(
transport, sk_hub, gek, "peer-setup-admin")
- # Store GEK bundle wrapped with AES-GCM (browser-compatible)
- node_bundle = wrap_gek_aes(new_gek, pk_x_raw)
+ # An ordinary member wraps a key of their choosing for the operator's public
+ # key and offers it to the node. The message that used to carry this no longer
+ # exists (the node wraps the GEK itself now), so it reaches no handler at all —
+ # a stronger outcome than the admin challenge this test used to assert.
+ node_bundle = wrap_gek_aes(attacker_gek, pk_x_raw)
ch_admin.send(_pack({
- "type": MNP.GEK_BUNDLE_STORE,
+ "type": "gek_bundle_store",
"v": MNP_VERSION,
"user_id": "node-operator",
"group_id": "g",
@@ -1351,12 +1505,12 @@ async def test_gek_auto_activate_on_node_bundle_store(sk_node, sk_hub, gek, shar
"nonce_b64": node_bundle["nonce_b64"],
"wrapped_b64": node_bundle["wrapped_b64"],
}))
- ack = await asyncio.wait_for(q_admin.get(), timeout=5.0)
- assert ack["type"] == "ack"
- await asyncio.sleep(0.2)
+ await asyncio.sleep(0.5)
+ assert q_admin.empty(), "the retired bundle message still gets a response"
- assert transport._ctx.get("gek") == new_gek
+ assert transport._ctx.get("gek") == gek, "group key was seized over MNP (C5b)"
+ assert await bundle_store.fetch("g", "node-operator") is None
await bundle_store.close()
await pc_admin.close()
@@ -1398,6 +1552,8 @@ async def test_webrtc_no_gek_connection_refused(sk_node, sk_hub, shared_dir):
token = _make_jwt(sk_hub)
channel.send(_pack({
"type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token,
+ "group_id": TEST_GROUP,
+ "nonce": base64.b64encode(os.urandom(NONCE_LEN)).decode(),
}))
msg = await asyncio.wait_for(received.get(), timeout=5.0)
@@ -1457,8 +1613,10 @@ async def test_gek_bundle_fetch_not_found(sk_node, sk_hub, gek, shared_dir, tmp_
await asyncio.wait_for(ready.wait(), timeout=5.0)
token = _make_jwt(sk_hub, groups=["g"])
+ nonce_c = os.urandom(NONCE_LEN)
channel.send(_pack({
"type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g",
+ "nonce": base64.b64encode(nonce_c).decode(),
}))
msg = await asyncio.wait_for(received.get(), timeout=5.0)
assert msg["type"] == MNP.HANDSHAKE_CHALLENGE