summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/api/revocation.py')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/revocation.py115
1 files changed, 101 insertions, 14 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
index c9c59c9..2e3323c 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
@@ -33,7 +33,7 @@ import time
import uuid
from typing import Any
-from fastapi import APIRouter, Depends, HTTPException, WebSocket, WebSocketDisconnect
+from fastapi import APIRouter, Depends, HTTPException, Request, WebSocket, WebSocketDisconnect
from pydantic import BaseModel
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
@@ -42,7 +42,7 @@ import jwt
from meshbay_hub.auth import hub_public_key_pem, decode_access_token
from meshbay_hub.api.deps import get_current_user, require_admin
from meshbay_hub.db.engine import get_db
-from meshbay_hub.db.models import Group, IPLog, User
+from meshbay_hub.db.models import Group, GroupMember, IPLog, Node, User
log = logging.getLogger(__name__)
@@ -128,38 +128,109 @@ async def _handle_chat_notify(group_id: str, sender_name: str, sender_user_id: s
log.warning("Chat notify failed: %s", e)
+async def _reject(ws: WebSocket, detail: str, code: int) -> None:
+ await ws.send_text(json.dumps({"type": "error", "detail": detail}))
+ await ws.close(code=code)
+
+
+async def _authorize_node_ws(token: str, claimed_id: str, claimed_groups) -> tuple:
+ """
+ Resolve a node WS registration against the database.
+
+ Returns (node_id, group_ids) on success, or (None, error_detail) on refusal.
+ Uses a short-lived session on purpose: a node WebSocket lives for hours, and a
+ request-scoped dependency would pin a PostgreSQL connection for its whole
+ lifetime, exhausting the pool once a handful of nodes connect.
+ """
+ from meshbay_hub.db.engine import get_session_factory
+
+ try:
+ decoded = decode_access_token(token)
+ except Exception as e:
+ return None, str(e)
+
+ if decoded.get("scope") != "node":
+ return None, "Node-scoped token required"
+
+ user_id = decoded.get("sub", "")
+ if not claimed_id:
+ return None, "node_id required"
+
+ async with get_session_factory()() as db:
+ node = await db.get(Node, claimed_id)
+ if node is None or node.user_id != user_id:
+ log.warning("Rejected WS registration for node %s by user %s",
+ claimed_id[:8], (user_id or "?")[:8])
+ return None, "node_id does not belong to this account"
+
+ user = await db.get(User, user_id)
+ if user is None or user.status != "active":
+ return None, "Account not active"
+
+ # Groups come from the database. The node may narrow the set to what it
+ # actually hosts, but it cannot widen it to groups it is not a member of —
+ # otherwise it could advertise itself as a source for any group on the hub.
+ result = await db.execute(
+ select(GroupMember.group_id).where(GroupMember.user_id == user_id))
+ authorized = {gid for (gid,) in result.all()}
+
+ claimed = set(claimed_groups or authorized)
+ return claimed_id, sorted(authorized & claimed)
+
+
@router.websocket("/v1/nodes/ws")
async def node_websocket(ws: WebSocket):
"""
Persistent WebSocket connection for nodes.
- Nodes authenticate with a JWT bearer in the first message.
- Hub sends revocation tokens as JSON messages.
+
+ Finding C2: this used to take `node_id` and `group_ids` straight from the
+ client's first message, with no check that the authenticated user owned that
+ node. Any registered user could connect with an ordinary browser token, claim a
+ victim node's id, and overwrite its entry in `_connected_nodes`. Every WebRTC
+ offer for that node was then relayed to the attacker, who answered with their
+ own SDP — a full node impersonation, and the DTLS channel binding does not help
+ because the attacker is the endpoint rather than a relay. The attacker received
+ the victim's encrypted keypair bundle, their chat, and their uploads.
+
+ Identity now comes from the token and the database, never from the message.
"""
await ws.accept()
node_id: str | None = None
try:
- # Auth: expect {"type": "auth", "token": "<jwt>"}
+ # Auth: expect {"type": "auth", "token": "<jwt>", "node_id": "..."}
raw = await ws.receive_text()
msg = json.loads(raw)
if msg.get("type") != "auth" or "token" not in msg:
- await ws.send_text(json.dumps({"type": "error", "detail": "Send auth first"}))
- await ws.close(code=4001)
+ await _reject(ws, "Send auth first", 4001)
return
try:
decoded = decode_access_token(msg["token"])
except Exception as e:
- await ws.send_text(json.dumps({"type": "error", "detail": str(e)}))
- await ws.close(code=4001)
+ await _reject(ws, str(e), 4001)
return
- node_id = msg.get("node_id") or decoded.get("sub", "unknown")
+ claimed_id = msg.get("node_id") or ""
+
+ # Refuse to displace a live registration rather than silently overwriting it.
+ if claimed_id and claimed_id in _connected_nodes:
+ await _reject(ws, "Node already connected", 4009)
+ return
+
+ resolved_id, result = await _authorize_node_ws(
+ msg["token"], claimed_id, msg.get("group_ids"))
+ if resolved_id is None:
+ await _reject(ws, result, 4003)
+ return
+ group_ids = result
+
+ user_id = decoded.get("sub", "")
+ node_id = resolved_id
_connected_nodes[node_id] = ws
- group_ids = msg.get("group_ids", [])
- if group_ids:
- _node_groups[node_id] = group_ids
- log.info("Node WS connected: %s (groups=%d)", node_id[:8], len(group_ids))
+ _node_groups[node_id] = group_ids
+ log.info("Node WS connected: %s (user=%s, groups=%d)",
+ node_id[:8], user_id[:8], len(group_ids))
await ws.send_text(json.dumps({"type": "auth_ok", "node_id": node_id}))
# Message loop — handle ping, punch_ready, etc.
@@ -203,12 +274,28 @@ class IncomingRequest(BaseModel):
async def notify_incoming(
node_id: str,
body: IncomingRequest,
+ request: Request,
current_user: User = Depends(get_current_user),
):
"""
Signal a node that a client wants to connect (NAT punch coordination).
Hub forwards the request via WebSocket; node punches NAT and replies punch_ready.
+
+ Finding H6: peer_ip was taken verbatim, so any authenticated user could make an
+ arbitrary node emit UDP packets to an address of their choosing — a small
+ reflection primitive using someone else's machine. The probe target must now be
+ the caller's own source address.
"""
+ from meshbay_hub.api.netutil import client_ip
+
+ caller_ip = client_ip(request)
+ if body.peer_ip != caller_ip:
+ raise HTTPException(
+ status_code=403,
+ detail="peer_ip must match the requesting address")
+ if not (1 <= body.peer_port <= 65535):
+ raise HTTPException(status_code=422, detail="Invalid peer_port")
+
ws = _connected_nodes.get(node_id)
if not ws:
raise HTTPException(status_code=404, detail="Node not connected")