diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-08-13 03:56:30 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-08-13 03:56:30 +0200 |
| commit | f0248975908ad670fa8a820f865bf22ea8d0172d (patch) | |
| tree | f4af64d36cacaccb4f6d13436e001aeb57e861e3 /packages | |
| parent | 35130e5528a52161630fd1c93572e1b2b7cd911b (diff) | |
| download | meshbay-f0248975908ad670fa8a820f865bf22ea8d0172d.tar.gz | |
feat: Phase 12 — P2P crypto material, password split, node Ed25519 auth
Baseline commit capturing in-progress Phase 12 work that was already present
in the working tree (uncommitted) before the Phase 11.5 security remediation
begins. Committed as-is, without review or modification, so that remediation
changes arrive as a separable diff.
Contents: BundleStore (P2P GEK + keypair bundles), password split
(auth_key / bundle_key), node Ed25519 auth (POST /v1/nodes/auth, node-scoped
JWT), GEK-HMAC handshake proof with DTLS channel binding, Ed25519 admin
challenge-response, node local admin UI rewrite, browser key persistence.
Not authored in this session — captured to establish a baseline.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Diffstat (limited to 'packages')
29 files changed, 3378 insertions, 772 deletions
diff --git a/packages/meshbay-common/src/meshbay_common/protocol.py b/packages/meshbay-common/src/meshbay_common/protocol.py index bf83906..55dcdde 100644 --- a/packages/meshbay-common/src/meshbay_common/protocol.py +++ b/packages/meshbay-common/src/meshbay_common/protocol.py @@ -40,6 +40,16 @@ class MNP: STREAM_DATA = "stream_data" # node sends encrypted fMP4 segment STREAM_END = "stream_end" # node signals end of stream EPHEMERAL_STREAM = "ephemeral_stream" # reserved — mobile live push + HANDSHAKE_CHALLENGE = "handshake_challenge" # node → client: GEK proof nonce + HANDSHAKE_RESPONSE = "handshake_response" # client → node: HMAC(GEK, nonce) + ADMIN_CHALLENGE = "admin_challenge" # node → client: Ed25519 sign challenge + ADMIN_RESPONSE = "admin_response" # client → node: Ed25519 signature + GEK_BUNDLE_STORE = "gek_bundle_store" # client → node: store wrapped GEK for a user + GEK_BUNDLE_FETCH = "gek_bundle_fetch" # client → node: request own wrapped GEK + GEK_BUNDLE_RESP = "gek_bundle_resp" # node → client: wrapped GEK bundle + KEYPAIR_BUNDLE_STORE = "keypair_bundle_store" # client → node: store encrypted keypair bundle + KEYPAIR_BUNDLE_FETCH = "keypair_bundle_fetch" # client → node: request own keypair bundle + KEYPAIR_BUNDLE_RESP = "keypair_bundle_resp" # node → client: encrypted keypair bundle # ── Index entry ─────────────────────────────────────────────────────────────── @@ -54,6 +64,8 @@ class IndexEntry: added_at: int # unix timestamp duration: int | None = None # seconds, for media thumb_hash: str | None = None # blake3 of thumbnail + uploader_id: str | None = None # user_id of who uploaded (None = pre-existing on disk) + uploader_pk: str | None = None # Ed25519 public key of uploader (base64 raw 32 bytes) @dataclass diff --git a/packages/meshbay-hub/src/meshbay_hub/api/admin.py b/packages/meshbay-hub/src/meshbay_hub/api/admin.py index 88e231a..164885d 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/admin.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/admin.py @@ -211,6 +211,7 @@ async def admin_list_groups( "name": g.name, "admin_id": g.admin_id, "visibility": g.visibility, + "description": g.description or "", "status": g.status, "created_at": g.created_at.isoformat(), "member_count": mc, @@ -265,25 +266,30 @@ async def admin_list_logs( offset: int = 0, limit: int = Query(default=50, le=200), ): - query = select(IPLog).order_by(IPLog.timestamp.desc()) + query = ( + select(IPLog, User.username) + .outerjoin(User, IPLog.user_id == User.id) + .order_by(IPLog.timestamp.desc()) + ) if user_id: query = query.where(IPLog.user_id == user_id) if event: query = query.where(IPLog.event == event) query = query.offset(offset).limit(limit) result = await db.execute(query) - logs = result.scalars().all() + rows = result.all() return { "logs": [ { "id": lg.id, "user_id": lg.user_id, + "username": uname or "", "event": lg.event, "ip_address": lg.ip_address, "detail": lg.detail, "timestamp": lg.timestamp.isoformat(), } - for lg in logs + for lg, uname in rows ], } diff --git a/packages/meshbay-hub/src/meshbay_hub/api/deps.py b/packages/meshbay-hub/src/meshbay_hub/api/deps.py index addba30..1bf57a4 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/deps.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/deps.py @@ -1,5 +1,10 @@ """ FastAPI shared dependencies — injected via Depends(). + +JWT scope enforcement: + - "user" scope (browser login): full access to all endpoints + - "node" scope (Ed25519 daemon auth): read-only group access + node operations + Node-scoped tokens CANNOT create/delete groups or manage membership. """ from fastapi import Depends, Header, HTTPException, status @@ -18,20 +23,13 @@ def set_admin_usernames(usernames: list[str]) -> None: _admin_usernames = set(usernames) -async def get_current_user( - authorization: str = Header(...), - db: AsyncSession = Depends(get_db), -) -> User: - """ - Verify the JWT bearer token and return the User from the database. - Node clients: verified locally with hub PK — no DB round-trip needed. - Hub API (web): must confirm user still exists and is active. - """ +async def _decode_token(authorization: str = Header(...)) -> dict: + """Decode and verify JWT bearer token. Returns full payload.""" try: scheme, token = authorization.split(None, 1) if scheme.lower() != "bearer": raise ValueError - payload = decode_access_token(token) + return decode_access_token(token) except Exception: raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, @@ -39,6 +37,15 @@ async def get_current_user( headers={"WWW-Authenticate": "Bearer"}, ) + +async def get_current_user( + payload: dict = Depends(_decode_token), + db: AsyncSession = Depends(get_db), +) -> User: + """ + Verify the JWT bearer token and return the User from the database. + Accepts both user-scoped and node-scoped tokens. + """ result = await db.execute( select(User).where(User.id == payload["sub"])) user = result.scalar_one_or_none() @@ -52,6 +59,19 @@ async def get_current_user( return user +async def require_user_scope( + payload: dict = Depends(_decode_token), + current_user: User = Depends(get_current_user), +) -> User: + """Reject node-scoped tokens — only browser (user-scope) can mutate groups.""" + if payload.get("scope") == "node": + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail="Node-scoped token cannot perform this operation — use browser", + ) + return current_user + + async def require_moderator( current_user: User = Depends(get_current_user), ) -> User: diff --git a/packages/meshbay-hub/src/meshbay_hub/api/groups.py b/packages/meshbay-hub/src/meshbay_hub/api/groups.py index e0ea016..88af764 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/groups.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/groups.py @@ -5,10 +5,10 @@ from pydantic import BaseModel from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession -from meshbay_hub.api.deps import get_current_user +from meshbay_hub.api.deps import get_current_user, require_user_scope from meshbay_hub.db.engine import get_db from meshbay_hub.db.models import ( - FederatedGroup, GEKBundle, Group, GroupMember, + FederatedGroup, Group, GroupMember, IPLog, SwarmSource, User, ) @@ -37,6 +37,7 @@ async def my_groups( "join_policy": g.join_policy, "created_at": g.created_at.isoformat(), "is_admin": g.admin_id == current_user.id, + "description": g.description or "", } for g in groups ] @@ -56,6 +57,8 @@ async def group_online_nodes( group = await db.get(Group, group_id) if not group: raise HTTPException(status_code=404, detail="Group not found") + if group.status != "active": + raise HTTPException(status_code=403, detail="Group is suspended") node_ids = get_online_nodes_for_group(group_id) nodes = [] @@ -85,6 +88,7 @@ async def list_public_groups( groups = [ { "id": g.id, "name": g.name, "join_policy": g.join_policy, + "description": g.description or "", "created_at": g.created_at.isoformat(), "source": "local", } for g in local @@ -193,7 +197,7 @@ async def group_members( async def join_group( group_id: str, request: Request, - current_user: User = Depends(get_current_user), + current_user: User = Depends(require_user_scope), db: AsyncSession = Depends(get_db), ): group = await db.get(Group, group_id) @@ -219,26 +223,23 @@ class GroupCreateRequest(BaseModel): name: str visibility: str = "private" # public|private join_policy: str = "invite" # open|request|invite - - -class GEKBundleRequest(BaseModel): - pk_eph_b64: str - nonce_b64: str - wrapped_b64: str + description: str | None = None @router.post("", status_code=201) async def create_group( body: GroupCreateRequest, request: Request, - current_user: User = Depends(get_current_user), + current_user: User = Depends(require_user_scope), db: AsyncSession = Depends(get_db), ): + desc = (body.description or "")[:512] if body.description else None group = Group( name=body.name, admin_id=current_user.id, visibility=body.visibility, join_policy=body.join_policy, + description=desc, ) db.add(group) await db.flush() # get group.id @@ -251,12 +252,11 @@ async def create_group( return {"group_id": group.id, "name": group.name} -@router.post("/{group_id}/members/{username}/gek", status_code=201) -async def store_gek_bundle( +@router.post("/{group_id}/members/{username}", status_code=201) +async def add_group_member( group_id: str, username: str, - body: GEKBundleRequest, - current_user: User = Depends(get_current_user), + current_user: User = Depends(require_user_scope), db: AsyncSession = Depends(get_db), ): group = await db.get(Group, group_id) @@ -270,26 +270,11 @@ async def store_gek_bundle( if not target: raise HTTPException(status_code=404, detail="User not found") - # Upsert GEK bundle - existing = await db.get(GEKBundle, (group_id, target.id)) new_member = False - if existing: - existing.pk_eph_b64 = body.pk_eph_b64 - existing.nonce_b64 = body.nonce_b64 - existing.wrapped_b64 = body.wrapped_b64 - else: - db.add(GEKBundle( - group_id=group_id, - user_id=target.id, - pk_eph_b64=body.pk_eph_b64, - nonce_b64=body.nonce_b64, - wrapped_b64=body.wrapped_b64, - )) - # Add member if not already in group - mem = await db.get(GroupMember, (group_id, target.id)) - if not mem: - db.add(GroupMember(group_id=group_id, user_id=target.id)) - new_member = True + mem = await db.get(GroupMember, (group_id, target.id)) + if not mem: + db.add(GroupMember(group_id=group_id, user_id=target.id)) + new_member = True if new_member: from meshbay_hub.api.notifications import create_notification @@ -303,26 +288,26 @@ async def store_gek_bundle( return {"status": "stored", "group_id": group_id, "username": username} -@router.get("/{group_id}/gek") -async def get_my_gek_bundle( +@router.delete("/{group_id}") +async def delete_group( group_id: str, - current_user: User = Depends(get_current_user), + request: Request, + current_user: User = Depends(require_user_scope), db: AsyncSession = Depends(get_db), ): group = await db.get(Group, group_id) if not group: raise HTTPException(status_code=404, detail="Group not found") + if group.admin_id != current_user.id: + raise HTTPException(status_code=403, detail="Only the group creator can delete") - bundle = await db.get(GEKBundle, (group_id, current_user.id)) - if not bundle: - raise HTTPException(status_code=404, detail="No GEK bundle for this user in this group") - - return { - "group_id": group_id, - "pk_eph_b64": bundle.pk_eph_b64, - "nonce_b64": bundle.nonce_b64, - "wrapped_b64": bundle.wrapped_b64, - } + from sqlalchemy import delete as sa_delete + await db.execute(sa_delete(GroupMember).where(GroupMember.group_id == group_id)) + db.add(IPLog(user_id=current_user.id, event="group_delete", + ip_address=_ip(request), detail=group.name)) + await db.delete(group) + await db.commit() + return {"status": "deleted", "group_id": group_id} def _ip(request: Request) -> str: diff --git a/packages/meshbay-hub/src/meshbay_hub/api/nodes.py b/packages/meshbay-hub/src/meshbay_hub/api/nodes.py index b970aa8..321e43c 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/nodes.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/nodes.py @@ -1,15 +1,84 @@ """Node endpoints — /v1/nodes/*""" +import base64 +import time + +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey +from cryptography.exceptions import InvalidSignature from fastapi import APIRouter, Depends, HTTPException, Request from pydantic import BaseModel +from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession +from meshbay_hub.auth import issue_access_token from meshbay_hub.api.deps import get_current_user +from meshbay_hub.api.middleware import limiter from meshbay_hub.db.engine import get_db -from meshbay_hub.db.models import IPLog, Node, User +from meshbay_hub.db.models import GroupMember, IPLog, Node, User router = APIRouter(prefix="/v1/nodes", tags=["nodes"]) +NODE_AUTH_TIMESTAMP_WINDOW = 60 # seconds + + +class NodeAuthRequest(BaseModel): + username: str + timestamp: int # unix epoch seconds + signature: str # base64 Ed25519 signature + + +@router.post("/auth") +@limiter.limit("10/minute") +async def node_auth( + body: NodeAuthRequest, + request: Request, + db: AsyncSession = Depends(get_db), +): + """Authenticate a node daemon via Ed25519 challenge-response. Returns node-scoped JWT.""" + now = int(time.time()) + if abs(now - body.timestamp) > NODE_AUTH_TIMESTAMP_WINDOW: + raise HTTPException(status_code=401, detail="Timestamp too old or too far in the future") + + result = await db.execute(select(User).where(User.username == body.username)) + user = result.scalar_one_or_none() + if not user: + raise HTTPException(status_code=401, detail="Invalid credentials") + if user.status != "active": + raise HTTPException(status_code=403, detail=f"Account {user.status}") + + if not user.pk_node_ed25519: + raise HTTPException( + status_code=401, + detail="No node key registered — link your node from the browser first", + ) + + message = f"meshbay:node_auth:{body.username}:{body.timestamp}".encode() + try: + pk_raw = base64.b64decode(user.pk_node_ed25519) + pk = Ed25519PublicKey.from_public_bytes(pk_raw) + sig = base64.b64decode(body.signature) + pk.verify(sig, message) + except (InvalidSignature, Exception): + db.add(IPLog(event="node_auth_fail", ip_address=_ip(request), detail=body.username)) + await db.commit() + raise HTTPException(status_code=401, detail="Invalid signature") + + memberships = await db.execute( + select(GroupMember.group_id).where(GroupMember.user_id == user.id)) + group_ids = [gid for (gid,) in memberships.all()] + + access_token = issue_access_token( + user.id, user.pk_node_ed25519, ttl=3600, groups=group_ids, scope="node") + + db.add(IPLog(user_id=user.id, event="node_auth", ip_address=_ip(request))) + await db.commit() + + return { + "access_token": access_token, + "token_type": "bearer", + "expires_in": 3600, + } + class NodeAnnounceRequest(BaseModel): pk_node: str diff --git a/packages/meshbay-hub/src/meshbay_hub/api/users.py b/packages/meshbay-hub/src/meshbay_hub/api/users.py index 2c2eede..53238de 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/users.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/users.py @@ -1,5 +1,6 @@ """User endpoints — /v1/users/*""" +import base64 import uuid from datetime import datetime, timezone, timedelta @@ -24,7 +25,7 @@ from meshbay_hub.api.middleware import limiter from meshbay_hub.config import HubConfig from meshbay_hub.db.engine import get_db from meshbay_hub.db.models import GroupMember, IPLog, RefreshToken, User -from meshbay_hub.api.deps import get_current_user +from meshbay_hub.api.deps import get_current_user, require_user_scope router = APIRouter(prefix="/v1/users", tags=["users"]) @@ -46,10 +47,10 @@ def _refresh_ttl() -> int: class RegisterRequest(BaseModel): username: str email: str - password: str + password: str | None = None # deprecated — legacy native clients + auth_key: str | None = None # PBKDF2-derived, new clients pk_user_ed25519: str # base64 raw 32B pk_user_x25519: str # base64 raw 32B - keypair_bundle: str | None = None # AES-GCM encrypted bundle (web clients) @field_validator("username") @classmethod @@ -61,17 +62,11 @@ class RegisterRequest(BaseModel): raise ValueError("username: only letters, digits, -, _, .") return v - @field_validator("password") - @classmethod - def password_strength(cls, v: str) -> str: - if len(v) < 8: - raise ValueError("password must be at least 8 characters") - return v - class LoginRequest(BaseModel): username: str - password: str + password: str | None = None # legacy (raw password) for migration + auth_key: str | None = None # PBKDF2-derived auth key (new scheme) class RefreshRequest(BaseModel): @@ -92,18 +87,23 @@ async def register( if existing.scalar_one_or_none(): raise HTTPException(status_code=409, detail="Username already taken") - pw_hash, pw_salt = hash_password(body.password) + credential = body.auth_key or body.password + if not credential: + raise HTTPException(status_code=400, detail="auth_key or password required") + + pw_hash, pw_salt = hash_password(credential) + # auth_key → pw_version 3 (password split); raw password → pw_version 2 (legacy) + pw_ver = current_pw_version() if body.auth_key else 2 hub_id = _cfg.identity.id if _cfg else "meshbay.org" user = User( username=body.username, email=encrypt_email(body.email), pw_hash=pw_hash, pw_salt=pw_salt, - pw_version=current_pw_version(), + pw_version=pw_ver, pk_ed25519=body.pk_user_ed25519, pk_x25519=body.pk_user_x25519, hub_id=hub_id, - keypair_bundle=body.keypair_bundle, ) db.add(user) db.add(IPLog( @@ -136,18 +136,52 @@ async def login( user = result.scalar_one_or_none() ip = _client_ip(request) - if not user or not verify_password( - body.password, user.pw_hash, user.pw_salt, version=user.pw_version - ): + + if not body.auth_key and not body.password: + raise HTTPException(status_code=401, detail="No credentials provided") + + if not user: db.add(IPLog(event="login_fail", ip_address=ip, detail=body.username)) await db.commit() raise HTTPException(status_code=401, detail="Invalid credentials") + if user.pw_version >= 3: + # New scheme: verify auth_key + if not body.auth_key or not verify_password( + body.auth_key, user.pw_hash, user.pw_salt, version=user.pw_version + ): + db.add(IPLog(event="login_fail", ip_address=ip, detail=body.username)) + await db.commit() + raise HTTPException(status_code=401, detail="Invalid credentials") + else: + # Legacy scheme: need raw password + if not body.password: + raise HTTPException(status_code=401, detail="auth_upgrade_required") + if not verify_password( + body.password, user.pw_hash, user.pw_salt, version=user.pw_version + ): + db.add(IPLog(event="login_fail", ip_address=ip, detail=body.username)) + await db.commit() + raise HTTPException(status_code=401, detail="Invalid credentials") + # Migrate to new scheme if auth_key provided alongside password + if body.auth_key: + new_hash, new_salt = hash_password(body.auth_key) + user.pw_hash = new_hash + user.pw_salt = new_salt + user.pw_version = current_pw_version() + elif user.pw_version < 2: + # Legacy rehash: upgrade Argon2 params within the password scheme (v1 -> v2) + new_hash, new_salt = hash_password(body.password) + user.pw_hash = new_hash + user.pw_salt = new_salt + user.pw_version = 2 + if user.status != "active": raise HTTPException(status_code=403, detail=f"Account {user.status}") - if pw_needs_rehash(user.pw_version): - new_hash, new_salt = hash_password(body.password) + # Rehash within the auth_key scheme if Argon2 params upgraded beyond v3 + if user.pw_version >= 3 and pw_needs_rehash(user.pw_version): + new_hash, new_salt = hash_password(body.auth_key) user.pw_hash = new_hash user.pw_salt = new_salt user.pw_version = current_pw_version() @@ -168,15 +202,12 @@ async def login( db.add(IPLog(user_id=user.id, event="login", ip_address=ip)) await db.commit() - resp = { + return { "access_token": access_token, "refresh_token": raw_rt, "token_type": "bearer", "expires_in": _ttl(), } - if user.keypair_bundle: - resp["keypair_bundle"] = user.keypair_bundle # encrypted, for web clients - return resp @router.post("/token/refresh") @@ -248,6 +279,64 @@ async def get_current_user_info( } +class NodeKeyRequest(BaseModel): + pk_node_ed25519: str # base64 raw 32B Ed25519 public key + + +@router.put("/me/node_key") +async def register_node_key( + body: NodeKeyRequest, + current_user: User = Depends(require_user_scope), + db: AsyncSession = Depends(get_db), +): + """Link a node daemon's Ed25519 public key to the operator's account.""" + try: + raw = base64.b64decode(body.pk_node_ed25519) + if len(raw) != 32: + raise ValueError + except Exception: + raise HTTPException(status_code=400, detail="Invalid Ed25519 public key (need 32 bytes base64)") + + current_user.pk_node_ed25519 = body.pk_node_ed25519 + await db.commit() + return {"status": "stored", "pk_node_ed25519": body.pk_node_ed25519} + + +class RotateKeysRequest(BaseModel): + pk_user_ed25519: str # base64 raw 32B + pk_user_x25519: str # base64 raw 32B + + +@router.put("/me/keys") +async def rotate_browser_keys( + body: RotateKeysRequest, + current_user: User = Depends(require_user_scope), + db: AsyncSession = Depends(get_db), +): + for field, label in [ + (body.pk_user_ed25519, "Ed25519"), + (body.pk_user_x25519, "X25519"), + ]: + try: + raw = base64.b64decode(field) + if len(raw) != 32: + raise ValueError + except Exception: + raise HTTPException( + status_code=400, + detail=f"Invalid {label} public key (need 32 bytes base64)", + ) + + current_user.pk_ed25519 = body.pk_user_ed25519 + current_user.pk_x25519 = body.pk_user_x25519 + await db.commit() + return { + "status": "updated", + "pk_ed25519": body.pk_user_ed25519, + "pk_x25519": body.pk_user_x25519, + } + + @router.get("/{username}/pubkeys") async def get_user_pubkeys( username: str, @@ -258,12 +347,15 @@ async def get_user_pubkeys( target = result.scalar_one_or_none() if not target: raise HTTPException(status_code=404, detail="User not found") - return { + resp = { "user_id": target.id, "username": target.username, "pk_ed25519": target.pk_ed25519, "pk_x25519": target.pk_x25519, } + if target.pk_node_ed25519: + resp["pk_node_ed25519"] = target.pk_node_ed25519 + return resp def _client_ip(request: Request) -> str: diff --git a/packages/meshbay-hub/src/meshbay_hub/auth.py b/packages/meshbay-hub/src/meshbay_hub/auth.py index 11ad112..563a1eb 100644 --- a/packages/meshbay-hub/src/meshbay_hub/auth.py +++ b/packages/meshbay-hub/src/meshbay_hub/auth.py @@ -30,8 +30,9 @@ _ARGON2_KEY_LEN = 32 _ARGON2_VERSIONS = { 1: {"iterations": 3, "memory_cost": 65536}, # 64 MB — initial 2: {"iterations": 3, "memory_cost": 262144}, # 256 MB — production target + 3: {"iterations": 3, "memory_cost": 262144}, # 256 MB — auth_key input (password split) } -_ARGON2_CURRENT_VERSION = 2 +_ARGON2_CURRENT_VERSION = 3 # Module-level hub keypair (loaded once at startup) _hub_sk_pem: bytes | None = None @@ -133,11 +134,13 @@ def issue_access_token( pk_user: str, ttl: int = 3600, groups: list[str] | None = None, + scope: str = "user", ) -> str: """ Issue a signed JWT access token. Includes jti (UUID4) — required to prevent replay and enable revocation. Includes groups — list of group_ids the user is a member of (node-side authz). + scope: "user" (browser, full access) or "node" (daemon, restricted). """ if _hub_sk_pem is None: raise RuntimeError("Hub keypair not loaded") @@ -151,6 +154,7 @@ def issue_access_token( "iat": now, "exp": now + ttl, "groups": groups or [], + "scope": scope, } return jwt.encode(payload, _hub_sk_pem, algorithm="EdDSA") diff --git a/packages/meshbay-hub/src/meshbay_hub/db/__init__.py b/packages/meshbay-hub/src/meshbay_hub/db/__init__.py index 5ef1d3c..62e5388 100644 --- a/packages/meshbay-hub/src/meshbay_hub/db/__init__.py +++ b/packages/meshbay-hub/src/meshbay_hub/db/__init__.py @@ -1,9 +1,9 @@ """Hub database layer.""" from .engine import init_db, close_db, get_db -from .models import Base, User, Node, Group, GroupMember, GEKBundle, RefreshToken, IPLog +from .models import Base, User, Node, Group, GroupMember, RefreshToken, IPLog __all__ = [ "init_db", "close_db", "get_db", "Base", "User", "Node", "Group", "GroupMember", - "GEKBundle", "RefreshToken", "IPLog", + "RefreshToken", "IPLog", ] diff --git a/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/2041a4060b3c_add_keypair_bundle_federated_groups_.py b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/2041a4060b3c_add_keypair_bundle_federated_groups_.py index 779efdc..a59a3d1 100644 --- a/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/2041a4060b3c_add_keypair_bundle_federated_groups_.py +++ b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/2041a4060b3c_add_keypair_bundle_federated_groups_.py @@ -63,14 +63,12 @@ def upgrade() -> None: ) op.create_index('ix_content_reports_group', 'content_reports', ['group_id'], unique=False) op.create_index('ix_content_reports_hash', 'content_reports', ['content_hash'], unique=False) - op.add_column('users', sa.Column('keypair_bundle', sa.Text(), nullable=True)) # ### end Alembic commands ### def downgrade() -> None: """Downgrade schema.""" # ### commands auto generated by Alembic - please adjust! ### - op.drop_column('users', 'keypair_bundle') op.drop_index('ix_content_reports_hash', table_name='content_reports') op.drop_index('ix_content_reports_group', table_name='content_reports') op.drop_table('content_reports') diff --git a/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/d28b9caf9f07_initial_schema.py b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/d28b9caf9f07_initial_schema.py index d4a9aa6..5192eb8 100644 --- a/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/d28b9caf9f07_initial_schema.py +++ b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/d28b9caf9f07_initial_schema.py @@ -84,17 +84,6 @@ def upgrade() -> None: sa.UniqueConstraint('token_hash') ) op.create_index('ix_refresh_tokens_hash', 'refresh_tokens', ['token_hash'], unique=False) - op.create_table('gek_bundles', - sa.Column('group_id', sa.String(length=36), nullable=False), - sa.Column('user_id', sa.String(length=36), nullable=False), - sa.Column('pk_eph_b64', sa.String(length=64), nullable=False), - sa.Column('nonce_b64', sa.String(length=32), nullable=False), - sa.Column('wrapped_b64', sa.String(length=128), nullable=False), - sa.Column('stored_at', sa.DateTime(timezone=True), nullable=False), - sa.ForeignKeyConstraint(['group_id'], ['groups.id'], ), - sa.ForeignKeyConstraint(['user_id'], ['users.id'], ), - sa.PrimaryKeyConstraint('group_id', 'user_id') - ) op.create_table('group_members', sa.Column('group_id', sa.String(length=36), nullable=False), sa.Column('user_id', sa.String(length=36), nullable=False), @@ -110,7 +99,6 @@ def downgrade() -> None: """Downgrade schema.""" # ### commands auto generated by Alembic - please adjust! ### op.drop_table('group_members') - op.drop_table('gek_bundles') op.drop_index('ix_refresh_tokens_hash', table_name='refresh_tokens') op.drop_table('refresh_tokens') op.drop_index('ix_nodes_user_id', table_name='nodes') diff --git a/packages/meshbay-hub/src/meshbay_hub/db/models.py b/packages/meshbay-hub/src/meshbay_hub/db/models.py index c420661..cdebd3c 100644 --- a/packages/meshbay-hub/src/meshbay_hub/db/models.py +++ b/packages/meshbay-hub/src/meshbay_hub/db/models.py @@ -6,7 +6,6 @@ Tables: nodes — node announcements groups — group registry group_members — group membership - gek_bundles — encrypted GEK per (group, user) refresh_tokens — hashed refresh tokens ip_logs — connection log for legal compliance (1-year retention) """ @@ -45,15 +44,14 @@ class User(Base): pw_version: Mapped[int] = mapped_column(Integer, default=1) pk_ed25519: Mapped[str] = mapped_column(String(64), nullable=False) # base64 raw 32B pk_x25519: Mapped[str] = mapped_column(String(64), nullable=False) # base64 raw 32B + pk_node_ed25519: Mapped[str | None] = mapped_column(String(64), nullable=True) # node daemon key hub_id: Mapped[str] = mapped_column(String(128), nullable=False) - keypair_bundle: Mapped[str | None] = mapped_column(Text) # AES-GCM encrypted, web clients only role: Mapped[str] = mapped_column(String(16), default="user") # user|moderator|admin status: Mapped[str] = mapped_column(String(16), default="active") # active|suspended|revoked created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_now) nodes: Mapped[list["Node"]] = relationship(back_populates="user") group_memberships: Mapped[list["GroupMember"]] = relationship(back_populates="user") - gek_bundles: Mapped[list["GEKBundle"]] = relationship(back_populates="user") refresh_tokens: Mapped[list["RefreshToken"]] = relationship(back_populates="user") ip_logs: Mapped[list["IPLog"]] = relationship(back_populates="user") @@ -89,11 +87,11 @@ class Group(Base): admin_id: Mapped[str] = mapped_column(ForeignKey("users.id"), nullable=False) visibility: Mapped[str] = mapped_column(String(16), default="private") # public|private join_policy: Mapped[str] = mapped_column(String(16), default="invite") # open|request|invite - status: Mapped[str] = mapped_column(String(16), default="active") # active|revoked + description: Mapped[str | None] = mapped_column(String(512)) + status: Mapped[str] = mapped_column(String(16), default="active") # active|suspended|revoked created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_now) members: Mapped[list["GroupMember"]] = relationship(back_populates="group") - gek_bundles: Mapped[list["GEKBundle"]] = relationship(back_populates="group") __table_args__ = (Index("ix_groups_name", "name"),) @@ -109,23 +107,6 @@ class GroupMember(Base): user: Mapped["User"] = relationship(back_populates="group_memberships") -# ── GEK bundles ─────────────────────────────────────────────────────────────── - -class GEKBundle(Base): - """Encrypted GEK bundle — opaque to the hub (hub cannot decrypt it).""" - __tablename__ = "gek_bundles" - - group_id: Mapped[str] = mapped_column(ForeignKey("groups.id"), primary_key=True) - user_id: Mapped[str] = mapped_column(ForeignKey("users.id"), primary_key=True) - pk_eph_b64: Mapped[str] = mapped_column(String(64), nullable=False) - nonce_b64: Mapped[str] = mapped_column(String(32), nullable=False) - wrapped_b64: Mapped[str] = mapped_column(String(128), nullable=False) - stored_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_now) - - group: Mapped["Group"] = relationship(back_populates="gek_bundles") - user: Mapped["User"] = relationship(back_populates="gek_bundles") - - # ── Refresh tokens ──────────────────────────────────────────────────────────── class RefreshToken(Base): diff --git a/packages/meshbay-hub/src/meshbay_hub/static/app.js b/packages/meshbay-hub/src/meshbay_hub/static/app.js index d8f9df5..ddff928 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/app.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/app.js @@ -66,6 +66,58 @@ async function getAllCachedIndexes() { // ── Auth persistence ───────────────────────────────────────────────────────── let _sessionKeys = null; +let _bundleKey = null; +let _pendingBundlePush = null; + +function _openKeyDB() { + return new Promise((resolve, reject) => { + const req = indexedDB.open('meshbay_keys', 1); + req.onupgradeneeded = () => req.result.createObjectStore('k'); + req.onsuccess = () => resolve(req.result); + req.onerror = () => reject(req.error); + }); +} +async function _storeBundleKey(key) { + try { + const db = await _openKeyDB(); + const tx = db.transaction('k', 'readwrite'); + tx.objectStore('k').put(key, 'bk'); + await new Promise(r => { tx.oncomplete = r; }); + db.close(); + } catch {} +} +async function _loadBundleKey() { + try { + const db = await _openKeyDB(); + const tx = db.transaction('k', 'readonly'); + const g = tx.objectStore('k').get('bk'); + const val = await new Promise(r => { g.onsuccess = () => r(g.result); }); + db.close(); + return val || null; + } catch { return null; } +} +async function _clearKeyDB() { + try { + const db = await _openKeyDB(); + const tx = db.transaction('k', 'readwrite'); + tx.objectStore('k').clear(); + await new Promise(r => { tx.oncomplete = r; }); + db.close(); + } catch {} +} +function _saveSessionKeys() { + try { + if (_sessionKeys) sessionStorage.setItem('meshbay_sk', JSON.stringify(_sessionKeys)); + } catch {} +} +function _restoreSessionKeys() { + try { + if (!_sessionKeys) { + const sk = sessionStorage.getItem('meshbay_sk'); + if (sk) _sessionKeys = JSON.parse(sk); + } + } catch {} +} function loadAuth() { try { @@ -81,6 +133,10 @@ function saveAuth(auth) { } else { localStorage.removeItem(AUTH_KEY); _sessionKeys = null; + _bundleKey = null; + _pendingBundlePush = null; + _clearKeyDB(); + try { sessionStorage.removeItem('meshbay_sk'); } catch {} } } @@ -139,9 +195,74 @@ function navigate(path) { const AuthContext = createContext(null); function useAuth() { return useContext(AuthContext); } +// ── User Menu ──────────────────────────────────────────────────────────────── + +function UserMenu({ user, theme, onThemeChange, onLogout }) { + const [open, setOpen] = useState(false); + const [langOpen, setLangOpen] = useState(false); + const ref = useRef(null); + + useEffect(() => { + if (!open) return; + const close = (e) => { + if (ref.current && !ref.current.contains(e.target)) setOpen(false); + }; + document.addEventListener('click', close); + return () => document.removeEventListener('click', close); + }, [open]); + + const resolved = resolveTheme(theme); + + return html` + <div class="user-menu-wrap" ref=${ref}> + <button class="user-menu-trigger" onClick=${() => setOpen(o => !o)}> + <span class="user-menu-avatar">${user.username[0].toUpperCase()}</span> + <span class="user-menu-name">${user.username}</span> + <span class="user-menu-caret">${open ? '▴' : '▾'}</span> + </button> + ${open && html` + <div class="user-menu-dropdown"> + <div class="user-menu-header"> + <span class="user-menu-avatar lg">${user.username[0].toUpperCase()}</span> + <div> + <div class="user-menu-uname">${user.username}</div> + <div class="user-menu-role">${user.role || 'user'}</div> + </div> + </div> + <div class="user-menu-divider"></div> + <button class="user-menu-item" onClick=${(e) => { e.stopPropagation(); setLangOpen(o => !o); }}> + <span class="umi-icon">${'\u{1F310}'}</span> ${t('usermenu.language')} + <span class="umi-arrow">${langOpen ? '▴' : '▾'}</span> + </button> + ${langOpen && LOCALES.map(l => html` + <button key=${l.code} class="user-menu-item user-menu-sub" + onClick=${() => { setLocale(l.code); window.location.reload(); }}> + <span class="umi-icon">${l.flag}</span> ${l.name} + ${getLocale() === l.code && html`<span class="umi-check">${'✓'}</span>`} + </button> + `)} + <button class="user-menu-item" onClick=${() => { setOpen(false); navigate('/settings'); }}> + <span class="umi-icon">${'⚙'}</span> ${t('usermenu.settings')} + </button> + <button class="user-menu-item" onClick=${() => { + onThemeChange(resolved === 'dark' ? 'light' : 'dark'); + }}> + <span class="umi-icon">${resolved === 'dark' ? '☀' : '☾'}</span> + ${' '}${resolved === 'dark' ? t('usermenu.theme_light') : t('usermenu.theme_dark')} + </button> + <div class="user-menu-divider"></div> + <button class="user-menu-item user-menu-logout" onClick=${onLogout}> + <span class="umi-icon">${'⏻'}</span> ${t('usermenu.logout')} + </button> + </div> + `} + </div> + `; +} + // ── Nav ────────────────────────────────────────────────────────────────────── -function Nav({ user, theme, onThemeToggle, onLogout, onMenuToggle, unreadCount }) { +function Nav({ user, theme, onThemeChange, onLogout, onMenuToggle, unreadCount }) { return html` <nav class="nav"> <div class="nav-left"> @@ -157,13 +278,9 @@ function Nav({ user, theme, onThemeToggle, onLogout, onMenuToggle, unreadCount } ${'\u{1F514}'}${unreadCount > 0 && html`<span class="notif-badge">${unreadCount}</span>`} </a> `} - <button class="nav-theme" onClick=${onThemeToggle} - aria-label="${t('nav.toggle_menu')}" title=${theme === 'dark' ? t('nav.light_mode') : t('nav.dark_mode')}> - ${theme === 'dark' ? '☀' : '☾'} - </button> ${user ? html` - <span class="nav-user">${user.username}</span> - <button class="nav-btn" onClick=${onLogout}>${t('nav.logout')}</button> + <${UserMenu} user=${user} theme=${theme} + onThemeChange=${onThemeChange} onLogout=${onLogout} /> ` : html` <a class="nav-btn" href="#/login">${t('nav.login')}</a> `} @@ -178,6 +295,22 @@ function Sidebar({ groups, route, menuOpen, role }) { const isStaff = role === 'moderator' || role === 'admin'; return html` <aside class="sidebar ${menuOpen ? 'open' : ''}"> + ${isStaff && html` + <div class="sidebar-section"> + <div class="sidebar-heading">${t('sidebar.admin')}</div> + <a class="sidebar-item ${route === '/admin' ? 'active' : ''}" + href="#/admin">${'\u{1F6E1}'} ${t('admin.title')}</a> + </div> + `} + <div class="sidebar-section"> + <div class="sidebar-heading">${t('sidebar.discover')}</div> + <a class="sidebar-item ${route === '/explore' ? 'active' : ''}" + href="#/explore">${t('sidebar.public_groups')}</a> + <a class="sidebar-item ${route === '/search' ? 'active' : ''}" + href="#/search">${t('sidebar.search')}</a> + <a class="sidebar-item ${route === '/create-group' ? 'active' : ''}" + href="#/create-group">${t('sidebar.create_group')}</a> + </div> <div class="sidebar-section"> <div class="sidebar-heading">${t('sidebar.my_groups')}</div> ${groups.length === 0 @@ -191,21 +324,6 @@ function Sidebar({ groups, route, menuOpen, role }) { `) } </div> - <div class="sidebar-section"> - <div class="sidebar-heading">${t('sidebar.discover')}</div> - <a class="sidebar-item ${route === '/explore' ? 'active' : ''}" - href="#/explore">${t('sidebar.public_groups')}</a> - <a class="sidebar-item ${route === '/search' ? 'active' : ''}" - href="#/search">${t('sidebar.search')}</a> - <a class="sidebar-item ${route === '/create-group' ? 'active' : ''}" - href="#/create-group">${t('sidebar.create_group')}</a> - <a class="sidebar-item ${route === '/settings' ? 'active' : ''}" - href="#/settings">${t('sidebar.settings')}</a> - ${isStaff && html` - <a class="sidebar-item ${route === '/admin' ? 'active' : ''}" - href="#/admin">${t('sidebar.admin')}</a> - `} - </div> </aside> `; } @@ -455,6 +573,7 @@ function ExplorePage({ token, myGroupIds }) { <div key=${g.id} class="group-card"> <a href="#/group/${g.id}" style="text-decoration:none;color:inherit"> <h3>${g.name}</h3> + ${g.description && html`<p class="group-card-desc">${g.description}</p>`} </a> <span class="badge">${g.join_policy}</span> ${g.source && g.source !== 'local' && html` @@ -484,6 +603,7 @@ function ExplorePage({ token, myGroupIds }) { function CreateGroupPage({ token, onCreated }) { const [name, setName] = useState(''); + const [description, setDescription] = useState(''); const [visibility, setVisibility] = useState('private'); const [joinPolicy, setJoinPolicy] = useState('invite'); const [error, setError] = useState(''); @@ -495,26 +615,12 @@ function CreateGroupPage({ token, onCreated }) { setLoading(true); setError(''); try { + const body = { name: name.trim(), visibility, join_policy: joinPolicy }; + if (description.trim()) body.description = description.trim().slice(0, 512); const data = await hubFetch('/v1/groups', { - method: 'POST', token, - body: { name: name.trim(), visibility, join_policy: joinPolicy }, + method: 'POST', token, body, }); - if (window.MeshBayCrypto) { - try { - const gek = window.MeshBayCrypto.generateGEK(); - const meResp = await hubFetch('/v1/users/me', { token }); - const pubkeys = await hubFetch(`/v1/users/${meResp.username}/pubkeys`); - const pkXBytes = Uint8Array.from(atob(pubkeys.pk_x25519), c => c.charCodeAt(0)); - const bundle = await window.MeshBayCrypto.wrapGEK(gek, pkXBytes); - await hubFetch(`/v1/groups/${data.group_id}/members/${meResp.username}/gek`, { - method: 'POST', token, body: bundle, - }); - } catch (e) { - console.warn('GEK wrap skipped:', e); - } - } - if (onCreated) onCreated(); navigate('/'); } catch (err) { @@ -525,31 +631,70 @@ function CreateGroupPage({ token, onCreated }) { }; return html` - <div class="page-center"> - <h2>${t('create_group.title')}</h2> - <p class="page-message" style="max-width:380px;text-align:center;margin-bottom:12px"> - ${t('create_group.hint')} - </p> - ${error && html`<p class="error-msg">${error}</p>`} - <form class="login-form" onSubmit=${onSubmit}> - <input type="text" placeholder="${t('create_group.name')}" - value=${name} onInput=${e => setName(e.target.value)} required /> - <label class="settings-label">${t('create_group.visibility')}</label> - <select class="settings-select" value=${visibility} - onChange=${e => setVisibility(e.target.value)}> - <option value="private">${t('create_group.private')}</option> - <option value="public">${t('create_group.public')}</option> - </select> - <label class="settings-label">${t('create_group.join_policy')}</label> - <select class="settings-select" value=${joinPolicy} - onChange=${e => setJoinPolicy(e.target.value)}> - <option value="invite">${t('create_group.invite')}</option> - <option value="open">${t('create_group.open')}</option> - </select> - <button type="submit" disabled=${loading}> - ${loading ? t('create_group.creating') : t('create_group.submit')} - </button> - </form> + <div class="create-group-page"> + <div class="create-group-card"> + <div class="create-group-header"> + <h2>${t('create_group.title')}</h2> + <p class="create-group-hint">${t('create_group.hint')}</p> + </div> + ${error && html`<div class="error-msg" style="margin-bottom:16px">${error}</div>`} + <form onSubmit=${onSubmit}> + <div class="form-field"> + <label class="form-label">${t('create_group.name')}</label> + <input type="text" placeholder="e.g. Family Photos, Project Files..." + value=${name} onInput=${e => setName(e.target.value)} required autofocus /> + </div> + + <div class="form-field"> + <label class="form-label">${t('create_group.description')}</label> + <textarea class="form-textarea" rows="3" maxlength="512" + placeholder="${t('create_group.description_hint')}" + value=${description} + onInput=${e => setDescription(e.target.value)} /> + <div class="form-char-count">${description.length}/512</div> + </div> + + <div class="form-field"> + <label class="form-label">${t('create_group.visibility')}</label> + <div class="option-cards"> + <button type="button" class="option-card ${visibility === 'private' ? 'selected' : ''}" + onClick=${() => setVisibility('private')}> + <span class="option-icon">🔒</span> + <span class="option-title">${t('create_group.private')}</span> + <span class="option-desc">Only invited members can see and access files</span> + </button> + <button type="button" class="option-card ${visibility === 'public' ? 'selected' : ''}" + onClick=${() => setVisibility('public')}> + <span class="option-icon">🌐</span> + <span class="option-title">${t('create_group.public')}</span> + <span class="option-desc">Anyone can discover and browse this group</span> + </button> + </div> + </div> + + <div class="form-field"> + <label class="form-label">${t('create_group.join_policy')}</label> + <div class="option-cards"> + <button type="button" class="option-card ${joinPolicy === 'invite' ? 'selected' : ''}" + onClick=${() => setJoinPolicy('invite')}> + <span class="option-icon">✉</span> + <span class="option-title">${t('create_group.invite')}</span> + <span class="option-desc">Members must be invited by an admin</span> + </button> + <button type="button" class="option-card ${joinPolicy === 'open' ? 'selected' : ''}" + onClick=${() => setJoinPolicy('open')}> + <span class="option-icon">🚪</span> + <span class="option-title">${t('create_group.open')}</span> + <span class="option-desc">Anyone can join without approval</span> + </button> + </div> + </div> + + <button class="create-group-submit" type="submit" disabled=${loading}> + ${loading ? t('create_group.creating') : t('create_group.submit')} + </button> + </form> + </div> </div> `; } @@ -617,7 +762,7 @@ async function pipelinedDownload(transport, gekKey, fileId, totalChunks, onChunk return results; } -function GroupPage({ groupId, group, token, username }) { +function GroupPage({ groupId, group, token, username, userId }) { const [status, setStatus] = useState('idle'); const [entries, setEntries] = useState([]); const [cached, setCached] = useState(false); @@ -632,10 +777,18 @@ function GroupPage({ groupId, group, token, username }) { const [tab, setTab] = useState('files'); const [uploading, setUploading] = useState(false); const [menuOpen, setMenuOpen] = useState(null); + const [isNodeAdmin, setIsNodeAdmin] = useState(false); const transportRef = useRef(null); const gekRef = useRef(null); useEffect(() => { + if (menuOpen === null) return; + const close = () => setMenuOpen(null); + document.addEventListener('click', close); + return () => document.removeEventListener('click', close); + }, [menuOpen]); + + useEffect(() => { let cancelled = false; getCachedGroupIndex(groupId).then(hit => { @@ -649,6 +802,8 @@ function GroupPage({ groupId, group, token, username }) { setStatus('discovering'); setError(''); gekRef.current = null; + if (!_bundleKey) _bundleKey = await _loadBundleKey(); + _restoreSessionKeys(); try { const nodesData = await hubFetch(`/v1/groups/${groupId}/nodes`, { token }); if (cancelled) return; @@ -657,15 +812,65 @@ function GroupPage({ groupId, group, token, username }) { return; } + // Session keys for P2P GEK bundle fetch (node delivers wrapped GEK) + const sessionKeys = _sessionKeys ? { + skXB64: _sessionKeys.skXB64, + pkXB64: _sessionKeys.pkXB64, + } : null; + setStatus('connecting'); const nodeId = nodesData.nodes[0].node_id; const transport = new window.MeshBayTransport('', token); transportRef.current = transport; - await transport.connect(nodeId, token, groupId); + const ack = await transport.connect( + nodeId, token, groupId, null, sessionKeys, _bundleKey, username); if (cancelled) return; + setIsNodeAdmin(!!ack.is_node_admin); + + // If transport recovered different session keys from node during handshake + if (transport.sessionKeys) { + const recovered = transport.sessionKeys; + if (!_sessionKeys || recovered.skXB64 !== _sessionKeys.skXB64) { + _sessionKeys = recovered; + if (!_sessionKeys.pkXB64) { + const pubkeys = await hubFetch( + `/v1/users/${username}/pubkeys`, { token }); + _sessionKeys.pkXB64 = pubkeys.pk_x25519; + } + _pendingBundlePush = null; + try { localStorage.removeItem(`meshbay_kp_${username}`); } catch {} + _saveSessionKeys(); + } + } + + // Push keypair bundle to node (new registration, localStorage → node) + if (_pendingBundlePush && transport.connected) { + try { + await transport.storeKeypairBundle(_pendingBundlePush); + try { localStorage.removeItem(`meshbay_kp_${username}`); } catch {} + _pendingBundlePush = null; + } catch (e) { + console.warn('[MeshBay] Bundle push to node deferred:', e.message); + } + } + + // Import GEK from transport (fetched from node during handshake) + if (transport.gekRaw && window.MeshBayCrypto) { + gekRef.current = await window.MeshBayCrypto.importGEK( + window.MeshBayCrypto.b64encode(transport.gekRaw)); + } + setStatus('fetching'); + transport.onIndexSync = (msg) => { + if (cancelled) return; + const synced = msg.entries || []; + setEntries(synced); + setCached(false); + cacheGroupIndex(groupId, group ? group.name : groupId, synced); + }; + const indexMsg = await transport.fetchIndex(); if (cancelled) return; const freshEntries = indexMsg.entries || []; @@ -692,6 +897,7 @@ function GroupPage({ groupId, group, token, username }) { return () => { cancelled = true; if (transportRef.current) { + transportRef.current.onIndexSync = null; transportRef.current.close(); transportRef.current = null; } @@ -705,11 +911,6 @@ function GroupPage({ groupId, group, token, username }) { setDlState({ fileId: entry.id, name: entry.name, progress: 0, total: entry.size }); try { - if (!gekRef.current && window.MeshBayCrypto) { - const gekB64 = await transport.fetchGEK(); - gekRef.current = await window.MeshBayCrypto.importGEK(gekB64); - } - const totalChunks = Math.ceil(entry.size / CHUNK_SIZE); let downloaded = 0; const onProgress = (bytes) => { @@ -764,8 +965,9 @@ function GroupPage({ groupId, group, token, username }) { const buf = new Uint8Array(await slice.arrayBuffer()); await transport.uploadChunk(file.name, i, totalChunks, buf); } + await new Promise(r => setTimeout(r, 2500)); const indexMsg = await transport.fetchIndex(); - setEntries(indexMsg.entries || []); + if (indexMsg.entries) setEntries(indexMsg.entries); } catch (err) { setError(err.message); } finally { @@ -777,7 +979,10 @@ function GroupPage({ groupId, group, token, username }) { const transport = transportRef.current; if (!transport || !transport.connected) return; try { - await transport.deleteFile(entry.id); + const signFn = (_sessionKeys && window.MeshBayKeys) + ? (challenge) => window.MeshBayKeys.signChallenge(_sessionKeys.skEdB64, challenge) + : null; + await transport.deleteFile(entry.id, signFn); const indexMsg = await transport.fetchIndex(); setEntries(indexMsg.entries || []); } catch (err) { @@ -846,12 +1051,30 @@ function GroupPage({ groupId, group, token, username }) { return html` <div> <div class="group-header"> - <h2>${group ? group.name : t('group.default_name')}</h2> + <div> + <h2 style="margin-bottom:${group && group.description ? '4px' : '0'}"> + ${group ? group.name : t('group.default_name')} + </h2> + ${group && group.description && html` + <p class="group-desc">${group.description}</p> + `} + </div> <span class="status-badge ${statusClass}"> ${(status === 'discovering' || status === 'connecting' || status === 'fetching') && html`<span class="spinner"></span>${' '}`} ${statusLabel} </span> + ${group && group.is_admin && html` + <button class="admin-btn danger" style="margin-left:auto" + onClick=${async () => { + if (!confirm(t('group.delete_group_confirm', { name: group.name }))) return; + try { + await hubFetch('/v1/groups/' + groupId, { method: 'DELETE', token }); + navigate('/'); + window.location.reload(); + } catch (err) { setError(err.message); } + }}>${t('group.delete_group')}</button> + `} </div> ${error && html`<div class="error-msg" style="margin-bottom:12px">${error}</div>`} ${dlState && html` @@ -954,21 +1177,21 @@ function GroupPage({ groupId, group, token, username }) { setMenuOpen(null); if (e.type === 'video') setVideoEntry(e); else setPreviewEntry(e); - }}>${t('group.view')}</button> + }}><span class="fmi">${'\u{1F441}'}</span> ${t('group.view')}</button> `} <button onClick=${() => { setMenuOpen(null); downloadFile(e); }}> - ${t('group.download')} + <span class="fmi">${'\u{2B07}'}</span> ${t('group.download')} </button> ${e.type === 'video' && html` <button onClick=${() => { setMenuOpen(null); setVideoEntry(e); }}> - ${t('group.play')} + <span class="fmi">${'\u{25B6}'}</span> ${t('group.play')} </button> `} - ${group && group.is_admin && status === 'connected' && html` + ${(isNodeAdmin || (userId && e.uploader_id === userId)) && status === 'connected' && html` <button class="danger" onClick=${() => { setMenuOpen(null); if (confirm(t('group.delete_confirm', { name: e.name }))) deleteFile(e); - }}>${t('group.delete')}</button> + }}><span class="fmi">${'\u{1F5D1}'}</span> ${t('group.delete')}</button> `} </div> `} @@ -987,14 +1210,19 @@ function GroupPage({ groupId, group, token, username }) { ${tab === 'chat' && status === 'connected' && html` <${ChatPanel} transportRef=${transportRef} username=${username} - entries=${entries} gekRef=${gekRef} onRefreshIndex=${refreshIndex} /> + entries=${entries} gekRef=${gekRef} onRefreshIndex=${refreshIndex} + onPreview=${(entry) => { + if (entry.type === 'video') setVideoEntry(entry); + else setPreviewEntry(entry); + }} /> `} ${tab === 'chat' && status !== 'connected' && html` <p class="page-message"><span class="spinner"></span>${' '}${t('status.connecting')}</p> `} ${tab === 'members' && html` - <${MembersPanel} groupId=${groupId} group=${group} token=${token} /> + <${MembersPanel} groupId=${groupId} group=${group} token=${token} + transportRef=${transportRef} gekRef=${gekRef} /> `} `} ${status === 'offline' && html` @@ -1046,10 +1274,6 @@ function FilePreview({ entry, transportRef, gekRef, onClose }) { return; } try { - if (!gekRef.current && window.MeshBayCrypto) { - const gekB64 = await transport.fetchGEK(); - gekRef.current = await window.MeshBayCrypto.importGEK(gekB64); - } const totalChunks = Math.ceil(entry.size / CHUNK_SIZE); let downloaded = 0; const chunks = await pipelinedDownload( @@ -1139,7 +1363,7 @@ function _b64ToU8(b64) { // ── Members Panel ──────────────────────────────────────────────────────── -function MembersPanel({ groupId, group, token }) { +function MembersPanel({ groupId, group, token, transportRef, gekRef }) { const [members, setMembers] = useState([]); const [adminId, setAdminId] = useState(''); const [loading, setLoading] = useState(true); @@ -1168,30 +1392,28 @@ function MembersPanel({ groupId, group, token }) { setInviting(true); setError(''); try { - const pubkeys = await hubFetch(`/v1/users/${inviteUser.trim()}/pubkeys`); + const transport = transportRef && transportRef.current; + const username = inviteUser.trim(); + + // Fetch invitee's public keys (hub = public key directory) + const pubkeys = await hubFetch(`/v1/users/${username}/pubkeys`, { token }); const pkXBytes = Uint8Array.from(atob(pubkeys.pk_x25519), c => c.charCodeAt(0)); - const gekB64 = await (async () => { - const transport = window._activeTransport; - if (transport && transport.connected) { - return await transport.fetchGEK(); - } - const bundleResp = await hubFetch(`/v1/groups/${groupId}/gek`, { token }); - if (!_sessionKeys) throw new Error('No session keys — log in via browser registration'); - const skXB64 = _sessionKeys.skXB64; - const skXRaw = Uint8Array.from(atob(skXB64), c => c.charCodeAt(0)); - const meResp = await hubFetch('/v1/users/me', { token }); - const myPubkeys = await hubFetch(`/v1/users/${meResp.username}/pubkeys`); - const myPkX = Uint8Array.from(atob(myPubkeys.pk_x25519), c => c.charCodeAt(0)); - const rawGek = await window.MeshBayCrypto.unwrapGEK(bundleResp, skXRaw, myPkX); - return btoa(String.fromCharCode(...rawGek)); - })(); + // Get raw GEK from the active transport connection + if (!transport || !transport.connected || !transport.gekRaw) { + throw new Error('Not connected to node or no GEK available'); + } + const gekBytes = transport.gekRaw; - const gekBytes = Uint8Array.from(atob(gekB64), c => c.charCodeAt(0)); + // Wrap GEK for invitee and store on node via P2P const bundle = await window.MeshBayCrypto.wrapGEK(gekBytes, pkXBytes); - await hubFetch(`/v1/groups/${groupId}/members/${inviteUser.trim()}/gek`, { - method: 'POST', token, body: bundle, + await transport.storeGekBundle(pubkeys.user_id, groupId, bundle); + + // Add member on hub (membership management only) + await hubFetch(`/v1/groups/${groupId}/members/${username}`, { + method: 'POST', token, body: {}, }); + setInviteUser(''); loadMembers(); } catch (err) { @@ -1263,19 +1485,17 @@ function _parsePayload(raw) { function ChatImage({ filename, entries, transportRef, gekRef }) { const [blobUrl, setBlobUrl] = useState(null); const [loading, setLoading] = useState(true); + const loadedRef = useRef(false); useEffect(() => { + if (loadedRef.current) return; let cancelled = false; const load = async () => { const transport = transportRef.current; - if (!transport || !transport.connected) { setLoading(false); return; } + if (!transport || !transport.connected) { setLoading(true); return; } const entry = entries.find(e => e.name === filename); - if (!entry) { setLoading(false); return; } + if (!entry) { setLoading(true); return; } try { - if (!gekRef.current && window.MeshBayCrypto) { - const gekB64 = await transport.fetchGEK(); - gekRef.current = await window.MeshBayCrypto.importGEK(gekB64); - } const totalChunks = Math.ceil(entry.size / CHUNK_SIZE); const chunks = await pipelinedDownload(transport, gekRef.current, entry.id, totalChunks); if (cancelled) return; @@ -1283,13 +1503,14 @@ function ChatImage({ filename, entries, transportRef, gekRef }) { const mime = ext === 'png' ? 'image/png' : ext === 'gif' ? 'image/gif' : ext === 'webp' ? 'image/webp' : ext === 'svg' ? 'image/svg+xml' : 'image/jpeg'; const blob = new Blob(chunks, { type: mime }); + loadedRef.current = true; setBlobUrl(URL.createObjectURL(blob)); } catch { /* ignore */ } if (!cancelled) setLoading(false); }; load(); return () => { cancelled = true; }; - }, [filename]); + }, [filename, entries.length]); useEffect(() => { return () => { if (blobUrl) URL.revokeObjectURL(blobUrl); }; @@ -1300,7 +1521,7 @@ function ChatImage({ filename, entries, transportRef, gekRef }) { return html`<img class="chat-att-thumb" src=${blobUrl} alt=${filename} />`; } -function ChatPanel({ transportRef, username, entries, gekRef, onRefreshIndex }) { +function ChatPanel({ transportRef, username, entries, gekRef, onRefreshIndex, onPreview }) { const [messages, setMessages] = useState([]); const [input, setInput] = useState(''); const [sending, setSending] = useState(false); @@ -1377,6 +1598,7 @@ function ChatPanel({ transportRef, username, entries, gekRef, onRefreshIndex }) const buf = new Uint8Array(await slice.arrayBuffer()); await transport.uploadChunk(file.name, i, totalChunks, buf); } + await new Promise(r => setTimeout(r, 2500)); if (onRefreshIndex) await onRefreshIndex(); const ext = file.name.split('.').pop().toLowerCase(); const ftype = ['jpg','jpeg','png','gif','webp','svg'].includes(ext) ? 'image' @@ -1423,7 +1645,11 @@ function ChatPanel({ transportRef, username, entries, gekRef, onRefreshIndex }) `} <div class="chat-bubble ${isOwn ? 'chat-bubble-own' : ''}"> ${att ? html` - <div class="chat-attachment"> + <div class="chat-attachment" style="cursor:pointer" onClick=${() => { + if (!onPreview) return; + const entry = entries.find(e => e.name === att.filename); + if (entry) onPreview(entry); + }}> ${att.type === 'image' ? html`<${ChatImage} filename=${att.filename} entries=${entries} transportRef=${transportRef} gekRef=${gekRef} />` @@ -1522,11 +1748,6 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { }; const startStream = async () => { - if (!gekRef.current && window.MeshBayCrypto) { - const gekB64 = await transport.fetchGEK(); - gekRef.current = await window.MeshBayCrypto.importGEK(gekB64); - } - transport.onStreamInit = (msg) => { if (cancelled) return; const mime = `video/mp4; codecs="${msg.codec}"`; @@ -1744,6 +1965,18 @@ function SettingsPage({ user, theme, onThemeChange, groups }) { try { return JSON.parse(localStorage.getItem('mb_muted') || '{}'); } catch { return {}; } }); + const [nodeKey, setNodeKey] = useState(''); + const [currentNodeKey, setCurrentNodeKey] = useState(null); + const [nodeKeyStatus, setNodeKeyStatus] = useState(''); + const [nodeKeyLoading, setNodeKeyLoading] = useState(false); + + useEffect(() => { + hubFetch(`/v1/users/${user.username}/pubkeys`, { token: user.token }) + .then(data => { + if (data.pk_node_ed25519) setCurrentNodeKey(data.pk_node_ed25519); + }) + .catch(() => {}); + }, [user.username, user.token]); const onLocaleChange = useCallback((e) => { const code = e.target.value; @@ -1764,6 +1997,26 @@ function SettingsPage({ user, theme, onThemeChange, groups }) { }); }, []); + const submitNodeKey = useCallback(async () => { + const key = nodeKey.trim(); + if (!key) return; + setNodeKeyLoading(true); + setNodeKeyStatus(''); + try { + await hubFetch('/v1/users/me/node_key', { + method: 'PUT', token: user.token, + body: { pk_node_ed25519: key }, + }); + setCurrentNodeKey(key); + setNodeKey(''); + setNodeKeyStatus(t('settings.node_key_success')); + } catch (e) { + setNodeKeyStatus(e.message); + } finally { + setNodeKeyLoading(false); + } + }, [nodeKey, user.token]); + return html` <div> <h2>${t('settings.title')}</h2> @@ -1781,6 +2034,32 @@ function SettingsPage({ user, theme, onThemeChange, groups }) { </div> <div class="settings-section"> + <h3 class="settings-heading">${t('settings.node_key')}</h3> + <p style="font-size:0.85em;color:var(--text-dim);margin-bottom:8px">${t('settings.node_key_desc')}</p> + ${currentNodeKey && html` + <div class="settings-row" style="margin-top:8px"> + <span class="settings-label">${t('settings.node_key_current')}</span> + <code class="settings-value" style="font-size:0.8em;word-break:break-all">${currentNodeKey}</code> + </div> + `} + <div style="display:flex;gap:8px;margin-top:10px;align-items:center"> + <input type="text" class="admin-search" style="flex:1;font-family:monospace;font-size:0.85em" + placeholder=${t('settings.node_key_placeholder')} + value=${nodeKey} onInput=${e => setNodeKey(e.target.value)} + onKeyDown=${e => e.key === 'Enter' && submitNodeKey()} /> + <button class="admin-btn" onClick=${submitNodeKey} + disabled=${nodeKeyLoading || !nodeKey.trim()}> + ${t('settings.node_key_submit')} + </button> + </div> + ${nodeKeyStatus && html` + <p style="margin-top:6px;font-size:0.85em;color:${nodeKeyStatus === t('settings.node_key_success') ? 'var(--green, #22c55e)' : 'var(--red, #ef4444)'}"> + ${nodeKeyStatus} + </p> + `} + </div> + + <div class="settings-section"> <h3 class="settings-heading">${t('settings.appearance')}</h3> <div class="settings-row"> <span class="settings-label">${t('settings.theme')}</span> @@ -1987,7 +2266,7 @@ function AdminPage({ token }) { <option value="admin">admin</option> </select> </td> - <td><span class="badge">${u.status}</span></td> + <td><span class="badge ${u.status === 'active' ? 'badge-ok' : u.status === 'suspended' ? 'badge-err' : ''}">${u.status}</span></td> <td>${new Date(u.created_at).toLocaleDateString()}</td> <td class="admin-actions"> <button class="admin-btn" onClick=${() => showUserDetail(u.id)}>${t('admin.btn_details')}</button> @@ -2024,7 +2303,7 @@ function AdminPage({ token }) { <td>${g.name}</td> <td><span class="badge">${g.visibility}</span></td> <td>${g.member_count}</td> - <td><span class="badge">${g.status}</span></td> + <td><span class="badge ${g.status === 'active' ? 'badge-ok' : g.status === 'suspended' ? 'badge-err' : ''}">${g.status}</span></td> <td>${new Date(g.created_at).toLocaleDateString()}</td> <td class="admin-actions"> ${g.status === 'active' @@ -2058,15 +2337,17 @@ function AdminPage({ token }) { <table class="admin-table"> <thead><tr> <th>${t('admin.col_time')}</th> + <th>${t('admin.col_user')}</th> <th>${t('admin.col_event')}</th> <th>${t('admin.col_ip')}</th> <th>${t('admin.col_detail')}</th> </tr></thead> <tbody> - ${logs.length === 0 && html`<tr><td colspan="4" class="admin-empty">${t('admin.no_logs')}</td></tr>`} + ${logs.length === 0 && html`<tr><td colspan="5" class="admin-empty">${t('admin.no_logs')}</td></tr>`} ${logs.map(lg => html` <tr key=${lg.id}> <td style="white-space:nowrap">${new Date(lg.timestamp).toLocaleString()}</td> + <td>${lg.username || ''}</td> <td><span class="badge">${lg.event}</span></td> <td>${lg.ip_address}</td> <td>${lg.detail || ''}</td> @@ -2208,8 +2489,8 @@ function App() { useEffect(() => { setMenuOpen(false); }, [route]); - const toggleTheme = useCallback(() => { - setTheme(prev => resolveTheme(prev) === 'dark' ? 'light' : 'dark'); + const changeTheme = useCallback((val) => { + setTheme(val); }, []); const authCtx = { @@ -2218,9 +2499,14 @@ function App() { let token, refreshToken; if (window.MeshBayKeys) { const data = await window.MeshBayKeys.loginAndRecover(username, password); - _sessionKeys = { skXB64: data.skXB64, skEdB64: data.skEdB64 }; token = data.accessToken; refreshToken = data.refreshToken; + _bundleKey = data.bundleKey; + await _storeBundleKey(_bundleKey); + if (data.skXB64) { + _sessionKeys = { skXB64: data.skXB64, skEdB64: data.skEdB64 }; + _pendingBundlePush = data.keypairBundleEnc; + } } else { const data = await hubFetch('/v1/users/login', { method: 'POST', @@ -2230,7 +2516,12 @@ function App() { refreshToken = data.refresh_token; } const me = await hubFetch('/v1/users/me', { token }); - const u = { username, token, refreshToken, role: me.role }; + if (_sessionKeys) { + const pubkeys = await hubFetch(`/v1/users/${username}/pubkeys`, { token }); + _sessionKeys.pkXB64 = pubkeys.pk_x25519; + _saveSessionKeys(); + } + const u = { username, userId: me.user_id, token, refreshToken, role: me.role }; setUser(u); saveAuth(u); }, @@ -2266,7 +2557,7 @@ function App() { const group = groups.find(g => g.id === groupId); page = html`<${GroupPage} groupId=${groupId} group=${group} token=${user.token} - username=${user.username} />`; + username=${user.username} userId=${user.userId} />`; } else if (route === '/admin') { page = (user.role === 'moderator' || user.role === 'admin') ? html`<${AdminPage} token=${user.token} />` @@ -2282,8 +2573,8 @@ function App() { <${AuthContext.Provider} value=${authCtx}> <${Nav} user=${user} - theme=${resolved} - onThemeToggle=${toggleTheme} + theme=${theme} + onThemeChange=${changeTheme} onLogout=${authCtx.logout} onMenuToggle=${() => setMenuOpen(o => !o)} unreadCount=${unreadCount} /> diff --git a/packages/meshbay-hub/src/meshbay_hub/static/crypto.js b/packages/meshbay-hub/src/meshbay_hub/static/crypto.js index eb96eef..5ebf624 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/crypto.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/crypto.js @@ -230,8 +230,24 @@ function b64encode(bytes) { return btoa(String.fromCharCode(...bytes)); } +// ── GEK proof (HMAC-SHA256 for handshake challenge) ───────────────────────── + +async function hmacGEK(gekRaw, nonceB64, offerFp, answerFp) { + const nonce = b64decode(nonceB64); + const data = concatBuffers([ + nonce, + offerFp || new Uint8Array(0), + answerFp || new Uint8Array(0), + ]); + const key = await crypto.subtle.importKey( + 'raw', gekRaw, { name: 'HMAC', hash: 'SHA-256' }, false, ['sign']); + const sig = await crypto.subtle.sign('HMAC', key, data); + return b64encode(new Uint8Array(sig)); +} + // Export for use in app.js window.MeshBayCrypto = { importGEK, deriveChunkKey, decryptChunk, decryptChunkBin, decryptFile, generateGEK, wrapGEK, unwrapGEK, encryptChunk, b64encode, b64decode, + hmacGEK, }; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/i18n.js b/packages/meshbay-hub/src/meshbay_hub/static/i18n.js index 2a91407..3450735 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/i18n.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/i18n.js @@ -89,6 +89,8 @@ const en = { 'group.view': 'View', 'group.delete': 'Delete', 'group.delete_confirm': 'Delete {name}?', + 'group.delete_group': 'Delete group', + 'group.delete_group_confirm': 'Permanently delete "{name}" and remove all members? Files on the node are preserved.', 'group.err_transport': 'Transport module not loaded', // Status @@ -130,12 +132,25 @@ const en = { 'settings.role': 'Role', 'settings.groups': 'Group notifications', 'settings.notifications': 'Notifications', + 'settings.node_key': 'Link Node', + 'settings.node_key_desc': 'If you operate a MeshBay node, paste its Ed25519 public key here to link it to your account. You can find this key in your node\'s admin dashboard.', + 'settings.node_key_placeholder': 'Paste node Ed25519 public key (base64)', + 'settings.node_key_submit': 'Link Node', + 'settings.node_key_success': 'Node linked successfully', + 'settings.node_key_current': 'Linked node key', // Sidebar 'sidebar.create_group': 'Create group', 'sidebar.settings': 'Settings', 'sidebar.admin': 'Admin', + // User menu + 'usermenu.language': 'Language', + 'usermenu.settings': 'Settings', + 'usermenu.theme_light': 'Light mode', + 'usermenu.theme_dark': 'Dark mode', + 'usermenu.logout': 'Logout', + // Admin panel 'admin.title': 'Administration', 'admin.tab_stats': 'Stats', @@ -205,6 +220,8 @@ const en = { 'create_group.join_policy': 'Join policy', 'create_group.invite': 'Invite only', 'create_group.open': 'Open (anyone can join)', + 'create_group.description': 'Description', + 'create_group.description_hint': 'What is this group about? (optional)', 'create_group.submit': 'Create', 'create_group.hint': 'A group needs a node to host files. You can create the group now and connect a node later.', 'create_group.creating': 'Creating...', @@ -245,7 +262,7 @@ const _strings = { en }; let _locale = 'en'; export const LOCALES = [ - { code: 'en', name: 'English' }, + { code: 'en', name: 'English', flag: '\u{1F1EC}\u{1F1E7}' }, ]; export function getLocale() { return _locale; } diff --git a/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js b/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js index 63baff5..ff3da33 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/keyderive.js @@ -23,6 +23,26 @@ const PBKDF2_ITERATIONS = 600000; // OWASP 2023 recommendation for PBKDF2-SHA512 const HUB = ''; // same origin +// ── Auth key derivation (password split) ────────────────────────────────────── + +/** + * Derive an auth key from password + username using PBKDF2-SHA512. + * This key is sent to the hub for authentication — the raw password never leaves the browser. + * Uses a different salt domain than deriveEncryptionKey (bundle key), so the two + * derived values are cryptographically independent. + */ +async function deriveAuthKey(password, username) { + const enc = new TextEncoder(); + const km = await crypto.subtle.importKey( + 'raw', enc.encode(password), 'PBKDF2', false, ['deriveBits']); + const salt = await crypto.subtle.digest( + 'SHA-256', enc.encode(`meshbay:auth:v1:${username}`)); + const bits = await crypto.subtle.deriveBits( + { name: 'PBKDF2', hash: 'SHA-512', salt, iterations: PBKDF2_ITERATIONS }, + km, 256); + return btoa(String.fromCharCode(...new Uint8Array(bits))); +} + // ── Key generation ──────────────────────────────────────────────────────────── /** @@ -105,20 +125,21 @@ async function decryptBundle(bundleB64, password, username) { * Full registration flow: * 1. Generate random keypairs * 2. Encrypt bundle with password - * 3. POST to hub (public keys + encrypted bundle) + * 3. POST to hub (public keys only — no keypair bundle) + * 4. Store encrypted bundle locally for backup to node on first connect * * Returns the raw private keys for immediate use after registration. */ async function registerUser(username, email, password) { const { skEdRaw, pkEdRaw, skXRaw, pkXRaw } = await generateKeypairs(); - // Convert SPKI public keys to raw 32-byte format expected by hub const pkEdCrypto = await crypto.subtle.importKey('spki', pkEdRaw, 'Ed25519', true, ['verify']); const pkXCrypto = await crypto.subtle.importKey('spki', pkXRaw, 'X25519', true, []); const pkEdBytes = new Uint8Array(await crypto.subtle.exportKey('raw', pkEdCrypto)); const pkXBytes = new Uint8Array(await crypto.subtle.exportKey('raw', pkXCrypto)); const encBundle = await encryptBundle(skEdRaw, skXRaw, password, username); + const authKey = await deriveAuthKey(password, username); const resp = await fetch(`${HUB}/v1/users/register`, { method: 'POST', @@ -126,39 +147,115 @@ async function registerUser(username, email, password) { body: JSON.stringify({ username, email, - password, + auth_key: authKey, pk_user_ed25519: btoa(String.fromCharCode(...pkEdBytes)), pk_user_x25519: btoa(String.fromCharCode(...pkXBytes)), - keypair_bundle: encBundle, // encrypted, hub stores but cannot read }), }); if (!resp.ok) throw new Error(`Registration failed: ${await resp.text()}`); - return { skEdRaw, skXRaw, pkEdBytes, pkXBytes }; + + // Store encrypted bundle locally — will be backed up to node on first group connect + try { localStorage.setItem(`meshbay_kp_${username}`, encBundle); } catch {} + + return { skEdRaw, skXRaw, pkEdBytes, pkXBytes, keypairBundleEnc: encBundle }; +} + +/** + * Decrypt a keypair bundle using a pre-derived AES-256 CryptoKey. + * Used when the bundle is fetched from the node (bundleKey was derived at login). + */ +async function decryptBundleWithKey(bundleB64, aesKey) { + const raw = Uint8Array.from(atob(bundleB64), c => c.charCodeAt(0)); + const nonce = raw.slice(0, 12); + const ct = raw.slice(12); + const plain = await crypto.subtle.decrypt({ name: 'AES-GCM', iv: nonce }, aesKey, ct); + return JSON.parse(new TextDecoder().decode(plain)); } /** - * Login and recover private keys from the encrypted bundle. + * Login and recover private keys. + * + * If localStorage has a keypair bundle (new registration, not yet pushed to node), + * decrypts it and returns the keys + encrypted bundle for push to node. + * Otherwise returns bundleKey so the caller can fetch from node during handshake. */ async function loginAndRecover(username, password) { + const authKey = await deriveAuthKey(password, username); + const resp = await fetch(`${HUB}/v1/users/login`, { method: 'POST', headers: { 'Content-Type': 'application/json' }, - body: JSON.stringify({ username, password }), + body: JSON.stringify({ username, auth_key: authKey }), }); + if (!resp.ok) throw new Error(`Login failed: ${await resp.text()}`); const data = await resp.json(); - const bundle = data.keypair_bundle; - if (!bundle) throw new Error('No keypair bundle in response — account may have been created via CLI'); - - const keys = await decryptBundle(bundle, password, username); - return { + const result = { accessToken: data.access_token, refreshToken: data.refresh_token, - skEdB64: keys.skEd, - skXB64: keys.skX, + bundleKey: await deriveEncryptionKey(password, username), }; + + // localStorage bundle = new registration, not yet pushed to node + const bundleEnc = (typeof localStorage !== 'undefined' + && localStorage.getItem(`meshbay_kp_${username}`)) || null; + + if (bundleEnc) { + const keys = await decryptBundle(bundleEnc, password, username); + result.skEdB64 = keys.skEd; + result.skXB64 = keys.skX; + result.keypairBundleEnc = bundleEnc; + } + + return result; +} + +async function regenerateKeys(token, username, password) { + const { skEdRaw, pkEdRaw, skXRaw, pkXRaw } = await generateKeypairs(); + + const pkEdCrypto = await crypto.subtle.importKey('spki', pkEdRaw, 'Ed25519', true, ['verify']); + const pkXCrypto = await crypto.subtle.importKey('spki', pkXRaw, 'X25519', true, []); + const pkEdBytes = new Uint8Array(await crypto.subtle.exportKey('raw', pkEdCrypto)); + const pkXBytes = new Uint8Array(await crypto.subtle.exportKey('raw', pkXCrypto)); + + const resp = await fetch(`${HUB}/v1/users/me/keys`, { + method: 'PUT', + headers: { + 'Content-Type': 'application/json', + 'Authorization': `Bearer ${token}`, + }, + body: JSON.stringify({ + pk_user_ed25519: btoa(String.fromCharCode(...pkEdBytes)), + pk_user_x25519: btoa(String.fromCharCode(...pkXBytes)), + }), + }); + + if (!resp.ok) throw new Error(`Key rotation failed: ${await resp.text()}`); + + const encBundle = await encryptBundle(skEdRaw, skXRaw, password, username); + try { localStorage.setItem(`meshbay_kp_${username}`, encBundle); } catch {} + + return { + skEdB64: btoa(String.fromCharCode(...new Uint8Array(skEdRaw))), + skXB64: btoa(String.fromCharCode(...new Uint8Array(skXRaw))), + pkEdB64: btoa(String.fromCharCode(...pkEdBytes)), + pkXB64: btoa(String.fromCharCode(...pkXBytes)), + keypairBundleEnc: encBundle, + }; +} + +async function signChallenge(skEdPkcs8B64, challengeB64) { + const skRaw = Uint8Array.from(atob(skEdPkcs8B64), c => c.charCodeAt(0)); + const sk = await crypto.subtle.importKey( + 'pkcs8', skRaw, { name: 'Ed25519' }, false, ['sign']); + const challenge = Uint8Array.from(atob(challengeB64), c => c.charCodeAt(0)); + const sig = await crypto.subtle.sign('Ed25519', sk, challenge); + return btoa(String.fromCharCode(...new Uint8Array(sig))); } -window.MeshBayKeys = { registerUser, loginAndRecover, generateKeypairs }; +window.MeshBayKeys = { + registerUser, loginAndRecover, regenerateKeys, generateKeypairs, signChallenge, + deriveAuthKey, decryptBundleWithKey, +}; diff --git a/packages/meshbay-hub/src/meshbay_hub/static/style.css b/packages/meshbay-hub/src/meshbay_hub/static/style.css index f3dcbf4..e430201 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/style.css +++ b/packages/meshbay-hub/src/meshbay_hub/static/style.css @@ -95,20 +95,81 @@ a:hover { text-decoration: underline; } .nav-user { color: var(--text-dim); font-size: 0.9em; } -.nav-theme { - background: none; - border: 1px solid rgba(255, 255, 255, 0.15); - color: var(--nav-text); - width: 32px; - height: 32px; - border-radius: 6px; - cursor: pointer; - font-size: 1.1em; +/* ── User Menu ──────────────────────────────────────────────────────────── */ + +.user-menu-wrap { position: relative; } + +.user-menu-trigger { display: flex; align-items: center; - justify-content: center; + gap: 8px; + background: rgba(255, 255, 255, 0.08); + border: 1px solid rgba(255, 255, 255, 0.12); + color: var(--nav-text); + padding: 4px 10px 4px 4px; + border-radius: 20px; + cursor: pointer; + font-size: 0.85em; + transition: background 0.12s; +} +.user-menu-trigger:hover { background: rgba(255, 255, 255, 0.15); } + +.user-menu-avatar { + width: 28px; height: 28px; + border-radius: 50%; + background: var(--accent); + color: #fff; + display: flex; align-items: center; justify-content: center; + font-weight: 700; font-size: 0.8em; + flex-shrink: 0; +} +.user-menu-avatar.lg { width: 36px; height: 36px; font-size: 0.95em; } + +.user-menu-name { max-width: 120px; overflow: hidden; text-overflow: ellipsis; white-space: nowrap; } +.user-menu-caret { font-size: 0.7em; color: rgba(255,255,255,0.5); } + +.user-menu-dropdown { + position: absolute; + right: 0; top: calc(100% + 6px); + background: var(--bg-surface); + border: 1px solid var(--border); + border-radius: 10px; + box-shadow: var(--shadow-lg); + min-width: 220px; + z-index: 120; + overflow: hidden; + padding: 4px 0; } -.nav-theme:hover { border-color: rgba(255, 255, 255, 0.3); background: none; } + +.user-menu-header { + display: flex; align-items: center; gap: 10px; + padding: 12px 14px 8px; +} +.user-menu-uname { font-weight: 600; font-size: 0.9em; color: var(--text); } +.user-menu-role { font-size: 0.75em; color: var(--text-dim); text-transform: capitalize; } + +.user-menu-divider { height: 1px; background: var(--border); margin: 4px 0; } + +.user-menu-section-label { + padding: 6px 14px 2px; + font-size: 0.7em; text-transform: uppercase; letter-spacing: 0.06em; + color: var(--text-dim); font-weight: 600; +} + +.user-menu-item { + display: flex; align-items: center; gap: 8px; + width: 100%; padding: 8px 14px; + background: none; border: none; border-radius: 0; + color: var(--text); font-size: 0.85em; + cursor: pointer; text-align: left; +} +.user-menu-item:hover { background: var(--bg-raised); } +.user-menu-logout { color: var(--error); } + +.umi-icon { width: 18px; text-align: center; font-size: 0.95em; flex-shrink: 0; } +.umi-arrow { margin-left: auto; font-size: 0.7em; color: var(--text-dim); } +.umi-check { margin-left: auto; color: var(--accent); font-size: 0.85em; } +.user-menu-sub { padding-left: 32px; } .nav-btn { background: rgba(255, 255, 255, 0.1); @@ -289,14 +350,17 @@ button:disabled { opacity: 0.5; cursor: not-allowed; } border-radius: 12px; font-size: 0.75em; } +.badge-ok { background: #16a34a20; color: var(--success); } +.badge-err { background: var(--error-bg); color: var(--error); } /* ── Group page header ───────────────────────────────────────────────────── */ .group-header { display: flex; - align-items: center; + align-items: flex-start; gap: 12px; margin-bottom: 16px; + flex-wrap: wrap; } .group-header h2 { margin-bottom: 0; } @@ -688,8 +752,10 @@ button:disabled { opacity: 0.5; cursor: not-allowed; } } .chat-bubble-own .chat-att-size { color: rgba(255, 255, 255, 0.6); } +.file-menu button { display: flex; align-items: center; gap: 8px; } .file-menu button.danger { color: var(--error); } .file-menu button.danger:hover { background: var(--error-bg); } +.fmi { font-size: 0.9em; width: 16px; text-align: center; flex-shrink: 0; } /* ── Download button + progress ──────────────────────────────────────────── */ @@ -1228,3 +1294,105 @@ button:disabled { opacity: 0.5; cursor: not-allowed; } color: var(--text-secondary); white-space: nowrap; } + +/* ── Create Group page ───────────────────────────────────────────────────── */ + +.create-group-page { + display: flex; + justify-content: center; + padding-top: 8px; +} + +.create-group-card { + width: 100%; + max-width: 520px; + background: var(--bg-surface); + border: 1px solid var(--border); + border-radius: 12px; + padding: 32px; + box-shadow: var(--shadow); +} + +.create-group-header { margin-bottom: 24px; } +.create-group-header h2 { margin-bottom: 8px; font-size: 1.3em; } +.create-group-hint { + color: var(--text-secondary); + font-size: 0.88em; + line-height: 1.5; +} + +.form-field { margin-bottom: 20px; } +.form-label { + display: block; + font-size: 0.82em; + font-weight: 600; + color: var(--text-secondary); + text-transform: uppercase; + letter-spacing: 0.04em; + margin-bottom: 8px; +} + +.option-cards { display: grid; grid-template-columns: 1fr 1fr; gap: 10px; } + +.option-card { + display: flex; flex-direction: column; align-items: center; gap: 4px; + padding: 14px 10px; + background: var(--bg-base); + border: 2px solid var(--border); + border-radius: 10px; + cursor: pointer; + text-align: center; + transition: border-color 0.12s, background 0.12s; +} +.option-card:hover { border-color: var(--accent); background: var(--bg-base); } +.option-card.selected { + border-color: var(--accent); + background: color-mix(in srgb, var(--accent) 8%, var(--bg-base)); +} + +.option-icon { font-size: 1.5em; margin-bottom: 2px; } +.option-title { font-weight: 600; font-size: 0.9em; color: var(--text); } +.option-desc { font-size: 0.75em; color: var(--text-dim); line-height: 1.3; } + +.create-group-submit { + width: 100%; + padding: 12px; + font-size: 1em; + font-weight: 600; + margin-top: 8px; + border-radius: 8px; +} + +.form-textarea { + width: 100%; + padding: 10px 12px; + border: 1px solid var(--border); + border-radius: 6px; + background: var(--bg-base); + color: var(--text); + font-size: 0.95em; + font-family: inherit; + resize: vertical; + min-height: 60px; + line-height: 1.5; +} +.form-textarea:focus { outline: none; border-color: var(--border-focus); } +.form-char-count { text-align: right; font-size: 0.75em; color: var(--text-dim); margin-top: 4px; } + +.group-desc { + font-size: 0.85em; + color: var(--text-secondary); + line-height: 1.4; + margin-bottom: 4px; +} + +.group-card-desc { + font-size: 0.8em; + color: var(--text-dim); + line-height: 1.3; + margin: 4px 0 6px; + display: -webkit-box; + -webkit-line-clamp: 2; + -webkit-box-orient: vertical; + overflow: hidden; +} diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport.js b/packages/meshbay-hub/src/meshbay_hub/static/transport.js index 1dc1ded..ca9c60e 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/transport.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/transport.js @@ -16,6 +16,16 @@ * transport.close(); */ +async function _pkFromSk(skPkcs8B64) { + const raw = Uint8Array.from(atob(skPkcs8B64), c => c.charCodeAt(0)); + const sk = await crypto.subtle.importKey('pkcs8', raw, { name: 'X25519' }, true, ['deriveBits']); + const jwk = await crypto.subtle.exportKey('jwk', sk); + const b64url = jwk.x; + const b64 = b64url.replace(/-/g, '+').replace(/_/g, '/'); + const pad = b64.length % 4; + return pad ? b64 + '='.repeat(4 - pad) : b64; +} + class MeshBayTransport { constructor(hubUrl, accessToken) { this._hubUrl = hubUrl; @@ -30,6 +40,7 @@ class MeshBayTransport { this._onStreamInit = null; this._onStreamData = null; this._onStreamEnd = null; + this._onIndexSync = null; } get connected() { return this._connected; } @@ -38,8 +49,15 @@ class MeshBayTransport { set onStreamInit(fn) { this._onStreamInit = fn; } set onStreamData(fn) { this._onStreamData = fn; } set onStreamEnd(fn) { this._onStreamEnd = fn; } + set onIndexSync(fn) { this._onIndexSync = fn; } + + get sessionKeys() { return this._sessionKeys; } - async connect(nodeId, jwtToken, groupId) { + async connect(nodeId, jwtToken, groupId, gekRaw, sessionKeys, bundleKey, username) { + this._gekRaw = gekRaw || null; + this._sessionKeys = sessionKeys || null; + this._bundleKey = bundleKey || null; + this._username = username || null; this._pc = new RTCPeerConnection({ iceServers: [{ urls: 'stun:stun.l.google.com:19302' }], }); @@ -106,22 +124,94 @@ class MeshBayTransport { } const answer = await resp.json(); + this._rawAnswerSdp = answer.sdp; await this._pc.setRemoteDescription({ type: 'answer', sdp: answer.sdp }); await channelReady; - const ack = await this._sendAndWait({ + const reply = await this._sendAndWait({ type: 'handshake', v: '0.1', token: jwtToken, group_id: groupId || '', }); - if (ack.type !== 'handshake_ack') { - throw new Error('MNP handshake rejected: ' + (ack.detail || JSON.stringify(ack))); + if (reply.type === 'handshake_challenge') { + if (!window.MeshBayCrypto) { + throw new Error('Node requires GEK proof but no crypto available'); + } + + // Recover session keys from node if not available locally (P2P keypair bundle) + if (!this._sessionKeys && this._bundleKey && window.MeshBayKeys) { + const kpResp = await this._sendAndWait({ + type: 'keypair_bundle_fetch', v: '0.1', + }); + if (kpResp.type === 'keypair_bundle_resp' && kpResp.found) { + const keys = await window.MeshBayKeys.decryptBundleWithKey( + kpResp.bundle_enc, this._bundleKey); + const pkXB64 = await _pkFromSk(keys.skX); + this._sessionKeys = { skXB64: keys.skX, skEdB64: keys.skEd, pkXB64 }; + } + } + + // Fetch wrapped GEK bundle from node (P2P only — hub never touches crypto) + if (!gekRaw && this._sessionKeys) { + const bundleResp = await this._sendAndWait({ + type: 'gek_bundle_fetch', v: '0.1', + }); + if (bundleResp.type === 'gek_bundle_resp' && bundleResp.found) { + const skXRaw = Uint8Array.from(atob(this._sessionKeys.skXB64), c => c.charCodeAt(0)); + const myPkX = Uint8Array.from(atob(this._sessionKeys.pkXB64), c => c.charCodeAt(0)); + try { + gekRaw = await window.MeshBayCrypto.unwrapGEK(bundleResp, skXRaw, myPkX); + this._gekRaw = gekRaw; + } catch (e) { + console.warn('[MeshBay] GEK unwrap failed with local keys, trying node keypair bundle'); + if (this._bundleKey && window.MeshBayKeys) { + const kpResp = await this._sendAndWait({ + type: 'keypair_bundle_fetch', v: '0.1', + }); + if (kpResp.type === 'keypair_bundle_resp' && kpResp.found) { + const keys = await window.MeshBayKeys.decryptBundleWithKey( + kpResp.bundle_enc, this._bundleKey); + const pkXB64 = await _pkFromSk(keys.skX); + this._sessionKeys = { skXB64: keys.skX, skEdB64: keys.skEd, pkXB64 }; + const skXRaw2 = Uint8Array.from(atob(keys.skX), c => c.charCodeAt(0)); + const myPkX2 = Uint8Array.from(atob(pkXB64), c => c.charCodeAt(0)); + gekRaw = await window.MeshBayCrypto.unwrapGEK(bundleResp, skXRaw2, myPkX2); + this._gekRaw = gekRaw; + } + } + } + } + } + + if (!gekRaw) { + throw new Error('Node requires GEK proof but no GEK available'); + } + + let proof = ''; + if (gekRaw) { + const offerFp = _extractDtlsFingerprint(this._pc.localDescription.sdp); + const answerFp = _extractDtlsFingerprint(this._rawAnswerSdp); + proof = await window.MeshBayCrypto.hmacGEK(gekRaw, reply.nonce, offerFp, answerFp); + } + const ack = await this._sendAndWait({ + type: 'handshake_response', + v: '0.1', + proof, + }); + if (ack.type !== 'handshake_ack') { + throw new Error('GEK proof rejected: ' + (ack.detail || JSON.stringify(ack))); + } + return ack; } - return ack; + if (reply.type !== 'handshake_ack') { + throw new Error('MNP handshake rejected: ' + (reply.detail || JSON.stringify(reply))); + } + + return reply; } async fetchIndex() { @@ -130,12 +220,6 @@ class MeshBayTransport { return msg; } - async fetchGEK() { - const msg = await this._sendAndWait({ type: 'gek_req', v: '0.1' }); - if (msg.type === 'error') throw new Error(msg.detail); - return msg.gek_b64; - } - async fetchChunk(fileId, chunkIndex) { const msg = await this._sendAndWait({ type: 'file_req', @@ -182,13 +266,25 @@ class MeshBayTransport { return msg; } - async deleteFile(fileId) { + async deleteFile(fileId, signFn) { const msg = await this._sendAndWait({ type: 'file_delete', v: '0.1', file_id: fileId, }); if (msg.type === 'error') throw new Error(msg.detail); + if (msg.type === 'admin_challenge') { + if (!signFn) throw new Error('Admin challenge received but no signing key available'); + const signature = await signFn(msg.challenge); + const ack = await this._sendAndWait({ + type: 'admin_response', + v: '0.1', + file_id: fileId, + signature, + }); + if (ack.type === 'error') throw new Error(ack.detail); + return ack; + } return msg; } @@ -208,6 +304,32 @@ class MeshBayTransport { return msg; } + async storeGekBundle(userId, groupId, bundle) { + const msg = await this._sendAndWait({ + type: 'gek_bundle_store', + v: '0.1', + user_id: userId, + group_id: groupId, + pk_eph_b64: bundle.pk_eph_b64, + nonce_b64: bundle.nonce_b64, + wrapped_b64: bundle.wrapped_b64, + }); + if (msg.type === 'error') throw new Error(msg.detail); + return msg; + } + + async storeKeypairBundle(bundleEnc) { + const msg = await this._sendAndWait({ + type: 'keypair_bundle_store', + v: '0.1', + bundle_enc: bundleEnc, + }); + if (msg.type === 'error') throw new Error(msg.detail); + return msg; + } + + get gekRaw() { return this._gekRaw; } + close() { if (this._channel) this._channel.close(); if (this._pc) this._pc.close(); @@ -226,6 +348,7 @@ class MeshBayTransport { reject(new Error('Response timeout')); }, 30000); this._pending.set(id, { + _reqType: obj.type, resolve: (msg) => { clearTimeout(timeout); this._pending.delete(id); resolve(msg); }, reject: (err) => { clearTimeout(timeout); this._pending.delete(id); reject(err); }, }); @@ -269,16 +392,25 @@ class MeshBayTransport { this._onChat(msg); return; } - if (msg.type === 'stream_init' && this._onStreamInit) { - this._onStreamInit(msg); + if (msg.type === 'stream_init') { + if (this._onStreamInit) this._onStreamInit(msg); + return; + } + if (msg.type === 'stream_data') { + if (this._onStreamData) this._onStreamData(msg); return; } - if (msg.type === 'stream_data' && this._onStreamData) { - this._onStreamData(msg); + if (msg.type === 'stream_end') { + if (this._onStreamEnd) this._onStreamEnd(msg); return; } - if (msg.type === 'stream_end' && this._onStreamEnd) { - this._onStreamEnd(msg); + + if (msg.type === 'index_sync' && msg.entries) { + if (this._onIndexSync) this._onIndexSync(msg); + const oldest = this._pending.entries().next(); + if (!oldest.done && oldest.value[1]._reqType === 'index_sync') { + oldest.value[1].resolve(msg); + } return; } @@ -485,5 +617,15 @@ function _b64decode(b64) { return bytes; } +function _extractDtlsFingerprint(sdp) { + const match = sdp.match(/a=fingerprint:sha-256 ([0-9A-Fa-f:]+)/); + if (!match) return new Uint8Array(0); + const hex = match[1].replace(/:/g, ''); + const bytes = new Uint8Array(hex.length / 2); + for (let i = 0; i < hex.length; i += 2) + bytes[i / 2] = parseInt(hex.substring(i, i + 2), 16); + return bytes; +} + // Export window.MeshBayTransport = MeshBayTransport; diff --git a/packages/meshbay-hub/tests/test_hub_api.py b/packages/meshbay-hub/tests/test_hub_api.py index f7c499e..a8232c1 100644 --- a/packages/meshbay-hub/tests/test_hub_api.py +++ b/packages/meshbay-hub/tests/test_hub_api.py @@ -10,7 +10,7 @@ from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey from cryptography.hazmat.primitives import serialization -from meshbay_common.crypto import generate_gek, pk_to_b64, wrap_gek +from meshbay_common.crypto import pk_to_b64 from meshbay_hub.api.deps import set_admin_usernames @@ -210,12 +210,8 @@ async def test_announce_and_get_node(client): # ── Groups + GEK bundles ────────────────────────────────────────────────────── @pytest.mark.asyncio -async def test_group_gek_roundtrip(client): - """Admin creates group, wraps GEK for member, member retrieves and can unwrap.""" - import jwt as pyjwt - from meshbay_common.crypto import unwrap_gek - - # Register admin (alice2) and member (bob2) +async def test_group_member_add(client): + """Admin creates group and adds member (GEK exchange happens P2P on node).""" pk_ed_a, pk_x_a, _ = _gen_user_keys() pk_ed_b, pk_x_b, sk_x_b = _gen_user_keys() @@ -227,46 +223,29 @@ async def test_group_gek_roundtrip(client): "username": uname, "email": email, "password": pwd, "pk_user_ed25519": pk_ed, "pk_user_x25519": pk_x}) - def _token(uname, pwd): - async def _inner(): - r = await client.post("/v1/users/login", - json={"username": uname, "password": pwd}) - return r.json()["access_token"] - return _inner - alice_token = (await client.post("/v1/users/login", json={"username": "alice2", "password": "alicepass99"})).json()["access_token"] - bob_token = (await client.post("/v1/users/login", - json={"username": "bob2", "password": "bobpass99"})).json()["access_token"] a_hdrs = {"Authorization": f"Bearer {alice_token}"} - b_hdrs = {"Authorization": f"Bearer {bob_token}"} - # Alice creates group r = await client.post("/v1/groups", json={"name": "mygroup"}, headers=a_hdrs) assert r.status_code == 201 group_id = r.json()["group_id"] - # Alice generates GEK and wraps it for bob - gek = generate_gek() - pk_bob_raw = base64.b64decode(pk_x_b) - bundle = wrap_gek(gek, pk_bob_raw) - - r = await client.post(f"/v1/groups/{group_id}/members/bob2/gek", - json=bundle, headers=a_hdrs) + # Add bob as member (hub handles membership only, GEK exchange is P2P) + r = await client.post(f"/v1/groups/{group_id}/members/bob2", + json={}, headers=a_hdrs) assert r.status_code == 201 - # Bob retrieves his bundle - r = await client.get(f"/v1/groups/{group_id}/gek", headers=b_hdrs) + # Verify bob is in the group + bob_token = (await client.post("/v1/users/login", + json={"username": "bob2", "password": "bobpass99"})).json()["access_token"] + b_hdrs = {"Authorization": f"Bearer {bob_token}"} + r = await client.get(f"/v1/groups/{group_id}/members", headers=b_hdrs) assert r.status_code == 200 - retrieved = r.json() - - # Bob unwraps — must recover original GEK - sk_b_raw = sk_x_b.private_bytes( - serialization.Encoding.Raw, serialization.PrivateFormat.Raw, - serialization.NoEncryption()) - recovered = unwrap_gek(retrieved, sk_b_raw, pk_bob_raw) - assert recovered == gek + members = [m["username"] for m in r.json()["members"]] + assert "alice2" in members + assert "bob2" in members @pytest.mark.asyncio @@ -291,12 +270,9 @@ async def test_non_admin_cannot_add_member(client): headers={"Authorization": f"Bearer {charlie_token}"}) group_id = r.json()["group_id"] - gek = generate_gek() - bundle = wrap_gek(gek, base64.b64decode(pk_x_b)) - # Dan (non-admin) tries to add a member → 403 - r = await client.post(f"/v1/groups/{group_id}/members/charlie/gek", - json=bundle, + r = await client.post(f"/v1/groups/{group_id}/members/charlie", + json={}, headers={"Authorization": f"Bearer {dan_token}"}) assert r.status_code == 403 @@ -331,10 +307,8 @@ async def test_jwt_contains_groups_claim(client): headers={"Authorization": f"Bearer {alice_token}"}) group_id = r.json()["group_id"] - gek = generate_gek() - bundle = wrap_gek(gek, base64.b64decode(pk_x_b)) - await client.post(f"/v1/groups/{group_id}/members/grp_bob/gek", - json=bundle, + await client.post(f"/v1/groups/{group_id}/members/grp_bob", + json={}, headers={"Authorization": f"Bearer {alice_token}"}) # Login again — groups should contain the new group @@ -381,10 +355,8 @@ async def test_my_groups(client): r = await client.post("/v1/groups", json={"name": "mg-group"}, headers={"Authorization": f"Bearer {alice_token}"}) group_id = r.json()["group_id"] - gek = generate_gek() - bundle = wrap_gek(gek, base64.b64decode(pk_x_b)) - await client.post(f"/v1/groups/{group_id}/members/mg_bob/gek", - json=bundle, + await client.post(f"/v1/groups/{group_id}/members/mg_bob", + json={}, headers={"Authorization": f"Bearer {alice_token}"}) # Re-login to get fresh token with group claims @@ -559,8 +531,7 @@ async def test_registered_email_not_plaintext(client, app): @pytest.mark.asyncio async def test_password_rehash_on_login(client, app): - """Users with pw_version=1 get rehashed to current version on login.""" - from meshbay_hub.auth import current_pw_version + """Users with pw_version=1 get rehashed to v2 on legacy password login.""" from meshbay_hub.db.engine import get_db from meshbay_hub.db.models import User from sqlalchemy import select @@ -590,16 +561,16 @@ async def test_password_rehash_on_login(client, app): await db.commit() break - # Login should succeed and trigger rehash + # Login should succeed and trigger legacy rehash (v1 -> v2) r = await client.post("/v1/users/login", json={ "username": "rehash_user", "password": "rehashpass9"}) assert r.status_code == 200 - # Verify pw_version is now current + # Verify pw_version is now 2 (legacy rehash stays within password scheme) async for db in get_db(): result = await db.execute(select(User).where(User.username == "rehash_user")) user = result.scalar_one() - assert user.pw_version == current_pw_version() + assert user.pw_version == 2 break # Login still works after rehash @@ -722,3 +693,205 @@ async def test_webapp_html_includes_scripts(client): assert 'type="module"' in html assert 'rel="stylesheet"' in html assert 'href="/style.css"' in html + + +# ── Password split (T1 fix) ───────────────────────────────────────────────── + +@pytest.mark.asyncio +async def test_register_with_auth_key(client, app): + """Registration with auth_key sets pw_version 3.""" + from meshbay_hub.db.engine import get_db + from meshbay_hub.db.models import User + from sqlalchemy import select + + pk_ed, pk_x, _ = _gen_user_keys() + r = await client.post("/v1/users/register", json={ + "username": "authuser", + "email": "auth@test.com", + "auth_key": "dGVzdGF1dGhrZXl0ZXN0YXV0aGtleXRlc3RhdXRo", + "pk_user_ed25519": pk_ed, + "pk_user_x25519": pk_x, + }) + assert r.status_code == 201 + assert "user_id" in r.json() + + async for db in get_db(): + result = await db.execute(select(User).where(User.username == "authuser")) + user = result.scalar_one() + assert user.pw_version == 3 + break + + +@pytest.mark.asyncio +async def test_register_with_password_sets_v2(client, app): + """Registration with raw password (legacy) sets pw_version 2.""" + from meshbay_hub.db.engine import get_db + from meshbay_hub.db.models import User + from sqlalchemy import select + + pk_ed, pk_x, _ = _gen_user_keys() + r = await client.post("/v1/users/register", json={ + "username": "legacyreg", + "email": "legacy@test.com", + "password": "legacypass99", + "pk_user_ed25519": pk_ed, + "pk_user_x25519": pk_x, + }) + assert r.status_code == 201 + + async for db in get_db(): + result = await db.execute(select(User).where(User.username == "legacyreg")) + user = result.scalar_one() + assert user.pw_version == 2 + break + + +@pytest.mark.asyncio +async def test_register_no_credentials_rejected(client): + """Registration without auth_key or password returns 400.""" + pk_ed, pk_x, _ = _gen_user_keys() + r = await client.post("/v1/users/register", json={ + "username": "nocred", + "email": "nocred@test.com", + "pk_user_ed25519": pk_ed, + "pk_user_x25519": pk_x, + }) + assert r.status_code == 400 + assert "auth_key or password required" in r.json()["detail"] + + +@pytest.mark.asyncio +async def test_login_with_auth_key(client, app): + """Login with auth_key for pw_version 3 account succeeds.""" + pk_ed, pk_x, _ = _gen_user_keys() + auth_key = "dGVzdGF1dGhrZXl0ZXN0YXV0aGtleXRlc3RhdXRo" + await client.post("/v1/users/register", json={ + "username": "authlogin", + "email": "authlogin@test.com", + "auth_key": auth_key, + "pk_user_ed25519": pk_ed, + "pk_user_x25519": pk_x, + }) + + r = await client.post("/v1/users/login", json={ + "username": "authlogin", "auth_key": auth_key}) + assert r.status_code == 200 + data = r.json() + assert "access_token" in data + assert "refresh_token" in data + + +@pytest.mark.asyncio +async def test_login_auth_key_wrong_rejected(client): + """Login with wrong auth_key returns 401.""" + pk_ed, pk_x, _ = _gen_user_keys() + await client.post("/v1/users/register", json={ + "username": "authwrong", + "email": "authwrong@test.com", + "auth_key": "dGVzdGF1dGhrZXl0ZXN0YXV0aGtleXRlc3RhdXRo", + "pk_user_ed25519": pk_ed, + "pk_user_x25519": pk_x, + }) + + r = await client.post("/v1/users/login", json={ + "username": "authwrong", "auth_key": "d3JvbmdrZXl3cm9uZ2tleXdyb25na2V5d3Jvbmc="}) + assert r.status_code == 401 + + +@pytest.mark.asyncio +async def test_login_v3_account_password_only_rejected(client): + """Login with raw password to a v3 (auth_key) account returns 401.""" + pk_ed, pk_x, _ = _gen_user_keys() + await client.post("/v1/users/register", json={ + "username": "v3nopw", + "email": "v3nopw@test.com", + "auth_key": "dGVzdGF1dGhrZXl0ZXN0YXV0aGtleXRlc3RhdXRo", + "pk_user_ed25519": pk_ed, + "pk_user_x25519": pk_x, + }) + + r = await client.post("/v1/users/login", json={ + "username": "v3nopw", "password": "somepassword"}) + assert r.status_code == 401 + + +@pytest.mark.asyncio +async def test_login_legacy_upgrade_required(client): + """Legacy account (pw_version 2) with auth_key only returns auth_upgrade_required.""" + pk_ed, pk_x, _ = _gen_user_keys() + await client.post("/v1/users/register", json={ + "username": "legacyupg", + "email": "legacyupg@test.com", + "password": "legacypass99", + "pk_user_ed25519": pk_ed, + "pk_user_x25519": pk_x, + }) + + r = await client.post("/v1/users/login", json={ + "username": "legacyupg", "auth_key": "dGVzdGF1dGhrZXl0ZXN0YXV0aGtleXRlc3RhdXRo"}) + assert r.status_code == 401 + assert r.json()["detail"] == "auth_upgrade_required" + + +@pytest.mark.asyncio +async def test_login_legacy_migration(client, app): + """Legacy account migrates to auth_key on login with both fields.""" + from meshbay_hub.auth import current_pw_version + from meshbay_hub.db.engine import get_db + from meshbay_hub.db.models import User + from sqlalchemy import select + + pk_ed, pk_x, _ = _gen_user_keys() + await client.post("/v1/users/register", json={ + "username": "migrateuser", + "email": "migrate@test.com", + "password": "migratepass9", + "pk_user_ed25519": pk_ed, + "pk_user_x25519": pk_x, + }) + + # Verify starts at pw_version 2 + async for db in get_db(): + result = await db.execute(select(User).where(User.username == "migrateuser")) + user = result.scalar_one() + assert user.pw_version == 2 + break + + auth_key = "bWlncmF0ZWF1dGhrZXltaWdyYXRlYXV0aGtleW1p" + + # Login with password + auth_key → should succeed and migrate + r = await client.post("/v1/users/login", json={ + "username": "migrateuser", + "password": "migratepass9", + "auth_key": auth_key, + }) + assert r.status_code == 200 + + # Verify pw_version is now 3 (migrated) + async for db in get_db(): + result = await db.execute(select(User).where(User.username == "migrateuser")) + user = result.scalar_one() + assert user.pw_version == current_pw_version() + assert user.pw_version == 3 + break + + # Login again with auth_key only → should succeed (migrated account) + r = await client.post("/v1/users/login", json={ + "username": "migrateuser", "auth_key": auth_key}) + assert r.status_code == 200 + assert "access_token" in r.json() + + # Old password no longer works (hash was replaced with auth_key hash) + r = await client.post("/v1/users/login", json={ + "username": "migrateuser", "password": "migratepass9"}) + assert r.status_code == 401 + + +@pytest.mark.asyncio +async def test_login_no_credentials_rejected(client): + """Login without auth_key or password returns 401.""" + r = await client.post("/v1/users/login", json={"username": "nobody"}) + assert r.status_code == 401 + assert "No credentials" in r.json()["detail"] + + diff --git a/packages/meshbay-hub/tests/test_node_auth.py b/packages/meshbay-hub/tests/test_node_auth.py new file mode 100644 index 0000000..e629d20 --- /dev/null +++ b/packages/meshbay-hub/tests/test_node_auth.py @@ -0,0 +1,281 @@ +"""Tests for node Ed25519 authentication and JWT scope enforcement. + +Node auth flow: register user (browser keys) → login → link node key → node auth. +The node's Ed25519 key is separate from the user's browser key. +""" + +import base64 +import time + +import pytest +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey +from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey +from cryptography.hazmat.primitives import serialization + + +def _gen_ed25519(): + sk = Ed25519PrivateKey.generate() + pk_raw = sk.public_key().public_bytes( + serialization.Encoding.Raw, serialization.PublicFormat.Raw) + return sk, base64.b64encode(pk_raw).decode() + + +def _gen_x25519(): + sk = X25519PrivateKey.generate() + pk_raw = sk.public_key().public_bytes( + serialization.Encoding.Raw, serialization.PublicFormat.Raw) + return sk, base64.b64encode(pk_raw).decode() + + +async def _register(client, username, pk_ed_b64, pk_x_b64): + r = await client.post("/v1/users/register", json={ + "username": username, "email": f"{username}@test.local", + "password": "testpass99", + "pk_user_ed25519": pk_ed_b64, "pk_user_x25519": pk_x_b64, + }) + assert r.status_code == 201 + return r.json()["user_id"] + + +async def _user_login(client, username, password="testpass99"): + r = await client.post("/v1/users/login", json={ + "username": username, "password": password, + }) + assert r.status_code == 200 + return r.json()["access_token"] + + +async def _link_node_key(client, user_token, pk_node_ed_b64): + r = await client.put("/v1/users/me/node_key", json={ + "pk_node_ed25519": pk_node_ed_b64, + }, headers={"Authorization": f"Bearer {user_token}"}) + assert r.status_code == 200 + return r + + +async def _node_auth(client, username, sk_node_ed): + timestamp = int(time.time()) + message = f"meshbay:node_auth:{username}:{timestamp}".encode() + signature = sk_node_ed.sign(message) + r = await client.post("/v1/nodes/auth", json={ + "username": username, + "timestamp": timestamp, + "signature": base64.b64encode(signature).decode(), + }) + return r + + +async def _setup_node_user(client, username): + """Register user with browser keys, login, link node key. Returns (node_sk, node_token).""" + _, pk_ed = _gen_ed25519() + _, pk_x = _gen_x25519() + await _register(client, username, pk_ed, pk_x) + + user_token = await _user_login(client, username) + + sk_node, pk_node = _gen_ed25519() + await _link_node_key(client, user_token, pk_node) + + return sk_node, user_token + + +# ── Core auth tests ────────────────────────────────────────────────────────── + + +@pytest.mark.asyncio +async def test_node_auth_success(client): + sk_node, _ = await _setup_node_user(client, "nodeuser") + + r = await _node_auth(client, "nodeuser", sk_node) + assert r.status_code == 200 + data = r.json() + assert "access_token" in data + assert data["token_type"] == "bearer" + + +@pytest.mark.asyncio +async def test_node_auth_wrong_key_rejected(client): + await _setup_node_user(client, "nodeuser2") + + wrong_sk, _ = _gen_ed25519() + r = await _node_auth(client, "nodeuser2", wrong_sk) + assert r.status_code == 401 + + +@pytest.mark.asyncio +async def test_node_auth_stale_timestamp(client): + sk_node, _ = await _setup_node_user(client, "nodeuser3") + + timestamp = int(time.time()) - 120 + message = f"meshbay:node_auth:nodeuser3:{timestamp}".encode() + signature = sk_node.sign(message) + r = await client.post("/v1/nodes/auth", json={ + "username": "nodeuser3", + "timestamp": timestamp, + "signature": base64.b64encode(signature).decode(), + }) + assert r.status_code == 401 + + +@pytest.mark.asyncio +async def test_node_auth_fails_without_node_key(client): + """Node auth must fail if no node key has been registered.""" + _, pk_ed = _gen_ed25519() + _, pk_x = _gen_x25519() + await _register(client, "nokey_user", pk_ed, pk_x) + + sk_any, _ = _gen_ed25519() + r = await _node_auth(client, "nokey_user", sk_any) + assert r.status_code == 401 + assert "No node key registered" in r.json()["detail"] + + +# ── Scope enforcement tests ────────────────────────────────────────────────── + + +@pytest.mark.asyncio +async def test_node_scope_blocks_group_create(client): + sk_node, _ = await _setup_node_user(client, "operator") + + r = await _node_auth(client, "operator", sk_node) + node_token = r.json()["access_token"] + + r = await client.post("/v1/groups", json={ + "name": "forbidden-group", "visibility": "private", + }, headers={"Authorization": f"Bearer {node_token}"}) + assert r.status_code == 403 + assert "Node-scoped" in r.json()["detail"] + + +@pytest.mark.asyncio +async def test_node_scope_blocks_add_member(client): + sk_node, user_token = await _setup_node_user(client, "op1") + + r = await client.post("/v1/groups", json={ + "name": "mygroup", "visibility": "private", "join_policy": "invite", + }, headers={"Authorization": f"Bearer {user_token}"}) + assert r.status_code == 201 + gid = r.json()["group_id"] + + _, pk2 = _gen_ed25519() + _, px2 = _gen_x25519() + await _register(client, "member1", pk2, px2) + + r = await _node_auth(client, "op1", sk_node) + node_token = r.json()["access_token"] + + r = await client.post(f"/v1/groups/{gid}/members/member1", + headers={"Authorization": f"Bearer {node_token}"}) + assert r.status_code == 403 + + +@pytest.mark.asyncio +async def test_node_scope_blocks_delete_group(client): + sk_node, user_token = await _setup_node_user(client, "op2") + + r = await client.post("/v1/groups", json={ + "name": "deleteme", "visibility": "public", "join_policy": "open", + }, headers={"Authorization": f"Bearer {user_token}"}) + gid = r.json()["group_id"] + + r = await _node_auth(client, "op2", sk_node) + node_token = r.json()["access_token"] + + r = await client.delete(f"/v1/groups/{gid}", + headers={"Authorization": f"Bearer {node_token}"}) + assert r.status_code == 403 + + +@pytest.mark.asyncio +async def test_node_scope_allows_read_members(client): + sk_node, user_token = await _setup_node_user(client, "op3") + + r = await client.post("/v1/groups", json={ + "name": "readgroup", "visibility": "private", "join_policy": "invite", + }, headers={"Authorization": f"Bearer {user_token}"}) + gid = r.json()["group_id"] + + r = await _node_auth(client, "op3", sk_node) + node_token = r.json()["access_token"] + + r = await client.get(f"/v1/groups/{gid}/members", + headers={"Authorization": f"Bearer {node_token}"}) + assert r.status_code == 200 + assert len(r.json()["members"]) == 1 + + +@pytest.mark.asyncio +async def test_node_scope_allows_pubkey_lookup(client): + sk_node, _ = await _setup_node_user(client, "op4") + + r = await _node_auth(client, "op4", sk_node) + node_token = r.json()["access_token"] + + r = await client.get("/v1/users/op4/pubkeys", + headers={"Authorization": f"Bearer {node_token}"}) + assert r.status_code == 200 + assert "pk_ed25519" in r.json() + assert r.json()["pk_node_ed25519"] is not None + + +@pytest.mark.asyncio +async def test_user_scope_still_works(client): + """Verify that user-scoped tokens (browser login) still have full access.""" + _, pk_ed = _gen_ed25519() + _, pk_x = _gen_x25519() + await _register(client, "webuser", pk_ed, pk_x) + + token = await _user_login(client, "webuser") + + r = await client.post("/v1/groups", json={ + "name": "browser-group", "visibility": "public", "join_policy": "open", + }, headers={"Authorization": f"Bearer {token}"}) + assert r.status_code == 201 + + +# ── Node key registration tests ────────────────────────────────────────────── + + +@pytest.mark.asyncio +async def test_link_node_key(client): + _, pk_ed = _gen_ed25519() + _, pk_x = _gen_x25519() + await _register(client, "linkuser", pk_ed, pk_x) + + token = await _user_login(client, "linkuser") + _, pk_node = _gen_ed25519() + + r = await client.put("/v1/users/me/node_key", json={ + "pk_node_ed25519": pk_node, + }, headers={"Authorization": f"Bearer {token}"}) + assert r.status_code == 200 + assert r.json()["pk_node_ed25519"] == pk_node + + +@pytest.mark.asyncio +async def test_link_node_key_invalid_format(client): + _, pk_ed = _gen_ed25519() + _, pk_x = _gen_x25519() + await _register(client, "badkey", pk_ed, pk_x) + + token = await _user_login(client, "badkey") + + r = await client.put("/v1/users/me/node_key", json={ + "pk_node_ed25519": "not-valid-base64!!!", + }, headers={"Authorization": f"Bearer {token}"}) + assert r.status_code == 400 + + +@pytest.mark.asyncio +async def test_link_node_key_blocked_for_node_scope(client): + """Node-scoped tokens must not be able to change the node key.""" + sk_node, _ = await _setup_node_user(client, "sneaky") + + r = await _node_auth(client, "sneaky", sk_node) + node_token = r.json()["access_token"] + + _, pk_evil = _gen_ed25519() + r = await client.put("/v1/users/me/node_key", json={ + "pk_node_ed25519": pk_evil, + }, headers={"Authorization": f"Bearer {node_token}"}) + assert r.status_code == 403 diff --git a/packages/meshbay-node/pyproject.toml b/packages/meshbay-node/pyproject.toml index 592de54..864fd04 100644 --- a/packages/meshbay-node/pyproject.toml +++ b/packages/meshbay-node/pyproject.toml @@ -17,6 +17,7 @@ dependencies = [ "aioquic>=1.0", # QUIC transport (MNP v2) — implemented in Phase 5 "websockets>=12.0", # hub→node revocation push "aiortc>=1.9", # WebRTC DataChannel for browser P2P (Phase 9) + "aiosqlite>=0.20", # async SQLite for chat, audit, bundle stores ] [project.optional-dependencies] diff --git a/packages/meshbay-node/src/meshbay_node/bundle_store.py b/packages/meshbay-node/src/meshbay_node/bundle_store.py new file mode 100644 index 0000000..e7c6981 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/bundle_store.py @@ -0,0 +1,105 @@ +""" +Bundle store — SQLite-backed storage for GEK bundles and keypair bundles. + +GEK bundles: ECIES-wrapped GEK targeted at a specific user's X25519 key. +Keypair bundles: AES-GCM encrypted (Ed25519 + X25519) private keys, encrypted +with the user's password-derived bundle_key. Opaque to the node. + +Both are stored and served over the P2P DataChannel during MNP handshake. +""" + +import logging +from pathlib import Path + +import aiosqlite + +log = logging.getLogger(__name__) + +_SCHEMA_GEK = """\ +CREATE TABLE IF NOT EXISTS gek_bundles ( + group_id TEXT NOT NULL, + user_id TEXT NOT NULL, + pk_eph_b64 TEXT NOT NULL, + nonce_b64 TEXT NOT NULL, + wrapped_b64 TEXT NOT NULL, + stored_at TEXT NOT NULL DEFAULT (datetime('now')), + PRIMARY KEY (group_id, user_id) +); +""" + +_SCHEMA_KEYPAIR = """\ +CREATE TABLE IF NOT EXISTS keypair_bundles ( + user_id TEXT PRIMARY KEY, + bundle_enc TEXT NOT NULL, + stored_at TEXT NOT NULL DEFAULT (datetime('now')) +); +""" + + +class BundleStore: + 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)) + await self._db.execute(_SCHEMA_GEK) + await self._db.execute(_SCHEMA_KEYPAIR) + await self._db.commit() + + async def store( + self, + group_id: str, + user_id: str, + pk_eph_b64: str, + nonce_b64: str, + wrapped_b64: str, + ) -> None: + assert self._db + await self._db.execute( + "INSERT OR REPLACE INTO gek_bundles " + "(group_id, user_id, pk_eph_b64, nonce_b64, wrapped_b64, stored_at) " + "VALUES (?, ?, ?, ?, ?, datetime('now'))", + (group_id, user_id, pk_eph_b64, nonce_b64, wrapped_b64), + ) + await self._db.commit() + + async def fetch(self, group_id: str, user_id: str) -> dict | None: + assert self._db + async with self._db.execute( + "SELECT pk_eph_b64, nonce_b64, wrapped_b64 FROM gek_bundles " + "WHERE group_id = ? AND user_id = ?", + (group_id, user_id), + ) as cursor: + row = await cursor.fetchone() + if not row: + return None + return { + "pk_eph_b64": row[0], + "nonce_b64": row[1], + "wrapped_b64": row[2], + } + + async def store_keypair(self, user_id: str, bundle_enc: str) -> None: + assert self._db + await self._db.execute( + "INSERT OR REPLACE INTO keypair_bundles " + "(user_id, bundle_enc, stored_at) VALUES (?, ?, datetime('now'))", + (user_id, bundle_enc), + ) + await self._db.commit() + + async def fetch_keypair(self, user_id: str) -> str | None: + assert self._db + async with self._db.execute( + "SELECT bundle_enc FROM keypair_bundles WHERE user_id = ?", + (user_id,), + ) as cursor: + row = await cursor.fetchone() + return row[0] if row else None + + async def close(self) -> None: + if self._db: + await self._db.close() + self._db = None diff --git a/packages/meshbay-node/src/meshbay_node/config.py b/packages/meshbay-node/src/meshbay_node/config.py index 9e6a391..a7a0785 100644 --- a/packages/meshbay-node/src/meshbay_node/config.py +++ b/packages/meshbay-node/src/meshbay_node/config.py @@ -52,6 +52,11 @@ visibility = "public" [keystore] # unlock_file = "~/.config/meshbay/unlock.key" # or set MESHBAY_UNLOCK_KEY env var + +# Node sovereignty: pin the operator's Ed25519 public key (base64, 32 bytes raw). +# Admin operations (file delete) require cryptographic proof of this key. +# Auto-pinned on first startup from the node operator's keystore. +# admin_pk_ed25519 = "base64-encoded-32-bytes" """ @@ -59,7 +64,6 @@ visibility = "public" class HubConfig: url: str = "https://meshbay.org" username: str = "" - password: str = "" # loaded from keystore or env; never written to TOML @dataclass @@ -94,6 +98,7 @@ class Config: groups: list[GroupConfig] = field(default_factory=list) keystore: KeystoreConfig = field(default_factory=KeystoreConfig) data_dir: Path = field(default_factory=lambda: Path.home() / ".local" / "share" / "meshbay") + admin_pk_ed25519: str = "" # base64 raw Ed25519 public key pinned locally # Back-compat: single-group access @property @@ -144,6 +149,9 @@ def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config: if "data_dir" in raw: cfg.data_dir = Path(raw["data_dir"]).expanduser().resolve() + if "admin_pk_ed25519" in raw: + cfg.admin_pk_ed25519 = raw["admin_pk_ed25519"] + ks = raw.get("keystore", {}) if "path" in ks: cfg.keystore.path = Path(ks["path"]).expanduser() @@ -155,8 +163,6 @@ 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 pwd := os.environ.get("MESHBAY_PASSWORD"): - cfg.hub.password = pwd if port := os.environ.get("MESHBAY_PORT"): cfg.node.port = int(port) diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index 5851b34..fe12909 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -23,6 +23,7 @@ Usage: """ import asyncio +import base64 import json import logging import signal @@ -30,12 +31,14 @@ import sys from pathlib import Path import uvicorn +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey from meshbay_common import MNP_VERSION from meshbay_common.protocol import MNP from meshbay_node.audit import AuditStore +from meshbay_node.bundle_store import BundleStore from meshbay_node.chat.store import ChatStore -from meshbay_node.config import Config, load_config, write_example_config +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 @@ -111,6 +114,7 @@ class NodeDaemon: self._denylist = Denylist() if Denylist else None self._chat_stores: dict[str, ChatStore] = {} self._audit_store: AuditStore | None = None + self._bundle_store: BundleStore | None = None self._indexers: list[DirectoryIndexer] = [] self._tasks: list[asyncio.Task] = [] self._hub: HubClient | None = None @@ -126,18 +130,46 @@ class NodeDaemon: ) log.info("Keys loaded: %s", keys.pk_ed25519_b64[:16]) - # 2. Hub connection + # 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 + from meshbay_node.ui import create_ui_app + ui_app = create_ui_app(self._state) + ui_cfg = uvicorn.Config( + ui_app, + host="127.0.0.1", + port=self._config.node.ui_port, + log_level="warning", + ) + 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) + + # 3. Hub connection (Ed25519 auth — retries until node key is linked) hub_cfg = HubConfig( hub_url=self._config.hub.url, username=self._config.hub.username, - password=self._config.hub.password, ) async with HubClient(hub_cfg, keys) as hub: self._hub = hub - session = await hub.startup(endpoint_hint=None) + session = await self._login_with_retry(hub) self._state["endpoint_hint"] = session.node_id - # 3. Build per-group contexts + # 4. Bundle store (P2P GEK bundles) + data_dir = self._config.data_dir + data_dir.mkdir(parents=True, exist_ok=True) + self._bundle_store = BundleStore(db_path=data_dir / "bundles.db") + await self._bundle_store.open() + log.info("Bundle store opened: %s", data_dir / "bundles.db") + + # X25519 key material for GEK unwrapping + from cryptography.hazmat.primitives import serialization + sk_x_raw = keys.sk_x25519.private_bytes( + serialization.Encoding.Raw, serialization.PrivateFormat.Raw, + serialization.NoEncryption()) + pk_x_raw = base64.b64decode(keys.pk_x25519_b64) + + # 4. Build per-group contexts groups_ctx: dict[str, dict] = {} for group_cfg in self._config.groups: if not group_cfg.id or not group_cfg.shared_dir: @@ -153,12 +185,13 @@ class NodeDaemon: gek = None if group_cfg.visibility == "private": - try: - gek = await hub.fetch_gek(group_cfg.id) + gek = await self._load_gek( + group_cfg.id, session.user_id, sk_x_raw, pk_x_raw) + if gek: log.info("GEK loaded for group %s", group_cfg.id[:8]) - except LookupError: - log.warning("No GEK for group %s — skipping", group_cfg.name) - continue + else: + log.info("No GEK yet for group %s — will accept first setup", + group_cfg.name) indexer = DirectoryIndexer( root=shared_root, @@ -183,9 +216,7 @@ class NodeDaemon: log.error("No valid groups configured — exiting") return - # 4. Chat stores (one SQLite DB per group) - data_dir = self._config.data_dir - data_dir.mkdir(parents=True, exist_ok=True) + # 5. Chat stores (one SQLite DB per group) for gid in groups_ctx: chat_db = data_dir / gid[:16] / "chat.db" store = ChatStore(db_path=chat_db) @@ -194,7 +225,7 @@ class NodeDaemon: groups_ctx[gid]["chat_store"] = store log.info("Chat stores opened: %d groups", len(self._chat_stores)) - # 4b. Audit store (legal compliance — IP + action logging) + # 6. Audit store (legal compliance — IP + action logging) audit_db = data_dir / "audit.db" self._audit_store = AuditStore(db_path=audit_db) await self._audit_store.open() @@ -219,6 +250,17 @@ class NodeDaemon: self._webrtc._ctx["hub_ws"] = _WsSender(hub) self._webrtc._ctx["node_user_id"] = session.user_id self._webrtc._ctx["audit_store"] = self._audit_store + self._webrtc._ctx["bundle_store"] = self._bundle_store + self._webrtc._ctx["sk_x25519_raw"] = sk_x_raw + 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) + if admin_pk: + self._webrtc._ctx["admin_pk_ed25519"] = admin_pk + log.info("Admin Ed25519 key pinned for node sovereignty") + else: + log.warning("No admin_pk_ed25519 — admin operations disabled") log.info("WebRTC transport ready") else: log.warning("WebRTC not available (aiortc not installed)") @@ -323,23 +365,13 @@ class NodeDaemon: log.info("HTTP API on port %d for group %s", group_cfg.http_port, group_cfg.name) - # 10. Local web UI + # 10. Update admin UI state (UI already running from step 2) self._state["groups_ctx"] = groups_ctx - self._state["config"] = self._config self._state["audit_store"] = self._audit_store + self._state["bundle_store"] = self._bundle_store self._state["webrtc"] = self._webrtc self._state["hub"] = hub - from meshbay_node.ui import create_ui_app - ui_app = create_ui_app(self._state) - ui_cfg = uvicorn.Config( - ui_app, - host="127.0.0.1", - port=self._config.node.ui_port, - log_level="warning", - ) - ui_server = uvicorn.Server(ui_cfg) - self._tasks.append(asyncio.create_task(ui_server.serve())) - log.info("Local UI at http://localhost:%d", self._config.node.ui_port) + self._state["pk_x25519_raw"] = pk_x_raw self._state["status"] = "running" log.info("Node ready — %d groups, WebRTC=%s, QUIC=%s", @@ -363,6 +395,75 @@ class NodeDaemon: await self._shutdown() + async def _login_with_retry(self, hub: HubClient): + """Login to hub, retrying if the node key hasn't been linked yet.""" + import httpx as _httpx + while True: + try: + 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, + ) + await asyncio.sleep(30) + else: + raise + except Exception as e: + log.warning("Hub login failed: %s — retrying in 10s", e) + await asyncio.sleep(10) + + async def _load_gek( + self, + group_id: str, + node_user_id: str, + sk_x_raw: bytes, + pk_x_raw: bytes, + ) -> bytes | None: + """Load GEK from local bundle store (node-only, hub never touches crypto).""" + from meshbay_common.crypto import unwrap_gek_aes + + if not self._bundle_store: + return None + + # Try node-specific bundle first (stored by init_gek for daemon reload), + # then fall back to operator's user bundle (legacy / pre-dual-key) + for user_key in [f"_node_{node_user_id}", node_user_id]: + bundle = await self._bundle_store.fetch(group_id, user_key) + if not bundle: + continue + try: + gek = unwrap_gek_aes(bundle, sk_x_raw, pk_x_raw) + log.info("GEK loaded from local bundle store for group %s (key=%s)", + group_id[:8], user_key[:16]) + return gek + except Exception as e: + log.debug("Failed to unwrap GEK bundle (key=%s): %s", user_key[:16], e) + + 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 + + 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 + async def _on_index_change(self, indexer: DirectoryIndexer) -> None: """Called when a DirectoryIndexer detects file changes.""" group_id = indexer.group_id @@ -429,6 +530,9 @@ class NodeDaemon: if self._audit_store: await self._audit_store.close() + if self._bundle_store: + await self._bundle_store.close() + for store in self._chat_stores.values(): await store.close() @@ -475,7 +579,7 @@ def main() -> None: calibrate_argon2() return - cfg = load_config(args.config) + 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") sys.exit(1) diff --git a/packages/meshbay-node/src/meshbay_node/hub_client.py b/packages/meshbay-node/src/meshbay_node/hub_client.py index ba9d3ff..432af0a 100644 --- a/packages/meshbay-node/src/meshbay_node/hub_client.py +++ b/packages/meshbay-node/src/meshbay_node/hub_client.py @@ -2,15 +2,15 @@ MeshBay Node — Hub client. Handles all communication from the node to a Mesh Hub: - - User registration (first run) - - Login → JWT (access token + refresh token) + - Ed25519 authentication (node-scoped JWT, no password material on node) - JWT offline verification and auto-refresh - Node announcement (endpoint_hint) - - GEK bundle retrieval for a group - User public key lookup (for GEK wrapping) + - Swarm hash registration -JWT verification is done locally using the hub's cached Ed25519 public key. -The hub is only contacted for login and refresh — not for every request. +The node authenticates via Ed25519 challenge-response (/v1/nodes/auth). +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 base64 @@ -23,10 +23,7 @@ from typing import Any, Callable import httpx import jwt -from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey -from cryptography.hazmat.primitives import serialization -from meshbay_common.crypto import pk_to_b64, unwrap_gek from meshbay_node.keystore import NodeKeys log = logging.getLogger(__name__) @@ -62,7 +59,6 @@ class HubSession: class HubConfig: hub_url: str username: str - password: str cache_dir: Path = field(default_factory=lambda: Path.home() / ".config" / "meshbay") @property @@ -110,80 +106,54 @@ class HubClient: log.info("Hub PK fetched and cached: %s", cache) return pem - # ── Registration ────────────────────────────────────────────────────────── - - async def register(self) -> str: - """Register this node's user on the hub. Returns user_id. Idempotent (409 ok).""" - r = await self._http.post("/v1/users/register", json={ - "username": self._config.username, - "password": self._config.password, - "pk_user_ed25519": self._keys.pk_ed25519_b64, - "pk_user_x25519": self._keys.pk_x25519_b64, - }) - if r.status_code == 201: - log.info("Registered user '%s' on hub", self._config.username) - return r.json()["user_id"] - if r.status_code == 409: - log.debug("User '%s' already registered", self._config.username) - return "" - r.raise_for_status() - return "" - - # ── Login ───────────────────────────────────────────────────────────────── + # ── Ed25519 authentication ─────────────────────────────────────────────── async def login(self) -> HubSession: - """Login, verify JWT offline, return HubSession.""" + """Authenticate via Ed25519 challenge-response. Returns node-scoped HubSession.""" hub_pk_pem = await self._fetch_hub_pk() - r = await self._http.post("/v1/users/login", json={ - "username": self._config.username, - "password": self._config.password, + timestamp = int(time.time()) + message = f"meshbay:node_auth:{self._config.username}:{timestamp}".encode() + signature = self._keys.sk_ed25519.sign(message) + + r = await self._http.post("/v1/nodes/auth", json={ + "username": self._config.username, + "timestamp": timestamp, + "signature": base64.b64encode(signature).decode(), }) r.raise_for_status() data = r.json() - access_token = data["access_token"] - refresh_token = data["refresh_token"] + access_token = data["access_token"] - # Verify offline — if this passes, the hub's identity is confirmed 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" assert "jti" in decoded, "Hub token missing jti — hub is outdated" + assert decoded.get("scope") == "node", \ + "Expected node-scoped token" - self._session = HubSession( - hub_url=self._config.hub_url, - username=self._config.username, - user_id=decoded["sub"], - access_token=access_token, - refresh_token=refresh_token, - hub_pk_pem=hub_pk_pem, - _token_exp=decoded["exp"], - ) + if self._session: + self._session.access_token = access_token + self._session._token_exp = decoded["exp"] + else: + self._session = HubSession( + hub_url=self._config.hub_url, + username=self._config.username, + user_id=decoded["sub"], + access_token=access_token, + refresh_token="", + hub_pk_pem=hub_pk_pem, + _token_exp=decoded["exp"], + ) log.info("Logged in as '%s' (exp in %ds)", self._config.username, self._session.token_expires_in) return self._session - async def refresh_token(self) -> None: - """Refresh the access token using the refresh token.""" - if self._session is None: - raise RuntimeError("Not logged in") - - r = await self._http.post("/v1/users/token/refresh", json={ - "refresh_token": self._session.refresh_token, - }) - r.raise_for_status() - new_token = r.json()["access_token"] - - decoded = jwt.decode(new_token, self._session.hub_pk_pem, algorithms=["EdDSA"]) - self._session.access_token = new_token - self._session._token_exp = decoded["exp"] - log.debug("Access token refreshed (exp in %ds)", self._session.token_expires_in) - async def ensure_fresh_token(self) -> None: - """Auto-refresh token if close to expiry.""" + """Re-authenticate with Ed25519 if token is close to expiry.""" if self._session and self._session.token_needs_refresh: - await self.refresh_token() + await self.login() # ── Node announcement ───────────────────────────────────────────────────── @@ -203,33 +173,6 @@ class HubClient: log.info("Node announced: %s (hint=%s)", node_id[:8], endpoint_hint) return node_id - # ── GEK retrieval ───────────────────────────────────────────────────────── - - async def fetch_gek(self, group_id: str) -> bytes: - """ - Fetch and unwrap the GEK bundle for a group. - Returns the raw GEK bytes. - """ - if self._session is None: - raise RuntimeError("Not logged in") - await self.ensure_fresh_token() - - r = await self._http.get(f"/v1/groups/{group_id}/gek", - headers=self._session.auth_headers) - if r.status_code == 404: - raise LookupError(f"No GEK bundle found for group {group_id!r}") - r.raise_for_status() - - bundle = r.json() - sk_x_raw = self._keys.sk_x25519.private_bytes( - serialization.Encoding.Raw, serialization.PrivateFormat.Raw, - serialization.NoEncryption()) - pk_x_raw = base64.b64decode(self._keys.pk_x25519_b64) - - gek = unwrap_gek(bundle, sk_x_raw, pk_x_raw) - log.info("GEK unwrapped for group %s", group_id[:8]) - return gek - # ── User pubkey lookup ──────────────────────────────────────────────────── async def get_user_pubkeys(self, username: str) -> dict: @@ -354,10 +297,9 @@ class HubClient: async def startup(self, endpoint_hint: str | None = None) -> HubSession: """ - Full startup sequence: register (idempotent) → login → announce node. - Returns an active HubSession. + Full startup sequence: Ed25519 login → announce node. + The operator must register separately (browser or setup script). """ - await self.register() session = await self.login() await self.announce_node(endpoint_hint) return session 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 e692c80..13e90c8 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -24,7 +24,10 @@ Signaling flow (handled externally by the hub): import asyncio import base64 +import hashlib +import hmac import logging +import os import struct from pathlib import Path from typing import Any @@ -32,7 +35,10 @@ from typing import Any import jwt import msgpack from aiortc import RTCPeerConnection, RTCSessionDescription, RTCDataChannel -from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey +from cryptography.hazmat.primitives.asymmetric.ed25519 import ( + Ed25519PrivateKey, + Ed25519PublicKey, +) from meshbay_common import MNP_VERSION from meshbay_common.crypto import pk_to_b64 @@ -46,6 +52,15 @@ CHUNK_SIZE = 1024 * 1024 MAX_MSG = 64 * 1024 * 1024 +def _extract_dtls_fingerprint(sdp: str) -> bytes: + """Extract the DTLS SHA-256 fingerprint from SDP as raw 32 bytes.""" + for line in sdp.splitlines(): + if line.startswith("a=fingerprint:sha-256 "): + hex_str = line.split(" ", 1)[1].replace(":", "") + return bytes.fromhex(hex_str) + return b"" + + STREAM_SEGMENT_SIZE = 256 * 1024 _H264_PROFILES = {"Baseline": "42", "Main": "4d", "High": "64", "High 10": "6e"} @@ -153,6 +168,9 @@ class WebRTCPeerSession: self._peer_id: str = peer_id self._remote_ip: str = "" self._username: str = "" + self._pk_user: str = "" + self._gek_challenge: bytes | None = None + self._admin_challenges: dict[str, bytes] = {} def _setup_channel(self, channel: RTCDataChannel) -> None: self._channel = channel @@ -171,6 +189,12 @@ class WebRTCPeerSession: try: if mtype == MNP.HANDSHAKE: 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 self._user_id is None: self._send({"type": "error", "detail": "Handshake required"}) elif mtype == MNP.INDEX_SYNC: @@ -179,8 +203,6 @@ class WebRTCPeerSession: self._do_file_request(msg) elif mtype == MNP.STREAM_SEGMENT: self._do_stream_segment(msg) - elif mtype == MNP.GEK_REQUEST: - self._do_gek_request() elif mtype == MNP.CHAT_MESSAGE: self._do_chat_message(msg) elif mtype == MNP.CHAT_HISTORY: @@ -189,6 +211,12 @@ class WebRTCPeerSession: self._do_file_upload(msg) elif mtype == MNP.FILE_DELETE: 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.KEYPAIR_BUNDLE_STORE: + asyncio.ensure_future(self._do_keypair_bundle_store(msg)) elif mtype == MNP.STREAM_REQUEST: asyncio.ensure_future(self._stream_video(msg)) else: @@ -234,23 +262,235 @@ class WebRTCPeerSession: self._send({"type": "error", "detail": "Group not hosted on this node"}) return - self._user_id = decoded["sub"] - self._group_id = group_id - self._username = decoded.get("username", "") + # 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", "") + + 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: + self._send({ + "type": "error", + "detail": "Group encryption not initialized — contact node operator", + }) + return + self._send(challenge) + + def _do_handshake_response(self, msg: dict) -> None: + if not self._gek_challenge or not hasattr(self, "_pending_sub"): + self._send({"type": "error", "detail": "No pending handshake challenge"}) + 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 + 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) + 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) + + data = self._gek_challenge + offer_fp + answer_fp + expected = hmac.new(gek, data, hashlib.sha256).digest() + if not hmac.compare_digest(proof_bytes, expected): + self._send({"type": "error", "detail": "GEK proof failed"}) + self._gek_challenge = None + self._audit_auth_failed(group_id, "GEK HMAC mismatch") + return + + self._gek_challenge = None + self._complete_handshake() + + def _complete_handshake(self) -> None: + self._user_id = self._pending_sub + self._group_id = self._pending_group + self._username = self._pending_username + self._pk_user = self._pending_pk_user peers = self._ctx.get("_peers") if peers is not None: peers[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], group_id[:8] if group_id else "none") - self._send({ + self._user_id[:8], + self._group_id[:8] if self._group_id else "none") + ack = { "type": MNP.HANDSHAKE_ACK, "v": MNP_VERSION, "node_pk": pk_to_b64(self._ctx["sk_node"].public_key()), - }) + "is_node_admin": bool(node_user_id and self._user_id == node_user_id), + } + if node_user_id: + ack["node_user_id"] = node_user_id + pk_x_b64 = self._ctx.get("pk_x25519_b64") + if pk_x_b64: + ack["node_pk_x25519"] = pk_x_b64 + self._send(ack) self._audit("handshake") + async def _do_gek_bundle_fetch(self) -> None: + """Serve the caller's wrapped GEK bundle during the handshake window.""" + bundle_store = self._ctx.get("bundle_store") + if not bundle_store: + self._send({"type": MNP.GEK_BUNDLE_RESP, "v": MNP_VERSION, "found": False}) + return + + group_id = getattr(self, "_pending_group", "") + user_id = getattr(self, "_pending_sub", "") + if not group_id or not user_id: + self._send({"type": "error", "detail": "No pending handshake"}) + return + + bundle = await bundle_store.fetch(group_id, user_id) + if bundle: + self._send({ + "type": MNP.GEK_BUNDLE_RESP, + "v": MNP_VERSION, + "found": True, + "pk_eph_b64": bundle["pk_eph_b64"], + "nonce_b64": bundle["nonce_b64"], + "wrapped_b64": bundle["wrapped_b64"], + }) + 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", "") + + 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"}) + 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: + return + + bundle = await bundle_store.fetch(group_id, user_id) + if not bundle: + 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) + 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)") + + async def _do_keypair_bundle_fetch(self) -> None: + """Serve the caller's encrypted keypair bundle during the handshake window.""" + bundle_store = self._ctx.get("bundle_store") + if not bundle_store: + self._send({"type": MNP.KEYPAIR_BUNDLE_RESP, "v": MNP_VERSION, "found": False}) + return + + user_id = getattr(self, "_pending_sub", "") + if not user_id: + self._send({"type": "error", "detail": "No pending handshake"}) + return + + bundle_enc = await bundle_store.fetch_keypair(user_id) + if bundle_enc: + self._send({ + "type": MNP.KEYPAIR_BUNDLE_RESP, + "v": MNP_VERSION, + "found": True, + "bundle_enc": bundle_enc, + }) + else: + self._send({"type": MNP.KEYPAIR_BUNDLE_RESP, "v": MNP_VERSION, "found": False}) + + async def _do_keypair_bundle_store(self, msg: dict) -> None: + """Store an encrypted keypair bundle (user backs up their own keys on node).""" + bundle_store = self._ctx.get("bundle_store") + if not bundle_store: + self._send({"type": "error", "detail": "Bundle store not available"}) + return + + bundle_enc = msg.get("bundle_enc", "") + if not bundle_enc: + self._send({"type": "error", "detail": "Missing bundle_enc"}) + return + + await bundle_store.store_keypair(self._user_id, bundle_enc) + log.info("Keypair bundle stored for user=%s", self._user_id[:8]) + self._audit("keypair_bundle_store") + self._send({ + "type": "ack", "v": MNP_VERSION, + "detail": "keypair_bundle_stored", + }) + def _audit_auth_failed(self, group_id: str, reason: str) -> None: audit = self._ctx.get("audit_store") if audit: @@ -275,6 +515,7 @@ class WebRTCPeerSession: { "id": e.id, "name": e.name, "path": e.path, "size": e.size, "type": e.type, "added_at": e.added_at, + "uploader_id": e.uploader_id, } for e in idx.entries ] @@ -286,18 +527,6 @@ class WebRTCPeerSession: "entries": entries, }) - def _do_gek_request(self) -> None: - ctx = self._group_ctx() - gek = ctx.get("gek") - if not gek: - self._send({"type": "error", "detail": "No GEK available"}) - return - self._send({ - "type": MNP.GEK_RESPONSE, - "v": MNP_VERSION, - "gek_b64": base64.b64encode(gek).decode(), - }) - def _do_file_request(self, msg: dict) -> None: ctx = self._group_ctx() file_id = msg["file_id"] @@ -378,7 +607,7 @@ class WebRTCPeerSession: if chat_store: raw = payload.encode() if isinstance(payload, str) else payload asyncio.ensure_future(chat_store.save_message( - sender_id=msg.get("sender_id", self._user_id), + sender_id=self._user_id, iteration=msg.get("iteration", 0), payload=raw, thread_id=msg.get("thread_id"), @@ -389,7 +618,7 @@ class WebRTCPeerSession: broadcast = { "type": MNP.CHAT_MESSAGE, "v": MNP_VERSION, - "sender_id": msg.get("sender_id", self._user_id), + "sender_id": self._user_id, "sender_name": sender_name, "payload": payload, "thread_id": msg.get("thread_id"), @@ -493,6 +722,18 @@ class WebRTCPeerSession: 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) + + def _register_uploader(self, ctx: dict, filename: str) -> None: + """Tag the index entry with the uploader's user_id after upload completes.""" + idx = ctx.get("index") + if not idx: + return + for entry in idx.entries: + if entry.name == filename and entry.path == "": + entry.uploader_id = self._user_id + entry.uploader_pk = self._pk_user + return def _do_file_delete(self, msg: dict) -> None: ctx = self._group_ctx() @@ -501,16 +742,76 @@ class WebRTCPeerSession: self._send({"type": "error", "detail": "Missing file_id"}) return - node_user_id = self._ctx.get("node_user_id") - if node_user_id and self._user_id != node_user_id: - self._send({"type": "error", "detail": "Only node admin can delete files"}) + entry = ctx["index"].get_entry(file_id) + if not entry: + 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: + self._send({"type": "error", "detail": "No authorized key for deletion"}) + return + + challenge = os.urandom(32) + self._admin_challenges[file_id] = challenge + self._send({ + "type": MNP.ADMIN_CHALLENGE, + "v": MNP_VERSION, + "challenge": base64.b64encode(challenge).decode(), + "file_id": file_id, + }) + + def _do_admin_response(self, msg: dict) -> None: + file_id = msg.get("file_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"}) + return + + try: + sig_bytes = base64.b64decode(sig_b64) + except Exception: + self._send({"type": "error", "detail": "Invalid signature encoding"}) return + 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: + try: + uploader_key = Ed25519PublicKey.from_public_bytes( + base64.b64decode(entry.uploader_pk)) + uploader_key.verify(sig_bytes, challenge) + verified = True + except Exception: + pass + + if not verified: + 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) + + 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(): file_path.unlink() diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py index 5e77ed8..b4885af 100644 --- a/packages/meshbay-node/src/meshbay_node/ui/app.py +++ b/packages/meshbay-node/src/meshbay_node/ui/app.py @@ -12,15 +12,17 @@ Served only on 127.0.0.1 — not exposed to the network. No authentication required (localhost only). """ +import base64 import json import logging import time from pathlib import Path from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Query -from fastapi.responses import HTMLResponse +from fastapi.responses import HTMLResponse, JSONResponse from meshbay_node import __version__ +from meshbay_common.crypto import generate_gek, wrap_gek_aes log = logging.getLogger(__name__) @@ -52,6 +54,7 @@ def create_ui_app(state: dict) -> FastAPI: "group_count": len(groups_ctx), "total_files": total_files, "webrtc_peers": webrtc.active_peers if webrtc else 0, + "pk_node_ed25519": state.get("pk_node_ed25519", ""), } @app.get("/api/groups") @@ -167,6 +170,95 @@ def create_ui_app(state: dict) -> FastAPI: ], } + # ── 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.""" + groups_ctx = state.get("groups_ctx", {}) + 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) + + bundle_store = state.get("bundle_store") + 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() + + 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) + + 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") + 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: + try: + node_bundle = wrap_gek_aes(gek, pk_x_node_raw) + await bundle_store.store( + group_id, f"_node_{node_user_id}", + node_bundle["pk_eph_b64"], node_bundle["nonce_b64"], + node_bundle["wrapped_b64"], + ) + log.info("GEK also wrapped for node keystore (daemon reload)") + except Exception as 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)) + + webrtc = state.get("webrtc") + if webrtc and "groups" in webrtc._ctx and group_id in webrtc._ctx["groups"]: + webrtc._ctx["groups"][group_id]["gek"] = gek + + return { + "status": "ok", + "group_id": group_id, + "wrapped_count": wrapped_count, + "total_members": len(members), + "errors": errors, + } + # ── Chat endpoints ─────────────────────────────────────────────────────── _chat_subscribers: list[WebSocket] = [] @@ -246,7 +338,10 @@ def _render_page(state: dict) -> str: webrtc = state.get("webrtc") total_files = sum(idx.count for idx in indexes.values()) peer_count = webrtc.active_peers if webrtc else 0 - status_color = {"running": "#22c55e", "error": "#ef4444"}.get(status, "#f59e0b") + status_color = { + "running": "#22c55e", "error": "#ef4444", + "waiting_for_node_key": "#f97316", + }.get(status, "#f59e0b") # Groups section groups_html = "" @@ -269,13 +364,34 @@ def _render_page(state: dict) -> str: f"<td>{_fmt_size(e.size)}</td><td>{e.path or '/'}</td></tr>" ) + has_gek = bool(ctx.get("gek")) + gek_badge = ( + '<span class="badge" style="background:#22c55e">GEK active</span>' + if has_gek + else '<span class="badge" style="background:#ef4444">No GEK</span>' + ) + gek_label = "Re-wrap GEK for all members" if has_gek else "Initialize GEK" + gek_color = "#3b82f6" if has_gek else "#22c55e" + gek_action = f""" + <div style="margin:10px 0"> + <button onclick="initGEK('{gid}')" + id="gek-btn-{gid[:8]}" + style="padding:8px 16px;background:{gek_color};color:#fff;border:none; + border-radius:6px;cursor:pointer;font-size:0.85em"> + {gek_label} + </button> + <span id="gek-status-{gid[:8]}" class="muted" style="margin-left:8px"></span> + </div>""" + groups_html += f""" <div class="card"> <h3>{name} <span class="badge" style="background:#6366f1">{vis}</span> + {gek_badge} </h3> <p><b>Directory:</b> <code>{shared}</code></p> <p><b>Files:</b> {fcount} — <b>Total:</b> {_fmt_size(total_size)}</p> + {gek_action} <p class="muted">ID: {gid}</p> <details><summary>File list</summary> <table> @@ -382,6 +498,21 @@ def _render_page(state: dict) -> str: <p><b>Node ID:</b> <code>{state.get("endpoint_hint") or "—"}</code></p> </div> + <h2>Link Node to Hub Account</h2> + <div class="card"> + <p>To connect to your group from a browser, link this node to your hub account. + Copy the key below and paste it in <b>Settings > Link Node</b> on the hub.</p> + <div style="margin:12px 0;display:flex;align-items:center;gap:8px"> + <code id="nodeKey" style="flex:1;padding:8px;word-break:break-all;background:var(--border); + border-radius:4px;font-size:0.9em;user-select:all">{state.get("pk_node_ed25519", "—")}</code> + <button onclick="navigator.clipboard.writeText(document.getElementById('nodeKey').textContent).then(()=>{{this.textContent='Copied!';setTimeout(()=>this.textContent='Copy',2000)}})" + style="padding:8px 16px;background:var(--accent);color:#fff;border:none;border-radius:6px; + cursor:pointer;font-size:0.85em;white-space:nowrap">Copy</button> + </div> + <p class="muted">This is the node's Ed25519 public key. It's safe to share — it identifies + this node but cannot be used to impersonate it.</p> + </div> + <div class="footer"> MeshBay Node v{__version__} — localhost only — <a href="/api/status">status</a> · @@ -392,7 +523,33 @@ def _render_page(state: dict) -> str: — auto-refresh 10s </div> </div> -<script>setTimeout(()=>location.reload(), 10000);</script> +<script> +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 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.style.color = '#22c55e'; + setTimeout(() => location.reload(), 2000); + }} else {{ + if (status) status.textContent = data.error || 'Failed'; + if (status) status.style.color = '#ef4444'; + if (btn) btn.disabled = false; + }} + }} catch (e) {{ + if (status) status.textContent = 'Error: ' + e.message; + if (status) status.style.color = '#ef4444'; + if (btn) btn.disabled = false; + }} +}} +setTimeout(()=>location.reload(), 10000); +</script> </body> </html>""" diff --git a/packages/meshbay-node/tests/test_daemon.py b/packages/meshbay-node/tests/test_daemon.py index faf12e3..1c5a07e 100644 --- a/packages/meshbay-node/tests/test_daemon.py +++ b/packages/meshbay-node/tests/test_daemon.py @@ -7,11 +7,13 @@ Hub interaction is mocked. """ import asyncio +import base64 import os 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 unittest.mock import AsyncMock, MagicMock, patch from meshbay_common.crypto import generate_gek @@ -20,6 +22,20 @@ 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() + pk_x_raw = sk_x.public_key().public_bytes( + serialization.Encoding.Raw, serialization.PublicFormat.Raw) + + mock_keys = MagicMock() + mock_keys.sk_ed25519 = sk_ed + mock_keys.pk_ed25519_b64 = "test" + mock_keys.sk_x25519 = sk_x + mock_keys.pk_x25519_b64 = base64.b64encode(pk_x_raw).decode() + return mock_keys + + @pytest.fixture def sk_hub(): return Ed25519PrivateKey.generate() @@ -48,7 +64,7 @@ def shared_dir(tmp_path): @pytest.fixture def node_config(tmp_path, shared_dir): return Config( - hub=HubConfig(url="http://localhost:9999", username="testuser", password="testpass"), + hub=HubConfig(url="http://localhost:9999", username="testuser"), node=NodeConfig(port=29000, quic_port=29010, http_port=29001, ui_port=28000), groups=[GroupConfig( id="g" * 32, @@ -70,10 +86,7 @@ async def test_daemon_creates_chat_store(tmp_path, node_config, gek, hub_pk_pem) daemon = NodeDaemon(node_config) sk_node = Ed25519PrivateKey.generate() - mock_keys = MagicMock() - mock_keys.sk_ed25519 = sk_node - mock_keys.pk_ed25519_b64 = "test" - mock_keys.pk_x25519_b64 = "test" + mock_keys = _mock_keystore_keys(sk_node) mock_session = MagicMock() mock_session.node_id = "node123" @@ -85,7 +98,6 @@ async def test_daemon_creates_chat_store(tmp_path, node_config, gek, hub_pk_pem) hub_instance = AsyncMock() hub_instance.startup = AsyncMock(return_value=mock_session) - hub_instance.fetch_gek = AsyncMock(return_value=gek) hub_instance.maintain_ws = AsyncMock() hub_instance.send_ws = AsyncMock() hub_instance._ws = None @@ -139,7 +151,7 @@ async def test_daemon_creates_chat_store(tmp_path, node_config, gek, hub_pk_pem) async def test_daemon_no_groups_exits(tmp_path): """Daemon with no valid groups exits cleanly.""" config = Config( - hub=HubConfig(url="http://localhost:9999", username="testuser", password="testpass"), + hub=HubConfig(url="http://localhost:9999", username="testuser"), node=NodeConfig(), groups=[GroupConfig(id="", name="empty", shared_dir="")], keystore=KeystoreConfig(path=tmp_path / "keystore.enc"), @@ -148,17 +160,22 @@ async def test_daemon_no_groups_exits(tmp_path): daemon = NodeDaemon(config) sk_node = Ed25519PrivateKey.generate() - mock_keys = MagicMock() - mock_keys.sk_ed25519 = sk_node - mock_keys.pk_ed25519_b64 = "test" + mock_keys = _mock_keystore_keys(sk_node) mock_session = MagicMock() mock_session.node_id = "node123" mock_session.user_id = "user123" mock_session.hub_pk_pem = b"pem" + mock_server = AsyncMock() + mock_server.serve = AsyncMock() + with patch("meshbay_node.daemon.load_or_create_keystore", return_value=mock_keys), \ - patch("meshbay_node.daemon.HubClient") as MockHub: + patch("meshbay_node.daemon.HubClient") as MockHub, \ + patch("meshbay_node.daemon.uvicorn") as mock_uvicorn: + + mock_uvicorn.Config = MagicMock() + mock_uvicorn.Server = MagicMock(return_value=mock_server) hub_instance = AsyncMock() hub_instance.startup = AsyncMock(return_value=mock_session) @@ -169,7 +186,6 @@ async def test_daemon_no_groups_exits(tmp_path): await daemon.run() - assert daemon._state["status"] == "starting" assert len(daemon._chat_stores) == 0 @@ -177,7 +193,7 @@ async def test_daemon_no_groups_exits(tmp_path): 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", password="testpass"), + hub=HubConfig(url="http://localhost:9999", username="testuser"), node=NodeConfig(port=29000, quic_port=29010, http_port=29001, ui_port=28000), groups=[GroupConfig( id="a" * 32, @@ -228,7 +244,7 @@ async def test_daemon_index_change_skips_other_group_peers( ): """Index change only pushes to peers in the same group.""" config = Config( - hub=HubConfig(url="http://localhost:9999", username="testuser", password="testpass"), + hub=HubConfig(url="http://localhost:9999", username="testuser"), node=NodeConfig(), groups=[], keystore=KeystoreConfig(path=tmp_path / "keystore.enc"), diff --git a/packages/meshbay-node/tests/test_hub_client.py b/packages/meshbay-node/tests/test_hub_client.py index fe8a2af..2975ee8 100644 --- a/packages/meshbay-node/tests/test_hub_client.py +++ b/packages/meshbay-node/tests/test_hub_client.py @@ -2,7 +2,6 @@ Tests for meshbay_node.hub_client — uses httpx.MockTransport to avoid network. """ -import base64 import json import os import time @@ -16,7 +15,7 @@ from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey from cryptography.hazmat.primitives import serialization -from meshbay_common.crypto import generate_gek, pk_to_b64, wrap_gek + from meshbay_node.hub_client import HubClient, HubConfig, HubSession from meshbay_node.keystore import NodeKeys @@ -51,16 +50,15 @@ def hub_config(tmp_path): return HubConfig( hub_url="http://fake-hub", username="testuser", - password="testpass99", cache_dir=tmp_path, ) -def make_token(sk_pem, user_id, pk_user_b64, hub_id="fake-hub", ttl=3600): +def make_node_token(sk_pem, user_id, pk_user_b64, hub_id="fake-hub", ttl=3600): now = int(time.time()) return jwt.encode({ "iss": hub_id, "sub": user_id, "pk_user": pk_user_b64, - "hub_id": hub_id, "jti": "test-jti", + "hub_id": hub_id, "jti": "test-jti", "scope": "node", "iat": now, "exp": now + ttl, }, sk_pem, algorithm="EdDSA") @@ -71,14 +69,14 @@ def make_token(sk_pem, user_id, pk_user_b64, hub_id="fake-hub", ttl=3600): async def test_login_verifies_jwt_offline(hub_keys, node_keys, hub_config): sk_hub, sk_hub_pem, pk_hub_pem = hub_keys user_id = "user-uuid-001" - token = make_token(sk_hub_pem, user_id, node_keys.pk_ed25519_b64) + token = make_node_token(sk_hub_pem, user_id, node_keys.pk_ed25519_b64) def handler(request): if request.url.path == "/v1/hub/pubkey": return httpx.Response(200, json={"pk_hub_pem": pk_hub_pem.decode()}) - if request.url.path == "/v1/users/login": + if request.url.path == "/v1/nodes/auth": return httpx.Response(200, json={ - "access_token": token, "refresh_token": "rt-abc", "expires_in": 3600}) + "access_token": token, "token_type": "bearer", "expires_in": 3600}) return httpx.Response(404) transport = httpx.MockTransport(handler) @@ -96,18 +94,18 @@ async def test_login_verifies_jwt_offline(hub_keys, node_keys, hub_config): @pytest.mark.asyncio async def test_login_rejects_missing_jti(hub_keys, node_keys, hub_config): sk_hub, sk_hub_pem, pk_hub_pem = hub_keys - # Token without jti bad_token = jwt.encode({ "iss": "fake-hub", "sub": "uid", "pk_user": node_keys.pk_ed25519_b64, - "hub_id": "fake-hub", "iat": int(time.time()), "exp": int(time.time()) + 3600, + "hub_id": "fake-hub", "scope": "node", + "iat": int(time.time()), "exp": int(time.time()) + 3600, }, sk_hub_pem, algorithm="EdDSA") def handler(request): if request.url.path == "/v1/hub/pubkey": return httpx.Response(200, json={"pk_hub_pem": pk_hub_pem.decode()}) - if request.url.path == "/v1/users/login": + if request.url.path == "/v1/nodes/auth": return httpx.Response(200, json={ - "access_token": bad_token, "refresh_token": "rt", "expires_in": 3600}) + "access_token": bad_token, "token_type": "bearer", "expires_in": 3600}) return httpx.Response(404) transport = httpx.MockTransport(handler) @@ -119,33 +117,16 @@ async def test_login_rejects_missing_jti(hub_keys, node_keys, hub_config): @pytest.mark.asyncio -async def test_register_idempotent(hub_keys, node_keys, hub_config): - def handler(request): - if request.url.path == "/v1/users/register": - return httpx.Response(409, json={"detail": "Username already taken"}) - return httpx.Response(404) - - transport = httpx.MockTransport(handler) - client = HubClient(hub_config, node_keys) - client._http = httpx.AsyncClient(transport=transport, base_url="http://fake-hub") - - # Should not raise on 409 - result = await client.register() - assert result == "" - - -@pytest.mark.asyncio async def test_token_needs_refresh(hub_keys, node_keys, hub_config): sk_hub, sk_hub_pem, pk_hub_pem = hub_keys - # Token expiring in 60s (< TOKEN_REFRESH_MARGIN of 300s) - short_token = make_token(sk_hub_pem, "uid", node_keys.pk_ed25519_b64, ttl=60) + short_token = make_node_token(sk_hub_pem, "uid", node_keys.pk_ed25519_b64, ttl=60) def handler(request): if request.url.path == "/v1/hub/pubkey": return httpx.Response(200, json={"pk_hub_pem": pk_hub_pem.decode()}) - if request.url.path == "/v1/users/login": + if request.url.path == "/v1/nodes/auth": return httpx.Response(200, json={ - "access_token": short_token, "refresh_token": "rt", "expires_in": 60}) + "access_token": short_token, "token_type": "bearer", "expires_in": 60}) return httpx.Response(404) transport = httpx.MockTransport(handler) @@ -157,37 +138,6 @@ async def test_token_needs_refresh(hub_keys, node_keys, hub_config): @pytest.mark.asyncio -async def test_fetch_gek(hub_keys, node_keys, hub_config): - """Admin wraps GEK for this node; client fetches and unwraps.""" - sk_hub, sk_hub_pem, pk_hub_pem = hub_keys - gek = generate_gek() - - # Simulate admin wrapping GEK for this node - pk_x_raw = base64.b64decode(node_keys.pk_x25519_b64) - bundle = wrap_gek(gek, pk_x_raw) - - token = make_token(sk_hub_pem, "uid", node_keys.pk_ed25519_b64) - - def handler(request): - if request.url.path == "/v1/hub/pubkey": - return httpx.Response(200, json={"pk_hub_pem": pk_hub_pem.decode()}) - if request.url.path == "/v1/users/login": - return httpx.Response(200, json={ - "access_token": token, "refresh_token": "rt", "expires_in": 3600}) - if "/v1/groups/" in request.url.path and request.url.path.endswith("/gek"): - return httpx.Response(200, json=bundle) - return httpx.Response(404) - - transport = httpx.MockTransport(handler) - client = HubClient(hub_config, node_keys) - client._http = httpx.AsyncClient(transport=transport, base_url="http://fake-hub") - - await client.login() - recovered = await client.fetch_gek("group-abc") - assert recovered == gek - - -@pytest.mark.asyncio async def test_hub_pk_cached(hub_keys, node_keys, hub_config, tmp_path): _, _, pk_hub_pem = hub_keys call_count = {"n": 0} diff --git a/packages/meshbay-node/tests/test_webrtc_transport.py b/packages/meshbay-node/tests/test_webrtc_transport.py index b6664e8..693a68b 100644 --- a/packages/meshbay-node/tests/test_webrtc_transport.py +++ b/packages/meshbay-node/tests/test_webrtc_transport.py @@ -9,6 +9,8 @@ Uses local loopback (no STUN/ICE needed for localhost). import asyncio import base64 +import hashlib +import hmac import os import struct import time @@ -24,9 +26,13 @@ from meshbay_common import MNP_VERSION from meshbay_common.crypto import ( generate_gek, pk_to_b64, + wrap_gek, + wrap_gek_aes, + unwrap_gek, ) from meshbay_common.webcrypto import chunk_key_aes, decrypt_chunk_aes from meshbay_common.protocol import MNP +from meshbay_node.bundle_store import BundleStore from meshbay_node.indexer import DirectoryIndexer from meshbay_node.transport.webrtc_server import WebRTCTransport @@ -42,6 +48,11 @@ def sk_hub(): @pytest.fixture +def sk_user(): + return Ed25519PrivateKey.generate() + + +@pytest.fixture def gek(): return generate_gek() @@ -60,7 +71,7 @@ def _hub_pk_pem(sk_hub): serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo) -def _make_jwt(sk_hub, groups=None): +def _make_jwt(sk_hub, groups=None, pk_user="test"): sk_pem = sk_hub.private_bytes( serialization.Encoding.PEM, serialization.PrivateFormat.PKCS8, @@ -69,7 +80,7 @@ def _make_jwt(sk_hub, groups=None): now = int(time.time()) return jwt.encode({ "iss": "test-hub", "sub": "user-001", - "pk_user": "test", "hub_id": "test-hub", + "pk_user": pk_user, "hub_id": "test-hub", "jti": "test-jti-webrtc", "iat": now, "exp": now + 3600, "groups": groups or [], }, sk_pem, algorithm="EdDSA") @@ -85,6 +96,108 @@ def _unpack(raw: bytes) -> dict: return msgpack.unpackb(raw[4:4 + length], raw=False) +def _extract_dtls_fp(sdp: str) -> bytes: + for line in sdp.splitlines(): + if line.startswith("a=fingerprint:sha-256 "): + return bytes.fromhex(line.split(" ", 1)[1].replace(":", "")) + 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) + channel.send(_pack({ + "type": MNP.HANDSHAKE, + "v": MNP_VERSION, + "token": token, + })) + 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) + 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).""" + pc = RTCPeerConnection() + q = asyncio.Queue() + buf = bytearray() + ch = pc.createDataChannel("mnp") + ready = asyncio.Event() + + @ch.on("open") + def on_open(): + ready.set() + + @ch.on("message") + def on_msg(message): + if isinstance(message, str): + message = message.encode() + buf.extend(message) + while len(buf) >= 4: + length = struct.unpack(">I", buf[:4])[0] + if len(buf) < 4 + length: + break + msg_bytes = bytes(buf[4:4 + length]) + del buf[:4 + length] + q.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) + + offer = await pc.createOffer() + await pc.setLocalDescription(offer) + 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) + + pk_user = "test" + if sk_user: + pk_user = base64.b64encode( + sk_user.public_key().public_bytes( + 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") + + 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) + assert msg["type"] == MNP.HANDSHAKE_ACK + return pc, ch, q + + @pytest.mark.asyncio async def test_webrtc_datachannel_handshake(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: browser sends MNP handshake, node responds with handshake_ack.""" @@ -120,15 +233,8 @@ async def test_webrtc_datachannel_handshake(sk_node, sk_hub, gek, shared_dir): await asyncio.sleep(0.5) - token = _make_jwt(sk_hub) - channel.send(_pack({ - "type": MNP.HANDSHAKE, - "v": MNP_VERSION, - "token": token, - })) - - msg = await asyncio.wait_for(received.get(), timeout=5.0) - assert msg["type"] == MNP.HANDSHAKE_ACK + msg = await _handshake_with_gek_proof(channel, received, sk_hub, gek, + browser_pc=browser_pc) assert msg["v"] == MNP_VERSION assert "node_pk" in msg @@ -153,15 +259,11 @@ async def test_webrtc_datachannel_file_transfer(sk_node, sk_hub, gek, shared_dir received = asyncio.Queue() channel = browser_pc.createDataChannel("mnp") + channel_ready = asyncio.Event() @channel.on("open") def on_open(): - token = _make_jwt(sk_hub) - channel.send(_pack({ - "type": MNP.HANDSHAKE, - "v": MNP_VERSION, - "token": token, - })) + channel_ready.set() buf = bytearray() @@ -186,8 +288,11 @@ async def test_webrtc_datachannel_file_transfer(sk_node, sk_hub, gek, shared_dir await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) - # 1) Handshake ack - ack = await asyncio.wait_for(received.get(), timeout=5.0) + await asyncio.wait_for(channel_ready.wait(), timeout=5.0) + + # 1) Handshake with GEK proof + ack = await _handshake_with_gek_proof(channel, received, sk_hub, gek, + browser_pc=browser_pc) assert ack["type"] == MNP.HANDSHAKE_ACK # 2) Request index @@ -220,13 +325,6 @@ async def test_webrtc_datachannel_file_transfer(sk_node, sk_hub, gek, shared_dir original = (shared_dir / "test.bin").read_bytes() assert plaintext == original - # 5) Request GEK over DataChannel - channel.send(_pack({"type": MNP.GEK_REQUEST, "v": MNP_VERSION})) - gek_msg = await asyncio.wait_for(received.get(), timeout=5.0) - assert gek_msg["type"] == MNP.GEK_RESPONSE - received_gek = base64.b64decode(gek_msg["gek_b64"]) - assert received_gek == gek - await browser_pc.close() await transport.close_all() @@ -342,44 +440,8 @@ async def test_webrtc_chat_send_and_history(sk_node, sk_hub, gek, shared_dir, tm ) transport._ctx["chat_store"] = chat_store - browser_pc = RTCPeerConnection() - received = asyncio.Queue() - buf = bytearray() - - channel = browser_pc.createDataChannel("mnp") - - @channel.on("open") - def on_open(): - token = _make_jwt(sk_hub) - channel.send(_pack({ - "type": MNP.HANDSHAKE, - "v": MNP_VERSION, - "token": token, - })) - - @channel.on("message") - def on_msg(message): - if isinstance(message, str): - message = message.encode() - buf.extend(message) - while len(buf) >= 4: - length = struct.unpack(">I", buf[:4])[0] - if len(buf) < 4 + length: - break - msg_bytes = bytes(buf[4:4 + length]) - del buf[:4 + length] - received.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) - - offer = await browser_pc.createOffer() - await browser_pc.setLocalDescription(offer) - - answer_sdp, _ = await transport.handle_offer( - browser_pc.localDescription.sdp, "peer-chat") - await browser_pc.setRemoteDescription( - RTCSessionDescription(sdp=answer_sdp, type="answer")) - - ack = await asyncio.wait_for(received.get(), timeout=5.0) - assert ack["type"] == MNP.HANDSHAKE_ACK + browser_pc, channel, received = await _setup_peer( + transport, sk_hub, gek, "peer-chat") channel.send(_pack({ "type": MNP.CHAT_MESSAGE, @@ -421,41 +483,8 @@ async def test_webrtc_chat_history_no_store(sk_node, sk_hub, gek, shared_dir): stun_servers=[], ) - browser_pc = RTCPeerConnection() - received = asyncio.Queue() - buf = bytearray() - - channel = browser_pc.createDataChannel("mnp") - - @channel.on("open") - def on_open(): - channel.send(_pack({ - "type": MNP.HANDSHAKE, "v": MNP_VERSION, - "token": _make_jwt(sk_hub), - })) - - @channel.on("message") - def on_msg(message): - if isinstance(message, str): - message = message.encode() - buf.extend(message) - while len(buf) >= 4: - length = struct.unpack(">I", buf[:4])[0] - if len(buf) < 4 + length: - break - msg_bytes = bytes(buf[4:4 + length]) - del buf[:4 + length] - received.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) - - offer = await browser_pc.createOffer() - await browser_pc.setLocalDescription(offer) - answer_sdp, _ = await transport.handle_offer( - browser_pc.localDescription.sdp, "peer-no-store") - await browser_pc.setRemoteDescription( - RTCSessionDescription(sdp=answer_sdp, type="answer")) - - ack = await asyncio.wait_for(received.get(), timeout=5.0) - assert ack["type"] == MNP.HANDSHAKE_ACK + browser_pc, channel, received = await _setup_peer( + transport, sk_hub, gek, "peer-no-store") channel.send(_pack({ "type": MNP.CHAT_HISTORY, "v": MNP_VERSION, "since": 0, "limit": 50, @@ -487,57 +516,8 @@ async def test_webrtc_chat_broadcast(sk_node, sk_hub, gek, shared_dir, tmp_path) ) transport._ctx["chat_store"] = chat_store - async def _connect_peer(peer_id, jwt_sub, groups=None): - pc = RTCPeerConnection() - q = asyncio.Queue() - b = bytearray() - ch = pc.createDataChannel("mnp") - - 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": "test", "hub_id": "test-hub", - "jti": f"jti-{peer_id}", "iat": now, "exp": now + 3600, - "groups": groups or [], - }, sk_h_pem, algorithm="EdDSA") - - @ch.on("open") - def on_open(): - ch.send(_pack({ - "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, - })) - - @ch.on("message") - def on_msg(message): - if isinstance(message, str): - message = message.encode() - b.extend(message) - while len(b) >= 4: - length = struct.unpack(">I", b[:4])[0] - if len(b) < 4 + length: - break - msg_bytes = bytes(b[4:4 + length]) - del b[:4 + length] - q.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) - - offer = await pc.createOffer() - await pc.setLocalDescription(offer) - answer_sdp, _ = await transport.handle_offer( - pc.localDescription.sdp, peer_id) - await pc.setRemoteDescription( - RTCSessionDescription(sdp=answer_sdp, type="answer")) - - ack = await asyncio.wait_for(q.get(), timeout=5.0) - assert ack["type"] == MNP.HANDSHAKE_ACK - return pc, ch, q - - pc_a, ch_a, q_a = await _connect_peer("peer-A", "user-A") - pc_b, ch_b, q_b = await _connect_peer("peer-B", "user-B") + pc_a, ch_a, q_a = await _setup_peer(transport, sk_hub, gek, "peer-A", "user-A") + pc_b, ch_b, q_b = await _setup_peer(transport, sk_hub, gek, "peer-B", "user-B") ch_a.send(_pack({ "type": MNP.CHAT_MESSAGE, "v": MNP_VERSION, "payload": "hi from A", @@ -617,18 +597,509 @@ async def test_webrtc_peer_cleanup_on_close(sk_node, sk_hub, gek, shared_dir): stun_servers=[], ) + browser_pc, channel, received = await _setup_peer( + transport, sk_hub, gek, "peer-cleanup") + + assert "user-001" in transport._ctx["_peers"] + assert transport.active_peers == 1 + + await transport.close_peer("peer-cleanup") + + assert "user-001" not in transport._ctx["_peers"] + assert transport.active_peers == 0 + + await browser_pc.close() + + +@pytest.mark.asyncio +async def test_webrtc_stream_segment_missing_file(sk_node, sk_hub, gek, shared_dir): + """WebRTC DataChannel: stream_segment for non-existent file returns error.""" + 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() + + transport = WebRTCTransport( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, + shared_root=shared_dir, index=indexer.index, + stun_servers=[], + ) + + browser_pc, channel, received = await _setup_peer( + transport, sk_hub, gek, "peer-stream") + + channel.send(_pack({ + "type": MNP.STREAM_SEGMENT, "v": MNP_VERSION, + "file_id": "nonexistent-file-id", + "segment_index": 0, "segment_duration": 4, + })) + + msg = await asyncio.wait_for(received.get(), timeout=5.0) + assert msg["type"] == "error" + assert "not found" in msg["detail"].lower() + + await browser_pc.close() + await transport.close_all() + + +@pytest.mark.asyncio +async def test_webrtc_wrong_gek_proof_rejected(sk_node, sk_hub, gek, shared_dir): + """WebRTC DataChannel: wrong GEK proof is rejected — hub admin can't fake membership.""" + 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() + + transport = WebRTCTransport( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, + shared_root=shared_dir, index=indexer.index, + stun_servers=[], + ) + browser_pc = RTCPeerConnection() received = asyncio.Queue() - buf = bytearray() + channel = browser_pc.createDataChannel("mnp") + + @channel.on("message") + def on_msg(message): + if isinstance(message, str): + message = message.encode() + received.put_nowait(_unpack(message)) + + offer = await browser_pc.createOffer() + await browser_pc.setLocalDescription(offer) + answer_sdp, _ = await transport.handle_offer( + browser_pc.localDescription.sdp, "peer-fake") + await browser_pc.setRemoteDescription( + RTCSessionDescription(sdp=answer_sdp, type="answer")) + + await asyncio.sleep(0.5) + + token = _make_jwt(sk_hub) + channel.send(_pack({ + "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, + })) + + challenge = await asyncio.wait_for(received.get(), timeout=5.0) + assert challenge["type"] == MNP.HANDSHAKE_CHALLENGE + + fake_gek = os.urandom(32) + nonce = base64.b64decode(challenge["nonce"]) + offer_fp = _extract_dtls_fp(browser_pc.localDescription.sdp) + answer_fp = _extract_dtls_fp(browser_pc.remoteDescription.sdp) + bad_proof = hmac.new(fake_gek, nonce + offer_fp + answer_fp, hashlib.sha256).digest() + channel.send(_pack({ + "type": MNP.HANDSHAKE_RESPONSE, "v": MNP_VERSION, + "proof": base64.b64encode(bad_proof).decode(), + })) + + msg = await asyncio.wait_for(received.get(), timeout=5.0) + assert msg["type"] == "error" + assert "GEK proof failed" in msg["detail"] + + await browser_pc.close() + await transport.close_all() + + +@pytest.mark.asyncio +async def test_webrtc_dtls_channel_binding_detects_mitm(sk_node, sk_hub, gek, shared_dir): + """WebRTC: DTLS channel binding detects fingerprint substitution (simulated MitM).""" + 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() + + transport = WebRTCTransport( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, + shared_root=shared_dir, index=indexer.index, + stun_servers=[], + ) + + browser_pc = RTCPeerConnection() + received = asyncio.Queue() + channel = browser_pc.createDataChannel("mnp") + + @channel.on("message") + def on_msg(message): + if isinstance(message, str): + message = message.encode() + received.put_nowait(_unpack(message)) + + offer = await browser_pc.createOffer() + await browser_pc.setLocalDescription(offer) + answer_sdp, _ = await transport.handle_offer( + browser_pc.localDescription.sdp, "peer-mitm") + await browser_pc.setRemoteDescription( + RTCSessionDescription(sdp=answer_sdp, type="answer")) + + await asyncio.sleep(0.5) + + token = _make_jwt(sk_hub) + channel.send(_pack({ + "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, + })) + + challenge = await asyncio.wait_for(received.get(), timeout=5.0) + assert challenge["type"] == MNP.HANDSHAKE_CHALLENGE + + nonce = base64.b64decode(challenge["nonce"]) + # Correct GEK but fake fingerprints — simulates MitM substituting DTLS certs + fake_fp = os.urandom(32) + proof = hmac.new(gek, nonce + fake_fp + fake_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) + assert msg["type"] == "error" + assert "GEK proof failed" in msg["detail"] + + await browser_pc.close() + await transport.close_all() + + +@pytest.mark.asyncio +async def test_webrtc_admin_challenge_response(sk_node, sk_hub, gek, shared_dir): + """WebRTC DataChannel: admin file delete requires Ed25519 challenge-response.""" + 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() + + sk_admin = Ed25519PrivateKey.generate() + + 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["admin_pk_ed25519"] = sk_admin.public_key() + transport._ctx["node_user_id"] = "user-001" + + browser_pc, channel, received = await _setup_peer( + transport, sk_hub, gek, "peer-admin") + + entry = indexer.index.entries[0] + channel.send(_pack({ + "type": MNP.FILE_DELETE, "v": MNP_VERSION, "file_id": entry.id, + })) + + 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 + + challenge = base64.b64decode(challenge_msg["challenge"]) + signature = sk_admin.sign(challenge) + channel.send(_pack({ + "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION, + "file_id": entry.id, + "signature": base64.b64encode(signature).decode(), + })) + + ack = await asyncio.wait_for(received.get(), timeout=5.0) + assert ack["type"] == MNP.FILE_DELETE_ACK + assert ack["file_id"] == entry.id + + assert indexer.index.get_entry(entry.id) is None + + await browser_pc.close() + await transport.close_all() + + +@pytest.mark.asyncio +async def test_webrtc_admin_bad_signature_rejected(sk_node, sk_hub, gek, shared_dir): + """WebRTC DataChannel: wrong Ed25519 signature is rejected — hub can't fake admin.""" + 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() + + sk_admin = Ed25519PrivateKey.generate() + sk_attacker = Ed25519PrivateKey.generate() + + 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["admin_pk_ed25519"] = sk_admin.public_key() + transport._ctx["node_user_id"] = "user-001" + + browser_pc, channel, received = await _setup_peer( + transport, sk_hub, gek, "peer-attacker") + + entry = indexer.index.entries[0] + channel.send(_pack({ + "type": MNP.FILE_DELETE, "v": MNP_VERSION, "file_id": entry.id, + })) + + 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) + channel.send(_pack({ + "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION, + "file_id": entry.id, + "signature": base64.b64encode(bad_sig).decode(), + })) + + msg = await asyncio.wait_for(received.get(), timeout=5.0) + assert msg["type"] == "error" + assert "signature" in msg["detail"].lower() or "verification" in msg["detail"].lower() + + assert indexer.index.get_entry(entry.id) is not None + + await browser_pc.close() + await transport.close_all() + + +@pytest.mark.asyncio +async def test_webrtc_stream_request_missing_file(sk_node, sk_hub, gek, shared_dir): + """WebRTC DataChannel: stream_request for non-existent file returns error.""" + 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() + + transport = WebRTCTransport( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, + shared_root=shared_dir, index=indexer.index, + stun_servers=[], + ) + + browser_pc, channel, received = await _setup_peer( + transport, sk_hub, gek, "peer-mse") + + channel.send(_pack({ + "type": MNP.STREAM_REQUEST, "v": MNP_VERSION, + "file_id": "nonexistent-file-id", + })) + + msg = await asyncio.wait_for(received.get(), timeout=5.0) + assert msg["type"] == "error" + assert "not found" in msg["detail"].lower() + + await browser_pc.close() + await transport.close_all() + + +@pytest.mark.asyncio +async def test_webrtc_uploader_delete_requires_challenge(sk_node, sk_hub, gek, shared_dir): + """Uploader must prove Ed25519 key ownership to delete — no uploader shortcut.""" + 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() + + sk_uploader = Ed25519PrivateKey.generate() + pk_uploader_b64 = base64.b64encode( + sk_uploader.public_key().public_bytes( + serialization.Encoding.Raw, serialization.PublicFormat.Raw) + ).decode() + + transport = WebRTCTransport( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, + shared_root=shared_dir, index=indexer.index, + stun_servers=[], + ) + # No admin_pk configured — only uploader_pk should authorize deletion + + browser_pc, channel, received = await _setup_peer( + transport, sk_hub, gek, "peer-uploader-del", sk_user=sk_uploader) + + # Tag an existing entry with the uploader's public key + entry = indexer.index.entries[0] + entry.uploader_id = "user-001" + entry.uploader_pk = pk_uploader_b64 + + # Request deletion — should get a challenge (no shortcut) + channel.send(_pack({ + "type": MNP.FILE_DELETE, "v": MNP_VERSION, "file_id": entry.id, + })) + + 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 + + # Sign with uploader's Ed25519 key + challenge = base64.b64decode(challenge_msg["challenge"]) + signature = sk_uploader.sign(challenge) + channel.send(_pack({ + "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION, + "file_id": entry.id, + "signature": base64.b64encode(signature).decode(), + })) + + ack = await asyncio.wait_for(received.get(), timeout=5.0) + assert ack["type"] == MNP.FILE_DELETE_ACK + assert ack["file_id"] == entry.id + + # Verify file was removed from index + assert indexer.index.get_entry(entry.id) is None + + await browser_pc.close() + await transport.close_all() + + +@pytest.mark.asyncio +async def test_webrtc_uploader_impersonation_blocked(sk_node, sk_hub, gek, shared_dir): + """Hub-forged JWT with same sub cannot delete — wrong Ed25519 key is rejected.""" + 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() + + # User A uploaded the file + sk_user_a = Ed25519PrivateKey.generate() + pk_a_b64 = base64.b64encode( + sk_user_a.public_key().public_bytes( + serialization.Encoding.Raw, serialization.PublicFormat.Raw) + ).decode() + + # User B is the attacker (different Ed25519 key, but hub forges JWT with same sub) + sk_user_b = Ed25519PrivateKey.generate() + + transport = WebRTCTransport( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, + shared_root=shared_dir, index=indexer.index, + stun_servers=[], + ) + # No admin_pk — only uploader_pk matters + + # Tag entry with user A's public key + entry = indexer.index.entries[0] + entry.uploader_id = "user-001" + entry.uploader_pk = pk_a_b64 + + # Connect as user B (same jwt_sub "user-001" via hub forgery, but B's Ed25519 key) + browser_pc, channel, received = await _setup_peer( + transport, sk_hub, gek, "peer-impersonator", + jwt_sub="user-001", sk_user=sk_user_b) + + # Request deletion — should get a challenge + channel.send(_pack({ + "type": MNP.FILE_DELETE, "v": MNP_VERSION, "file_id": entry.id, + })) + + challenge_msg = await asyncio.wait_for(received.get(), timeout=5.0) + 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) + channel.send(_pack({ + "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION, + "file_id": entry.id, + "signature": base64.b64encode(bad_sig).decode(), + })) + + msg = await asyncio.wait_for(received.get(), timeout=5.0) + assert msg["type"] == "error" + assert "verification" in msg["detail"].lower() or "signature" in msg["detail"].lower() + + # File must still exist in the index + assert indexer.index.get_entry(entry.id) is not None + + await browser_pc.close() + await transport.close_all() + + +# ── GEK bundle P2P exchange tests ────────────────────────────────────────── + + +@pytest.fixture +def x25519_keypair(): + from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey + sk = X25519PrivateKey.generate() + sk_raw = sk.private_bytes( + serialization.Encoding.Raw, serialization.PrivateFormat.Raw, + serialization.NoEncryption()) + pk_raw = sk.public_key().public_bytes( + serialization.Encoding.Raw, serialization.PublicFormat.Raw) + return sk_raw, pk_raw + + +@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.""" + 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() + 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 + + # 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) + + 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"], + })) + ack = await asyncio.wait_for(q_admin.get(), timeout=5.0) + assert ack["type"] == "ack" + assert ack["detail"] == "gek_bundle_stored" + + # 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"] + + # Unwrap to verify it's correct + recovered = unwrap_gek(stored, sk_x_raw, pk_x_raw) + assert recovered == gek + + await bundle_store.close() + await pc_admin.close() + await transport.close_all() + + +@pytest.mark.asyncio +async def test_gek_bundle_fetch_during_handshake(sk_node, sk_hub, gek, shared_dir, + tmp_path, x25519_keypair): + """Browser fetches GEK bundle from node during the handshake challenge window.""" + 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() + + sk_x_raw, pk_x_raw = x25519_keypair + bundle_store = BundleStore(db_path=tmp_path / "bundles.db") + await bundle_store.open() + + # Pre-populate a bundle for user-001 in group "g" + bundle = wrap_gek(gek, pk_x_raw) + await bundle_store.store("g", "user-001", + bundle["pk_eph_b64"], bundle["nonce_b64"], + bundle["wrapped_b64"]) + + 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 + + # Connect manually: handshake → challenge → gek_bundle_fetch → response + browser_pc = RTCPeerConnection() + received = asyncio.Queue() + buf = bytearray() channel = browser_pc.createDataChannel("mnp") + ready = asyncio.Event() @channel.on("open") def on_open(): - channel.send(_pack({ - "type": MNP.HANDSHAKE, "v": MNP_VERSION, - "token": _make_jwt(sk_hub), - })) + ready.set() @channel.on("message") def on_msg(message): @@ -646,49 +1117,162 @@ async def test_webrtc_peer_cleanup_on_close(sk_node, sk_hub, gek, shared_dir): offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( - browser_pc.localDescription.sdp, "peer-cleanup") + browser_pc.localDescription.sdp, "peer-fetch") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) + await asyncio.wait_for(ready.wait(), timeout=5.0) + + # Step 1: Send handshake with group_id so _pending_group is set + token = _make_jwt(sk_hub, groups=["g"]) + channel.send(_pack({ + "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g", + })) + msg = await asyncio.wait_for(received.get(), timeout=5.0) + assert msg["type"] == MNP.HANDSHAKE_CHALLENGE + + # Step 2: Fetch GEK bundle from node (during challenge window) + channel.send(_pack({"type": MNP.GEK_BUNDLE_FETCH, "v": MNP_VERSION})) + bundle_resp = await asyncio.wait_for(received.get(), timeout=5.0) + assert bundle_resp["type"] == MNP.GEK_BUNDLE_RESP + assert bundle_resp["found"] is True + + # Step 3: Unwrap GEK and compute HMAC proof + recovered_gek = unwrap_gek(bundle_resp, sk_x_raw, pk_x_raw) + 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() + + # Step 4: Complete handshake + 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) assert ack["type"] == MNP.HANDSHAKE_ACK - assert "user-001" in transport._ctx["_peers"] - assert transport.active_peers == 1 + await bundle_store.close() + await browser_pc.close() + await transport.close_all() - await transport.close_peer("peer-cleanup") - assert "user-001" not in transport._ctx["_peers"] - assert transport.active_peers == 0 +# ── Keypair bundle P2P tests ───────────────────────────────────────────────── + + +@pytest.mark.asyncio +async def test_keypair_bundle_store_and_fetch(sk_node, sk_hub, gek, shared_dir, tmp_path): + """Keypair bundle stored on node, then fetched during handshake window.""" + 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() + + 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 + + # Connect and store a keypair bundle + pc1, ch1, q1 = await _setup_peer(transport, sk_hub, gek, "peer-kp-store") + ch1.send(_pack({ + "type": MNP.KEYPAIR_BUNDLE_STORE, + "v": MNP_VERSION, + "bundle_enc": "encrypted-keypair-data-base64", + })) + ack = await asyncio.wait_for(q1.get(), timeout=5.0) + assert ack["type"] == "ack" + assert ack["detail"] == "keypair_bundle_stored" + + # Verify in DB + stored = await bundle_store.fetch_keypair("user-001") + assert stored == "encrypted-keypair-data-base64" + + await pc1.close() + + # New connection: fetch during handshake window + browser_pc = RTCPeerConnection() + received = asyncio.Queue() + buf = bytearray() + channel = browser_pc.createDataChannel("mnp") + ready = asyncio.Event() + + @channel.on("open") + def on_open(): + ready.set() + + @channel.on("message") + def on_msg(message): + if isinstance(message, str): + message = message.encode() + buf.extend(message) + while len(buf) >= 4: + length = struct.unpack(">I", buf[:4])[0] + if len(buf) < 4 + length: + break + msg_bytes = bytes(buf[4:4 + length]) + del buf[:4 + length] + received.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) + + offer = await browser_pc.createOffer() + await browser_pc.setLocalDescription(offer) + answer_sdp, _ = await transport.handle_offer( + browser_pc.localDescription.sdp, "peer-kp-fetch") + await browser_pc.setRemoteDescription( + RTCSessionDescription(sdp=answer_sdp, type="answer")) + await asyncio.wait_for(ready.wait(), timeout=5.0) + + token = _make_jwt(sk_hub, groups=["g"]) + channel.send(_pack({ + "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g", + })) + msg = await asyncio.wait_for(received.get(), timeout=5.0) + assert msg["type"] == MNP.HANDSHAKE_CHALLENGE + + # Fetch keypair bundle during challenge window + channel.send(_pack({"type": MNP.KEYPAIR_BUNDLE_FETCH, "v": MNP_VERSION})) + kp_resp = await asyncio.wait_for(received.get(), timeout=5.0) + assert kp_resp["type"] == MNP.KEYPAIR_BUNDLE_RESP + assert kp_resp["found"] is True + assert kp_resp["bundle_enc"] == "encrypted-keypair-data-base64" + + await bundle_store.close() await browser_pc.close() + await transport.close_all() @pytest.mark.asyncio -async def test_webrtc_stream_segment_missing_file(sk_node, sk_hub, gek, shared_dir): - """WebRTC DataChannel: stream_segment for non-existent file returns error.""" +async def test_keypair_bundle_fetch_not_found(sk_node, sk_hub, gek, shared_dir, tmp_path): + """Keypair bundle fetch returns found=false when no bundle exists.""" 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() + 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 browser_pc = RTCPeerConnection() received = asyncio.Queue() buf = bytearray() - channel = browser_pc.createDataChannel("mnp") + ready = asyncio.Event() @channel.on("open") def on_open(): - channel.send(_pack({ - "type": MNP.HANDSHAKE, "v": MNP_VERSION, - "token": _make_jwt(sk_hub), - })) + ready.set() @channel.on("message") def on_msg(message): @@ -706,52 +1290,150 @@ async def test_webrtc_stream_segment_missing_file(sk_node, sk_hub, gek, shared_d offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( - browser_pc.localDescription.sdp, "peer-stream") + browser_pc.localDescription.sdp, "peer-kp-none") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) + await asyncio.wait_for(ready.wait(), timeout=5.0) - ack = await asyncio.wait_for(received.get(), timeout=5.0) - assert ack["type"] == MNP.HANDSHAKE_ACK + token = _make_jwt(sk_hub, groups=["g"]) + channel.send(_pack({ + "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g", + })) + msg = await asyncio.wait_for(received.get(), timeout=5.0) + assert msg["type"] == MNP.HANDSHAKE_CHALLENGE + channel.send(_pack({"type": MNP.KEYPAIR_BUNDLE_FETCH, "v": MNP_VERSION})) + resp = await asyncio.wait_for(received.get(), timeout=5.0) + assert resp["type"] == MNP.KEYPAIR_BUNDLE_RESP + assert resp["found"] is False + + await bundle_store.close() + await browser_pc.close() + await transport.close_all() + + +@pytest.mark.asyncio +async def test_gek_auto_activate_on_node_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).""" + 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() + + sk_x_raw, pk_x_raw = x25519_keypair + bundle_store = BundleStore(db_path=tmp_path / "bundles.db") + await bundle_store.open() + + new_gek = generate_gek() + + 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["node_user_id"] = "node-operator" + 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() + + 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) + ch_admin.send(_pack({ + "type": MNP.GEK_BUNDLE_STORE, + "v": MNP_VERSION, + "user_id": "node-operator", + "group_id": "g", + "pk_eph_b64": node_bundle["pk_eph_b64"], + "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) + + assert transport._ctx.get("gek") == new_gek + + await bundle_store.close() + await pc_admin.close() + await transport.close_all() + + +@pytest.mark.asyncio +async def test_webrtc_no_gek_connection_refused(sk_node, sk_hub, shared_dir): + """WebRTC DataChannel: connection refused when GEK is not initialized.""" + hub_pk_pem = _hub_pk_pem(sk_hub) + indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None) + await indexer.initial_scan() + + transport = WebRTCTransport( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=None, + shared_root=shared_dir, index=indexer.index, + stun_servers=[], + ) + + browser_pc = RTCPeerConnection() + received = asyncio.Queue() + channel = browser_pc.createDataChannel("mnp") + + @channel.on("message") + def on_msg(message): + if isinstance(message, str): + message = message.encode() + received.put_nowait(_unpack(message)) + + offer = await browser_pc.createOffer() + await browser_pc.setLocalDescription(offer) + answer_sdp, _ = await transport.handle_offer( + browser_pc.localDescription.sdp, "peer-no-gek") + await browser_pc.setRemoteDescription( + RTCSessionDescription(sdp=answer_sdp, type="answer")) + + await asyncio.sleep(0.5) + + token = _make_jwt(sk_hub) channel.send(_pack({ - "type": MNP.STREAM_SEGMENT, "v": MNP_VERSION, - "file_id": "nonexistent-file-id", - "segment_index": 0, "segment_duration": 4, + "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" - assert "not found" in msg["detail"].lower() + assert "not initialized" in msg["detail"].lower() await browser_pc.close() await transport.close_all() @pytest.mark.asyncio -async def test_webrtc_stream_request_missing_file(sk_node, sk_hub, gek, shared_dir): - """WebRTC DataChannel: stream_request for non-existent file returns error.""" +async def test_gek_bundle_fetch_not_found(sk_node, sk_hub, gek, shared_dir, tmp_path): + """GEK bundle fetch returns found=false when no bundle exists.""" 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() + 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 browser_pc = RTCPeerConnection() received = asyncio.Queue() buf = bytearray() - channel = browser_pc.createDataChannel("mnp") + ready = asyncio.Event() @channel.on("open") def on_open(): - channel.send(_pack({ - "type": MNP.HANDSHAKE, "v": MNP_VERSION, - "token": _make_jwt(sk_hub), - })) + ready.set() @channel.on("message") def on_msg(message): @@ -769,21 +1451,23 @@ async def test_webrtc_stream_request_missing_file(sk_node, sk_hub, gek, shared_d offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( - browser_pc.localDescription.sdp, "peer-mse") + browser_pc.localDescription.sdp, "peer-nofound") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) + await asyncio.wait_for(ready.wait(), timeout=5.0) - ack = await asyncio.wait_for(received.get(), timeout=5.0) - assert ack["type"] == MNP.HANDSHAKE_ACK - + token = _make_jwt(sk_hub, groups=["g"]) channel.send(_pack({ - "type": MNP.STREAM_REQUEST, "v": MNP_VERSION, - "file_id": "nonexistent-file-id", + "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g", })) - msg = await asyncio.wait_for(received.get(), timeout=5.0) - assert msg["type"] == "error" - assert "not found" in msg["detail"].lower() + assert msg["type"] == MNP.HANDSHAKE_CHALLENGE + + channel.send(_pack({"type": MNP.GEK_BUNDLE_FETCH, "v": MNP_VERSION})) + resp = await asyncio.wait_for(received.get(), timeout=5.0) + assert resp["type"] == MNP.GEK_BUNDLE_RESP + assert resp["found"] is False + await bundle_store.close() await browser_pc.close() await transport.close_all() |