aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-hub/src')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/groups.py130
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/revocation.py34
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/daemon.py61
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/3dc91cd4ea52_group_hosted_at.py37
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/db/models.py5
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/tasks/cleanup.py47
6 files changed, 305 insertions, 9 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/groups.py b/packages/meshbay-hub/src/meshbay_hub/api/groups.py
index 8e3197c..8283276 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/groups.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/groups.py
@@ -2,7 +2,7 @@
from fastapi import APIRouter, Depends, HTTPException, Request
from pydantic import BaseModel
-from sqlalchemy import select
+from sqlalchemy import func, or_, select
from sqlalchemy.ext.asyncio import AsyncSession
from meshbay_hub.api.deps import get_current_user, require_user_scope
@@ -26,10 +26,16 @@ async def my_groups(
db: AsyncSession = Depends(get_db),
):
"""List groups the current user belongs to."""
+ from meshbay_hub.api.revocation import get_online_nodes_for_group
+
result = await db.execute(
select(Group)
.join(GroupMember, Group.id == GroupMember.group_id)
- .where(GroupMember.user_id == current_user.id, Group.status == "active")
+ .where(GroupMember.user_id == current_user.id, Group.status == "active",
+ # A group nobody hosts yet is the owner's business alone. Someone
+ # added to it before a node exists would see a name they cannot
+ # open and cannot be told why.
+ or_(Group.hosted_at.is_not(None), Group.admin_id == current_user.id))
.order_by(Group.name)
)
groups = result.scalars().all()
@@ -48,6 +54,15 @@ async def my_groups(
"created_at": g.created_at.isoformat(),
"is_admin": g.admin_id == current_user.id,
"description": g.description or "",
+ # Presence, from the socket registry the hub already keeps for
+ # signaling — so the sidebar gets it on the request it already
+ # makes, with no poll and no timer. It says a node serving this
+ # group is connected *to the hub*; it does not promise this
+ # browser can reach it, and a hub is free to lie about it. The
+ # client downgrades to offline on its own failed connection,
+ # which is the evidence that actually concerns the user.
+ "node_online": bool(get_online_nodes_for_group(g.id)),
+ "hosted": g.hosted_at is not None,
}
for g in groups
]
@@ -88,7 +103,12 @@ async def list_public_groups(
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")
+ # Unhosted groups are absent from the directory: until a node announces it,
+ # a group has no files, no key and nothing to connect to, so listing it only
+ # produces a dead end. Its owner still sees it in /mine while they set it up.
+ query = select(Group).where(Group.visibility == "public",
+ Group.status == "active",
+ Group.hosted_at.is_not(None))
if q:
query = query.where(Group.name.ilike(f"%{q}%"))
result = await db.execute(
@@ -246,10 +266,50 @@ async def join_group(
class GroupCreateRequest(BaseModel):
name: str
visibility: str = "private" # public|private
- join_policy: str = "invite" # open|request|invite
+ join_policy: str = "invite" # open (public groups) | invite (private)
description: str | None = None
+MAX_PUBLIC_GROUPS = 10
+
+
+async def _check_public_group_quota(db: AsyncSession, user: User) -> None:
+ """Refuse an eleventh live public group from the same owner.
+
+ Public groups are the ones that cost other people something: they appear in
+ Discover and anyone may join them, so a script that opens hundreds fills the
+ directory for everybody. Private groups are invisible to anyone not invited
+ and are not capped.
+
+ Counted: public, still active, owned by this user. A group suspended by
+ moderation does not hold a slot — the owner is already being dealt with, and
+ keeping the slot occupied would punish them twice. Deleting one frees a slot,
+ since the row is gone.
+
+ Creation is the only place this can be checked, and deliberately so: PATCH
+ refuses to change visibility at all, so a private group cannot be flipped
+ public behind the cap. **If visibility ever becomes editable, this check has
+ to move with it.**
+ """
+ if user.role in ("admin", "moderator"):
+ return # the cap is an anti-spam measure, not a rule about operating an instance
+
+ count = (await db.execute(
+ select(func.count())
+ .select_from(Group)
+ .where(Group.admin_id == user.id,
+ Group.visibility == "public",
+ Group.status == "active")
+ )).scalar_one()
+
+ if count >= MAX_PUBLIC_GROUPS:
+ raise HTTPException(
+ status_code=409,
+ detail=f"You already run {count} public groups, which is the limit of "
+ f"{MAX_PUBLIC_GROUPS}. Delete one you no longer use, or create "
+ f"this one as private — private groups are not limited.")
+
+
@router.post("", status_code=201)
async def create_group(
body: GroupCreateRequest,
@@ -257,6 +317,19 @@ async def create_group(
current_user: User = Depends(require_user_scope),
db: AsyncSession = Depends(get_db),
):
+ if body.visibility == "public":
+ # A public group that admits nobody is a contradiction: it is listed in
+ # the directory, so people find it and then discover they cannot get in.
+ # Admission by request was considered and dropped — between strangers the
+ # only channel is the hub, so the one-time code would travel through the
+ # very party it exists to keep out, and would protect nothing.
+ if body.join_policy != "open":
+ raise HTTPException(
+ status_code=422,
+ detail="A public group is open to join. Make it private if you "
+ "want to choose who comes in.")
+ await _check_public_group_quota(db, current_user)
+
desc = (body.description or "")[:512] if body.description else None
group = Group(
name=body.name,
@@ -324,6 +397,55 @@ async def remove_group_member(
return {"status": "removed", "group_id": group_id, "username": username}
+@router.post("/{group_id}/leave")
+async def leave_group(
+ group_id: str,
+ request: Request,
+ current_user: User = Depends(require_user_scope),
+ db: AsyncSession = Depends(get_db),
+):
+ """
+ Leave a group you are a member of.
+
+ Deliberately separate from `DELETE /{group_id}/members/{username}`, which is
+ the owner removing somebody else and is refused to everyone else. Reusing it
+ would have meant relaxing that check for the self case, and an authorization
+ rule with an exception in it is the kind that gets read wrong later.
+
+ The owner cannot leave: the group would be left with no one able to admit a
+ member, edit it or delete it. That is the same answer the removal endpoint
+ already gives, and the same one account deletion gives while you still own
+ groups — hand the group over (not yet possible) or delete it.
+
+ This is only the hub's half, exactly as for removal: membership is gone, so
+ signaling will not reach a node and the next token will not name this group.
+ The node keeps what it holds — the identity it pinned, the keypair bundle,
+ and the files uploaded — until its operator unpins them, and whoever left
+ still holds the group key they were served, so the operator should rotate it
+ (`meshbay-node gek-init`) if that matters.
+ """
+ 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=409,
+ detail="You own this group, so you cannot leave it — it would be left "
+ "with nobody able to manage it. Delete the group instead.")
+
+ membership = await db.get(GroupMember, (group_id, current_user.id))
+ if not membership:
+ raise HTTPException(status_code=404, detail="You are not a member of this group")
+
+ await db.delete(membership)
+ db.add(IPLog(user_id=current_user.id, event="group_leave",
+ ip_address=client_ip(request),
+ detail=f"left {group.name}"))
+ await db.commit()
+ return {"status": "left", "group_id": group_id}
+
+
class GroupUpdateRequest(BaseModel):
description: str | None = None
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
index 1f6d7f0..d555f9f 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
@@ -31,11 +31,12 @@ import json
import logging
import time
import uuid
+from datetime import datetime, timezone
from typing import Any
from fastapi import APIRouter, Depends, HTTPException, Request, WebSocket, WebSocketDisconnect
from pydantic import BaseModel
-from sqlalchemy import select
+from sqlalchemy import select, update
from sqlalchemy.ext.asyncio import AsyncSession
import jwt
@@ -67,6 +68,36 @@ def get_online_nodes_for_group(group_id: str) -> list[str]:
return [nid for nid, gids in _node_groups.items() if group_id in gids]
+
+async def _mark_hosted(group_ids: list[str]) -> None:
+ """Stamp the first time a node announced it hosts each of these groups.
+
+ `group_ids` is already narrowed to what this node may claim — the caller
+ derives it from the database and a node can only shrink the set, never widen
+ it (finding C2) — so being announced here is evidence the group has a host.
+
+ Set once. A node going offline does not un-host a group, and re-stamping on
+ every reconnection would make `hosted_at` a "last seen" field, which is what
+ the in-memory registry is already for.
+ """
+ from meshbay_hub.db.engine import get_session_factory
+ from meshbay_hub.db.models import Group
+
+ if not group_ids:
+ return
+ try:
+ async with get_session_factory()() as db:
+ await db.execute(
+ update(Group)
+ .where(Group.id.in_(group_ids), Group.hosted_at.is_(None))
+ .values(hosted_at=datetime.now(timezone.utc)))
+ await db.commit()
+ except Exception as e:
+ # A group that stays unhosted in the table is visible to its owner and
+ # collected later; failing the socket over it would take the node down.
+ log.warning("Could not mark groups hosted: %s", e)
+
+
async def broadcast_revocation(token: str) -> int:
"""Push a signed revocation token to all connected nodes. Returns count sent."""
payload = json.dumps({"type": "revocation", "token": token})
@@ -236,6 +267,7 @@ async def node_websocket(ws: WebSocket):
node_id = resolved_id
_connected_nodes[node_id] = ws
_node_groups[node_id] = group_ids
+ await _mark_hosted(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}))
diff --git a/packages/meshbay-hub/src/meshbay_hub/daemon.py b/packages/meshbay-hub/src/meshbay_hub/daemon.py
index ff66fee..4af26f5 100644
--- a/packages/meshbay-hub/src/meshbay_hub/daemon.py
+++ b/packages/meshbay-hub/src/meshbay_hub/daemon.py
@@ -1,6 +1,7 @@
-"""Entry point for the meshbay-hub systemd service."""
+"""Entry point for the meshbay-hub systemd service, and its few CLI chores."""
import argparse
+import asyncio
import logging
import sys
from pathlib import Path
@@ -15,6 +16,25 @@ def main() -> None:
parser.add_argument("--config", type=Path, default=None)
parser.add_argument("--log-level", default="INFO",
choices=["DEBUG", "INFO", "WARNING", "ERROR"])
+ # Optional on purpose: the systemd unit runs `meshbay-hub --config …` with no
+ # subcommand and must go on starting the server.
+ sub = parser.add_subparsers(dest="command")
+
+ prune = sub.add_parser(
+ "prune-groups",
+ help="delete groups no node ever hosted (meant for cron)")
+ prune.add_argument("--days", type=int, default=None,
+ help="grace period since creation (default 7)")
+ prune.add_argument("--dry-run", action="store_true",
+ help="list what would go, delete nothing")
+ # Accepted after the subcommand too. Everyone writes the cron line as
+ # `prune-groups --config …`, and argparse only takes an option before the
+ # subcommand unless the subparser declares it as well. SUPPRESS so that
+ # leaving it out here does not overwrite a value given before it.
+ prune.add_argument("--config", type=Path, default=argparse.SUPPRESS)
+ prune.add_argument("--log-level", default=argparse.SUPPRESS,
+ choices=["DEBUG", "INFO", "WARNING", "ERROR"])
+
args = parser.parse_args()
logging.basicConfig(
@@ -24,6 +44,9 @@ def main() -> None:
cfg = load_config(args.config)
+ if args.command == "prune-groups":
+ sys.exit(asyncio.run(_prune_groups(cfg, args.days, args.dry_run)))
+
uvicorn.run(
"meshbay_hub.app:create_app",
factory=True,
@@ -34,5 +57,41 @@ def main() -> None:
)
+async def _prune_groups(cfg, days: int | None, dry_run: bool) -> int:
+ """Collect groups that were created and never given a node.
+
+ A group with no host has no files, no key and nothing to connect to, and is
+ invisible to everyone but its owner — so it is litter rather than data. The
+ grace period is counted from creation and `hosted_at` is never cleared, so a
+ node being offline today cannot make a live group look abandoned.
+
+ Deliberately a command rather than a loop inside the server: deleting other
+ people's groups on a timer nobody asked for is the kind of thing an operator
+ should schedule knowingly, and `--dry-run` lets them see the list first.
+ """
+ from meshbay_hub.db.engine import close_db, get_session_factory, init_db
+ from meshbay_hub.tasks.cleanup import UNHOSTED_GRACE_DAYS, prune_unhosted_groups
+
+ grace = UNHOSTED_GRACE_DAYS if days is None else days
+ # init_db also runs create_all. That is what the server does at startup, so
+ # it adds no risk here — and it creates missing *tables* only, never a
+ # missing column, which is why deploys run alembic separately.
+ await init_db(cfg.db.url)
+ try:
+ async with get_session_factory()() as db:
+ gone = await prune_unhosted_groups(db, grace_days=grace, dry_run=dry_run)
+ finally:
+ await close_db()
+
+ verb = "would delete" if dry_run else "deleted"
+ if not gone:
+ print(f"No group older than {grace} day(s) is still unhosted.")
+ return 0
+ print(f"{verb} {len(gone)} group(s) unhosted for more than {grace} day(s):")
+ for gid, name in gone:
+ print(f" {gid} {name}")
+ return 0
+
+
if __name__ == "__main__":
main()
diff --git a/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/3dc91cd4ea52_group_hosted_at.py b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/3dc91cd4ea52_group_hosted_at.py
new file mode 100644
index 0000000..57aa9a7
--- /dev/null
+++ b/packages/meshbay-hub/src/meshbay_hub/db/migrations/versions/3dc91cd4ea52_group_hosted_at.py
@@ -0,0 +1,37 @@
+"""group hosted_at
+
+Records the first time a node announced that it hosts a group. Groups without
+it are shown to their owner only, and `meshbay-hub prune-groups` collects the
+ones that never got a node.
+
+Backfilled to created_at for every existing group: they predate the rule, and
+starting the clock on them retroactively would delete live groups on the first
+run of the reaper.
+
+Revision ID: 3dc91cd4ea52
+Revises: d1f47a90c3b2
+Create Date: 2026-08-16 03:20:00.575130
+
+"""
+from typing import Sequence, Union
+
+from alembic import op
+import sqlalchemy as sa
+
+
+# revision identifiers, used by Alembic.
+revision: str = '3dc91cd4ea52'
+down_revision: Union[str, Sequence[str], None] = 'd1f47a90c3b2'
+branch_labels: Union[str, Sequence[str], None] = None
+depends_on: Union[str, Sequence[str], None] = None
+
+
+def upgrade() -> None:
+ op.add_column("groups",
+ sa.Column("hosted_at", sa.DateTime(timezone=True), nullable=True))
+ # Existing groups are grandfathered in rather than left to expire.
+ op.execute("UPDATE groups SET hosted_at = created_at WHERE hosted_at IS NULL")
+
+
+def downgrade() -> None:
+ op.drop_column("groups", "hosted_at")
diff --git a/packages/meshbay-hub/src/meshbay_hub/db/models.py b/packages/meshbay-hub/src/meshbay_hub/db/models.py
index 7d5a3e3..a220a64 100644
--- a/packages/meshbay-hub/src/meshbay_hub/db/models.py
+++ b/packages/meshbay-hub/src/meshbay_hub/db/models.py
@@ -99,6 +99,11 @@ class Group(Base):
description: Mapped[str | None] = mapped_column(String(512))
status: Mapped[str] = mapped_column(String(16), default="active") # active|suspended|revoked
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=_now)
+ # First time a node registered on /v1/nodes/ws announcing that it hosts this
+ # group. Until then the group has no files, no key and nobody to serve it, so
+ # it is shown to its owner only and is what `prune-groups` collects. Set once
+ # and never cleared: a node going offline does not un-host a group.
+ hosted_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True))
members: Mapped[list["GroupMember"]] = relationship(back_populates="group")
diff --git a/packages/meshbay-hub/src/meshbay_hub/tasks/cleanup.py b/packages/meshbay-hub/src/meshbay_hub/tasks/cleanup.py
index 7674c6e..c387100 100644
--- a/packages/meshbay-hub/src/meshbay_hub/tasks/cleanup.py
+++ b/packages/meshbay-hub/src/meshbay_hub/tasks/cleanup.py
@@ -1,13 +1,13 @@
-"""Scheduled cleanup tasks — IP log purge (1-year retention)."""
+"""Scheduled cleanup tasks — IP log purge, and unhosted group collection."""
import asyncio
import logging
from datetime import datetime, timedelta, timezone
-from sqlalchemy import delete
+from sqlalchemy import delete, select
from sqlalchemy.ext.asyncio import AsyncSession
-from meshbay_hub.db.models import IPLog
+from meshbay_hub.db.models import Group, GroupMember, IPLog
log = logging.getLogger(__name__)
@@ -38,3 +38,44 @@ async def cleanup_loop(get_session):
await asyncio.sleep(CLEANUP_INTERVAL_HOURS * 3600)
except asyncio.CancelledError:
return
+
+
+# ── Groups that never got a node ──────────────────────────────────────────────
+
+UNHOSTED_GRACE_DAYS = 7
+
+
+async def find_unhosted_groups(db: AsyncSession, grace_days: int = UNHOSTED_GRACE_DAYS):
+ """Groups created more than `grace_days` ago that no node has ever announced.
+
+ `hosted_at` is set the first time a node registers claiming the group and is
+ never cleared, so this finds groups that were created and then abandoned —
+ not ones whose node happens to be offline today. That distinction is the
+ whole reason the column exists rather than a check against the live socket
+ registry, which would delete every group during a hub restart.
+ """
+ cutoff = datetime.now(timezone.utc) - timedelta(days=grace_days)
+ result = await db.execute(
+ select(Group).where(Group.hosted_at.is_(None), Group.created_at < cutoff))
+ return list(result.scalars().all())
+
+
+async def prune_unhosted_groups(db: AsyncSession, grace_days: int = UNHOSTED_GRACE_DAYS,
+ dry_run: bool = False) -> list[tuple[str, str]]:
+ """Delete abandoned groups. Returns [(id, name)] of what was (or would be) removed.
+
+ Memberships go with the group — there is no cascade configured, and leaving
+ orphan rows behind would keep the group in everyone's /mine query through the
+ join. Nothing on a node is touched: the hub does not command those machines,
+ and by definition no node ever claimed this group anyway.
+ """
+ doomed = await find_unhosted_groups(db, grace_days)
+ if not doomed or dry_run:
+ return [(g.id, g.name) for g in doomed]
+
+ ids = [g.id for g in doomed]
+ await db.execute(delete(GroupMember).where(GroupMember.group_id.in_(ids)))
+ await db.execute(delete(Group).where(Group.id.in_(ids)))
+ await db.commit()
+ log.info("Pruned %d group(s) that no node ever hosted", len(doomed))
+ return [(g.id, g.name) for g in doomed]