diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/ops')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/__init__.py | 167 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/apps.py | 295 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/chat.py | 254 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/core.py | 57 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/files.py | 114 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/groups.py | 248 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/members.py | 314 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/node_toml.py | 204 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/roots.py | 285 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/settings.py | 340 |
10 files changed, 2278 insertions, 0 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/ops/__init__.py b/packages/meshbay-node/src/meshbay_node/ops/__init__.py new file mode 100644 index 0000000..bfdfd7b --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/ops/__init__.py @@ -0,0 +1,167 @@ +""" +Operator operations — one implementation, several front doors. + +Three things ask this node to act: the CLI (over the loopback admin API), the +local admin UI, and — from Stage B3 — signed MNP messages from a paired client. +They must agree, and the way to make them agree is not to write the operation +three times and hope. + +**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 +the same shape one size down. So each operation lives here once, takes the +daemon's `state`, and knows nothing about HTTP, argv or MNP. The adapters +translate: `ui/app.py` turns `OpError` into a JSON response, the CLI prints it, +the MNP handler sends an error frame. + +**Authorization is not here.** Reaching this module already means the caller got +past its adapter's check — the loopback session token (11.5.3) for the API, an +Ed25519 signature verified against the roster for MNP. These functions do what +they are told; deciding who may tell them is the adapter's job and stays visible +in the adapter. +""" + +from meshbay_node.ops.apps import ( + _validate_app_dirs, + rematch_video, + set_app_directories, + set_app_directory, + set_chat_directory, + set_chat_link_preview, + set_enabled_apps, + set_musicbrainz_enabled, + set_search_listed, + set_tmdb_config, + set_tmdb_enabled, +) +from meshbay_node.ops.chat import ( + _wrap_for_node, + chat_epoch_keys, + chat_status, + encrypt_chat_history, + ensure_chat_epoch, + open_chat_epoch, + prune_chat, +) +from meshbay_node.ops.core import ( + OpError, + _config, + _group_ctx, + _hub, + _roster, +) +from meshbay_node.ops.files import ( + delete_file, + index_cache_stats, + prune_index_cache, +) +from meshbay_node.ops.groups import ( + attach_group, + detach_group, + list_groups, + set_gek, +) +from meshbay_node.ops.members import ( + _invite_url, + cancel_invite, + cancel_link_invitation, + create_invite, + create_link_invitation, + create_link_invite, + pair_operator, + read_roster, + resolve_user, + revoke_member, + unpin_member, +) +from meshbay_node.ops.node_toml import ( + _find_group_range, + _insert_roots_block, + _remove_roots_block, + _update_node_toml, + _update_root_field, +) +from meshbay_node.ops.roots import ( + add_root, + eject_root, + plug_root, + remove_root, + update_root, +) +from meshbay_node.ops.settings import ( + NODE_SETTING_WRITERS, + _group_limits, + clear_denylist, + get_node_settings, + list_transfers, + read_denylist, + reload_config, + set_node_settings, + set_scan_settings, + set_transfer_limits, + start_reload, +) + +__all__ = [ + "NODE_SETTING_WRITERS", + "OpError", + "_config", + "_find_group_range", + "_group_ctx", + "_group_limits", + "_hub", + "_insert_roots_block", + "_invite_url", + "_remove_roots_block", + "_roster", + "_update_node_toml", + "_update_root_field", + "_validate_app_dirs", + "_wrap_for_node", + "add_root", + "attach_group", + "cancel_invite", + "cancel_link_invitation", + "chat_epoch_keys", + "chat_status", + "clear_denylist", + "create_invite", + "create_link_invitation", + "create_link_invite", + "delete_file", + "detach_group", + "eject_root", + "encrypt_chat_history", + "ensure_chat_epoch", + "get_node_settings", + "index_cache_stats", + "list_groups", + "list_transfers", + "open_chat_epoch", + "pair_operator", + "plug_root", + "prune_chat", + "prune_index_cache", + "read_denylist", + "read_roster", + "reload_config", + "rematch_video", + "remove_root", + "resolve_user", + "revoke_member", + "set_app_directories", + "set_app_directory", + "set_chat_directory", + "set_chat_link_preview", + "set_enabled_apps", + "set_gek", + "set_musicbrainz_enabled", + "set_node_settings", + "set_scan_settings", + "set_search_listed", + "set_tmdb_config", + "set_tmdb_enabled", + "set_transfer_limits", + "start_reload", + "unpin_member", + "update_root", +] diff --git a/packages/meshbay-node/src/meshbay_node/ops/apps.py b/packages/meshbay-node/src/meshbay_node/ops/apps.py new file mode 100644 index 0000000..0e4f6dd --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/ops/apps.py @@ -0,0 +1,295 @@ +"""Per-group application settings: which apps, their directories, TMDB, MusicBrainz, chat.""" + +from __future__ import annotations + +import logging + +from meshbay_common.background import spawn + +from meshbay_node.ops.core import OpError, _group_ctx, _roster +from meshbay_node.roots import RootSet + +log = logging.getLogger("meshbay_node.ops") + + +# ── Applications ───────────────────────────────────────────────────────────── + +async def set_enabled_apps(state: dict, group_id: str, apps: list[str]) -> dict: + """ + Which group "applications" (Chat, Files, ...) are shown to members. + + Same shape as other signed ops: lives on the node (roster.db), takes + effect without a restart, and is signed by the operator (webrtc_server.py + checks the caller's own admin-authority allow-list before this runs). + """ + roster = _roster(state) + ctx = _group_ctx(state, group_id) + # See the same guard in webrtc/group_ops.py _do_apps_enabled: Files cannot be + # turned off, and both writers put it at the front so the two agree. + if "files" not in apps: + apps = ["files"] + list(apps) + await roster.set_enabled_apps(group_id, apps, + set_by=state.get("node_user_id", "")) + ctx["enabled_apps"] = apps + log.info("Enabled apps for group %s: %s", group_id[:8], ",".join(sorted(apps))) + return {"apps": apps, "group_id": group_id} + + +# ── TMDB config (Videos app) ───────────────────────────────────────────────── + +async def set_tmdb_config(state: dict, token: str | None = None, + language: str | None = None) -> dict: + """ + Whether the node uses a custom API token instead of the shipped default, + and in what language it queries TMDB (docs/MESHBAY_DESIGN.md §9.7). + + Node-wide (roster.py group_settings, group_id="") rather than per-group + like set_enabled_apps: the token and the shared-cache + language are one operator's budget and one credential, not a per-group + or per-viewer concern. Whether TMDB is used *at all* is the per-group + decision set_tmdb_enabled below makes instead. `token=""` explicitly + clears a previously-set custom token (reverts to the shipped default); + `token=None` leaves whatever was there unchanged. Same discipline for + `language`. + """ + roster = _roster(state) + await roster.set_tmdb_config(token, language, set_by=state.get("node_user_id", "")) + # `token=None` means "leave whatever was there" (§ set_tmdb_config's own + # docstring) — so the customized flag only changes when a value (a real + # token, or "" to clear one) was actually given. + if token is not None: + state["tmdb_token_customized"] = bool(token) + if language is not None: + state["tmdb_language"] = language + log.info("TMDB config: custom_token=%s language=%s", + bool(token), language or state.get("tmdb_language", "")) + return { + "token_customized": state.get("tmdb_token_customized", False), + "language": state.get("tmdb_language", ""), + } + + +async def set_tmdb_enabled(state: dict, group_id: str, enabled: bool) -> dict: + """ + Whether TMDB lookups run for this group at all (docs/MESHBAY_DESIGN.md + §9.7) — per-group, unlike set_tmdb_config above: an operator running a + real media library alongside test/demo groups on one node wants + outbound TMDB traffic (and API quota) spent for the one that needs it, + not all of them just because one process serves both. + """ + roster = _roster(state) + ctx = _group_ctx(state, group_id) + await roster.set_tmdb_enabled(group_id, enabled, set_by=state.get("node_user_id", "")) + ctx["tmdb_enabled"] = enabled + log.info("TMDB enabled for group %s: %s", group_id[:8], enabled) + return {"enabled": enabled, "group_id": group_id} + + +# ── MusicBrainz config (Music app) ─────────────────────────────────────────── + +# set_musicbrainz_config removed — MusicBrainz contact is now the owner's +# hub email, resolved at login (daemon.py / musicbrainz.py). + +async def set_musicbrainz_enabled(state: dict, group_id: str, enabled: bool) -> dict: + """ + Whether MusicBrainz lookups run for this group at all + (docs/MESHBAY_DESIGN.md §9.8) — per-group from the start, same reasoning as + set_tmdb_enabled: a real media-library group and a test/demo group on + one node need not share the decision to make outbound requests. + """ + roster = _roster(state) + ctx = _group_ctx(state, group_id) + await roster.set_musicbrainz_enabled(group_id, enabled, set_by=state.get("node_user_id", "")) + ctx["musicbrainz_enabled"] = enabled + log.info("MusicBrainz enabled for group %s: %s", group_id[:8], enabled) + return {"enabled": enabled, "group_id": group_id} + + +# ── App directories ────────────────────────────────────────────────────────── + +def _validate_app_dirs(state: dict, group_id: str, paths: list[str], *, + require_writable: bool) -> list[str]: + """ + Every path an app is pointed at must live inside one of the group's roots. + + The per-app setters this replaces validated nothing: a typo, or a path left + behind by a root that was removed, was stored and then quietly matched no + entry — an app showing an empty tab with no way to tell "misconfigured" + from "no files yet". Refusing at the point of setting is the only moment + the operator is present to be told. + + Not `RootSet.resolve()`, deliberately: that also refuses a directory whose + root is currently *unavailable*, and an operator must be able to configure + a library on a drive they have unplugged. What is checked here is the + shape — inside a named root, no traversal — which does not change with + what happens to be mounted. + """ + roots: RootSet | None = _group_ctx(state, group_id).get("roots") + if roots is None: + raise OpError("Group has no roots", status=503) + + clean: list[str] = [] + for raw in paths: + path = str(raw or "").strip().strip("/") + if not path: + continue + if ".." in path.split("/"): + raise OpError(f"{path!r} is not a directory inside this group", + status=400) + found = roots.split(path) + if found is None: + raise OpError( + f"{path!r} is not inside any of this group's shared " + f"directories", status=400, + extra={"available": roots.names}) + root, _tail = found + if require_writable and not root.writable: + raise OpError( + f"{root.name!r} is read-only, and this setting needs a " + f"directory that accepts uploads", status=400) + clean.append(path) + return sorted(set(clean)) + + +async def set_app_directories(state: dict, group_id: str, app_key: str, + paths: list[str], *, + require_writable: bool = False) -> dict: + """ + Which folder(s) inside the group's shared roots an application works over. + + One function for every app, keyed by the app's own name: adding an + application is a registry entry and a settings component, not another + near-identical op here — one per app differing only in the key it wrote + and whether it took a string or a list. + + Empty means nothing configured, which every app reads as "show nothing + until an operator has chosen" — never "the whole group index". Pointing an + app at the whole library is a decision, not a default nobody made. + + A change always fires (never awaits) a sweep of what the new directories + already contain: the ordinary per-change enrichment path only looks at + entries new since the last broadcast, so files already sitting in a folder + when it was chosen would otherwise never be picked up. + """ + roster = _roster(state) + ctx = _group_ctx(state, group_id) + clean = _validate_app_dirs(state, group_id, paths, + require_writable=require_writable) + await roster.set_app_directories(group_id, app_key, clean, + set_by=state.get("node_user_id", "")) + ctx[f"{app_key}_directories"] = clean + # An app whose directories are also published under a second name (chat's + # single destination) has that name re-derived here: leaving it behind + # would make the two disagree within a single run, and only until a restart + # — the shape of bug that reads as "it works after a restart". + from meshbay_node.roster import Roster + alias = Roster.ctx_alias(app_key, clean) + if alias: + ctx[alias[0]] = alias[1] + log.info("%s directories for group %s: %s", app_key, group_id[:8], + ", ".join(clean) or "(none)") + + enrich = (state.get("enrich_app_dirs_fns") or {}).get(app_key) + if enrich: + spawn(enrich(group_id)) + return {"app": app_key, "directories": clean, "group_id": group_id} + + +async def set_app_directory(state: dict, group_id: str, app_key: str, + path: str, *, + require_writable: bool = False) -> dict: + """ + The single-directory form, for an app that only ever wants one. + + Stored as a one-element list like every other app, because two storage + shapes for one idea is what made `video_root` (scalar) and `photo_roots` + (list) need separate ops, separate MNP messages and separate widgets to + say the same thing. `path=""` clears it. + """ + result = await set_app_directories( + state, group_id, app_key, [path] if path else [], + require_writable=require_writable) + dirs = result["directories"] + return {**result, "path": dirs[0] if dirs else ""} + + +# ── Chat ───────────────────────────────────────────────────────────────────── + +async def set_chat_directory(state: dict, group_id: str, path: str) -> dict: + """ + Where chat attachments are written. + + `require_writable`, unlike every other app directory: this one is a + *destination*, not a view. Pointing it at a read-only root would produce an + attachment button that fails at the moment somebody uses it, which is the + failure mode the RO/RW model exists to move earlier. + """ + return await set_app_directory(state, group_id, "chat", path, + require_writable=True) + + +async def set_chat_link_preview(state: dict, group_id: str, + enabled: bool) -> dict: + """ + Whether the node fetches a page's title and image when a member posts a + link. + + Outbound third-party traffic on the operator's connection, caused by a + message they did not write and pointing at a URL they did not choose — so + it is theirs to switch off, on the same reasoning as the per-group TMDB + switch. Absent means on, because that is what the node did before this + existed. + """ + roster = _roster(state) + ctx = _group_ctx(state, group_id) + await roster.set_chat_link_preview(group_id, enabled, + set_by=state.get("node_user_id", "")) + ctx["chat_link_preview"] = enabled + log.info("Chat link previews for group %s: %s", group_id[:8], + "on" if enabled else "off") + return {"enabled": enabled, "group_id": group_id} + + +async def set_search_listed(state: dict, group_id: str, listed: bool) -> dict: + """ + Whether this group's files appear in members' cross-group Search. + + A presentation choice, and it must never be described as more: a member + still lists the whole group by opening it, the node serves the index + exactly as before, and a client that ignores the flag lists the group in + Search too. What it buys is a family album not turning up in the middle of + a film library. Absent means listed. + """ + roster = _roster(state) + ctx = _group_ctx(state, group_id) + await roster.set_search_listed(group_id, listed, + set_by=state.get("node_user_id", "")) + ctx["search_listed"] = listed + log.info("Search listing for group %s: %s", group_id[:8], + "on" if listed else "off") + return {"listed": listed, "group_id": group_id} + + +# ── Videos: force TMDB re-matching ─────────────────────────────────────────── +# +# `media_cache.file_tmdb` is keyed by a file's content hash and is otherwise +# only pruned on deletion, so a fixed matcher/parser never dislodges a match +# already in cache. This drops a group's *auto-resolved* mappings so the +# next `media_meta_req` for each poster tile re-resolves against the current +# code. Re-resolution is lazy and calls TMDB once per unique title — real +# API budget — so this is an explicit operator action, never a background job. +# Manual "Fix match" corrections (media_cache.tmdb_override) are kept. + +async def rematch_video(state: dict, group_id: str) -> dict: + media_cache = state.get("media_cache") + if media_cache is None: + raise OpError("No media cache in this process", status=503) + indexer = (state.get("indexers") or {}).get(group_id) + if indexer is None: + raise OpError("Unknown group", status=404) + file_ids = [e.id for e in indexer.index.entries if e.type == "video"] + removed = await media_cache.clear_tmdb_matches(file_ids) + log.info("Video rematch for group %s: %d auto match(es) cleared across %d video file(s)", + group_id[:8], removed, len(file_ids)) + return {"status": "cleared", "removed": removed, "videos": len(file_ids), + "group_id": group_id} diff --git a/packages/meshbay-node/src/meshbay_node/ops/chat.py b/packages/meshbay-node/src/meshbay_node/ops/chat.py new file mode 100644 index 0000000..469e907 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/ops/chat.py @@ -0,0 +1,254 @@ +"""A group's chat keys (epochs) and its stored history.""" + +from __future__ import annotations + +import logging +import time as _time + +from meshbay_common.chatbox import new_epoch_key +from meshbay_common.crypto import unwrap_gek_aes, wrap_gek_aes + +from meshbay_node.ops.core import OpError, _group_ctx + +log = logging.getLogger("meshbay_node.ops") + + +# ── Chat epoch keys ────────────────────────────────────────────────────────── +# +# The key a group's chat archive is encrypted under. Generated here, by the +# node, and never by a member — the C5b rule is about key material arriving from +# outside, and this is the same rule that lets `gek_rotate` be a signed +# instruction rather than a delivery. +# +# An *epoch* rather than a rotation, and the distinction is the whole design: +# opening a new one stops a departing member reading what comes next, while +# every earlier epoch is kept and still delivered to current members, so the +# history they could already read stays readable. Rotating instead — replacing +# the key, as `set_gek` does — would make every message anyone ever sent +# permanently unreadable to everybody, which is what a plain GEK-derived +# archive key would have done on the very first `member unpin` +# (finding F4, docs/MESHBAY_DESIGN.md §13.6). + + +async def _wrap_for_node(state: dict, key: bytes) -> dict: + """ + Wrap a key to the node's own X25519 key, the way `set_gek` does for the GEK. + + Wrapped, not raw: the claim chat encryption makes is against someone who + obtains the node's storage *without the keystore password*, and the node's + X25519 private key is what the keystore protects. A raw key in SQLite would + leave nothing behind that claim. + """ + pk_x_node_raw = state.get("pk_x25519_raw") + if not pk_x_node_raw: + raise OpError("Node identity not available", status=503) + return wrap_gek_aes(key, pk_x_node_raw) + + +async def chat_epoch_keys(state: dict, group_id: str) -> list[dict]: + """ + Every chat epoch key this group has, oldest first, in the clear *in memory*. + + Cached on the group context: unwrapping is an ECIES operation per epoch and + this is on the path of every member connecting to a group with chat on. + """ + ctx = _group_ctx(state, group_id) + cached = ctx.get("chat_epoch_keys") + if cached is not None: + return cached + + bundle_store = state.get("bundle_store") + if not bundle_store: + raise OpError("Bundle store not available", status=503) + sk_x_raw = state.get("sk_x25519_raw") + pk_x_raw = state.get("pk_x25519_raw") + if not (sk_x_raw and pk_x_raw): + raise OpError("Node identity not available", status=503) + + keys: list[dict] = [] + for row in await bundle_store.fetch_chat_epochs(group_id): + try: + keys.append({"epoch": row["epoch"], + "key": unwrap_gek_aes(row, sk_x_raw, pk_x_raw)}) + except Exception as e: + # Loud, and not fatal: one unreadable epoch must not take the + # readable ones with it. The messages of that epoch are lost, which + # is a thing the operator needs told rather than a thing to hide. + log.error("chat: epoch %d of group %s will not unwrap (%s) — " + "its messages are unreadable", row["epoch"], + group_id[:8], e) + ctx["chat_epoch_keys"] = keys + return keys + + +async def open_chat_epoch(state: dict, group_id: str) -> dict: + """ + Open a new chat epoch. Idempotent only in the sense that it always adds one. + + Called when the set of devices that may read *future* messages shrinks: a + member removed, a device revoked or unpinned, the group key rotated, or the + operator asking directly. Never on a schedule — an epoch nobody needed is an + epoch key the node has to keep for ever. + """ + bundle_store = state.get("bundle_store") + if not bundle_store: + raise OpError("Bundle store not available", status=503) + + epoch = await bundle_store.latest_chat_epoch(group_id) + 1 + key = new_epoch_key() + wrapped = await _wrap_for_node(state, key) + await bundle_store.store_chat_epoch( + group_id, epoch, wrapped["pk_eph_b64"], wrapped["nonce_b64"], + wrapped["wrapped_b64"]) + + # Tolerant of a group context that does not exist yet: the daemon opens the + # first epoch **while it is building** `groups_ctx`, before publishing it on + # the state, because a group with no epoch key is a group nobody can speak + # in. Insisting on the context here would make start-up the one moment this + # cannot be called. + ctx = (state.get("groups_ctx") or {}).get(group_id) + if ctx is not None: + cached = ctx.get("chat_epoch_keys") + if cached is not None: + cached.append({"epoch": epoch, "key": key}) + ctx["chat_epoch"] = epoch + + # The transports hold their own view of the group, exactly as `set_gek` + # notes: an epoch that did not reach them would have members sealing under + # a key the node no longer thinks is current. + for transport_key in ("webrtc", "quic_server"): + transport = state.get(transport_key) + groups = getattr(transport, "_ctx", {}).get("groups") if transport else None + if groups and group_id in groups: + groups[group_id]["chat_epoch"] = epoch + groups[group_id].pop("chat_epoch_keys", None) + + log.info("Chat epoch %d opened for group %s", epoch, group_id[:8]) + return {"epoch": epoch} + + +async def ensure_chat_epoch(state: dict, group_id: str) -> int: + """The current epoch, opening the first one if the group has none.""" + bundle_store = state.get("bundle_store") + if not bundle_store: + raise OpError("Bundle store not available", status=503) + epoch = await bundle_store.latest_chat_epoch(group_id) + if epoch: + return epoch + return (await open_chat_epoch(state, group_id))["epoch"] + + +async def chat_status(state: dict, group_id: str) -> dict: + """What the operator needs to decide anything about this group's chat.""" + ctx = _group_ctx(state, group_id) + bundle_store = state.get("bundle_store") + store = ctx.get("chat_store") + plain = sealed = 0 + if store is not None: + plain, sealed = await store.count_by_format() + return { + "group_id": group_id, + "epoch": (await bundle_store.latest_chat_epoch(group_id) + if bundle_store else 0), + # Rows written before MNP 2.0. Not a state the node can be *in* — chat + # is always encrypted now — but a state its disk can be in until + # `chat encrypt-history` has run, and the operator has to be told, + # because those messages are the ones still readable off a stolen disk. + "plaintext_messages": plain, + "encrypted_messages": sealed, + } + + +async def encrypt_chat_history(state: dict, group_id: str) -> dict: + """ + Re-encrypt the messages written before this group turned encryption on. + + Deliberately **not** done by the switch. It rewrites the only copy of a + conversation, and a toggle that does that is one somebody flips twice; this + is an explicit command, it copies the database first, and it runs in one + transaction. + + The node can do this at all only because it holds those rows in plaintext — + it is the last moment at which anyone can. Afterwards nothing on this + machine can read them without an epoch key. + + Messages are sealed under a **synthetic device** belonging to the node, not + under the original sender's key: the node does not hold anyone's signing key + and must not pretend to. They are marked as such, so a reader is told these + carry the node's word for who wrote them — which is all they ever carried, + since they were written before signing existed. + """ + import shutil + + from meshbay_common.chatbox import seal + + ctx = _group_ctx(state, group_id) + store = ctx.get("chat_store") + if store is None: + raise OpError("This group has no chat store", status=404) + + epoch = await ensure_chat_epoch(state, group_id) + keys = {k["epoch"]: k["key"] for k in await chat_epoch_keys(state, group_id)} + key = keys.get(epoch) + if not key: + raise OpError("No chat key for this group", status=503) + + sk_node = state.get("sk_node") + if sk_node is None: + raise OpError("Node identity not available", status=503) + from cryptography.hazmat.primitives import serialization + + device_raw = sk_node.public_key().public_bytes( + serialization.Encoding.Raw, serialization.PublicFormat.Raw) + import base64 as _b64 + + device_b64 = _b64.b64encode(device_raw).decode() + + backup = store.db_path.with_name( + f"{store.db_path.name}.bak-{int(_time.time())}") + shutil.copy2(store.db_path, backup) + + converted = 0 + for row in await store.all_plaintext(): + text = (row.payload.decode("utf-8", errors="replace") + if isinstance(row.payload, bytes) else str(row.payload)) + env = seal(key, group_id, epoch, device_b64, device_raw, sk_node, { + "text": text, + "thread_id": row.thread_id, + "sender_name": row.sender_name, + "sent_at": int(row.timestamp), + # The node sealed this after the fact; it did not witness it being + # signed. Said in the payload rather than inferred from the device. + "migrated": True, + }) + await store.reseal(row.id, epoch=epoch, device=device_raw, + nonce=env["nonce"], ct=env["ct"], sig=env["sig"]) + converted += 1 + await store.commit() + + log.info("Chat history re-encrypted for group %s: %d message(s), backup %s", + group_id[:8], converted, backup.name) + return {"group_id": group_id, "converted": converted, + "backup": str(backup), "epoch": epoch} + + +async def prune_chat(state: dict, group_id: str, max_age_days: int) -> dict: + """ + Delete messages older than `max_age_days`. Epoch keys are never touched. + + An epoch whose messages have all aged out costs 32 bytes and keeps the + operation reversible in the only direction that matters: nothing that is + still stored becomes unreadable because something else was deleted. + """ + ctx = _group_ctx(state, group_id) + store = ctx.get("chat_store") + if store is None: + raise OpError("This group has no chat store", status=404) + if max_age_days < 1: + raise OpError("max_age_days must be at least 1", status=400) + removed = await store.delete_older_than( + _time.time() - max_age_days * 86400) + log.info("Chat retention for group %s: %d message(s) removed", + group_id[:8], removed) + return {"group_id": group_id, "removed": removed, + "max_age_days": max_age_days} diff --git a/packages/meshbay-node/src/meshbay_node/ops/core.py b/packages/meshbay-node/src/meshbay_node/ops/core.py new file mode 100644 index 0000000..7c19d78 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/ops/core.py @@ -0,0 +1,57 @@ +"""What every operation shares: the refusal it raises and the lookups into `state`.""" + +from __future__ import annotations + +from typing import Any + + +class OpError(Exception): + """ + An operation refused, with enough for any adapter to report it. + + `status` is an HTTP code because one adapter needs one; the others ignore it. + `extra` carries the "here is what would have worked" payload — a bare "no + such group" leaves an operator guessing at a UUID. + """ + + def __init__(self, message: str, *, status: int = 400, + extra: dict[str, Any] | None = None): + super().__init__(message) + self.message = message + self.status = status + self.extra = extra or {} + + def as_dict(self) -> dict: + return {"error": self.message, **self.extra} + + +# ── Shared lookups ─────────────────────────────────────────────────────────── + +def _roster(state: dict): + roster = state.get("roster") + if not roster: + raise OpError("Roster not available", status=503) + return roster + + +def _hub(state: dict): + hub = state.get("hub") + if not hub or not hub._session: + raise OpError("Hub not connected", status=503) + return hub + + +def _group_ctx(state: dict, group_id: str) -> dict: + groups_ctx = state.get("groups_ctx", {}) + if group_id not in groups_ctx: + raise OpError("Group not hosted on this node", status=404, + extra={"available": [ + {"id": gid} for gid in groups_ctx]}) + return groups_ctx[group_id] + + +def _config(state: dict): + config = state.get("config") + if not config: + raise OpError("No config loaded", status=503) + return config diff --git a/packages/meshbay-node/src/meshbay_node/ops/files.py b/packages/meshbay-node/src/meshbay_node/ops/files.py new file mode 100644 index 0000000..4609627 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/ops/files.py @@ -0,0 +1,114 @@ +"""Files on disk and the index cache.""" + +from __future__ import annotations + +import asyncio +import logging +from pathlib import Path + +from meshbay_node.ops.core import OpError, _group_ctx + +log = logging.getLogger("meshbay_node.ops") + + +# ── Files ──────────────────────────────────────────────────────────────────── + +async def delete_file(state: dict, group_id: str, file_id: str) -> dict: + """ + Remove a file from a group. Milestone 14.11 — the last operator action that + needed a browser. + + Authorization happened in the adapter. On the loopback path that is the + session token, which means physical or SSH access to the machine hosting the + files — an operator who can run this can also `rm` the file, so the check is + not weaker than the alternative. + """ + ctx = _group_ctx(state, group_id) + index = ctx.get("index") + roots = ctx.get("roots") + if not index or not roots: + raise OpError("Group has no index", status=503) + + entry = index.get_entry(file_id) + if not entry: + raise OpError("No such file in this group", status=404) + + from meshbay_node.roots import entry_abs_path + path = entry_abs_path(roots, entry) + if path is None: + raise OpError( + f"{entry.name!r} is in root {entry.path.split('/')[0]!r}, which is " + f"not readable right now — the file is frozen, not gone", status=409) + + try: + path.unlink() + except FileNotFoundError: + # Already gone from disk; drop the stale entry rather than refusing. + log.warning("Index named a file that is not on disk: %s", path) + except OSError as e: + raise OpError(f"Cannot delete {entry.name!r}: {e}", status=500) from e + + index.remove_entry(file_id) + log.info("File deleted by operator: %s/%s", entry.path, entry.name) + return {"status": "deleted", "name": entry.name, "path": entry.path, + "group_id": group_id} + + +# ── Index cache maintenance ─────────────────────────────────────────────────── +# +# The (path, size, mtime) -> hash accelerator (indexer/cache.py) is node-wide +# and grows for as long as a path was ever seen — a folder an operator later +# stops sharing (root removed, or every group hosting it is deleted) leaves +# its rows behind forever otherwise. Nothing about correctness needs this: +# a stale row just sits unused (lookup() keys on the live path string, so a +# path nothing scans any more is never looked up). This is disk space +# hygiene the operator can run when they want it, not a background job. + +async def index_cache_stats(state: dict) -> dict: + """Row count only — cheap, safe to call on every dashboard render. + The actual staleness check (prune_index_cache) is not this cheap and + must never run implicitly.""" + cache = state.get("index_cache") + return {"count": await cache.count() if cache else 0} + + +async def prune_index_cache(state: dict) -> dict: + """ + Drop cache rows that cannot be right for anything any more: the path is + not under any group's root at all, or it is under a root that is + available right now and the file is genuinely gone from disk. + + Deliberately leaves alone anything under a root that is currently + *unavailable* (a disconnected drive) — indexer.py's own rule is that + such a root freezes rather than empties, precisely so it does not pay a + full rehash the moment it comes back. Pruning through an unavailable + root here would reintroduce exactly that cost via a different door, so + an owning-but-unavailable root wins over "the file isn't there right + now" every time, unconditionally. + + A row lost here costs one rehash the next time that path is scanned, + never a wrong answer: lookup() (cache.py) always re-validates size and + mtime against a live stat() before trusting a cached hash. + """ + cache = state.get("index_cache") + if cache is None: + raise OpError("No index cache in this process", status=503) + + indexers = list((state.get("indexers") or {}).values()) + roots = [root for indexer in indexers for root in indexer.roots] + + def _is_stale(path_str: str) -> bool: + path = Path(path_str) + owning = [r for r in roots if r.path in path.parents] + if not owning: + return True + if any(not r.available for r in owning): + return False + return not path.exists() + + paths = await cache.all_paths() + stale = await asyncio.to_thread(lambda: [p for p in paths if _is_stale(p)]) + await cache.remove_many(stale) + log.info("Index cache pruned: %d stale row(s) removed, %d kept", + len(stale), len(paths) - len(stale)) + return {"status": "pruned", "removed": len(stale), "kept": len(paths) - len(stale)} diff --git a/packages/meshbay-node/src/meshbay_node/ops/groups.py b/packages/meshbay-node/src/meshbay_node/ops/groups.py new file mode 100644 index 0000000..905aa44 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/ops/groups.py @@ -0,0 +1,248 @@ +"""The groups this node hosts and their keys.""" + +from __future__ import annotations + +import logging +from pathlib import Path + +from meshbay_common.crypto import generate_gek, wrap_gek_aes + +from meshbay_node.config import DEFAULT_CONFIG_PATH +from meshbay_node.ops.core import OpError, _config, _group_ctx, _hub +from meshbay_node.ops.node_toml import _find_group_range + +log = logging.getLogger("meshbay_node.ops") + + +# ── Group keys ─────────────────────────────────────────────────────────────── + +async def set_gek(state: dict, group_id: str, *, rotate: bool = False) -> dict: + """ + Generate the group key and activate it, or rotate an existing one. + + Nothing is pre-wrapped for members. Each member's copy is produced when they + connect, for a key they proved they hold (`join_request`) — pre-wrapping used + to fetch public keys from the hub, which is H3 with the node as the victim + instead of the inviter. Only the node's own copy is stored, so the daemon can + reload the key across restarts without the operator's browser. + + **`rotate` generates a fresh key even when one exists.** That is the point of + it: after a revocation the ex-member still holds the current key, and nothing + else takes it away from them. Without `rotate` an existing key is kept, so + running this twice is not destructive by accident. + """ + ctx = _group_ctx(state, group_id) + hub = _hub(state) + + if rotate and ctx.get("visibility") == "public": + raise OpError( + "Key rotation is not available for public groups", status=400) + + bundle_store = state.get("bundle_store") + if not bundle_store: + raise OpError("Bundle store not available", status=503) + + existing = ctx.get("gek") + gek = generate_gek() if (rotate or not existing) else existing + rotated = bool(existing) and gek is not existing + errors: list[str] = [] + + roster = state.get("roster") + authorized = len(await roster.list_members(group_id)) if roster else 0 + + 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) + + ctx["gek"] = gek + log.info("GEK %s for group %s — %d authorized member(s) will receive it " + "on connect", "rotated" if rotated else "initialized", + group_id[:8], authorized) + + # The transport holds its own view of the group; a rotation that did not + # reach it would keep serving the old key until the daemon restarted. + for transport_key in ("webrtc", "quic_server"): + transport = state.get(transport_key) + groups = getattr(transport, "_ctx", {}).get("groups") if transport else None + if groups and group_id in groups: + groups[group_id]["gek"] = gek + + indexes = state.get("indexes") or {} + index = indexes.get(group_id) + if index is not None: + # The index is encrypted under the GEK; leaving the old key on it would + # serve members a listing they cannot open. + index.gek = gek + + return { + "status": "rotated" if rotated else "ok", + "group_id": group_id, + "rotated": rotated, + "authorized_members": authorized, + "errors": errors, + } + + +# ── Groups and roots ───────────────────────────────────────────────────────── + +async def list_groups(state: dict) -> dict: + """What this node hosts, with live status. Milestone 14.2.""" + config = state.get("config") + groups_ctx = state.get("groups_ctx", {}) + peers = state.get("peers") or {} + out = [] + for gid, ctx in groups_ctx.items(): + cfg = next((g for g in config.groups if g.id == gid), None) if config else None + idx = ctx.get("index") + roots = ctx.get("roots") + out.append({ + "id": gid, + "name": cfg.name if cfg else gid[:8], + "visibility": cfg.visibility if cfg else "private", + "join_policy": cfg.join_policy if cfg else "invite", + "has_gek": bool(ctx.get("gek")), + "file_count": idx.count if idx else 0, + "index_version": idx.version if idx else 0, + # With paths: this answers the loopback API, which is the + # operator's own channel. `meshbay-node root list` printed "?" for + # every directory without it — it was reading a field the member + # form of this deliberately omits. + "roots": roots.describe(with_paths=True) if roots else [], + "peers": sum(1 for p in peers.values() if p.get("group_id") == gid), + }) + roster = state.get("roster") + has_operator = False + if roster: + members = await roster.list_members() + has_operator = any(m["role"] == "operator" and m["status"] == "active" + for m in members) + from meshbay_node.config import node_settings_defaults + # No config (a test, an unconfigured node) falls back to NodeConfig()'s own + # values rather than to numbers repeated here, which is the copy this used + # to be: it was missing three settings and reported them as null. + defaults = node_settings_defaults(config.node if config else None) + if roster: + settings = await roster.node_settings(defaults) + else: + settings = defaults + return {"groups": out, "operator_paired": has_operator, "settings": settings} + + +async def attach_group(state: dict, name: str, shared_dir: str, + writable: bool = True) -> dict: + """ + Write a new [[groups]] block into node.toml. + + 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. + """ + if not name or not shared_dir: + raise OpError("name and shared_dir are required") + config = _config(state) + hub = _hub(state) + try: + mine = await hub.list_my_groups() + except Exception as e: + raise OpError(f"Could not list groups: {e}", status=502) from e + + match = [g for g in mine if g["id"] == name or g["name"] == name] + if not match: + raise OpError(f"No group of yours is called {name!r}", status=404, + extra={"available": [{"name": g["name"], "id": g["id"]} + for g in mine]}) + if len(match) > 1: + raise OpError(f"Several of your groups are called {name!r} — use the id", + status=409, + extra={"available": [{"name": g["name"], "id": g["id"]} + for g in match]}) + group = match[0] + + if any(g.id == group["id"] for g in config.groups): + raise OpError(f"{group['name']!r} is already hosted by this node", status=409) + + path = Path(shared_dir).expanduser() + try: + path.mkdir(parents=True, exist_ok=True) + except OSError as e: + raise OpError(f"Cannot create {path}: {e}") from e + + conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) + join_policy = group.get("join_policy", "invite") + block = (f'\n[[groups]]\n' + f'id = "{group["id"]}"\n' + f'name = "{group["name"]}"\n' + f'visibility = "{group.get("visibility", "private")}"\n' + f'join_policy = "{join_policy}"\n') + # No `upload_dir` here. `GroupConfig.__post_init__` still *reads* it, so an + # existing node.toml keeps working — but what it does on read is force every + # other root read-only and append that path as the one writable one, which + # is the model this refactor replaced. Writing it into a group created + # today would mean two mechanisms deciding the same thing, one of them + # invisible: `group add --dir X --writable --upload-dir Y` silently made X + # read-only. A second writable directory is `root add <path> --writable`. + block += (f'\n [[groups.roots]]\n' + # Forward slashes: a Windows path in a TOML basic string is a + # parse error (`\U`, `\a`, ... are escapes). pathlib reads `/`. + f' path = "{path.as_posix()}"\n' + f' writable = {"true" if writable else "false"}\n') + try: + with conf_path.open("a", encoding="utf-8", newline="\n") as f: + f.write(block) + except OSError as e: + raise OpError(f"Cannot write {conf_path}: {e}", status=500) from e + + result = {"group_id": group["id"], "name": group["name"], + "shared_dir": str(path), "config": str(conf_path), + "writable": writable, + "note": "restart the node to pick it up"} + return result + + +async def detach_group(state: dict, name: str) -> dict: + """ + Remove a [[groups]] block from node.toml. + + Does not touch the hub — only stops this node from hosting the group + after the next reload or restart. + """ + if not name: + raise OpError("group name or id is required") + config = _config(state) + + match = [g for g in config.groups if g.id == name or g.name == name] + if not match: + raise OpError(f"No hosted group matches {name!r}", status=404, + extra={"available": [{"name": g.name, "id": g.id} + for g in config.groups]}) + group = match[0] + + conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) + text = conf_path.read_text(encoding="utf-8") + lines = text.split("\n") + + rng = _find_group_range(lines, group.id) + if rng is None: + raise OpError(f"Group {group.id[:8]} not found in {conf_path}") + + start, end = rng + while end < len(lines) and lines[end].strip() == "": + end += 1 + + new_lines = lines[:start] + lines[end:] + conf_path.write_text("\n".join(new_lines), encoding="utf-8", newline="\n") + log.info("Group detached: %s (%s) removed from %s", group.name, group.id[:8], conf_path) + + return {"group_id": group.id, "name": group.name, "config": str(conf_path), + "note": "restart the node to stop hosting it"} diff --git a/packages/meshbay-node/src/meshbay_node/ops/members.py b/packages/meshbay-node/src/meshbay_node/ops/members.py new file mode 100644 index 0000000..2fd3899 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/ops/members.py @@ -0,0 +1,314 @@ +"""Who may use the node: pairing, invitations, revocation.""" + +from __future__ import annotations + +import logging +from urllib.parse import urlsplit + +from meshbay_common.crypto import pk_to_b64 +from meshbay_common.join import ROLE_MEMBER, ROLE_OPERATOR + +from meshbay_node.ops.core import OpError, _group_ctx, _hub, _roster +from meshbay_node.roster import LinkInviteLimit + +log = logging.getLogger("meshbay_node.ops") + + +# ── Roster ─────────────────────────────────────────────────────────────────── + +async def read_roster(state: dict, group_id: str = "") -> dict: + roster = state.get("roster") + if not roster: + return {"identities": [], "members": [], "invites": []} + members = await roster.list_members(group_id or None) + members = [m for m in members if m.get("pk_ed25519") is not None] + return { + "identities": await roster.list_identities(), + "members": members, + "invites": await roster.list_invites(), + } + + +async def resolve_user(state: dict, username: str) -> dict: + """ + Map a username to an account id. + + 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 there. + """ + 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 + raise OpError(f"Unknown user {username!r}", status=404) + + +async def pair_operator(state: dict) -> dict: + """ + 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). Returned + once and stored only as a hash. + """ + roster = _roster(state) + user_id = state.get("node_user_id") + if not user_id: + raise OpError("Node not connected to hub yet", status=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} + + +async def create_invite(state: dict, group_id: str, username: str, *, + user_id: str = "", + created_by: str = "local-cli") -> dict: + """ + Issue an invitation code. + + 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. + + When ``user_id`` is supplied directly (MNP path), the hub lookup is skipped. + """ + roster = _roster(state) + _group_ctx(state, group_id) + if not user_id: + hub = _hub(state) + try: + account = await hub.get_user_pubkeys(username) + except Exception as e: + raise OpError(f"Unknown user {username!r}: {e}", status=404) from e + user_id = account["user_id"] + + # Hub membership first, and fatal if it fails. + # + # `/v1/groups/mine` joins `GroupMember`, so someone who was never registered + # does not see the group at all and can never redeem the code. Creating the + # invite first and tolerating a failed registration — which is what this did + # — hands the operator a code that cannot work, and says nothing. Worse, an + # unreachable hub raised *after* the roster write, leaving a valid code + # nobody was ever given; every retry left another. + # + # Registering before the roster write means a failure costs nothing: no code + # exists to be orphaned. A membership row without an invite is harmless — + # without the code there is still no group key. + # + # The endpoint is idempotent (`if not mem: db.add(...)`, no 409), so the SPA + # registering the same membership again right after `createInvite` + # (group-settings.js) costs nothing either. + # + # Skipped only when there is no username to register with: the MNP path + # allows an empty one (`username || ''` in transport.js), and there the SPA + # is the one that registers. + if username: + try: + await _hub(state).add_group_member(group_id, username) + except Exception as e: + raise OpError( + f"Could not register {username!r} on the hub, so the invite " + f"could not be redeemed: {e}", status=502) from e + + 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=user_id, + role=ROLE_MEMBER, + created_by=created_by, + ttl=ttl, + username=username, + ) + + invites = await roster.list_invites() + expires = next((i["expires_at"] for i in invites + if i["user_id"] == user_id + and i["group_id"] == group_id), "") + return {"code": code, "expires_at": expires, + "username": username, "user_id": user_id} + + +async def create_link_invite(state: dict, group_id: str, *, + created_by: str = "local-cli") -> dict: + """ + Issue a code bound to no account, for an invitation link. + + Nothing is registered on the hub here, unlike `create_invite`: there is no + account to register yet. The hub half is a ticket the inviter's client asks + the hub for, bound to the invitee's address (docs/MESHBAY_DESIGN.md §7.3). + """ + roster = _roster(state) + _group_ctx(state, group_id) + config = state.get("config") + ttl = (config.node.invite_ttl_hours if config else 168) * 3600 + try: + code, invite_id, expires = await roster.create_link_invite( + group_id, created_by, ttl=ttl) + except LinkInviteLimit as e: + raise OpError(str(e), status=429) from e + log.info("Invitation link issued: group=%s invite=%s", group_id[:8], invite_id[:8]) + return {"code": code, "invite_id": invite_id, "expires_at": expires, + "group_id": group_id} + + +async def cancel_invite(state: dict, group_id: str, invite_id: str) -> dict: + """Take back an unredeemed invitation link. Unknown or spent is a refusal, + so a mistyped handle does not read as success.""" + roster = _roster(state) + _group_ctx(state, group_id) + if not await roster.cancel_invite(group_id, invite_id): + raise OpError("No unredeemed invitation link with that id in this group", + status=404) + log.info("Invitation link cancelled: group=%s invite=%s", group_id[:8], invite_id[:8]) + return {"cancelled": True, "invite_id": invite_id, "group_id": group_id} + + +def _invite_url(hub_url: str, group_id: str, ticket: str, node_pk_b64: str, code: str) -> str: + """ + An invitation link, in the one shape the hub and the interface also write + (docs/MESHBAY_DESIGN.md §3.4): everything after `#`, and the node key + URL-safe and unpadded. `test_invite_link_client.py` (hub) holds it to the hub's. + """ + parts = urlsplit(hub_url) + origin = f"{parts.scheme}://{parts.netloc}" + n = node_pk_b64.replace("+", "-").replace("/", "_").rstrip("=") + return f"{origin}/#/invite?v=1&g={group_id}&t={ticket}&n={n}&c={code}" + + +async def create_link_invitation(state: dict, group_id: str, email: str, *, + created_by: str = "local-cli") -> dict: + """ + A whole invitation link, from the operator's own machine: the node's code, + then the hub's ticket bound to `email`, then the link. + + In that order because the ticket names the code's handle. A ticket the hub + refuses takes the code back with it — a code nobody can reach the node with + would only hold one of the group's places. The hub is never asked to mail: + the operator sends the link. + """ + email = (email or "").strip() + if "@" not in email: + raise OpError("An invitation link is bound to an e-mail address", status=422) + hub = _hub(state) + sk_node = state.get("sk_node") + if sk_node is None: + raise OpError("Node key not loaded", status=503) + node = await create_link_invite(state, group_id, created_by=created_by) + try: + ticket = await hub.create_invite_link( + group_id, email, node["expires_at"], node["invite_id"]) + except Exception as e: + await _roster(state).cancel_invite(group_id, node["invite_id"]) + raise OpError(f"The hub refused the link, so none was made: {e}", + status=502) from e + return { + "link": _invite_url(hub.hub_url, group_id, ticket["ticket"], + pk_to_b64(sk_node.public_key()), node["code"]), + "expires_at": ticket["expires_at"], + "invite_id": node["invite_id"], + "email": email, + } + + +async def cancel_link_invitation(state: dict, group_id: str, invite_id: str) -> dict: + """ + Take a link back, both halves: the node's code first, which is what stops + anyone joining, then the hub's ticket — attempted even when the first half + finds nothing to cancel, so neither is left behind (the member-removal rule). + """ + roster = _roster(state) + _group_ctx(state, group_id) + node_cancelled = await roster.cancel_invite(group_id, invite_id) + hub_cancelled = False + hub = state.get("hub") + if hub and hub._session: + try: + for link in await hub.list_invite_links(group_id): + if link.get("node_invite_id") == invite_id and link.get("status") == "pending": + await hub.delete_invite_link(group_id, link["link_id"]) + hub_cancelled = True + except Exception as e: + log.warning("Invitation link %s: the hub half was not cancelled: %s", + invite_id[:8], e) + if not node_cancelled and not hub_cancelled: + raise OpError("No unredeemed invitation link with that id in this group", + status=404) + log.info("Invitation link cancelled: group=%s invite=%s node=%s hub=%s", + group_id[:8], invite_id[:8], node_cancelled, hub_cancelled) + return {"cancelled": True, "invite_id": invite_id, + "node": node_cancelled, "hub": hub_cancelled} + + +async def revoke_member(state: dict, user_id: str, group_id: str) -> dict: + """ + 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. + + **An unredeemed invite is a membership that has not happened yet**, so it is + revoked here too, and on its own it is enough for this to be a removal. A + member row appears only when a code is consumed: somebody invited to the + wrong group has none, this refused them with "no such member", and the + browser's removal — node half first, deliberately — died on that refusal + before it reached the hub half. They stayed a member on the hub, with a live + code, and the interface offered no other way to take either back. + """ + roster = _roster(state) + revoked = await roster.set_status(group_id, user_id, "revoked") + dropped = await roster.drop_invites(group_id, user_id) + if not revoked and not dropped: + raise OpError("No such member in that group", status=404) + log.info("Member revoked: user=%s group=%s member=%s invites_dropped=%d", + user_id[:8], group_id[:8], revoked, dropped) + return {"status": "revoked", "user_id": user_id, "group_id": group_id, + "was_member": revoked, "invites_dropped": dropped, + # Only what is true: somebody who never redeemed a code never held + # the key, and telling an operator to rotate it teaches them that + # the advice is noise. + "reminder": ("rotate the group key: meshbay-node gek rotate" + if revoked else "")} + + +async def unpin_member(state: dict, user_id: str) -> dict: + """Forget a pinned identity, so the person can pair again with a new key.""" + roster = _roster(state) + if not await roster.unpin(user_id): + raise OpError("No such pinned identity", status=404) + # Drop the stored keypair bundle too. Left behind, it is served to the next + # connection, which then cannot open it (the passphrase may have changed + # since) and dies in the identity step before it ever reaches the join the + # unpin was meant to enable. + bundle_store = state.get("bundle_store") + if bundle_store: + try: + await bundle_store.delete_keypair(user_id) + except Exception: + log.warning("unpin: could not drop keypair bundle for %s", user_id[:8]) + log.info("Identity unpinned: user=%s", user_id[:8]) + return {"status": "unpinned", "user_id": user_id} diff --git a/packages/meshbay-node/src/meshbay_node/ops/node_toml.py b/packages/meshbay-node/src/meshbay_node/ops/node_toml.py new file mode 100644 index 0000000..f711f26 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/ops/node_toml.py @@ -0,0 +1,204 @@ +"""Edits to node.toml made as text, so the operator's comments survive them.""" + +from __future__ import annotations + +import re +from pathlib import Path + +from meshbay_node.ops.core import OpError + + +def _find_group_range(lines: list[str], group_id: str) -> tuple[int, int] | None: + """Line range of a [[groups]] block by id: (start, end_exclusive).""" + id_re = re.compile(r'^\s*id\s*=\s*"([^"]*)"') + block_starts: list[int] = [] + for i, line in enumerate(lines): + if line.strip() == "[[groups]]": + block_starts.append(i) + + for j, start in enumerate(block_starts): + boundary = block_starts[j + 1] if j + 1 < len(block_starts) else len(lines) + for k in range(start + 1, boundary): + s = lines[k].strip() + if s.startswith("[") and s != "[[groups.roots]]": + boundary = k + break + for k in range(start + 1, boundary): + m = id_re.match(lines[k]) + if m and m.group(1) == group_id: + return (start, boundary) + return None + + +def _update_node_toml(conf_path: Path, updates: dict) -> None: + """Write changed [node] settings back to node.toml without disturbing comments. + + For each key, if the line exists (commented or not) it is replaced in place; + otherwise the key is appended to the end of the [node] section. + """ + if not conf_path.exists(): + return + text = conf_path.read_text(encoding="utf-8") + lines = text.split("\n") + + node_start = None + node_end = len(lines) + for i, line in enumerate(lines): + stripped = line.strip() + if stripped == "[node]": + node_start = i + elif node_start is not None and re.match(r'^\[', stripped): + node_end = i + break + + if node_start is None: + lines.append("") + lines.append("[node]") + node_start = len(lines) - 1 + node_end = len(lines) + + def _format_value(key, value): + if isinstance(value, bool): + return f"{key} = {'true' if value else 'false'}" + if isinstance(value, list): + items = ", ".join(f'"{v}"' for v in value) + return f"{key} = [{items}]" + return f"{key} = {value}" + + remaining = dict(updates) + for i in range(node_start + 1, node_end): + for key in list(remaining): + pattern = re.compile( + r'^(\s*#?\s*)' + re.escape(key) + r'\s*=\s*.*$') + if pattern.match(lines[i]): + value = remaining.pop(key) + lines[i] = _format_value(key, value) + break + + for key, value in remaining.items(): + lines.insert(node_end, _format_value(key, value)) + node_end += 1 + + conf_path.write_text("\n".join(lines), encoding="utf-8", newline="\n") + + +def _insert_roots_block(conf_path: Path, group_id: str, + root_block: str) -> None: + """Append a [[groups.roots]] block inside the matching [[groups]] section.""" + text = conf_path.read_text(encoding="utf-8") + lines = text.split("\n") + + rng = _find_group_range(lines, group_id) + if rng is None: + raise OpError(f"Group {group_id[:8]} not found in {conf_path}") + + _start, end = rng + insert_at = end + while insert_at > _start + 1 and lines[insert_at - 1].strip() == "": + insert_at -= 1 + + new_lines = (lines[:insert_at] + + [""] + + root_block.rstrip("\n").split("\n") + + lines[insert_at:]) + conf_path.write_text("\n".join(new_lines), encoding="utf-8", newline="\n") + + +def _remove_roots_block(conf_path: Path, group_id: str, + resolved_path: str) -> None: + """Remove a [[groups.roots]] block whose resolved path matches.""" + text = conf_path.read_text(encoding="utf-8") + lines = text.split("\n") + + rng = _find_group_range(lines, group_id) + if rng is None: + raise OpError(f"Group {group_id[:8]} not found in {conf_path}") + + start, end = rng + path_re = re.compile(r'^\s*path\s*=\s*"([^"]*)"') + roots_starts: list[int] = [] + for i in range(start + 1, end): + if lines[i].strip() == "[[groups.roots]]": + roots_starts.append(i) + + for j, rs in enumerate(roots_starts): + rs_end = roots_starts[j + 1] if j + 1 < len(roots_starts) else end + for k in range(rs, rs_end): + m = path_re.match(lines[k]) + if m: + try: + p = str(Path(m.group(1)).expanduser().resolve()) + except OSError: + continue + if p == resolved_path: + rm_start = rs + if rm_start > 0 and lines[rm_start - 1].strip() == "": + rm_start -= 1 + new_lines = lines[:rm_start] + lines[rs_end:] + conf_path.write_text("\n".join(new_lines), encoding="utf-8", newline="\n") + return + + raise OpError("Root path not found in config", status=404) + + +def _update_root_field(conf_path: Path, group_id: str, + resolved_path: str, *, + writable: bool, removable: bool) -> None: + """Update writable/removable fields on a root in node.toml.""" + text = conf_path.read_text(encoding="utf-8") + lines = text.split("\n") + + rng = _find_group_range(lines, group_id) + if rng is None: + raise OpError(f"Group {group_id[:8]} not found in {conf_path}") + + start, end = rng + path_re = re.compile(r'^\s*path\s*=\s*"([^"]*)"') + writable_re = re.compile(r'^\s*(writable|upload)\s*=') + removable_re = re.compile(r'^\s*removable\s*=') + roots_starts: list[int] = [] + for i in range(start + 1, end): + if lines[i].strip() == "[[groups.roots]]": + roots_starts.append(i) + + for j, rs in enumerate(roots_starts): + rs_end = roots_starts[j + 1] if j + 1 < len(roots_starts) else end + found_path = False + for k in range(rs, rs_end): + m = path_re.match(lines[k]) + if m: + try: + p = str(Path(m.group(1)).expanduser().resolve()) + except OSError: + continue + if p == resolved_path: + found_path = True + break + if not found_path: + continue + + writable_idx = None + removable_idx = None + for k in range(rs, rs_end): + if writable_re.match(lines[k]): + writable_idx = k + if removable_re.match(lines[k]): + removable_idx = k + + if writable_idx is not None: + lines[writable_idx] = f" writable = {'true' if writable else 'false'}" + else: + lines.insert(rs_end, f" writable = {'true' if writable else 'false'}") + if removable_idx is not None and removable_idx >= rs_end: + removable_idx += 1 + rs_end += 1 + + if removable_idx is not None: + lines[removable_idx] = f" removable = {'true' if removable else 'false'}" + else: + lines.insert(rs_end, f" removable = {'true' if removable else 'false'}") + + conf_path.write_text("\n".join(lines), encoding="utf-8", newline="\n") + return + + raise OpError("Root path not found in config", status=404) diff --git a/packages/meshbay-node/src/meshbay_node/ops/roots.py b/packages/meshbay-node/src/meshbay_node/ops/roots.py new file mode 100644 index 0000000..e3e2781 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/ops/roots.py @@ -0,0 +1,285 @@ +"""A group's directories: adding, removing, changing and ejecting them.""" + +from __future__ import annotations + +import logging +from dataclasses import asdict +from pathlib import Path + +from meshbay_node.config import DEFAULT_CONFIG_PATH +from meshbay_node.ops.core import OpError, _config, _group_ctx, _roster +from meshbay_node.ops.node_toml import _insert_roots_block, _remove_roots_block, _update_root_field +from meshbay_node.roots import RootError, RootSet, off_disk + +log = logging.getLogger("meshbay_node.ops") + + +async def add_root(state: dict, group_id: str, path: str, *, + name: str = "", kind: str = "generic", + writable: bool = False, + removable: bool = False) -> dict: + """ + Add a directory to a group, refusing anything ambiguous. + + Validated against the group's existing roots *before* being written, so a + config that would be refused at startup is refused here instead — where the + operator is watching and can fix it. + """ + config = _config(state) + cfg = next((g for g in config.groups if g.id == group_id), None) + if cfg is None: + raise OpError("Group not configured on this node", status=404) + + specs = [asdict(r) for r in cfg.roots] + specs.append({"path": path, "name": name, "kind": kind, + "writable": writable, "removable": removable}) + try: + built = RootSet.build(specs) + except RootError as e: + raise OpError(str(e)) from e + + added = built.roots[-1] + + try: + added.path.mkdir(parents=True, exist_ok=True) + except OSError as e: + raise OpError(f"Cannot create {added.path}: {e}") from e + + conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) + root_block = f' [[groups.roots]]\n path = "{added.path.as_posix()}"' + if name: + root_block += f'\n name = "{added.name}"' + if kind != "generic": + root_block += f'\n kind = "{added.kind}"' + if writable: + root_block += '\n writable = true' + if removable: + root_block += '\n removable = true' + _insert_roots_block(conf_path, group_id, root_block) + + from meshbay_node.config import RootSpec + cfg.roots.append(RootSpec( + path=str(added.path), name=added.name, kind=added.kind, + writable=added.writable, removable=added.removable)) + + # Deliberately *not* mutating the live RootSet in place. + # + # `DirectoryIndexer.retarget` decides what to scan by diffing the names it + # already has against the ones it is given — so handing it the same object, + # edited, means the new root is in both sides of the comparison and is + # never scanned. It would appear in the table and stay permanently empty. + # `_reload_config_inner` diffs the same way and would likewise conclude + # nothing changed. The caller reloads instead, which builds a fresh set + # from the file this just wrote. + # + # `built` is that set, computed here only to validate and to answer with; + # what the node serves comes from the reload. + log.info("Root added: %s → group %s", added.name, group_id[:8]) + return {"status": "added", "name": added.name, "path": str(added.path), + "group_id": group_id, "roots": built.describe()} + + +async def remove_root(state: dict, group_id: str, root_name: str) -> dict: + """Remove a named root from a group. At least one root must remain.""" + config = _config(state) + cfg = next((g for g in config.groups if g.id == group_id), None) + if cfg is None: + raise OpError("Group not configured on this node", status=404) + + from meshbay_common.paths import fold + + from meshbay_node.roots import derive_name + target = fold(root_name) + match_idx = None + for i, r in enumerate(cfg.roots): + try: + rname = r.name or derive_name(Path(r.path).expanduser().resolve()) + except Exception: + continue + if fold(rname) == target: + match_idx = i + break + + if match_idx is None: + raise OpError(f"No root named {root_name!r} in this group", status=404) + if len(cfg.roots) < 2: + raise OpError("Cannot remove the only root", status=400) + + removed = cfg.roots[match_idx] + resolved = str(Path(removed.path).expanduser().resolve()) + + conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) + _remove_roots_block(conf_path, group_id, resolved) + + cfg.roots.pop(match_idx) + + # Not mutating the live set here either — see `add_root`. Dropping the + # root from it would leave `retarget` unable to tell that its entries + # should go, so the removed directory's files would stay in the index. + # + # Built from the config this just edited, and never returned empty: an + # empty list is a *valid answer* meaning "this group has no directories", + # which the client cannot tell from "the node could not say" — it would + # blank the operator's table on an op that succeeded. + result_roots = RootSet.build([asdict(r) for r in cfg.roots]).describe() + + log.info("Root removed: %s from group %s", root_name, group_id[:8]) + return {"status": "removed", "name": root_name, "group_id": group_id, + "roots": result_roots} + + +async def update_root(state: dict, group_id: str, root_name: str, *, + writable: bool | None = None, + removable: bool | None = None) -> dict: + """Toggle writable/removable on an existing root without removing it.""" + config = _config(state) + cfg = next((g for g in config.groups if g.id == group_id), None) + if cfg is None: + raise OpError("Group not configured on this node", status=404) + + from meshbay_common.paths import fold + + from meshbay_node.roots import RootSet + target = fold(root_name) + match = None + for r in cfg.roots: + rname = r.name or str(Path(r.path).name) + if fold(rname) == target: + match = r + break + if match is None: + raise OpError(f"No root named {root_name!r} in this group", status=404) + + changed = False + if writable is not None and match.writable != writable: + match.writable = writable + changed = True + if removable is not None and match.removable != removable: + match.removable = removable + changed = True + + if not changed: + specs = [asdict(r) for r in cfg.roots] + built = RootSet.build(specs) + return {"status": "unchanged", "name": root_name, "group_id": group_id, + "roots": built.describe()} + + conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) + _update_root_field(conf_path, group_id, str(Path(match.path).expanduser().resolve()), + writable=match.writable, removable=match.removable) + + # Update the live RootSet so GET /api/groups returns correct data + # immediately, without waiting for the async reload to finish. + live_roots: RootSet | None = state.get("groups_ctx", {}).get( + group_id, {}).get("roots") + if live_roots: + for lr in live_roots.roots: + lr_name = lr.name or str(Path(lr.path).name) + if fold(lr_name) == target: + if writable is not None: + lr.writable = writable + if removable is not None: + lr.removable = removable + break + + # Built from config when there is no live set, never returned empty: an + # empty list is a *valid answer* meaning "this group has no directories", + # and the client cannot tell it from "the node could not say". It would + # blank the operator's table on an op that succeeded. + result_roots = (live_roots.describe() if live_roots + else RootSet.build([asdict(r) for r in cfg.roots]).describe()) + + log.info("Root updated: %s (writable=%s, removable=%s) in group %s", + root_name, match.writable, match.removable, group_id[:8]) + return {"status": "updated", "name": root_name, "group_id": group_id, + "roots": result_roots} + + +async def eject_root(state: dict, group_id: str, root_name: str) -> dict: + """Mark a removable root as ejected so the operator can safely unplug.""" + config = _config(state) + cfg = next((g for g in config.groups if g.id == group_id), None) + if cfg is None: + raise OpError("Group not configured on this node", status=404) + + from meshbay_common.paths import fold + target = fold(root_name) + ctx = _group_ctx(state, group_id) + roots: RootSet | None = ctx.get("roots") + if not roots: + raise OpError("Group has no roots", status=503) + + root = None + for r in roots: + if fold(r.name) == target: + root = r + break + if root is None: + raise OpError(f"No root named {root_name!r} in this group", status=404) + if not root.removable: + raise OpError(f"Root {root_name!r} is not marked as removable", status=400) + if root.ejected: + return {"status": "already_ejected", "name": root_name, + "group_id": group_id, "roots": roots.describe()} + + # The indexer stops its watchdog and freezes the entries; it holds the same + # RootSet object, but the flags are set here too so a context whose indexer + # was replaced by a retarget cannot be left disagreeing with the roster. + indexer = state.get("indexers", {}).get(group_id) + if indexer: + indexer.eject_root(root_name) + root.ejected = True + root.available = False + + await _roster(state).set_root_ejected( + group_id, root_name, True, set_by=state.get("node_user_id", "")) + + log.info("Root ejected: %s from group %s", root_name, group_id[:8]) + return {"status": "ejected", "name": root_name, "group_id": group_id, + "roots": roots.describe()} + + +async def plug_root(state: dict, group_id: str, root_name: str) -> dict: + """Re-enable an ejected root after the device is plugged back in.""" + config = _config(state) + cfg = next((g for g in config.groups if g.id == group_id), None) + if cfg is None: + raise OpError("Group not configured on this node", status=404) + + from meshbay_common.paths import fold + target = fold(root_name) + ctx = _group_ctx(state, group_id) + roots: RootSet | None = ctx.get("roots") + if not roots: + raise OpError("Group has no roots", status=503) + + root = None + for r in roots: + if fold(r.name) == target: + root = r + break + if root is None: + raise OpError(f"No root named {root_name!r} in this group", status=404) + if not root.ejected: + return {"status": "already_plugged", "name": root_name, + "group_id": group_id, "roots": roots.describe()} + if not await off_disk(roots, root.is_live): + raise OpError( + f"Directory not found: {root.path}. Is the device connected?", + status=409) + + # Persisted before the rescan, which can take minutes on a large library: + # a crash halfway through must leave the root plugged, not ejected with + # entries half rebuilt. + await _roster(state).set_root_ejected( + group_id, root_name, False, set_by=state.get("node_user_id", "")) + + indexer = state.get("indexers", {}).get(group_id) + if indexer: + await indexer.plug_root(root_name) + root.ejected = False + root.available = await off_disk(roots, root.is_live) + + log.info("Root plugged: %s in group %s", root_name, group_id[:8]) + return {"status": "plugged", "name": root_name, "group_id": group_id, + "roots": roots.describe()} diff --git a/packages/meshbay-node/src/meshbay_node/ops/settings.py b/packages/meshbay-node/src/meshbay_node/ops/settings.py new file mode 100644 index 0000000..eade6ca --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/ops/settings.py @@ -0,0 +1,340 @@ +"""Node-wide settings: denylist, node.toml values, transfer caps, scan timing, reload.""" + +from __future__ import annotations + +import logging +from pathlib import Path + +from meshbay_common.background import spawn + +from meshbay_node.config import DEFAULT_CONFIG_PATH +from meshbay_node.ops.core import OpError, _config, _group_ctx, _roster +from meshbay_node.ops.node_toml import _update_node_toml +from meshbay_node.roster import Roster + +log = logging.getLogger("meshbay_node.ops") + + +# ── Revocation denylist ────────────────────────────────────────────────────── + +async def read_denylist(state: dict) -> dict: + """Milestone 14.10 — what the node is currently refusing.""" + denylist = state.get("denylist") + if not denylist: + return {"users": [], "groups": [], "jtis": [], "count": 0} + entries = denylist.entries() + return {**entries, "count": sum(len(v) for v in entries.values())} + + +async def clear_denylist(state: dict, *, subject: str = "") -> dict: + """ + Drop denylist entries — all of them, or one identifier. + + Deliberately not silent: a cleared denylist re-admits whoever it was keeping + out, and the count is what tells the operator whether they undid one + revocation or all of them. + """ + denylist = state.get("denylist") + if not denylist: + raise OpError("No denylist in this process", status=503) + removed = denylist.clear(subject) + log.warning("Denylist cleared (%s): %d entr(y/ies) removed", + subject or "all", removed) + return {"status": "cleared", "removed": removed, "subject": subject or "all"} + + +# ── Node settings ──────────────────────────────────────────────────────────── + +# What `set_node_settings` accepts, and how each value is validated. A module +# constant so a test can hold its key set against `Roster.node_setting_keys()`: +# this is the third list of the same settings, and the first two had already +# drifted apart once — the reader's defaults covered fewer settings than the +# resolver answered for, which is how node.toml's transfer pools came to be +# parsed and then ignored. The kinds here are the *writer's* validation and +# deliberately not the resolver's coercions. +NODE_SETTING_WRITERS: dict[str, tuple[str, str]] = { + "invite_ttl_hours": ("int", Roster.SETTING_INVITE_TTL), + "pair_ttl_hours": ("int", Roster.SETTING_PAIR_TTL), + "device_request_ttl_minutes": ("int", Roster.SETTING_DEVICE_TTL), + "max_concurrent_streams": ("int", Roster.SETTING_MAX_STREAMS), + "max_concurrent_downloads": ("int", Roster.SETTING_MAX_DOWNLOADS), + "max_concurrent_uploads": ("int", Roster.SETTING_MAX_UPLOADS), + "max_upload_gb": ("size", Roster.SETTING_MAX_UPLOAD_GB), + "transcode_incompatible_video": ("bool", Roster.SETTING_TRANSCODE), + "stun_servers": ("stun_list", Roster.SETTING_STUN_SERVERS), + "ice_interfaces": ("list", Roster.SETTING_ICE_INTERFACES), +} + + +async def get_node_settings(state: dict) -> dict: + """Return current effective node settings.""" + from meshbay_node.config import node_settings_defaults + roster = _roster(state) + config = _config(state) + defaults = node_settings_defaults(config.node) + if roster: + return await roster.node_settings(defaults) + return defaults + +async def set_node_settings(state: dict, settings: dict) -> dict: + """Update node-level daemon settings. Writes to both roster.db and node.toml.""" + roster = _roster(state) + config = _config(state) + nd = config.node + conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) + + allowed_keys = NODE_SETTING_WRITERS + + set_by = state.get("node_user_id", "") + updated = {} + for key, value in settings.items(): + if key not in allowed_keys: + continue + kind, setting_key = allowed_keys[key] + if kind == "int": + try: + v = int(value) + except (TypeError, ValueError): + raise OpError(f"{key} must be an integer") + if v < 1: + raise OpError(f"{key} must be positive") + setattr(nd, key, v) + await roster.set_node_setting(setting_key, str(v), set_by) + updated[key] = v + elif kind == "size": + # A quantity, not a count: half a gigabyte is a legitimate ceiling + # on a small disk, so this one is not run through the `int` branch + # above, whose floor of 1 would round it to "refuse everything". + # bool before float, as config.py does it: `true` is not 1 GB. + if isinstance(value, bool): + raise OpError(f"{key} must be a number") + try: + fv = float(value) + except (TypeError, ValueError): + raise OpError(f"{key} must be a number") + if fv <= 0: + raise OpError(f"{key} must be greater than zero") + setattr(nd, key, fv) + await roster.set_node_setting(setting_key, repr(fv), set_by) + updated[key] = fv + elif kind == "bool": + v = bool(value) + setattr(nd, key, v) + await roster.set_node_setting(setting_key, "1" if v else "0", set_by) + updated[key] = v + elif kind in ("list", "stun_list"): + import json as _json + if not isinstance(value, list): + raise OpError(f"{key} must be a list") + v = [str(s) for s in value] + if kind == "stun_list": + for s in v: + if not s.startswith("stun:"): + raise OpError(f"Invalid STUN server: {s} (must start with stun:)") + setattr(nd, key, v) + await roster.set_node_setting(setting_key, _json.dumps(v), set_by) + updated[key] = v + + if updated: + _update_node_toml(conf_path, updated) + if "max_concurrent_streams" in updated: + webrtc = state.get("webrtc") + # `webrtc._stream_sem` was assigned here for months. That attribute + # has never existed -- the pool is `ctx["_transcode_sem"]` -- so the + # `hasattr` guard was always False and the setting only ever took + # effect on a restart, which docs/MESHBAY_DESIGN.md §6.8 says it + # does not need. + if webrtc is not None: + webrtc.set_capacity( + max_concurrent_streams=updated["max_concurrent_streams"]) + if ("max_concurrent_downloads" in updated + or "max_concurrent_uploads" in updated): + webrtc = state.get("webrtc") + if webrtc is not None: + webrtc.set_capacity( + max_concurrent_downloads=updated.get( + "max_concurrent_downloads"), + max_concurrent_uploads=updated.get( + "max_concurrent_uploads")) + if "max_upload_gb" in updated: + webrtc = state.get("webrtc") + if webrtc is not None: + webrtc.set_capacity(max_upload_gb=updated["max_upload_gb"]) + if "stun_servers" in updated: + webrtc = state.get("webrtc") + if webrtc and hasattr(webrtc, '_stun'): + webrtc._stun = updated["stun_servers"] + from meshbay_node.transport.stun_multi import set_servers as _set_stun + _set_stun(updated["stun_servers"]) + if "ice_interfaces" in updated: + from meshbay_node.transport.ice_filter import install as install_ice_filter + install_ice_filter(updated["ice_interfaces"] or None) + + log.info("Node settings updated: %s", updated) + return {"updated": updated} + + +# ── Transfers ──────────────────────────────────────────────────────────────── + +async def set_transfer_limits(state: dict, group_id: str, + downloads: int, uploads: int) -> dict: + """How many transfers one member may run at once in this group. + + Same shape as every other operator setting: lives on the node (roster.db, + not the hub and not node.toml, for the reason change 5 gives — a hub that + decided this would have authority over someone else's machine), signed + (webrtc_server checks the caller's admin authority before this runs), and + live, so the pools are updated in place rather than at the next restart. + """ + roster = _roster(state) + ctx = _group_ctx(state, group_id) + limits = await roster.set_transfer_limits( + group_id, {"download": downloads, "upload": uploads}, + set_by=state.get("node_user_id", "")) + ctx["transfer_limits"] = limits + webrtc = state.get("webrtc") + slots = getattr(webrtc, "_ctx", {}).get("_transfer_slots") if webrtc else None + granted = slots.set_group_limits(group_id, limits) if slots else [] + # And **tell them**. The node-wide path (`WebRTCTransport.set_capacity`) + # does this and this one did not: the leases were granted in the pool and + # the peers waiting on them were never told, so a cap raised from 2 to 4 + # left both transfers sitting at "waiting" until the client's own watchdog + # re-asked a minute later. That is §5.2's first row — "node granted a slot, + # the push was lost" — reached by writing the grant and forgetting the send, + # which is the same omission as the missing `touch()` one layer up. + for lease in granted: + webrtc._notify_granted(lease) + log.info("Transfer limits for group %s: %s (%d started at once)", + group_id[:8], limits, len(granted)) + return {"group_id": group_id, "limits": limits, + "started": [x.tr for x in granted]} + +async def list_transfers(state: dict) -> dict: + """Live transfer leases and queue depth. + + The operator's window into "is anything actually holding a slot". When + somebody reports a transfer stuck at waiting, this is the only thing that + says whether the node ever had them in a queue — the alternative is reading + a log for a line that, by definition, is not being printed. + + Carries no filename and no path: a lease holds neither, and this is exactly + where it would be tempting to add one. + """ + webrtc = state.get("webrtc") + ctx = getattr(webrtc, "_ctx", {}) if webrtc else {} + slots = ctx.get("_transfer_slots") + if slots is None: + from meshbay_node.transfers import DEFAULT_MAX_CONCURRENT, DEFAULT_MAX_PER_MEMBER, KINDS + # No pool built means nothing has transferred since the daemon started, + # which is a real answer and not an error. + # + # The caps still have to be the operator's own. Reporting the module + # defaults here was worse than reporting nothing: `transfers set 2 2` + # answered "applied now", and `transfers show` immediately said 0/8 — + # a setting written, acknowledged and displayed wrong, which reads + # exactly like the hot-swap that did nothing for months. Found by + # running it, not by a test: the test asserted the defaults and so + # agreed with the bug. + return {"pools": { + k: {"in_use": 0, + "cap": int(ctx.get(f"max_concurrent_{k}s") + or DEFAULT_MAX_CONCURRENT), + "per_member": DEFAULT_MAX_PER_MEMBER, + "queued": 0} + for k in KINDS}, "leases": [], "groups": _group_limits(state)} + out = slots.snapshot() + out["groups"] = _group_limits(state) + return out + + +def _group_limits(state: dict) -> list[dict]: + """Each group's per-member caps, as the operator set them. + + Reported because `transfers show` used to print only the node's default and + an operator reading "2 per member" had no way to tell whether that was this + group's setting or the fallback — and no way to change it either, since the + signed op had no door but MNP. Both were the same bug wearing two faces. + """ + from meshbay_node.transfers import DEFAULT_MAX_PER_MEMBER + + config = state.get("config") + groups_ctx = state.get("groups_ctx") or {} + out = [] + for group in (getattr(config, "groups", None) or []): + limits = (groups_ctx.get(group.id) or {}).get("transfer_limits") or {} + out.append({ + "group_id": group.id, + "name": group.name, + "download": int(limits.get("download") or DEFAULT_MAX_PER_MEMBER), + "upload": int(limits.get("upload") or DEFAULT_MAX_PER_MEMBER), + "set": bool(limits), + }) + return out + + +# ── Scan settings ──────────────────────────────────────────────────────────── + +async def set_scan_settings(state: dict, group_id: str, reconcile_interval_secs: float, + debounce_secs: float) -> dict: + """ + How often the indexer's reconciliation backstop runs, and how long a + changed file is left alone before being hashed (indexer.py + DirectoryIndexer). Persisted like set_enabled_apps — + but there is also a *live* DirectoryIndexer object to update, since it + reads these once at construction and runs its own background loop with + them rather than consulting groups_ctx on every use. + """ + roster = _roster(state) + await roster.set_scan_settings(group_id, reconcile_interval_secs, debounce_secs, + set_by=state.get("node_user_id", "")) + indexer = state.get("indexers", {}).get(group_id) + if indexer: + indexer.reconcile_secs = reconcile_interval_secs + indexer.debounce_secs = debounce_secs + # Apply the new interval now rather than after whatever backoff had + # already stretched the wait to. + indexer.note_activity() + # Optional, unlike _group_ctx(): a group can be persisted here before it + # is hot-loaded (or in a test that only cares about the roster/indexer + # side), and that must not turn a successful write into a 404. + ctx = state.get("groups_ctx", {}).get(group_id) + if ctx is not None: + ctx["reconcile_interval_secs"] = reconcile_interval_secs + ctx["debounce_secs"] = debounce_secs + log.info("Scan settings for group %s: reconcile=%.0fs debounce=%.0fs", + group_id[:8], reconcile_interval_secs, debounce_secs) + return {"reconcile_interval_secs": reconcile_interval_secs, + "debounce_secs": debounce_secs, "group_id": group_id} + + +# ── Reload ────────────────────────────────────────────────────────────────── + +async def reload_config(state: dict) -> dict: + """Hot-reload node.toml without dropping connections. Blocks until the + reload actually finishes — see start_reload for why the loopback route + uses that instead.""" + reload_fn = state.get("reload_fn") + if not reload_fn: + raise OpError("Reload not available", status=503) + await reload_fn() + return {"status": "reloaded"} + + +async def start_reload(state: dict) -> dict: + """ + Same as reload_config, but does not wait for the reload to finish. + + The loopback route uses this one: the Electron bridge caps every call at + a fixed 30s (main.js node:call), and hot-loading a brand-new group runs + its full initial scan synchronously inside _reload_config_inner() + (daemon.py) before that coroutine returns — minutes, not seconds, on a + real library (found against a 45 GB group on the same slow disk the + StarWars benchmark used). The reload keeps running on the daemon's own + event loop either way; add_root/remove_root below already fire it the + same way for exactly this reason. + """ + reload_fn = state.get("reload_fn") + if not reload_fn: + raise OpError("Reload not available", status=503) + spawn(reload_fn()) + return {"status": "reloading"} |