summaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-08-13 03:56:30 +0200
committerChristophe Besson <cbesson@gmail.com>2026-08-13 03:56:30 +0200
commitf0248975908ad670fa8a820f865bf22ea8d0172d (patch)
treef4af64d36cacaccb4f6d13436e001aeb57e861e3 /packages
parent35130e5528a52161630fd1c93572e1b2b7cd911b (diff)
downloadmeshbay-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')
-rw-r--r--packages/meshbay-common/src/meshbay_common/protocol.py12
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/admin.py12
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/deps.py40
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/groups.py77
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/nodes.py71
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/users.py140
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/auth.py6
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/db/__init__.py4
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/2041a4060b3c_add_keypair_bundle_federated_groups_.py2
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/d28b9caf9f07_initial_schema.py12
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/db/models.py25
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/app.js549
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/crypto.js16
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/i18n.js19
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/keyderive.js127
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/style.css192
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/transport.js178
-rw-r--r--packages/meshbay-hub/tests/test_hub_api.py279
-rw-r--r--packages/meshbay-hub/tests/test_node_auth.py281
-rw-r--r--packages/meshbay-node/pyproject.toml1
-rw-r--r--packages/meshbay-node/src/meshbay_node/bundle_store.py105
-rw-r--r--packages/meshbay-node/src/meshbay_node/config.py12
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py160
-rw-r--r--packages/meshbay-node/src/meshbay_node/hub_client.py128
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py353
-rw-r--r--packages/meshbay-node/src/meshbay_node/ui/app.py163
-rw-r--r--packages/meshbay-node/tests/test_daemon.py44
-rw-r--r--packages/meshbay-node/tests/test_hub_client.py76
-rw-r--r--packages/meshbay-node/tests/test_webrtc_transport.py1066
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} &mdash; <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 &gt; 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__} &mdash; localhost only &mdash;
<a href="/api/status">status</a> &middot;
@@ -392,7 +523,33 @@ def _render_page(state: dict) -> str:
&mdash; 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()