aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src')
-rw-r--r--packages/meshbay-node/src/meshbay_node/audit.py14
-rw-r--r--packages/meshbay-node/src/meshbay_node/blocklist.py83
-rw-r--r--packages/meshbay-node/src/meshbay_node/config.py7
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py90
-rw-r--r--packages/meshbay-node/src/meshbay_node/hub_client.py43
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/group_index.py151
-rw-r--r--packages/meshbay-node/src/meshbay_node/modules/__init__.py0
-rw-r--r--packages/meshbay-node/src/meshbay_node/musicbrainz.py4
-rw-r--r--packages/meshbay-node/src/meshbay_node/platform.py8
-rw-r--r--packages/meshbay-node/src/meshbay_node/replication.py140
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/quic_client.py23
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py27
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py19
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/video_meta.py13
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py19
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py12
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/handshake.py14
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py31
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/wire.py25
-rw-r--r--packages/meshbay-node/src/meshbay_node/ui/app.py45
23 files changed, 317 insertions, 457 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/audit.py b/packages/meshbay-node/src/meshbay_node/audit.py
index b382107..81ee600 100644
--- a/packages/meshbay-node/src/meshbay_node/audit.py
+++ b/packages/meshbay-node/src/meshbay_node/audit.py
@@ -33,20 +33,6 @@ CREATE INDEX IF NOT EXISTS idx_audit_user ON audit_log(user_id);
CREATE INDEX IF NOT EXISTS idx_audit_event ON audit_log(event);
"""
-EVENTS = {
- "connect",
- "disconnect",
- "handshake",
- "file_download",
- "file_upload",
- "file_delete",
- "stream_video",
- "chat_message",
- "chat_history",
- "index_sync",
- "auth_failed",
-}
-
RETENTION_DAYS = 365
diff --git a/packages/meshbay-node/src/meshbay_node/blocklist.py b/packages/meshbay-node/src/meshbay_node/blocklist.py
new file mode 100644
index 0000000..29d8666
--- /dev/null
+++ b/packages/meshbay-node/src/meshbay_node/blocklist.py
@@ -0,0 +1,83 @@
+"""
+The hub's content blocklist, as this node applies it (docs/MESHBAY_DESIGN.md §7.5).
+
+Hashes of public content that the hub's moderation has blocked. A node hosting a
+public group stops serving those files there: they leave the index members are
+sent, and a request for one is refused. Nothing is deleted from the operator's
+disk — the node stops serving, and what the operator keeps is theirs to decide.
+
+Private groups are untouched. The hub never learns what a private group holds,
+so there is nothing it could have blocked in one, and nothing here reaches them.
+
+Kept on disk as well as in memory, so a node that restarts while the hub is
+unreachable does not serve again, meanwhile, what it had already stopped serving.
+The hub's list is the authority: a full sync replaces this one.
+
+It is an exact match on the content id. A file changed by one byte is another
+id, and files over the partial-hash threshold are identified by a sample of their
+bytes (§6.3) — a moderation tool, not a guarantee.
+"""
+
+import json
+import logging
+import os
+import re
+from pathlib import Path
+
+log = logging.getLogger(__name__)
+
+_HASH = re.compile(r"^[0-9a-f]{64}$")
+
+
+class ContentBlocklist:
+ def __init__(self, path: Path | None = None):
+ self._path = path
+ self._hashes: set[str] = set()
+ if path is not None:
+ try:
+ data = json.loads(path.read_text(encoding="utf-8"))
+ self._hashes = {h for h in data.get("hashes", []) if _HASH.match(h)}
+ except FileNotFoundError:
+ pass
+ except (OSError, ValueError) as e:
+ # A damaged file is not a reason to refuse to start; the next
+ # sync with the hub rewrites it.
+ log.warning("Content blocklist %s unreadable: %s", path, e)
+
+ def __contains__(self, content_hash: object) -> bool:
+ return content_hash in self._hashes
+
+ def __len__(self) -> int:
+ return len(self._hashes)
+
+ def __iter__(self):
+ return iter(self._hashes)
+
+ def replace(self, hashes) -> bool:
+ """The hub's full list. True when it differs from what was applied."""
+ new = {h for h in hashes if isinstance(h, str) and _HASH.match(h)}
+ if new == self._hashes:
+ return False
+ self._hashes = new
+ self._save()
+ return True
+
+ def apply(self, add=(), remove=()) -> bool:
+ """One pushed change. True when it changed anything."""
+ before = set(self._hashes)
+ self._hashes |= {h for h in add if isinstance(h, str) and _HASH.match(h)}
+ self._hashes -= {h for h in remove if isinstance(h, str)}
+ if self._hashes == before:
+ return False
+ self._save()
+ return True
+
+ def _save(self) -> None:
+ if self._path is None:
+ return
+ tmp = self._path.with_suffix(".tmp")
+ try:
+ tmp.write_text(json.dumps({"hashes": sorted(self._hashes)}), encoding="utf-8")
+ os.replace(tmp, self._path)
+ except OSError as e:
+ log.warning("Content blocklist %s not saved: %s", self._path, e)
diff --git a/packages/meshbay-node/src/meshbay_node/config.py b/packages/meshbay-node/src/meshbay_node/config.py
index 563e395..2cd8a8f 100644
--- a/packages/meshbay-node/src/meshbay_node/config.py
+++ b/packages/meshbay-node/src/meshbay_node/config.py
@@ -505,10 +505,3 @@ def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config:
"MESHBAY_MAX_CONCURRENT_STREAMS")
return cfg
-
-
-def write_example_config(path: Path = DEFAULT_CONFIG_PATH) -> None:
- """Write an example config file if none exists."""
- if not path.exists():
- path.parent.mkdir(parents=True, exist_ok=True)
- path.write_text(EXAMPLE_CONFIG, encoding="utf-8", newline="\n")
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py
index 220a908..45a6b9a 100644
--- a/packages/meshbay-node/src/meshbay_node/daemon.py
+++ b/packages/meshbay-node/src/meshbay_node/daemon.py
@@ -44,6 +44,7 @@ from meshbay_common.protocol import MNP
from meshbay_node import uploads as uploads_mod
from meshbay_node.audit import RETENTION_DAYS as AUDIT_RETENTION_DAYS
from meshbay_node.audit import AuditStore
+from meshbay_node.blocklist import ContentBlocklist
from meshbay_node.bundle_store import BundleStore
from meshbay_node.chat.store import ChatStore
from meshbay_node.cli.dispatch import run, start
@@ -114,6 +115,10 @@ class NodeDaemon(EnrichmentMixin):
# Persisted so a restart does not silently un-revoke everyone (H4)
self._denylist = (
Denylist(path=config.data_dir / "denylist.json") if Denylist else None)
+ # The hub's content blocklist, applied in public groups only (§7.5).
+ # Persisted for the same reason: a restart while the hub is unreachable
+ # must not serve again what had stopped being served.
+ self._blocklist = ContentBlocklist(config.data_dir / "blocklist.json")
self._chat_stores: dict[str, ChatStore] = {}
# One instance, shared by every group's DirectoryIndexer — see
# indexer/cache.py's docstring for why this stopped being per-group.
@@ -525,6 +530,7 @@ class NodeDaemon(EnrichmentMixin):
self._webrtc._ctx["pk_x25519_b64"] = keys.pk_x25519_b64
self._webrtc._ctx["roster"] = self._roster
+ self._webrtc._ctx["blocklist"] = self._blocklist
# The MNP adapter calls the same operations as the loopback API
# (meshbay_node.ops), and those take the daemon's state. Handing
# the transport a second set of lookups is how two paths to one
@@ -611,11 +617,16 @@ class NodeDaemon(EnrichmentMixin):
except Exception as e:
log.warning("Invalid revocation token: %s", e)
+ async def on_connected():
+ await self._sync_blocklist(hub)
+
ws_task = asyncio.create_task(hub.maintain_ws(
on_incoming=on_incoming,
on_revocation=on_revocation,
on_webrtc_offer=on_webrtc_offer,
group_ids=lambda: list((self._state.get("groups_ctx") or {}).keys()),
+ on_connected=on_connected,
+ on_blocklist=self._on_blocklist_update,
))
self._tasks.append(ws_task)
log.info("Hub WS task started")
@@ -667,12 +678,6 @@ class NodeDaemon(EnrichmentMixin):
await indexer.initial_scan()
log.info("Background scan complete for %s: %d files",
name, indexer.index.count)
- # Swarm registration for public groups (after files are known).
- if gctx.get("visibility") == "public":
- endpoint = f"webrtc:{self._config.node.quic_port}"
- hashes = [e.id for e in gctx["index"].entries]
- if hashes:
- await self._register_swarm(hashes, endpoint)
# initial_scan() itself never calls on_change (it predates
# the concept — every existing caller only cared about the
# scan finishing, not about notifying anyone) — but Videos
@@ -1454,9 +1459,10 @@ class NodeDaemon(EnrichmentMixin):
# node refuses every handshake while the GEK is None (NS8) — so this is
# "nobody is listening", not a case to send in clear for.
if peers and idx.gek:
- msg = (index_delta_message(idx, delta, indexer.roots)
+ hidden = self._blocklist if self._is_public(group_id) else ()
+ msg = (index_delta_message(idx, delta, indexer.roots, hidden)
if delta is not None
- else index_sync_message(idx, indexer.roots))
+ else index_sync_message(idx, indexer.roots, hidden))
pushed = 0
for session in peers:
try:
@@ -1468,20 +1474,53 @@ class NodeDaemon(EnrichmentMixin):
log.info("Index %s pushed to %d WebRTC peers",
"delta" if delta is not None else "sync", pushed)
- # 11.9 — Register file hashes with hub swarm table (public groups only, H7)
- group_cfg = next(
- (g for g in self._config.groups if g.id == group_id), None)
- if (self._hub and self._state.get("endpoint_hint")
- and group_cfg and group_cfg.visibility == "public"):
- # Only the newly added hashes once there is a delta to know them
- # from — registering the whole library again on every change is
- # the same O(changes x library size) cost the delta above exists
- # to avoid.
- hashes = ([e.id for e in delta.additions] if delta is not None
- else [e.id for e in idx.entries])
- if hashes:
- endpoint = f"webrtc:{self._config.node.quic_port}"
- spawn(self._register_swarm(hashes, endpoint))
+ # ── Content blocklist (public groups, §7.5) ──────────────────────────────
+
+ def _is_public(self, group_id: str) -> bool:
+ return any(g.id == group_id and g.visibility == "public"
+ for g in self._config.groups)
+
+ def _hosts_public_group(self) -> bool:
+ return any(g.visibility == "public" for g in self._config.groups)
+
+ async def _sync_blocklist(self, hub) -> None:
+ """The hub's whole list, on every connection. Only a node hosting a
+ public group asks: nothing else here could have been blocked."""
+ if not self._hosts_public_group():
+ return
+ try:
+ hashes = await hub.fetch_blocklist()
+ except Exception as e:
+ # Keep applying the list on disk; the next connection tries again.
+ log.warning("Content blocklist sync failed: %s", e)
+ return
+ if self._blocklist.replace(hashes):
+ log.info("Content blocklist: %d hashes", len(self._blocklist))
+ self._push_public_indexes()
+
+ def _on_blocklist_update(self, add: list, remove: list) -> None:
+ if not self._hosts_public_group():
+ return
+ if self._blocklist.apply(add, remove):
+ log.info("Content blocklist updated: +%d -%d", len(add), len(remove))
+ self._push_public_indexes()
+
+ def _push_public_indexes(self) -> None:
+ """Resend each public group's whole index, so what the list now hides
+ leaves every connected member's view, and what it released comes back."""
+ if not self._webrtc:
+ return
+ for indexer in self._indexers:
+ idx = indexer.index
+ if not self._is_public(idx.group_id) or not idx.gek:
+ continue
+ msg = index_sync_message(idx, indexer.roots, self._blocklist)
+ for session in list(self._webrtc._sessions.values()):
+ if session._group_id == idx.group_id:
+ try:
+ session._send(msg)
+ except Exception:
+ pass
def _drop_group_sessions(self, group_id: str) -> None:
"""Close live sessions for a revoked group (H4)."""
@@ -1492,13 +1531,6 @@ class NodeDaemon(EnrichmentMixin):
spawn(session.close())
log.info("Dropped session for revoked group %s", group_id[:8])
- async def _register_swarm(self, hashes: list[str], endpoint: str) -> None:
- try:
- n = await self._hub.register_swarm(hashes, endpoint)
- log.info("Swarm: registered %d/%d hashes", n, len(hashes))
- except Exception as e:
- log.warning("Swarm registration failed: %s", e)
-
async def _shutdown(self) -> None:
log.info("Shutting down...")
self._state["status"] = "stopping"
diff --git a/packages/meshbay-node/src/meshbay_node/hub_client.py b/packages/meshbay-node/src/meshbay_node/hub_client.py
index 346a4cd..2ce8993 100644
--- a/packages/meshbay-node/src/meshbay_node/hub_client.py
+++ b/packages/meshbay-node/src/meshbay_node/hub_client.py
@@ -6,7 +6,6 @@ Handles all communication from the node to a Mesh Hub:
- JWT offline verification and auto-refresh
- Node announcement (endpoint_hint)
- User public key lookup (for GEK wrapping)
- - Swarm hash registration
The node authenticates via Ed25519 challenge-response (/v1/nodes/auth).
No auth_key or password is ever stored on or transmitted from the node.
@@ -335,6 +334,8 @@ class HubClient:
on_revocation: Any = None,
on_webrtc_offer: Any = None,
group_ids: list[str] | None = None, # static list or callable returning one
+ on_connected: Any = None,
+ on_blocklist: Any = None,
) -> None:
"""
Maintain a persistent WebSocket connection to the hub.
@@ -394,6 +395,12 @@ class HubClient:
self._ws = ws
log.info("Hub WS connected")
+ if on_connected:
+ # Off the read loop, for the reason offers are: a slow
+ # fetch must not stop this socket being read.
+ task = asyncio.create_task(on_connected())
+ pending.add(task)
+ task.add_done_callback(pending.discard)
async for raw in ws:
msg = json.loads(raw)
@@ -406,6 +413,9 @@ class HubClient:
elif mtype == "revocation" and on_revocation:
on_revocation(msg.get("token", ""))
+ elif mtype == "blocklist_update" and on_blocklist:
+ on_blocklist(msg.get("add") or [], msg.get("remove") or [])
+
elif mtype == "webrtc_offer" and on_webrtc_offer:
# Answered off the read loop on purpose. Awaiting the
# handler here meant one slow negotiation stopped the
@@ -461,26 +471,25 @@ class HubClient:
log.warning("Could not deliver WebRTC answer to %s: %s",
str(msg.get("peer_id"))[:8], e)
- # ── Swarm registration ─────────────────────────────────────────────────
+ # ── Content blocklist (public groups) ─────────────────────────────────
- async def register_swarm(self, content_hashes: list[str], endpoint: str) -> int:
- """Register file hashes in the hub swarm table. Returns count registered."""
+ async def fetch_blocklist(self, max_pages: int = 100) -> set[str]:
+ """The hub's content blocklist, every page of it (docs/MESHBAY_DESIGN.md §7.5)."""
if self._session is None:
raise RuntimeError("Not logged in")
await self.ensure_fresh_token()
-
- registered = 0
- for h in content_hashes:
- try:
- r = await self._http.post("/v1/swarm/register", json={
- "content_hash": h,
- "endpoint": endpoint,
- }, headers=self._session.auth_headers)
- if r.status_code in (201, 200):
- registered += 1
- except Exception:
- pass
- return registered
+ hashes: set[str] = set()
+ after = ""
+ for _ in range(max_pages):
+ r = await self._http.get("/v1/blocklist", params={"after": after},
+ headers=self._session.auth_headers)
+ r.raise_for_status()
+ page = r.json()
+ hashes.update(page.get("hashes") or [])
+ after = page.get("next") or ""
+ if not after:
+ return hashes
+ raise RuntimeError("Content blocklist longer than expected; not applied")
# ── Convenience: full startup sequence ───────────────────────────────────
diff --git a/packages/meshbay-node/src/meshbay_node/indexer/group_index.py b/packages/meshbay-node/src/meshbay_node/indexer/group_index.py
index 9e4e780..85bc13e 100644
--- a/packages/meshbay-node/src/meshbay_node/indexer/group_index.py
+++ b/packages/meshbay-node/src/meshbay_node/indexer/group_index.py
@@ -1,67 +1,41 @@
"""
-Mesh Group Index — encrypted file listing for a group.
+Mesh Group Index — the file listing for one group, and the deltas between versions.
-Wire format (private group):
- msgpack({entries, version, group_id}) → zstd compress → GEK ChaCha20 encrypt → sign
-
-Wire format (public group):
- msgpack({entries, version, group_id}) → sign (no encryption)
+What travels on the wire is built from it by `transport/wire.py` and sealed under
+the group key (`meshbay_common.groupbox`); this module holds no encoding of its own.
Delta format:
{base_version, version, additions: [...], deletions: [id, ...]}
"""
-import base64
import logging
-from dataclasses import asdict, dataclass, field
+from dataclasses import dataclass, field
-import blake3
-import msgpack
-import zstandard as zstd
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
-from meshbay_common.crypto import (
- pk_to_b64,
- sign_chunk,
- verify_chunk_signature,
-)
from meshbay_common.protocol import IndexDelta, IndexEntry
-from meshbay_common.webcrypto import (
- chunk_key_aes as derive_chunk_key,
-)
-from meshbay_common.webcrypto import (
- decrypt_chunk_aes as decrypt_chunk,
-)
-from meshbay_common.webcrypto import (
- encrypt_chunk_aes as encrypt_chunk,
-)
log = logging.getLogger(__name__)
-ZSTD_LEVEL = 3 # fast compression
-INDEX_CHUNK = 0 # the index itself is treated as chunk 0 of a virtual "index file"
-
@dataclass
class GroupIndex:
"""
- Encrypted, signed Mesh Group Index for one group.
+ Mesh Group Index for one group.
Usage:
idx = GroupIndex(group_id="...", sk_node=sk, gek=gek_bytes)
idx.add_entry(entry)
- wire_bytes = idx.serialize() # for sending to members
- recovered = GroupIndex.deserialize(wire_bytes, sk_node=sk, gek=gek_bytes)
"""
group_id: str
sk_node: Ed25519PrivateKey
gek: bytes | None = None # None → public group (no encryption)
version: int = 1
- # The group's roots and whether each is readable right now. Travels inside
- # the encrypted payload because it names the operator's directories, and a
- # member needs it to tell "temporarily unavailable" from "deleted" — a
- # distinction the entries alone cannot carry, since an unavailable root's
- # files are still listed. Absent in an index written before roots existed.
+ # The group's roots and whether each is readable right now, as the indexer
+ # last described them. A member needs this to tell "temporarily unavailable"
+ # from "deleted" — a distinction the entries alone cannot carry, since an
+ # unavailable root's files are still listed. What members receive is built
+ # from the RootSet by `transport/wire.py`, not from this copy.
roots: list = field(default_factory=list)
_entries: dict = field(default_factory=dict, repr=False) # id → IndexEntry
@@ -101,111 +75,6 @@ class GroupIndex:
idx._entries = dict(entries_by_id)
return idx
- # ── Serialisation ─────────────────────────────────────────────────────────
-
- def serialize(self) -> bytes:
- """
- Produce a signed index envelope:
- msgpack → zstd → [GEK encrypt if private] → sign → length-prefixed envelope
-
- **This is not an MNP message.** It was the payload of `index_sync` on the QUIC
- transport, while WebRTC sent plain entries under the same type — one message
- type, two encodings (2026-09-03). Both transports now build `index_sync` from
- `transport/wire.py`. This stays as a correct at-rest/interchange format, and
- as the only thing that signs and encrypts a whole index; read it as that, not
- as a wire contract.
- """
- payload = msgpack.packb({
- "group_id": self.group_id,
- "version": self.version,
- "roots": list(self.roots),
- "entries": [asdict(e) for e in self.entries],
- }, use_bin_type=True)
-
- compressed = zstd.compress(payload, level=ZSTD_LEVEL)
-
- if self.gek is not None:
- # Private group: encrypt with GEK-derived key
- idx_hash = blake3.blake3(compressed).digest()
- ckey = derive_chunk_key(self.gek, idx_hash, INDEX_CHUNK)
- nonce, ct = encrypt_chunk(ckey, compressed)
- sig = sign_chunk(self.sk_node, INDEX_CHUNK, nonce, blake3.blake3(ct).digest())
- envelope = msgpack.packb({
- "type": "index",
- "encrypted": True,
- "version": self.version,
- "group_id": self.group_id,
- "idx_hash_b64": base64.b64encode(idx_hash).decode(),
- "nonce_b64": base64.b64encode(nonce).decode(),
- "ct_b64": base64.b64encode(ct).decode(),
- "sig_b64": base64.b64encode(sig).decode(),
- "pk_node_b64": pk_to_b64(self.sk_node.public_key()),
- }, use_bin_type=True)
- else:
- # Public group: just sign the compressed payload
- payload_hash = blake3.blake3(compressed).digest()
- sig = sign_chunk(self.sk_node, INDEX_CHUNK,
- bytes(12), # zero nonce for plaintext
- payload_hash)
- envelope = msgpack.packb({
- "type": "index",
- "encrypted": False,
- "version": self.version,
- "group_id": self.group_id,
- "data_b64": base64.b64encode(compressed).decode(),
- "hash_b64": base64.b64encode(payload_hash).decode(),
- "sig_b64": base64.b64encode(sig).decode(),
- "pk_node_b64": pk_to_b64(self.sk_node.public_key()),
- }, use_bin_type=True)
-
- return envelope
-
- @classmethod
- def deserialize(
- cls,
- data: bytes,
- sk_node: Ed25519PrivateKey,
- gek: bytes | None = None,
- ) -> "GroupIndex":
- """Deserialize, verify signature, and decrypt (if private)."""
- from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey
-
- envelope = msgpack.unpackb(data, raw=False)
-
- pk_node_raw = base64.b64decode(envelope["pk_node_b64"])
- pk_node = Ed25519PublicKey.from_public_bytes(pk_node_raw)
- sig = base64.b64decode(envelope["sig_b64"])
-
- if envelope["encrypted"]:
- if gek is None:
- raise ValueError("GEK required to decrypt private group index")
- ct = base64.b64decode(envelope["ct_b64"])
- nonce = base64.b64decode(envelope["nonce_b64"])
- ct_hash = blake3.blake3(ct).digest()
- verify_chunk_signature(pk_node, INDEX_CHUNK, nonce, ct_hash, sig)
-
- idx_hash = base64.b64decode(envelope["idx_hash_b64"])
- ckey = derive_chunk_key(gek, idx_hash, INDEX_CHUNK)
- compressed = decrypt_chunk(ckey, nonce, ct)
- else:
- compressed = base64.b64decode(envelope["data_b64"])
- payload_hash = base64.b64decode(envelope["hash_b64"])
- verify_chunk_signature(pk_node, INDEX_CHUNK, bytes(12), payload_hash, sig)
-
- payload = msgpack.unpackb(zstd.decompress(compressed), raw=False)
- idx = cls(
- group_id=payload["group_id"],
- sk_node=sk_node,
- gek=gek,
- version=payload["version"],
- # Absent from an index written before roots existed; an empty list
- # reads as "nothing known about availability", not "no roots".
- roots=payload.get("roots") or [],
- )
- for e in payload["entries"]:
- idx.add_entry(IndexEntry(**e))
- return idx
-
# ── Delta ─────────────────────────────────────────────────────────────────
def diff(self, previous: "GroupIndex") -> IndexDelta:
diff --git a/packages/meshbay-node/src/meshbay_node/modules/__init__.py b/packages/meshbay-node/src/meshbay_node/modules/__init__.py
deleted file mode 100644
index e69de29..0000000
--- a/packages/meshbay-node/src/meshbay_node/modules/__init__.py
+++ /dev/null
diff --git a/packages/meshbay-node/src/meshbay_node/musicbrainz.py b/packages/meshbay-node/src/meshbay_node/musicbrainz.py
index ce2d233..e3681ea 100644
--- a/packages/meshbay-node/src/meshbay_node/musicbrainz.py
+++ b/packages/meshbay-node/src/meshbay_node/musicbrainz.py
@@ -183,10 +183,6 @@ class MusicBrainzClient:
results = (data or {}).get("releases", [])
return _best_match_release(artist, album, results)
- async def release_details(self, mbid: str) -> dict | None:
- """Full release details, including recordings (tracklist)."""
- return await self._get(_BASE_URL + f"release/{mbid}", {"inc": "recordings+artist-credits"})
-
async def fetch_cover_art(self, mbid: str) -> bytes | None:
"""
The release's front cover, or None if Cover Art Archive has nothing
diff --git a/packages/meshbay-node/src/meshbay_node/platform.py b/packages/meshbay-node/src/meshbay_node/platform.py
index 09397a2..2255036 100644
--- a/packages/meshbay-node/src/meshbay_node/platform.py
+++ b/packages/meshbay-node/src/meshbay_node/platform.py
@@ -226,6 +226,14 @@ def _startup_vbs() -> Path:
def _pid_alive(pid: int) -> bool:
"""Whether a process with this pid exists, in any session (tasklist lists
session 0 too, where a service-mode daemon runs; opening it would not)."""
+ if sys.platform != "win32":
+ # A zombie counts as gone: it has exited, and only its parent has not
+ # reaped it yet -- os.kill(pid, 0) would still find it.
+ try:
+ stat = Path(f"/proc/{pid}/stat").read_text(encoding="ascii")
+ except OSError:
+ return False
+ return stat.rsplit(")", 1)[1].split()[0] != "Z"
r = subprocess.run(["tasklist", "/FI", f"PID eq {pid}", "/NH", "/FO", "CSV"],
capture_output=True, text=True)
return f'"{pid}"' in r.stdout
diff --git a/packages/meshbay-node/src/meshbay_node/replication.py b/packages/meshbay-node/src/meshbay_node/replication.py
deleted file mode 100644
index 297ed0b..0000000
--- a/packages/meshbay-node/src/meshbay_node/replication.py
+++ /dev/null
@@ -1,140 +0,0 @@
-"""
-MeshBay Node — content replication (node-to-node, admin-authorized).
-
-A replication node downloads files from a source node and stores them
-locally, then registers itself as an additional swarm source in the hub.
-This provides redundancy and improves availability for public content.
-
-Only public content is replicated (no GEK needed).
-Private content replication requires the GEK and is admin-controlled.
-
-Usage:
- replicator = ContentReplicator(
- hub_url=..., access_token=..., source_endpoint=...,
- local_dir=Path("/data/replicated"), node_pk_b64=...,
- )
- await replicator.replicate_file(file_id, file_name, file_size)
-"""
-
-import logging
-from pathlib import Path
-
-import blake3
-import httpx
-
-log = logging.getLogger(__name__)
-
-CHUNK_SIZE = 1024 * 1024 # 1 MB
-
-
-class ContentReplicator:
- """
- Downloads public files from a source node and registers as swarm source.
- """
-
- def __init__(
- self,
- hub_url: str,
- access_token: str,
- source_endpoint: str, # "http://ip:port" of source node HTTP API
- local_dir: Path,
- node_pk_b64: str,
- ):
- self._hub_url = hub_url.rstrip("/")
- self._access_token = access_token
- self._source_endpoint = source_endpoint.rstrip("/")
- self._local_dir = local_dir
- self._node_pk_b64 = node_pk_b64
- self._local_dir.mkdir(parents=True, exist_ok=True)
-
- @property
- def _auth_headers(self) -> dict:
- return {"Authorization": f"Bearer {self._access_token}"}
-
- async def fetch_index(self) -> list[dict]:
- """Fetch the public Mesh Group Index from the source node."""
- async with httpx.AsyncClient(timeout=30) as c:
- r = await c.get(f"{self._source_endpoint}/index")
- r.raise_for_status()
- return r.json()["entries"]
-
- async def replicate_file(
- self,
- file_id: str,
- file_name: str,
- file_size: int,
- progress_cb=None,
- ) -> Path:
- """
- Download a public file from the source node, verify integrity,
- save locally, and register as swarm source in the hub.
- Returns the local file path.
- """
- local_path = self._local_dir / file_name
- if local_path.exists():
- # Verify hash
- existing_hash = blake3.blake3(local_path.read_bytes()).hexdigest()
- if existing_hash == file_id:
- log.info("Already have %s, skipping", file_name)
- await self._register_swarm(file_id, local_path)
- return local_path
-
- log.info("Replicating %s (%d bytes) from %s", file_name, file_size, self._source_endpoint)
-
- # Stream download chunk by chunk
- n_chunks = max(1, (file_size + CHUNK_SIZE - 1) // CHUNK_SIZE)
- with open(local_path, "wb") as f:
- async with httpx.AsyncClient(timeout=60) as c:
- for chunk_idx in range(n_chunks):
- # Download full file (simpler for public content)
- if chunk_idx == 0:
- r = await c.get(
- f"{self._source_endpoint}/file/{file_id}",
- headers=self._auth_headers,
- )
- r.raise_for_status()
- f.write(r.content)
- if progress_cb:
- progress_cb(len(r.content), file_size)
- break # full file downloaded in one request
-
- # Verify hash
- actual_hash = blake3.blake3(local_path.read_bytes()).hexdigest()
- if actual_hash != file_id:
- local_path.unlink(missing_ok=True)
- raise ValueError(f"Hash mismatch: expected {file_id[:16]}, got {actual_hash[:16]}")
-
- log.info("Replicated %s (hash OK)", file_name)
- await self._register_swarm(file_id, local_path)
- return local_path
-
- async def _register_swarm(self, content_hash: str, local_path: Path) -> None:
- """Register this node as a swarm source for the content hash in the hub."""
- try:
- async with httpx.AsyncClient(timeout=10) as c:
- r = await c.post(
- f"{self._hub_url}/v1/swarm/register",
- json={"content_hash": content_hash, "endpoint": self._source_endpoint},
- headers=self._auth_headers,
- )
- r.raise_for_status()
- log.debug("Registered as swarm source for %s", content_hash[:16])
- except Exception as e:
- log.warning("Failed to register swarm source: %s", e)
-
- async def replicate_all(self, progress_cb=None) -> list[Path]:
- """Replicate all public files from the source node."""
- entries = await self.fetch_index()
- results = []
- for entry in entries:
- try:
- path = await self.replicate_file(
- file_id=entry["id"],
- file_name=entry["name"],
- file_size=entry["size"],
- progress_cb=progress_cb,
- )
- results.append(path)
- except Exception as e:
- log.error("Failed to replicate %s: %s", entry["name"], e)
- return results
diff --git a/packages/meshbay-node/src/meshbay_node/transport/quic_client.py b/packages/meshbay-node/src/meshbay_node/transport/quic_client.py
index 8f3fb74..d548103 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/quic_client.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/quic_client.py
@@ -212,16 +212,19 @@ class QuicChunkClient:
"session — refusing to handshake without channel binding")
binding = quic_binding(self._peer_cert_der)
- # MNP 3.4: a signed challenge proves the node key before we send anything
- # else. A wrong signature is refused; an absent one is an older node.
- if reply.get("sig"):
- try:
- Ed25519PublicKey.from_public_bytes(
- base64.b64decode(reply.get("node_pk", ""))
- ).verify(base64.b64decode(reply["sig"]), challenge_transcript(
- self._group_id, nonce_c, nonce_s, binding))
- except Exception as exc:
- raise ConnectionError(f"Node challenge signature invalid: {exc}") from exc
+ # A signed challenge proves the node key before we send anything else.
+ # Every node we can reach signs (the floor is 4.0, signing is 3.4), and
+ # one without a certificate to bind could not complete the proof anyway,
+ # so a missing signature is refused like a wrong one.
+ if not reply.get("sig"):
+ raise ConnectionError("Node challenge is not signed")
+ try:
+ Ed25519PublicKey.from_public_bytes(
+ base64.b64decode(reply.get("node_pk", ""))
+ ).verify(base64.b64decode(reply["sig"]), challenge_transcript(
+ self._group_id, nonce_c, nonce_s, binding))
+ except Exception as exc:
+ raise ConnectionError(f"Node challenge signature invalid: {exc}") from exc
self._proto._send(self._ctrl_stream, {
"type": MNP.HANDSHAKE_RESPONSE,
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py
index f42d8e8..daaee62 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py
@@ -5,6 +5,7 @@ import base64
import os
import time
+import msgpack
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey
from meshbay_common import MNP_VERSION
from meshbay_common.adminop import (
@@ -42,6 +43,16 @@ from meshbay_common.adminop import (
from meshbay_common.crypto import pk_to_b64
from meshbay_common.protocol import MNP
+# What one connection may have waiting for a signature. Anyone authenticated can
+# ask for a challenge — the signature is what is checked, and it comes later — so
+# without a bound a member who never answers makes the node keep every request,
+# payload and all, for the life of the connection (§13.5b). A person signs one
+# operation at a time; a handful covers a settings page saving several at once.
+MAX_PENDING_ADMIN_OPS = 8
+# The subject and payload of one pending operation, packed. A path is at most a
+# few KiB, and the largest field a legitimate request carries is a directory list.
+MAX_ADMIN_OP_BYTES = 64 * 1024
+
# Which executor runs each signed operation once its signature has been
# checked. Every one runs as a task of the session.
_ADMIN_EXECUTORS = {
@@ -98,8 +109,22 @@ class AdminMixin:
(e.g. root management from a NodePage connection).
"""
gid = group_id if group_id is not None else (self._group_id or "")
+ now = time.time()
+ for op_id, pending in list(self._admin_ops.items()):
+ if now - pending["ts"] > ADMIN_CHALLENGE_TTL:
+ del self._admin_ops[op_id]
+ if len(self._admin_ops) >= MAX_PENDING_ADMIN_OPS:
+ self._send({"type": "error", "detail": "Too many operations waiting for a "
+ "signature", "code": "too_many_pending"})
+ self._audit("admin_pending_flood", op)
+ return
+ if len(msgpack.packb([subject, payload or {}], use_bin_type=True)) \
+ > MAX_ADMIN_OP_BYTES:
+ self._send({"type": "error", "detail": "Request too large",
+ "code": "too_large"})
+ return
nonce = os.urandom(32)
- ts = int(time.time())
+ ts = int(now)
op_id = base64.b64encode(os.urandom(16)).decode()
self._admin_ops[op_id] = {
"op": op, "subject": subject, "nonce": nonce, "ts": ts,
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py
index b02e59a..e678a13 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admission.py
@@ -8,7 +8,12 @@ import time
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey
from meshbay_common import MNP_VERSION
-from meshbay_common.adminop import OP_INVITE_CANCEL, OP_INVITE_CREATE, OP_INVITE_LINK_CREATE
+from meshbay_common.adminop import (
+ OP_INVITE_CANCEL,
+ OP_INVITE_CREATE,
+ OP_INVITE_LINK_CREATE,
+ invite_create_subject,
+)
from meshbay_common.crypto import wrap_gek_aes
from meshbay_common.device import (
DEVICE_TTL,
@@ -71,11 +76,13 @@ class AdmissionMixin:
})
return
- self._issue_admin_challenge(OP_INVITE_CREATE, invitee_id, {
- "group_id": group_id,
- "user_id": invitee_id,
- "username": str(msg.get("username", ""))[:64],
- })
+ username = str(msg.get("username", ""))[:64]
+ self._issue_admin_challenge(
+ OP_INVITE_CREATE, invite_create_subject(invitee_id, username), {
+ "group_id": group_id,
+ "user_id": invitee_id,
+ "username": username,
+ })
def _do_invite_link_create(self, msg: dict) -> None:
"""
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py
index f2b6fa5..857db33 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py
@@ -64,6 +64,8 @@ class MusicMixin:
"""
ctx = self._group_ctx()
file_id = msg.get("file_id", "")
+ if self._refuse_blocked(file_id):
+ return
entry = ctx["index"].get_entry(file_id)
if not entry:
self._send({"type": "error", "detail": "File not found"})
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py
index a425c65..4337e24 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py
@@ -228,6 +228,8 @@ class StreamingMixin:
async def _stream_video_inner(self, msg: dict) -> None:
ctx = self._group_ctx()
file_id = msg.get("file_id", "")
+ if self._refuse_blocked(file_id):
+ return
entry = ctx["index"].get_entry(file_id)
if not entry:
self._send({"type": "error", "detail": "File not found"})
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py
index bd2ab52..70781eb 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py
@@ -39,6 +39,8 @@ class SubtitlesMixin:
"""
ctx = self._group_ctx()
file_id = msg.get("file_id", "")
+ if self._refuse_blocked(file_id):
+ return
entry = ctx["index"].get_entry(file_id)
if not entry:
self._send({"type": "error", "detail": "File not found"})
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/video_meta.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/video_meta.py
index 834bac4..a9b35dc 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/video_meta.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/video_meta.py
@@ -11,6 +11,7 @@ from meshbay_common.adminop import (
OP_TMDB_ENABLED,
OP_TMDB_OVERRIDE,
OP_TMDB_REMATCH,
+ tmdb_config_subject,
)
from meshbay_common.protocol import MNP
@@ -60,12 +61,12 @@ class VideoMetaMixin:
if not self._has_admin_authority():
self._send({"type": "error", "detail": "No authorized key for this"})
return
- # The subject is the signed, audited, human-shown string — it must
- # never contain the token itself (it would end up in the audit log
- # in plaintext). The actual token travels only in `payload`, which
- # is node-side context, never re-sent or re-verified from the wire.
- # The language is not a secret, so it travels in the subject itself.
- subject = f"custom_token={'yes' if token else 'no'},language={language or 'default'}"
+ # The subject is the signed, audited string, so it must never contain
+ # the token itself (it would end up in the audit log in plaintext); it
+ # carries the token's SHA-256 instead, which binds the signature to this
+ # token without writing it down. `None` (unchanged) and `""` (clear)
+ # stay distinct, for the token and the language alike.
+ subject = tmdb_config_subject(token, language)
self._issue_admin_challenge(
OP_TMDB_CONFIG, subject,
payload={"token": token, "language": language},
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
index c3fb623..882ddbc 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
@@ -246,6 +246,25 @@ class SessionCore:
return self._ctx["groups"].get(self._group_id) or {}
return self._ctx
+ def _hidden_ids(self):
+ """The content blocklist, in a public group; nothing anywhere else.
+
+ A private group's content never reaches the hub, so nothing in it can have
+ been blocked there (docs/MESHBAY_DESIGN.md §7.5, `meshbay_node.blocklist`).
+ """
+ if self._group_ctx().get("visibility") != "public":
+ return ()
+ return self._ctx.get("blocklist") or ()
+
+ def _refuse_blocked(self, file_id) -> bool:
+ """Refuse a file the blocklist names, in a public group. True if refused."""
+ if not isinstance(file_id, str) or file_id not in self._hidden_ids():
+ return False
+ self._send({"type": "error", "detail": "This file is not available here.",
+ "code": "content_blocked", "file_id": file_id})
+ self._audit("content_blocked", file_id[:16])
+ return True
+
def _register_peer(self) -> None:
"""Add this connection to its group's peer set.
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py
index e76e23e..f629563 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/files.py
@@ -179,7 +179,7 @@ class FilesMixin:
def _do_index_sync(self) -> None:
ctx = self._group_ctx()
- self._send(index_sync_message(ctx["index"], ctx.get("roots")))
+ self._send(index_sync_message(ctx["index"], ctx.get("roots"), self._hidden_ids()))
async def _try_serve_thumbnail(
self, thumb_hash: str, chunk_index: int, gek: bytes | None,
@@ -236,8 +236,18 @@ class FilesMixin:
self._note_unleased(tr)
file_id = msg["file_id"]
chunk_index = msg["chunk_index"]
+ if self._refuse_blocked(file_id):
+ return
entry = ctx["index"].get_entry(file_id)
if not entry:
+ # A blocked file's own thumbnail is a preview of it.
+ hidden = self._hidden_ids()
+ if hidden and file_id in {
+ e.thumb_hash for e in map(ctx["index"].get_entry, hidden)
+ if e is not None and e.thumb_hash}:
+ self._send({"type": "error", "detail": "This file is not available here.",
+ "code": "content_blocked", "file_id": file_id})
+ return
thumb = await self._try_serve_thumbnail(file_id, chunk_index, ctx.get("gek"))
if thumb is not None:
log.debug("file_req file_id=%s chunk=%s: served as thumbnail",
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/handshake.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/handshake.py
index eceb2b4..f701072 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/handshake.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/handshake.py
@@ -24,7 +24,6 @@ from meshbay_common.handshake import (
)
from meshbay_common.protocol import MNP
-from meshbay_node import transfers as transfers_mod
from meshbay_node.indexer.indexer import DirectoryIndexer
from meshbay_node.transport.webrtc.channel import _extract_dtls_fingerprint, _get_remote_ip
from meshbay_node.transport.webrtc.limits import MAX_MSG
@@ -266,19 +265,6 @@ class HandshakeMixin:
# No `chat_encrypted` beside it: there is no switch. A peer that
# reached this point speaks MNP 2.0, and 2.0 has no plaintext chat.
"chat_epoch": int(self._group_ctx().get("chat_epoch", 0) or 0),
- # This member's own transfer caps in this group, so the interface
- # can say "2 of 2 of your slots are busy" rather than draw a bare
- # spinner. Absent reads as "no limit known" and the hint is simply
- # not drawn — never as "unlimited", which would have the interface
- # contradicting the node.
- "transfer_limits": {
- "download": self._slots().member_cap(
- transfers_mod.DOWNLOAD,
- (self._group_id or "", self._user_id or "")),
- "upload": self._slots().member_cap(
- transfers_mod.UPLOAD,
- (self._group_id or "", self._user_id or "")),
- },
# So a client that connects mid-scan shows the indexing state
# immediately, instead of waiting for the next periodic
# INDEX_PROGRESS push. Never a path or filename — see
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py
index 9d58d8c..f7bbbfa 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py
@@ -14,6 +14,8 @@ from meshbay_common.adminop import (
OP_ROOT_UPDATE,
OP_SET_SCAN_SETTINGS,
OP_TRANSFER_LIMITS,
+ group_attach_subject,
+ root_add_subject,
)
from meshbay_common.protocol import MNP
@@ -246,10 +248,11 @@ class NodeOpsMixin:
# is ignored rather than obeyed: on load it forces every other root
# read-only, which is the model the RO/RW one replaced. A second
# writable directory is `root_add` with `writable`.
+ writable = bool(msg.get("writable", True))
+ # The directory being exposed is signed, not only the group's name.
self._issue_admin_challenge(
- OP_GROUP_ATTACH, name,
- payload={"name": name, "shared_dir": shared_dir,
- "writable": bool(msg.get("writable", True))},
+ OP_GROUP_ATTACH, group_attach_subject(name, shared_dir, writable),
+ payload={"name": name, "shared_dir": shared_dir, "writable": writable},
group_id="")
async def _admin_exec_group_attach(
@@ -342,16 +345,20 @@ class NodeOpsMixin:
if not self._has_admin_authority():
self._send({"type": "error", "detail": "No authorized key for this"})
return
+ payload = {
+ "group_id": target_group, "path": path,
+ "name": str(msg.get("name", ""))[:128],
+ "kind": str(msg.get("kind", "generic"))[:16],
+ "writable": bool(msg.get("writable", msg.get("upload", False))),
+ "removable": bool(msg.get("removable", False)),
+ }
+ # Everything the executor acts on is signed — `writable` decides whether
+ # every member may write there. The group is in the transcript itself.
self._issue_admin_challenge(
- OP_ROOT_ADD, path,
- payload={
- "group_id": target_group, "path": path,
- "name": str(msg.get("name", ""))[:128],
- "kind": str(msg.get("kind", "generic"))[:16],
- "writable": bool(msg.get("writable", msg.get("upload", False))),
- "removable": bool(msg.get("removable", False)),
- },
- group_id=target_group)
+ OP_ROOT_ADD,
+ root_add_subject(path, payload["name"], payload["kind"],
+ payload["writable"], payload["removable"]),
+ payload=payload, group_id=target_group)
async def _admin_exec_root_add(
self, pending: dict, transcript: bytes, sig: bytes,
diff --git a/packages/meshbay-node/src/meshbay_node/transport/wire.py b/packages/meshbay-node/src/meshbay_node/transport/wire.py
index 6986b01..258edb5 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/wire.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/wire.py
@@ -12,11 +12,10 @@ envelope produced by `GroupIndex.serialize()`. Same message type, two encodings,
consumer each and nothing asserting they matched. Same failure mode as the two
`file_chunk` encoders, and the same fix: one builder, used by both.
-`GroupIndex.serialize()`/`deserialize()` are unchanged and still tested — they remain
-a correct signed index envelope — but they no longer describe any MNP message. Read
-them as an at-rest/interchange format, not as a wire contract. It is also not a
-candidate for reuse below: it compresses with zstd, which no browser can decompress
-(`DecompressionStream` offers gzip and deflate only).
+That envelope, and `GroupIndex.serialize()`/`deserialize()` which produced it, are
+gone: once both transports built `index_sync` here, nothing stored or exchanged it.
+It was never a candidate for reuse below either — it compressed with zstd, which no
+browser can decompress (`DecompressionStream` offers gzip and deflate only).
Since MNP 1.0 both messages carry their payload **sealed under a GEK-derived subkey**
(`meshbay_common.groupbox`). Only the routing fields — `type`, `v`, `group_id` — stay
@@ -66,7 +65,7 @@ def list_dirs(roots: RootSet | None) -> list[str]:
return sorted(out)[:MAX_DIRS]
-def index_sync_message(index, roots: RootSet | None) -> dict:
+def index_sync_message(index, roots: RootSet | None, hidden=()) -> dict:
"""
The full `index_sync` message for one group.
@@ -74,10 +73,13 @@ def index_sync_message(index, roots: RootSet | None) -> dict:
without them a folder someone just created, or one they emptied, does not exist as
far as a client is concerned, and a member cannot tell "the drive is unplugged"
from "it is all still there".
+
+ `hidden` is the content blocklist in a public group (`meshbay_node.blocklist`):
+ those entries are not sent at all.
"""
payload = {
"version": index.version,
- "entries": [index_entry_wire(e) for e in index.entries],
+ "entries": [index_entry_wire(e) for e in index.entries if e.id not in hidden],
"dirs": list_dirs(roots),
"roots": roots.describe() if roots else [],
}
@@ -89,7 +91,7 @@ def index_sync_message(index, roots: RootSet | None) -> dict:
}
-def index_delta_message(index, delta, roots=None) -> dict:
+def index_delta_message(index, delta, roots=None, hidden=()) -> dict:
"""
One `index_delta` — what changed since the last thing this node broadcast.
@@ -104,13 +106,16 @@ def index_delta_message(index, delta, roots=None) -> dict:
page: the delta that told them something had changed was the one message
that could not say what. It is a handful of dicts, bounded by the number of
directories a group has, and it is sealed with the rest.
+
+ `hidden`, as for `index_sync_message`: a blocked entry is never added or
+ updated; its deletion still goes out.
"""
payload = {
"base_version": delta.base_version,
"version": delta.version,
- "additions": [index_entry_wire(e) for e in delta.additions],
+ "additions": [index_entry_wire(e) for e in delta.additions if e.id not in hidden],
"deletions": list(delta.deletions),
- "updates": [index_entry_wire(e) for e in delta.updates],
+ "updates": [index_entry_wire(e) for e in delta.updates if e.id not in hidden],
}
if roots is not None:
payload["roots"] = roots.describe()
diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py
index 50beb26..bbc4649 100644
--- a/packages/meshbay-node/src/meshbay_node/ui/app.py
+++ b/packages/meshbay-node/src/meshbay_node/ui/app.py
@@ -20,7 +20,6 @@ from fastapi.responses import JSONResponse
from meshbay_common.background import spawn
from meshbay_node import __version__, ops
-from meshbay_node.indexer.indexer import DirectoryIndexer
log = logging.getLogger(__name__)
@@ -493,50 +492,6 @@ def create_ui_app(state: dict) -> FastAPI:
raise HTTPException(400, "apps must be a non-empty list")
return await _op(lambda: ops.set_enabled_apps(state, group_id, apps))
- # ── App directories (operator only, localhost) ────────────────────────
- #
- # The loopback twin of the `app_directories` MNP op. One endpoint for every
- # application, keyed by the app's own name, so adding one needs no route
- # here — the same reason the op is generic. `ALLOWED_APPS` is checked on
- # the MNP path; here the caller is already on localhost holding the run
- # token, and `ops` refuses a directory outside the group's roots either
- # way, so an unknown key writes one unread settings row and nothing else.
-
- @app.put("/api/groups/{group_id}/app-directories/{app_key}")
- async def set_app_directories(group_id: str, app_key: str, payload: dict):
- dirs = payload.get("directories")
- if not isinstance(dirs, list):
- raise HTTPException(400, "directories must be a list")
- return await _op(lambda: ops.set_app_directories(
- state, group_id, app_key, [str(d) for d in dirs]))
-
- @app.put("/api/groups/{group_id}/chat-directory")
- async def set_chat_directory(group_id: str, payload: dict):
- return await _op(lambda: ops.set_chat_directory(
- state, group_id, str(payload.get("path") or "")))
-
- @app.put("/api/groups/{group_id}/chat-link-preview")
- async def set_chat_link_preview(group_id: str, payload: dict):
- return await _op(lambda: ops.set_chat_link_preview(
- state, group_id, bool(payload.get("enabled", True))))
-
- @app.put("/api/groups/{group_id}/search-listed")
- async def set_search_listed(group_id: str, payload: dict):
- return await _op(lambda: ops.set_search_listed(
- state, group_id, bool(payload.get("listed", True))))
-
- # ── Scan settings (operator only, localhost) ──────────────────────────
-
- @app.put("/api/groups/{group_id}/scan-settings")
- async def set_scan_settings(group_id: str, payload: dict):
- return await _op(lambda: ops.set_scan_settings(
- state, group_id,
- float(payload.get("reconcile_interval_secs",
- DirectoryIndexer.DEFAULT_RECONCILE_SECS)),
- float(payload.get("debounce_secs",
- DirectoryIndexer.DEFAULT_DEBOUNCE_SECS)),
- ))
-
# ── Reload config ────────────────────────────────────────────────────
@app.post("/api/reload")