diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-08-09 04:39:34 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-08-09 04:39:34 +0200 |
| commit | fb91c4545c757711e1b5fd354ca4b311c89fd2c0 (patch) | |
| tree | eb1aee6cc0fb5fc020eed2763009dea5a32eb8ee /packages/meshbay-hub/src/meshbay_hub/api | |
| parent | 77d76421829161df6b1ef628b4e6e051a2c3c2ee (diff) | |
| download | meshbay-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.py | 47 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/groups.py | 118 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/hub.py | 26 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/middleware.py | 13 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/nodes.py | 66 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/users.py | 199 |
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" |