"""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 ( FederatedGroup, GEKBundle, Group, GroupMember, IPLog, SwarmSource, User, ) router = APIRouter(prefix="/v1/groups", tags=["groups"]) @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() return { "groups": [ { "id": g.id, "name": g.name, "visibility": g.visibility, "join_policy": g.join_policy, "created_at": g.created_at.isoformat(), "is_admin": g.admin_id == current_user.id, } 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") 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, "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" @router.post("/v1/swarm/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 content hash (public swarm).""" 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} @router.get("/v1/swarm/{content_hash}") async def swarm_sources( content_hash: str, db: AsyncSession = Depends(get_db), ): """Return list of nodes that can serve a content hash.""" 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], } 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)) 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 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}", ) 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")