diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-08-18 02:15:02 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-08-18 02:15:02 +0200 |
| commit | e9d5e979fdab9a1cc3c729d602e6f27207b9480c (patch) | |
| tree | b5993f2c81b760ba56f251457edf84dd91ad63dc /packages/meshbay-node/src/meshbay_node/ui/app.py | |
| parent | 50ebb4f2e620dad8e1fbca8307b97c5e10e7e6c0 (diff) | |
| download | meshbay-e9d5e979fdab9a1cc3c729d602e6f27207b9480c.tar.gz | |
feat(node): several named roots per group, and one implementation per operation
Stage A — a group's content is a set of named roots
---------------------------------------------------
`shared_dir` becomes a list of {name, path, kind}. The name is the directory's
basename, derived once at add time and *stored*: recomputing it would
re-identify a whole library the day someone renames a folder on disk. Duplicate
names are refused case-insensitively and no root may contain another — both
compared with NFC folding, because most of these directories live on exFAT or
NTFS where `Films` and `films` are one directory.
Every index path carries its root name, in a one-root group as much as in a
five-root one. One path shape has to be got right once; two have to be kept
right for ever.
**A root that goes away freezes; it never empties.** Unmounting a volume makes
watchdog report every file under it as deleted, or presents an empty directory
to the next scan. Acting on either propagates deletions for a whole library to
every member, as though the owner had erased it. So a deletion is acted on only
once its root is confirmed readable, and availability is tracked per root — one
unplugged drive leaves the others serving. 12 tests, verified to fail against an
indexer without the check.
Events are not trusted to be complete either: ReadDirectoryChangesW drops them
under load and inotify on a FUSE mount misses changes made outside it. A
periodic reconciliation sweep is the only thing that recovers a missed event.
MNP 0.2 → 0.3 (additive). The hub needs no change: SwarmSource carries a content
hash, a node id and an endpoint — no paths, no filenames — and private groups
register nothing (H7).
Stage B — one implementation behind every front door
----------------------------------------------------
C1 and C6 were both "a second path into the node with its own weaker
handshake". Two implementations of `revoke` with two authorization checks is
that shape one size down. `meshbay_node/ops.py` holds each operation once,
takes the daemon state, and knows nothing about HTTP, argv or MNP. The loopback
API is one `_op(...)` line per endpoint; the MNP handlers call the same
functions. test_ops.py asserts the shape rather than trusting it.
Phase 14 is finished on top of it — `group list`, `gek init|rotate`, `reload`
(SIGHUP), `denylist show|clear`, `file list|rm`. **No operator action requires a
browser any more.** Plus `gek_rotate` and `member_unpin` as operator-signed MNP
operations: rotation is the half of revocation that revocation cannot do, since
the ex-member holds the current key, and the node generates the replacement
with its own CSPRNG — no key material crosses the wire, which is what the C5b
rule is actually about.
Two bugs found by running it rather than by testing it
------------------------------------------------------
GroupIndex is keyed by **content hash**, so the same bytes at two paths are one
entry — which is also why a scan reports ten files and indexes nine.
Reconciliation compared paths, so it decided the second path was a missed event
every 60 s, rewrote the entry and pushed an index update to every connected
peer. Seen in a live node's log.
`meshbay-node reload` crashed on first use with `subprocess` unimported: the
module compiles fine, which is the "syntax, not names" trap already recorded for
the SPA. test_cli_dispatch.py now walks every verb and refuses to let one be
added to the parser without an entry there.
Also corrected: protocol.py declared a second MNP_VERSION of "0.1" while the
wire carried "0.2" — harmless only because nothing imported it. And
_do_dir_create/_do_dir_delete referenced an undefined `filename` on their error
path.
740 tests pass; QE/deploy/e2e.py passes end to end against the live deployment.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/ui/app.py')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ui/app.py | 318 |
1 files changed, 57 insertions, 261 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py index b78f78d..050829e 100644 --- a/packages/meshbay-node/src/meshbay_node/ui/app.py +++ b/packages/meshbay-node/src/meshbay_node/ui/app.py @@ -23,6 +23,7 @@ from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Query from fastapi.responses import HTMLResponse, JSONResponse from meshbay_node import __version__ +from meshbay_node import ops from meshbay_node.config import DEFAULT_CONFIG_PATH from meshbay_common.crypto import generate_gek, wrap_gek_aes from meshbay_common.join import ROLE_MEMBER, ROLE_OPERATOR @@ -30,6 +31,24 @@ from meshbay_common.join import ROLE_MEMBER, ROLE_OPERATOR log = logging.getLogger(__name__) +def _op(coro): + """ + Run an operation and translate its refusal into a JSON response. + + The operations live in `meshbay_node.ops` and know nothing about HTTP. This + is the whole of the HTTP adapter: without it each handler would carry its + own status codes, and the MNP handler in Stage B3 would carry a second set + that slowly stopped agreeing. + """ + async def run(): + try: + return await coro() + except ops.OpError as e: + return JSONResponse(e.as_dict(), e.status) + return run() + + + def create_ui_app(state: dict) -> FastAPI: app = FastAPI( title="MeshBay Node Admin", @@ -110,90 +129,27 @@ def create_ui_app(state: dict) -> FastAPI: @app.get("/api/groups") async def api_groups(): - groups_ctx = state.get("groups_ctx", {}) - config = state.get("config") - result = [] - for gid, ctx in groups_ctx.items(): - cfg = None - if config: - cfg = next((g for g in config.groups if g.id == gid), None) - idx = ctx.get("index") - result.append({ - "id": gid, - "name": cfg.name if cfg else gid[:8], - "shared_dir": str(ctx.get("shared_root", "")), - "visibility": cfg.visibility if cfg else "private", - "file_count": idx.count if idx else 0, - "index_version": idx.version if idx else 0, - }) - return {"groups": result} - + return await _op(lambda: ops.list_groups(state)) @app.post("/api/groups/attach") async def attach_group(payload: dict): - """ - Write a new [[groups]] block into node.toml. + return await _op(lambda: ops.attach_group( + state, + (payload.get("name") or "").strip(), + (payload.get("shared_dir") or "").strip(), + )) - The name-to-id lookup happens here because this process is the one logged - into the hub. Nothing is created on the hub: the group already exists, - this only tells the node to host it. - """ - name = (payload.get("name") or "").strip() - shared_dir = (payload.get("shared_dir") or "").strip() - if not name or not shared_dir: - return JSONResponse({"error": "name and shared_dir are required"}, 400) + @app.delete("/api/groups/{group_id}/files/{file_id}") + async def delete_file(group_id: str, file_id: str): + """Milestone 14.11 — the last operator action that needed a browser.""" + return await _op(lambda: ops.delete_file(state, group_id, file_id)) - config = state.get("config") - if not config: - return JSONResponse({"error": "No config loaded"}, 503) + @app.get("/api/denylist") + async def api_denylist(): + return await _op(lambda: ops.read_denylist(state)) - hub = state.get("hub") - if not hub or not hub._session: - return JSONResponse({"error": "Hub not connected"}, 503) - try: - mine = await hub.list_my_groups() - except Exception as e: - return JSONResponse({"error": f"Could not list groups: {e}"}, 502) - - match = [g for g in mine if g["id"] == name or g["name"] == name] - if not match: - return JSONResponse({ - "error": f"No group of yours is called {name!r}", - "available": [{"name": g["name"], "id": g["id"]} for g in mine], - }, 404) - if len(match) > 1: - return JSONResponse({ - "error": f"Several of your groups are called {name!r} — use the id", - "available": [{"name": g["name"], "id": g["id"]} for g in match], - }, 409) - group = match[0] - - if any(g.id == group["id"] for g in config.groups): - return JSONResponse( - {"error": f"{group['name']!r} is already hosted by this node"}, 409) - - path = Path(shared_dir).expanduser() - try: - path.mkdir(parents=True, exist_ok=True) - except OSError as e: - return JSONResponse({"error": f"Cannot create {path}: {e}"}, 400) - - # Appended as text rather than re-serialised: node.toml is hand-written - # and full of comments explaining decisions, and a round trip through a - # TOML writer would throw all of that away. - conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) - block = (f'\n[[groups]]\n' - f'id = "{group["id"]}"\n' - f'name = "{group["name"]}"\n' - f'shared_dir = "{path}"\n' - f'visibility = "{group.get("visibility", "private")}"\n') - try: - with conf_path.open("a") as f: - f.write(block) - except OSError as e: - return JSONResponse({"error": f"Cannot write {conf_path}: {e}"}, 500) - - return {"group_id": group["id"], "name": group["name"], - "shared_dir": str(path), "config": str(conf_path)} + @app.post("/api/denylist/clear") + async def api_denylist_clear(subject: str = ""): + return await _op(lambda: ops.clear_denylist(state, subject=subject)) @app.get("/api/groups/{group_id}/files") async def api_group_files(group_id: str): @@ -279,7 +235,11 @@ def create_ui_app(state: dict) -> FastAPI: { "id": g.id, "name": g.name, - "shared_dir": g.shared_dir, + "roots": [ + {"path": r.path, "name": r.name, "kind": r.kind, + "upload": r.upload} + for r in g.roots + ], "visibility": g.visibility, } for g in config.groups @@ -290,205 +250,33 @@ def create_ui_app(state: dict) -> FastAPI: @app.post("/api/operator/pair") async def operator_pair(): - """ - Issue a one-time code that pairs a browser as this node's operator. - - The code is the whole point: it binds the operator's browser identity key - to their account without asking the hub, which is what stops a hub from - naming itself node administrator (M3, and the same substitution as H3). - It is returned once and stored only as a hash. - """ - roster = state.get("roster") - user_id = state.get("node_user_id") - if not roster or not user_id: - return JSONResponse({"error": "Node not connected to hub yet"}, 503) - - config = state.get("config") - ttl = (config.node.pair_ttl_hours if config else 24) * 3600 - code = await roster.create_invite( - group_id="", # operator authority is node-wide - user_id=user_id, - role=ROLE_OPERATOR, - created_by="local-cli", - ttl=ttl, - username=(config.hub.username if config else ""), - ) - invites = await roster.list_invites() - expires = next((i["expires_at"] for i in invites - if i["user_id"] == user_id and i["role"] == ROLE_OPERATOR), "") - return {"code": code, "expires_at": expires, "user_id": user_id} + return await _op(lambda: ops.pair_operator(state)) @app.get("/api/roster") async def api_roster(group_id: str = ""): - roster = state.get("roster") - if not roster: - return {"identities": [], "members": [], "invites": []} - return { - "identities": await roster.list_identities(), - "members": await roster.list_members(group_id or None), - "invites": await roster.list_invites(), - } + return await _op(lambda: ops.read_roster(state, group_id)) @app.post("/api/groups/{group_id}/invites") async def create_invite(group_id: str, username: str): - """ - Issue an invitation code from the CLI, without a browser. - - The hub is asked for the account id and nothing else — never for a key. - A hub that answered with the wrong account would produce an invite whose - code it never learns, since the code goes to a human out of band. - """ - roster = state.get("roster") - groups_ctx = state.get("groups_ctx", {}) - if not roster: - return JSONResponse({"error": "Roster not available"}, 503) - if group_id not in groups_ctx: - return JSONResponse({"error": "Group not hosted on this node"}, 404) - - hub = state.get("hub") - if not hub or not hub._session: - return JSONResponse({"error": "Hub not connected"}, 503) - try: - account = await hub.get_user_pubkeys(username) - except Exception as e: - return JSONResponse({"error": f"Unknown user {username!r}: {e}"}, 404) - - config = state.get("config") - ttl = (config.node.invite_ttl_hours if config else 168) * 3600 - code = await roster.create_invite( - group_id=group_id, - user_id=account["user_id"], - role=ROLE_MEMBER, - created_by="local-cli", - ttl=ttl, - username=username, - ) - invites = await roster.list_invites() - expires = next((i["expires_at"] for i in invites - if i["user_id"] == account["user_id"] - and i["group_id"] == group_id), "") - return {"code": code, "expires_at": expires, - "username": username, "user_id": account["user_id"]} + return await _op(lambda: ops.create_invite(state, group_id, username)) @app.get("/api/resolve") async def resolve_user(username: str): - """ - Map a username to an account id for the CLI. - - The roster answers first — it is the node's own record. The hub is the - fallback for identities pinned before invitations carried a name, and for - people admitted through an open-join group. Only an account id comes back; - no key is ever taken from here. - """ - roster = state.get("roster") - if roster: - for ident in await roster.list_identities(): - if ident["username"] == username: - return {"user_id": ident["user_id"], "source": "roster"} - hub = state.get("hub") - if hub and hub._session: - try: - account = await hub.get_user_pubkeys(username) - return {"user_id": account["user_id"], "source": "hub"} - except Exception: - pass - return JSONResponse({"error": f"Unknown user {username!r}"}, 404) + return await _op(lambda: ops.resolve_user(state, username)) @app.post("/api/members/{user_id}/revoke") async def revoke_member(user_id: str, group_id: str): - """ - Stop serving the group key to someone. - - Takes effect on their next connection: the key is wrapped on demand, so - there is no stored bundle left behind that would outlive this. Rotating - the group key is still required — they hold the current one. - """ - roster = state.get("roster") - if not roster: - return JSONResponse({"error": "Roster not available"}, 503) - if not await roster.set_status(group_id, user_id, "revoked"): - return JSONResponse({"error": "No such member in that group"}, 404) - log.info("Member revoked: user=%s group=%s", user_id[:8], group_id[:8]) - return {"status": "revoked", "user_id": user_id, "group_id": group_id, - "reminder": "rotate the group key: meshbay-node gek-init"} + return await _op(lambda: ops.revoke_member(state, user_id, group_id)) @app.post("/api/members/{user_id}/unpin") async def unpin_member(user_id: str): - """Forget a pinned identity, so the person can pair again with a new key.""" - roster = state.get("roster") - if not roster: - return JSONResponse({"error": "Roster not available"}, 503) - if not await roster.unpin(user_id): - return JSONResponse({"error": "No such pinned identity"}, 404) - log.info("Identity unpinned: user=%s", user_id[:8]) - return {"status": "unpinned", "user_id": user_id} + return await _op(lambda: ops.unpin_member(state, user_id)) # ── GEK initialization (operator only, localhost) ────────────────────── @app.post("/api/groups/{group_id}/gek") - async def init_gek(group_id: str): - """ - Generate the group key and activate it. - - It used to be wrapped here for every member, using public keys fetched from - the hub — which is H3 with the node as the victim instead of the inviter: a - hub answering with its own key was handed the group key by the node itself. - - Nothing is pre-wrapped for members now. Each member's copy is produced when - they connect, for a key they proved they hold (`join_request`). Only the - node's own copy is stored, so the daemon can reload the key across restarts - without the operator's browser. - """ - groups_ctx = state.get("groups_ctx", {}) - if group_id not in groups_ctx: - return JSONResponse({"error": "Group not hosted on this node"}, 404) - - hub = state.get("hub") - if not hub or not hub._session: - return JSONResponse({"error": "Hub not connected"}, 503) - - bundle_store = state.get("bundle_store") - if not bundle_store: - return JSONResponse({"error": "Bundle store not available"}, 503) - - existing_gek = groups_ctx[group_id].get("gek") - gek = existing_gek or generate_gek() - errors: list[str] = [] - - roster = state.get("roster") - authorized = len(await roster.list_members(group_id)) if roster else 0 - - # Store a copy wrapped for the node keystore X25519 key so the daemon can - # reload the GEK on restart without the operator's browser keys. - node_user_id = hub._session.user_id if hub._session else None - pk_x_node_raw = state.get("pk_x25519_raw") - if pk_x_node_raw and node_user_id: - try: - node_bundle = wrap_gek_aes(gek, pk_x_node_raw) - await bundle_store.store( - group_id, f"_node_{node_user_id}", - node_bundle["pk_eph_b64"], node_bundle["nonce_b64"], - node_bundle["wrapped_b64"], - ) - log.info("GEK wrapped for node keystore (daemon reload)") - except Exception as e: - errors.append(f"node keystore: {e}") - log.warning("Failed to wrap GEK for node keystore: %s", e) - - groups_ctx[group_id]["gek"] = gek - log.info("GEK initialized for group %s — %d authorized member(s) will " - "receive it on connect", group_id[:8], authorized) - - webrtc = state.get("webrtc") - if webrtc and "groups" in webrtc._ctx and group_id in webrtc._ctx["groups"]: - webrtc._ctx["groups"][group_id]["gek"] = gek - - return { - "status": "ok", - "group_id": group_id, - "authorized_members": authorized, - "errors": errors, - } + async def init_gek(group_id: str, rotate: bool = False): + return await _op(lambda: ops.set_gek(state, group_id, rotate=rotate)) # ── Chat endpoints ─────────────────────────────────────────────────────── @@ -656,7 +444,15 @@ def _render_page(state: dict, roster_view: dict | None = None) -> str: cfg = next((g for g in config.groups if g.id == gid), None) idx = ctx.get("index") name = cfg.name if cfg else gid[:8] - shared = ctx.get("shared_root", "") + roots = ctx.get("roots") + # An unavailable root is shown as such rather than hidden: its files are + # still listed and still in the index, and hiding the root would make a + # frozen library look deleted — the exact confusion this is meant to + # prevent. + shared = ", ".join( + f"{r.name} → {r.path}" + ("" if r.available else " [UNAVAILABLE]") + for r in roots + ) if roots else "" vis = cfg.visibility if cfg else "private" fcount = idx.count if idx else 0 total_size = sum(e.size for e in idx.entries) if idx else 0 |