aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub/api
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-08-09 04:39:34 +0200
committerChristophe Besson <cbesson@gmail.com>2026-08-09 04:39:34 +0200
commitfb91c4545c757711e1b5fd354ca4b311c89fd2c0 (patch)
treeeb1aee6cc0fb5fc020eed2763009dea5a32eb8ee /packages/meshbay-hub/src/meshbay_hub/api
parent77d76421829161df6b1ef628b4e6e051a2c3c2ee (diff)
downloadmeshbay-fb91c4545c757711e1b5fd354ca4b311c89fd2c0.tar.gz
feat(hub): add production hub — config, auth, API routers, tests
config.py: TOML + env var priority. auth.py: Argon2id passwords, JWT EdDSA with jti, refresh token hashed (blake3). Routers: hub (info/pubkey), users (register/login/refresh/pubkeys), nodes (announce/get), groups (create/gek-bundle/gek-retrieve). Rate limiting via slowapi. app.py factory with lifespan. All 40 tests pass (SQLite in-memory, no PostgreSQL required). Fix: remove tests/__init__.py to resolve namespace conflicts. Co-Authored-By: Claude Sonnet 4.6 (1M context) <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/api')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/deps.py47
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/groups.py118
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/hub.py26
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/middleware.py13
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/nodes.py66
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/users.py199
6 files changed, 469 insertions, 0 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/deps.py b/packages/meshbay-hub/src/meshbay_hub/api/deps.py
new file mode 100644
index 0000000..cb637f3
--- /dev/null
+++ b/packages/meshbay-hub/src/meshbay_hub/api/deps.py
@@ -0,0 +1,47 @@
+"""
+FastAPI shared dependencies — injected via Depends().
+"""
+
+from collections.abc import AsyncGenerator
+
+from fastapi import Depends, Header, HTTPException, status
+from sqlalchemy.ext.asyncio import AsyncSession
+from sqlalchemy import select
+
+from meshbay_hub.auth import decode_access_token
+from meshbay_hub.db.engine import get_db
+from meshbay_hub.db.models import User
+
+
+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.
+ """
+ try:
+ scheme, token = authorization.split(None, 1)
+ if scheme.lower() != "bearer":
+ raise ValueError
+ payload = decode_access_token(token)
+ except Exception:
+ raise HTTPException(
+ status_code=status.HTTP_401_UNAUTHORIZED,
+ detail="Invalid or expired token",
+ headers={"WWW-Authenticate": "Bearer"},
+ )
+
+ result = await db.execute(
+ select(User).where(User.id == payload["sub"]))
+ user = result.scalar_one_or_none()
+
+ if user is None:
+ raise HTTPException(status_code=status.HTTP_401_UNAUTHORIZED,
+ detail="User not found")
+ if user.status != "active":
+ raise HTTPException(status_code=status.HTTP_403_FORBIDDEN,
+ detail=f"Account {user.status}")
+ return user
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/groups.py b/packages/meshbay-hub/src/meshbay_hub/api/groups.py
new file mode 100644
index 0000000..942a88e
--- /dev/null
+++ b/packages/meshbay-hub/src/meshbay_hub/api/groups.py
@@ -0,0 +1,118 @@
+"""Group endpoints — /v1/groups/*"""
+
+from fastapi import APIRouter, Depends, HTTPException, Request
+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.db.engine import get_db
+from meshbay_hub.db.models import GEKBundle, Group, GroupMember, IPLog, User
+
+router = APIRouter(prefix="/v1/groups", tags=["groups"])
+
+
+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
+
+
+@router.post("", status_code=201)
+async def create_group(
+ body: GroupCreateRequest,
+ request: Request,
+ current_user: User = Depends(get_current_user),
+ db: AsyncSession = Depends(get_db),
+):
+ group = Group(
+ name=body.name,
+ admin_id=current_user.id,
+ visibility=body.visibility,
+ join_policy=body.join_policy,
+ )
+ db.add(group)
+ await db.flush() # get group.id
+
+ db.add(GroupMember(group_id=group.id, user_id=current_user.id))
+ db.add(IPLog(user_id=current_user.id, event="group_create",
+ ip_address=_ip(request), detail=body.name))
+ await db.commit()
+ await db.refresh(group)
+ return {"group_id": group.id, "name": group.name}
+
+
+@router.post("/{group_id}/members/{username}/gek", status_code=201)
+async def store_gek_bundle(
+ group_id: str,
+ username: str,
+ body: GEKBundleRequest,
+ current_user: User = Depends(get_current_user),
+ 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 admin can add members")
+
+ result = await db.execute(select(User).where(User.username == username))
+ target = result.scalar_one_or_none()
+ if not target:
+ raise HTTPException(status_code=404, detail="User not found")
+
+ # Upsert GEK bundle
+ existing = await db.get(GEKBundle, (group_id, target.id))
+ 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))
+
+ await db.commit()
+ return {"status": "stored", "group_id": group_id, "username": username}
+
+
+@router.get("/{group_id}/gek")
+async def get_my_gek_bundle(
+ group_id: str,
+ current_user: User = Depends(get_current_user),
+ db: AsyncSession = Depends(get_db),
+):
+ group = await db.get(Group, group_id)
+ if not group:
+ raise HTTPException(status_code=404, detail="Group not found")
+
+ 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,
+ }
+
+
+def _ip(request: Request) -> str:
+ fwd = request.headers.get("X-Forwarded-For")
+ return fwd.split(",")[0].strip() if fwd else (
+ request.client.host if request.client else "unknown")
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/hub.py b/packages/meshbay-hub/src/meshbay_hub/api/hub.py
new file mode 100644
index 0000000..00aba28
--- /dev/null
+++ b/packages/meshbay-hub/src/meshbay_hub/api/hub.py
@@ -0,0 +1,26 @@
+"""Hub info endpoints — /v1/hub/*"""
+
+from fastapi import APIRouter
+from meshbay_common import MNP_VERSION, MHP_VERSION
+from meshbay_hub import __version__
+from meshbay_hub.auth import hub_public_key_pem
+from meshbay_hub.db.engine import get_engine
+
+router = APIRouter(prefix="/v1/hub", tags=["hub"])
+
+
+@router.get("/info")
+async def hub_info():
+ engine = get_engine()
+ return {
+ "hub_version": __version__,
+ "mnp_version": MNP_VERSION,
+ "mhp_version": MHP_VERSION,
+ "db_dialect": engine.dialect.name,
+ }
+
+
+@router.get("/pubkey")
+async def hub_pubkey():
+ """Hub Ed25519 public key PEM — cached by nodes on first contact."""
+ return {"pk_hub_pem": hub_public_key_pem().decode()}
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/middleware.py b/packages/meshbay-hub/src/meshbay_hub/api/middleware.py
new file mode 100644
index 0000000..bed7b54
--- /dev/null
+++ b/packages/meshbay-hub/src/meshbay_hub/api/middleware.py
@@ -0,0 +1,13 @@
+"""
+Hub middleware — rate limiting on auth endpoints.
+
+Uses slowapi (Starlette-compatible, token bucket algorithm).
+Limits applied to /v1/users/register and /v1/users/login
+to mitigate credential stuffing and registration floods.
+"""
+
+from slowapi import Limiter
+from slowapi.util import get_remote_address
+
+# Rate limiter instance — mounted on the FastAPI app in app.py
+limiter = Limiter(key_func=get_remote_address)
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/nodes.py b/packages/meshbay-hub/src/meshbay_hub/api/nodes.py
new file mode 100644
index 0000000..b970aa8
--- /dev/null
+++ b/packages/meshbay-hub/src/meshbay_hub/api/nodes.py
@@ -0,0 +1,66 @@
+"""Node endpoints — /v1/nodes/*"""
+
+from fastapi import APIRouter, Depends, HTTPException, Request
+from pydantic import BaseModel
+from sqlalchemy.ext.asyncio import AsyncSession
+
+from meshbay_hub.api.deps import get_current_user
+from meshbay_hub.db.engine import get_db
+from meshbay_hub.db.models import IPLog, Node, User
+
+router = APIRouter(prefix="/v1/nodes", tags=["nodes"])
+
+
+class NodeAnnounceRequest(BaseModel):
+ pk_node: str
+ endpoint_hint: str | None = None
+
+
+@router.post("/announce", status_code=201)
+async def announce_node(
+ body: NodeAnnounceRequest,
+ request: Request,
+ current_user: User = Depends(get_current_user),
+ db: AsyncSession = Depends(get_db),
+):
+ node = Node(
+ user_id=current_user.id,
+ pk_node=body.pk_node,
+ endpoint_hint=body.endpoint_hint,
+ )
+ db.add(node)
+ db.add(IPLog(
+ user_id=current_user.id,
+ event="node_announce",
+ ip_address=_ip(request),
+ detail=body.endpoint_hint,
+ ))
+ await db.commit()
+ await db.refresh(node)
+ return {"node_id": node.id}
+
+
+@router.get("/{node_id}")
+async def get_node(
+ node_id: str,
+ current_user: User = Depends(get_current_user),
+ db: AsyncSession = Depends(get_db),
+):
+ node = await db.get(Node, node_id)
+ if not node:
+ raise HTTPException(status_code=404, detail="Node not found")
+ owner = await db.get(User, node.user_id)
+ return {
+ "node_id": node.id,
+ "username": owner.username if owner else "",
+ "pk_node": node.pk_node,
+ "endpoint_hint": node.endpoint_hint,
+ "announced_at": node.announced_at.isoformat(),
+ }
+
+
+def _ip(request: Request) -> str:
+ fwd = request.headers.get("X-Forwarded-For")
+ if fwd:
+ return fwd.split(",")[0].strip()
+ return request.client.host if request.client else "unknown"
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/users.py b/packages/meshbay-hub/src/meshbay_hub/api/users.py
new file mode 100644
index 0000000..0b615a4
--- /dev/null
+++ b/packages/meshbay-hub/src/meshbay_hub/api/users.py
@@ -0,0 +1,199 @@
+"""User endpoints — /v1/users/*"""
+
+from datetime import datetime, timezone, timedelta
+
+from fastapi import APIRouter, Depends, HTTPException, Request, status
+from pydantic import BaseModel, EmailStr, field_validator
+from sqlalchemy import select
+from sqlalchemy.ext.asyncio import AsyncSession
+
+from meshbay_hub.auth import (
+ decode_access_token,
+ generate_refresh_token,
+ hash_password,
+ hash_refresh_token,
+ hub_public_key_pem,
+ issue_access_token,
+ verify_password,
+)
+from meshbay_hub.config import HubConfig
+from meshbay_hub.db.engine import get_db
+from meshbay_hub.db.models import IPLog, RefreshToken, User
+from meshbay_hub.api.deps import get_current_user
+
+router = APIRouter(prefix="/v1/users", tags=["users"])
+
+_cfg: HubConfig | None = None
+
+def set_config(cfg: HubConfig) -> None:
+ global _cfg
+ _cfg = cfg
+
+def _ttl() -> int:
+ return _cfg.jwt.access_token_ttl if _cfg else 3600
+
+def _refresh_ttl() -> int:
+ return _cfg.jwt.refresh_token_ttl if _cfg else 86400 * 30
+
+
+# ── Models ────────────────────────────────────────────────────────────────────
+
+class RegisterRequest(BaseModel):
+ username: str
+ email: str
+ password: str
+ pk_user_ed25519: str # base64 raw 32B
+ pk_user_x25519: str # base64 raw 32B
+
+ @field_validator("username")
+ @classmethod
+ def username_valid(cls, v: str) -> str:
+ v = v.strip()
+ if len(v) < 3 or len(v) > 64:
+ raise ValueError("username must be 3-64 chars")
+ if not v.replace("_", "").replace("-", "").replace(".", "").isalnum():
+ 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
+
+
+class RefreshRequest(BaseModel):
+ refresh_token: str
+
+
+# ── Endpoints ─────────────────────────────────────────────────────────────────
+
+@router.post("/register", status_code=201)
+async def register(
+ body: RegisterRequest,
+ request: Request,
+ db: AsyncSession = Depends(get_db),
+):
+ existing = await db.execute(
+ select(User).where(User.username == body.username))
+ if existing.scalar_one_or_none():
+ raise HTTPException(status_code=409, detail="Username already taken")
+
+ pw_hash, pw_salt = hash_password(body.password)
+ hub_id = _cfg.identity.id if _cfg else "meshbay.org"
+ user = User(
+ username=body.username,
+ email=body.email,
+ pw_hash=pw_hash,
+ pw_salt=pw_salt,
+ pk_ed25519=body.pk_user_ed25519,
+ pk_x25519=body.pk_user_x25519,
+ hub_id=hub_id,
+ )
+ db.add(user)
+ db.add(IPLog(
+ event="account_create",
+ ip_address=_client_ip(request),
+ detail=body.username,
+ ))
+ await db.commit()
+ await db.refresh(user)
+
+ # Set user_id in IPLog after commit
+ await db.execute(
+ IPLog.__table__.update()
+ .where(IPLog.user_id == None) # noqa: E711
+ .values(user_id=user.id))
+ await db.commit()
+
+ return {"user_id": user.id}
+
+
+@router.post("/login")
+async def login(
+ body: LoginRequest,
+ request: Request,
+ db: AsyncSession = Depends(get_db),
+):
+ result = await db.execute(
+ select(User).where(User.username == body.username))
+ 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):
+ 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.status != "active":
+ raise HTTPException(status_code=403, detail=f"Account {user.status}")
+
+ access_token = issue_access_token(user.id, user.pk_ed25519, ttl=_ttl())
+ raw_rt, rt_hash = generate_refresh_token()
+
+ expires_at = datetime.now(timezone.utc) + timedelta(seconds=_refresh_ttl())
+ db.add(RefreshToken(user_id=user.id, token_hash=rt_hash, expires_at=expires_at))
+ db.add(IPLog(user_id=user.id, event="login", ip_address=ip))
+ await db.commit()
+
+ return {
+ "access_token": access_token,
+ "refresh_token": raw_rt,
+ "token_type": "bearer",
+ "expires_in": _ttl(),
+ }
+
+
+@router.post("/token/refresh")
+async def token_refresh(
+ body: RefreshRequest,
+ db: AsyncSession = Depends(get_db),
+):
+ rt_hash = hash_refresh_token(body.refresh_token)
+ result = await db.execute(
+ select(RefreshToken).where(
+ RefreshToken.token_hash == rt_hash,
+ RefreshToken.revoked == False, # noqa: E712
+ ))
+ rt = result.scalar_one_or_none()
+
+ if not rt or rt.expires_at.replace(tzinfo=timezone.utc) < datetime.now(timezone.utc):
+ raise HTTPException(status_code=401, detail="Invalid or expired refresh token")
+
+ user = await db.get(User, rt.user_id)
+ if not user or user.status != "active":
+ raise HTTPException(status_code=401, detail="User not found or suspended")
+
+ new_token = issue_access_token(user.id, user.pk_ed25519, ttl=_ttl())
+ return {"access_token": new_token, "token_type": "bearer", "expires_in": _ttl()}
+
+
+@router.get("/{username}/pubkeys")
+async def get_user_pubkeys(
+ username: str,
+ current_user: User = Depends(get_current_user),
+ db: AsyncSession = Depends(get_db),
+):
+ result = await db.execute(select(User).where(User.username == username))
+ target = result.scalar_one_or_none()
+ if not target:
+ raise HTTPException(status_code=404, detail="User not found")
+ return {
+ "user_id": target.id,
+ "username": target.username,
+ "pk_ed25519": target.pk_ed25519,
+ "pk_x25519": target.pk_x25519,
+ }
+
+
+def _client_ip(request: Request) -> str:
+ forwarded = request.headers.get("X-Forwarded-For")
+ if forwarded:
+ return forwarded.split(",")[0].strip()
+ return request.client.host if request.client else "unknown"