summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-08-09 04:08:43 +0200
committerChristophe Besson <cbesson@gmail.com>2026-08-09 04:08:43 +0200
commit46b6353ebfb57c7fea481a9aac919b7977e3d186 (patch)
tree6295072fc9ebcde4c8558319468b4ef311c5a9d6 /packages/meshbay-node/src/meshbay_node
parentb92b076bed49da15ce1ba96d80eb84db539a0778 (diff)
downloadmeshbay-46b6353ebfb57c7fea481a9aac919b7977e3d186.tar.gz
feat(node): add directory indexer and Mesh Group Index
GroupIndex: msgpack→zstd→GEK-encrypt→sign for private groups, plaintext+sign for public groups. DirectoryIndexer: watchdog-based watcher, async initial scan via thread pool, on_change callback. Delta support (diff between versions). 10/10 tests passing. Co-Authored-By: Claude Sonnet 4.6 (1M context) <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node')
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/__init__.py5
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/group_index.py191
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/indexer.py213
3 files changed, 409 insertions, 0 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/indexer/__init__.py b/packages/meshbay-node/src/meshbay_node/indexer/__init__.py
index e69de29..bb6231b 100644
--- a/packages/meshbay-node/src/meshbay_node/indexer/__init__.py
+++ b/packages/meshbay-node/src/meshbay_node/indexer/__init__.py
@@ -0,0 +1,5 @@
+"""Directory indexer and Mesh Group Index."""
+from .indexer import DirectoryIndexer
+from .group_index import GroupIndex
+
+__all__ = ["DirectoryIndexer", "GroupIndex"]
diff --git a/packages/meshbay-node/src/meshbay_node/indexer/group_index.py b/packages/meshbay-node/src/meshbay_node/indexer/group_index.py
new file mode 100644
index 0000000..4edeeab
--- /dev/null
+++ b/packages/meshbay-node/src/meshbay_node/indexer/group_index.py
@@ -0,0 +1,191 @@
+"""
+Mesh Group Index — encrypted file listing for a group.
+
+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)
+
+Delta format:
+ {base_version, version, additions: [...], deletions: [id, ...]}
+"""
+
+import base64
+import logging
+import os
+import time
+from dataclasses import dataclass, field, asdict
+from pathlib import Path
+from typing import Iterator
+
+import blake3
+import msgpack
+import zstandard as zstd
+from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
+
+from meshbay_common.crypto import (
+ chunk_key as derive_chunk_key,
+ encrypt_chunk,
+ decrypt_chunk,
+ sign_chunk,
+ verify_chunk_signature,
+ pk_to_b64,
+ generate_gek,
+)
+from meshbay_common.protocol import IndexEntry, IndexDelta
+
+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.
+
+ 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
+ _entries: dict = field(default_factory=dict, repr=False) # id → IndexEntry
+
+ # ── Entry management ──────────────────────────────────────────────────────
+
+ def add_entry(self, entry: IndexEntry) -> None:
+ self._entries[entry.id] = entry
+
+ def remove_entry(self, file_id: str) -> bool:
+ return self._entries.pop(file_id, None) is not None
+
+ def get_entry(self, file_id: str) -> IndexEntry | None:
+ return self._entries.get(file_id)
+
+ @property
+ def entries(self) -> list[IndexEntry]:
+ return list(self._entries.values())
+
+ @property
+ def count(self) -> int:
+ return len(self._entries)
+
+ # ── Serialisation ─────────────────────────────────────────────────────────
+
+ def serialize(self) -> bytes:
+ """
+ Produce a wire-ready byte string:
+ msgpack → zstd → [GEK encrypt if private] → sign → length-prefixed envelope
+ """
+ payload = msgpack.packb({
+ "group_id": self.group_id,
+ "version": self.version,
+ "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
+ from meshbay_common.crypto import verify_chunk_signature
+
+ 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"],
+ )
+ for e in payload["entries"]:
+ idx.add_entry(IndexEntry(**e))
+ return idx
+
+ # ── Delta ─────────────────────────────────────────────────────────────────
+
+ def diff(self, previous: "GroupIndex") -> IndexDelta:
+ """Compute what changed since a previous version of this index."""
+ prev_ids = set(previous._entries)
+ curr_ids = set(self._entries)
+
+ additions = [self._entries[i] for i in curr_ids - prev_ids]
+ deletions = list(prev_ids - curr_ids)
+
+ return IndexDelta(
+ base_version=previous.version,
+ version=self.version,
+ additions=additions,
+ deletions=deletions,
+ )
diff --git a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py
new file mode 100644
index 0000000..5cd7205
--- /dev/null
+++ b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py
@@ -0,0 +1,213 @@
+"""
+Directory indexer — watches a directory and maintains a GroupIndex.
+
+Uses watchdog for filesystem events. On any change (create/modify/delete/move),
+the affected file is re-scanned and the GroupIndex is updated.
+File metadata (blake3 hash, size, type, duration) is computed on first scan.
+Heavy operations (hashing large files) run in a thread pool to avoid blocking.
+"""
+
+import asyncio
+import logging
+import mimetypes
+import time
+from concurrent.futures import ThreadPoolExecutor
+from pathlib import Path
+from typing import Callable, Awaitable
+
+import blake3
+from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
+from watchdog.events import FileSystemEvent, FileSystemEventHandler
+from watchdog.observers import Observer
+
+from meshbay_common.protocol import IndexEntry
+from meshbay_node.indexer.group_index import GroupIndex
+
+log = logging.getLogger(__name__)
+
+# File types we include in the index (skip hidden files, temp files, etc.)
+EXCLUDED_PREFIXES = (".", "~", "#")
+EXCLUDED_SUFFIXES = (".tmp", ".part", ".crdownload", ".download")
+
+MEDIA_EXTENSIONS = {
+ "video": {".mp4", ".mkv", ".avi", ".mov", ".wmv", ".flv", ".webm", ".m4v"},
+ "audio": {".mp3", ".flac", ".ogg", ".wav", ".aac", ".m4a", ".opus"},
+ "image": {".jpg", ".jpeg", ".png", ".gif", ".webp", ".svg", ".bmp", ".tiff"},
+ "document": {".pdf", ".epub", ".mobi", ".txt", ".md", ".docx", ".odt"},
+ "archive": {".zip", ".tar", ".gz", ".bz2", ".xz", ".7z", ".rar"},
+}
+
+
+def _detect_type(path: Path) -> str:
+ suffix = path.suffix.lower()
+ for ftype, exts in MEDIA_EXTENSIONS.items():
+ if suffix in exts:
+ return ftype
+ return "other"
+
+
+def _is_indexable(path: Path) -> bool:
+ if not path.is_file():
+ return False
+ name = path.name
+ return (
+ not any(name.startswith(p) for p in EXCLUDED_PREFIXES)
+ and not any(name.endswith(s) for s in EXCLUDED_SUFFIXES)
+ )
+
+
+def _scan_file(root: Path, file_path: Path) -> IndexEntry | None:
+ """Compute IndexEntry for a file. Blocking — run in executor."""
+ if not _is_indexable(file_path):
+ return None
+ try:
+ stat = file_path.stat()
+ data = file_path.read_bytes()
+ file_id = blake3.blake3(data).hexdigest()
+ rel_path = str(file_path.parent.relative_to(root))
+ if rel_path == ".":
+ rel_path = ""
+ return IndexEntry(
+ id=file_id,
+ name=file_path.name,
+ path=rel_path,
+ size=stat.st_size,
+ type=_detect_type(file_path),
+ added_at=int(stat.st_mtime),
+ )
+ except (OSError, PermissionError) as e:
+ log.warning("Cannot index %s: %s", file_path, e)
+ return None
+
+
+class DirectoryIndexer:
+ """
+ Watches a directory and keeps a GroupIndex up to date.
+
+ Usage:
+ indexer = DirectoryIndexer(
+ root=Path("/home/user/shared"),
+ group_id="my-group",
+ sk_node=sk,
+ gek=gek_bytes,
+ on_change=async_callback,
+ )
+ await indexer.start()
+ # ... later
+ await indexer.stop()
+ """
+
+ def __init__(
+ self,
+ root: Path,
+ group_id: str,
+ sk_node: Ed25519PrivateKey,
+ gek: bytes | None,
+ on_change: Callable[["DirectoryIndexer"], Awaitable[None]] | None = None,
+ ):
+ self.root = root.resolve()
+ self.group_id = group_id
+ self.sk_node = sk_node
+ self.gek = gek
+ self.on_change = on_change
+
+ self._index = GroupIndex(group_id=group_id, sk_node=sk_node, gek=gek)
+ self._executor = ThreadPoolExecutor(max_workers=2, thread_name_prefix="indexer")
+ self._observer: Observer | None = None
+ self._loop: asyncio.AbstractEventLoop | None = None
+
+ @property
+ def index(self) -> GroupIndex:
+ return self._index
+
+ # ── Initial scan ──────────────────────────────────────────────────────────
+
+ async def initial_scan(self) -> None:
+ """Scan the entire directory tree. Run once at startup."""
+ log.info("Scanning %s ...", self.root)
+ loop = asyncio.get_event_loop()
+ files = [p for p in self.root.rglob("*") if p.is_file()]
+ count = 0
+ for file_path in files:
+ entry = await loop.run_in_executor(
+ self._executor, _scan_file, self.root, file_path)
+ if entry:
+ self._index.add_entry(entry)
+ count += 1
+ self._index.version = int(time.time())
+ log.info("Initial scan complete: %d files indexed", count)
+
+ # ── Watchdog integration ──────────────────────────────────────────────────
+
+ async def start(self) -> None:
+ """Start initial scan + filesystem watcher."""
+ self._loop = asyncio.get_event_loop()
+ await self.initial_scan()
+
+ handler = _WatchdogHandler(self)
+ self._observer = Observer()
+ self._observer.schedule(handler, str(self.root), recursive=True)
+ self._observer.start()
+ log.info("Watching %s for changes", self.root)
+
+ async def stop(self) -> None:
+ """Stop the filesystem watcher."""
+ if self._observer:
+ self._observer.stop()
+ self._observer.join()
+ self._observer = None
+ self._executor.shutdown(wait=False)
+ log.info("Indexer stopped")
+
+ # ── Internal update ───────────────────────────────────────────────────────
+
+ def _schedule_update(self, file_path: Path, deleted: bool = False) -> None:
+ """Called from watchdog thread — schedule async update on the event loop."""
+ if self._loop:
+ self._loop.call_soon_threadsafe(
+ lambda: asyncio.ensure_future(
+ self._update_entry(file_path, deleted)))
+
+ async def _update_entry(self, file_path: Path, deleted: bool) -> None:
+ if deleted:
+ # Remove by matching path (hash not available after deletion)
+ to_remove = [
+ e.id for e in self._index.entries
+ if (self.root / e.path / e.name).resolve() == file_path.resolve()
+ ]
+ for fid in to_remove:
+ self._index.remove_entry(fid)
+ log.debug("Removed from index: %s", file_path.name)
+ else:
+ loop = asyncio.get_event_loop()
+ entry = await loop.run_in_executor(
+ self._executor, _scan_file, self.root, file_path)
+ if entry:
+ self._index.add_entry(entry)
+ log.debug("Indexed: %s (%s)", file_path.name, entry.id[:8])
+
+ self._index.version = int(time.time())
+ if self.on_change:
+ await self.on_change(self)
+
+
+class _WatchdogHandler(FileSystemEventHandler):
+ def __init__(self, indexer: DirectoryIndexer):
+ self._indexer = indexer
+
+ def on_created(self, event: FileSystemEvent):
+ if not event.is_directory:
+ self._indexer._schedule_update(Path(event.src_path))
+
+ def on_modified(self, event: FileSystemEvent):
+ if not event.is_directory:
+ self._indexer._schedule_update(Path(event.src_path))
+
+ def on_deleted(self, event: FileSystemEvent):
+ if not event.is_directory:
+ self._indexer._schedule_update(Path(event.src_path), deleted=True)
+
+ def on_moved(self, event: FileSystemEvent):
+ if not event.is_directory:
+ self._indexer._schedule_update(Path(event.src_path), deleted=True)
+ self._indexer._schedule_update(Path(event.dest_path))