"""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, require_user_scope from meshbay_hub.api.netutil import client_ip from meshbay_hub.db.engine import get_db from meshbay_hub.db.models import ( FederatedGroup, Group, GroupMember, IPLog, SwarmSource, User, ) router = APIRouter(prefix="/v1/groups", tags=["groups"]) # Swarm endpoints live at /v1/swarm/*. They were previously declared on the groups # router with a full path, which mounted them at /v1/groups/v1/swarm/* (H7). swarm_router = APIRouter(prefix="/v1/swarm", tags=["swarm"]) @router.get("/mine") async def my_groups( current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db), ): """List groups the current user belongs to.""" result = await db.execute( select(Group) .join(GroupMember, Group.id == GroupMember.group_id) .where(GroupMember.user_id == current_user.id, Group.status == "active") .order_by(Group.name) ) groups = result.scalars().all() muted_rows = await db.execute( select(GroupMember.group_id, GroupMember.muted) .where(GroupMember.user_id == current_user.id)) muted_map = {gid: bool(m) for gid, m in muted_rows.all()} return { "groups": [ { "id": g.id, "name": g.name, "visibility": g.visibility, "muted": muted_map.get(g.id, False), "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 ] } @router.get("/{group_id}/nodes") async def group_online_nodes( group_id: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db), ): """Return online nodes that serve a group (for WebRTC connection).""" from meshbay_hub.api.revocation import get_online_nodes_for_group from meshbay_hub.db.models import Node 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 = [] for nid in node_ids: node = await db.get(Node, nid) if node: nodes.append({"node_id": nid, "pk_node": node.pk_node}) return {"nodes": nodes} @router.get("") async def list_public_groups( db: AsyncSession = Depends(get_db), q: str = "", limit: int = 50, offset: int = 0, include_federated: bool = True, ): """List/search public groups — local and optionally federated. No auth required.""" query = select(Group).where(Group.visibility == "public", Group.status == "active") if q: query = query.where(Group.name.ilike(f"%{q}%")) result = await db.execute( query.order_by(Group.created_at.desc()).limit(limit).offset(offset) ) local = result.scalars().all() 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 ] if include_federated: fed_query = select(FederatedGroup) if q: fed_query = fed_query.where(FederatedGroup.name.ilike(f"%{q}%")) fed_result = await db.execute( fed_query.order_by(FederatedGroup.updated_at.desc()).limit(limit) ) for fg in fed_result.scalars().all(): groups.append({ "id": fg.id, "name": fg.name, "join_policy": fg.join_policy, "updated_at": fg.updated_at.isoformat(), "source": fg.source_hub, }) return {"groups": groups, "total": len(groups)} # ── Swarm (content replication) ─────────────────────────────────────────────── class SwarmRegisterRequest(BaseModel): content_hash: str # blake3 hex endpoint: str # "ip:port" @swarm_router.post("/register", status_code=201) async def swarm_register( body: SwarmRegisterRequest, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db), ): """ Node registers itself as a source for a PUBLIC content hash. Finding H7: the node registered hashes for every group it hosted, private ones included, and this route was mounted at /v1/groups/v1/swarm/register — so the node's calls 404'd and the leak was masked by a routing bug rather than prevented. Nodes now filter by group visibility before calling, and the path is correct, so the filter has to be right. """ from meshbay_hub.csam import check_content_hash if check_content_hash(body.content_hash): raise HTTPException(status_code=451, detail="Content blocked") from datetime import datetime, timezone existing = await db.get(SwarmSource, (body.content_hash, current_user.id)) now = datetime.now(timezone.utc) if existing: existing.endpoint = body.endpoint existing.last_seen = now else: db.add(SwarmSource( content_hash=body.content_hash, node_id=current_user.id, endpoint=body.endpoint, )) await db.commit() return {"status": "registered", "hash": body.content_hash} @swarm_router.get("/{content_hash}") async def swarm_sources( content_hash: str, current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db), ): """ Return nodes that can serve a content hash. Authenticated (H7): an open endpoint lets anyone probe whether a given file exists anywhere in the network and which node holds it. """ from datetime import datetime, timezone, timedelta cutoff = datetime.now(timezone.utc) - timedelta(minutes=30) result = await db.execute( select(SwarmSource) .where( SwarmSource.content_hash == content_hash, SwarmSource.last_seen > cutoff, ) ) sources = result.scalars().all() return { "hash": content_hash, "sources": [{"node_id": s.node_id, "endpoint": s.endpoint} for s in sources], } @router.get("/{group_id}/members") async def group_members( 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") mem = await db.get(GroupMember, (group_id, current_user.id)) if not mem: raise HTTPException(status_code=403, detail="Not a member") result = await db.execute( select(User.id, User.username) .join(GroupMember, User.id == GroupMember.user_id) .where(GroupMember.group_id == group_id) ) members = [{"user_id": uid, "username": uname} for uid, uname in result.all()] return { "group_id": group_id, "admin_id": group.admin_id, "members": members, } @router.post("/{group_id}/join") async def join_group( group_id: str, 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.status != "active": raise HTTPException(status_code=403, detail="Group is not active") if group.join_policy != "open": raise HTTPException(status_code=403, detail="Group does not allow open joining") existing = await db.get(GroupMember, (group_id, current_user.id)) if existing: raise HTTPException(status_code=409, detail="Already a member") db.add(GroupMember(group_id=group_id, user_id=current_user.id)) db.add(IPLog(user_id=current_user.id, event="group_join", ip_address=client_ip(request), detail=group.name)) await db.commit() return {"status": "joined", "group_id": group_id, "name": group.name} class GroupCreateRequest(BaseModel): name: str visibility: str = "private" # public|private join_policy: str = "invite" # open|request|invite description: str | None = None @router.post("", status_code=201) async def create_group( body: GroupCreateRequest, request: Request, 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 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=client_ip(request), detail=body.name)) await db.commit() await db.refresh(group) return {"group_id": group.id, "name": group.name} class GroupUpdateRequest(BaseModel): description: str | None = None @router.patch("/{group_id}") async def update_group( group_id: str, body: GroupUpdateRequest, current_user: User = Depends(require_user_scope), db: AsyncSession = Depends(get_db), ): """ Change the group's description. Owner only. Only the description: name, visibility and join policy are what members joined on the strength of, and a group that can quietly become public is a different thing from the one they agreed to. Those need a decision about who is told, not a PATCH. """ 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 owner can edit it") if body.description is not None: desc = body.description.strip()[:512] group.description = desc or None await db.commit() return {"group_id": group.id, "description": group.description or ""} @router.post("/{group_id}/members/{username}", status_code=201) async def add_group_member( group_id: str, username: str, 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 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") new_member = False 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 await create_notification( db, target.id, "group_invite", f"You were added to {group.name}", link=f"#/group/{group_id}", group_id=group_id, ) await db.commit() return {"status": "stored", "group_id": group_id, "username": username} class MuteRequest(BaseModel): muted: bool @router.post("/{group_id}/mute") async def set_group_mute( group_id: str, body: MuteRequest, current_user: User = Depends(require_user_scope), db: AsyncSession = Depends(get_db), ): """ Turn this group's notifications on or off, for this account. Server-side on purpose: it used to be a checkbox in the browser's localStorage that nothing ever read, so turning notifications off for a group had no effect anywhere. Now nothing is created in the first place. """ membership = await db.get(GroupMember, (group_id, current_user.id)) if not membership: raise HTTPException(status_code=404, detail="Not a member of this group") membership.muted = body.muted await db.commit() return {"status": "ok", "group_id": group_id, "muted": body.muted} @router.delete("/{group_id}") async def delete_group( group_id: str, 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") 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=client_ip(request), detail=group.name)) await db.delete(group) await db.commit() return {"status": "deleted", "group_id": group_id}