summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/daemon.py
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/daemon.py')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py60
1 files changed, 59 insertions, 1 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py
index 6a8bbc9..5b6e770 100644
--- a/packages/meshbay-node/src/meshbay_node/daemon.py
+++ b/packages/meshbay-node/src/meshbay_node/daemon.py
@@ -31,6 +31,8 @@ from pathlib import Path
import uvicorn
+from meshbay_common import MNP_VERSION
+from meshbay_common.protocol import MNP
from meshbay_node.chat.store import ChatStore
from meshbay_node.config import Config, load_config, write_example_config
from meshbay_node.hub_client import HubClient, HubConfig
@@ -161,6 +163,7 @@ class NodeDaemon:
group_id=group_cfg.id,
sk_node=keys.sk_ed25519,
gek=gek,
+ on_change=self._on_index_change,
)
await indexer.start()
self._indexers.append(indexer)
@@ -330,7 +333,14 @@ class NodeDaemon:
"yes" if self._webrtc else "no",
"yes" if self._quic_server else "no")
- # 11. Wait for shutdown
+ # 11. Initial swarm registration
+ endpoint = f"webrtc:{self._config.node.quic_port}"
+ for gctx in groups_ctx.values():
+ hashes = [e.id for e in gctx["index"].entries]
+ if hashes:
+ asyncio.ensure_future(self._register_swarm(hashes, endpoint))
+
+ # 12. Wait for shutdown
stop_event = asyncio.Event()
loop = asyncio.get_event_loop()
for sig in (signal.SIGINT, signal.SIGTERM):
@@ -339,6 +349,54 @@ class NodeDaemon:
await self._shutdown()
+ async def _on_index_change(self, indexer: DirectoryIndexer) -> None:
+ """Called when a DirectoryIndexer detects file changes."""
+ group_id = indexer.group_id
+ idx = indexer.index
+ log.info("Index changed for group %s: %d files (v%d)",
+ group_id[:8], idx.count, idx.version)
+
+ # 11.5 — Push updated index to connected WebRTC peers in this group
+ if self._webrtc:
+ entries = [
+ {
+ "id": e.id, "name": e.name, "path": e.path,
+ "size": e.size, "type": e.type, "added_at": e.added_at,
+ }
+ for e in idx.entries
+ ]
+ sync_msg = {
+ "type": MNP.INDEX_SYNC,
+ "v": MNP_VERSION,
+ "group_id": idx.group_id,
+ "version": idx.version,
+ "entries": entries,
+ }
+ pushed = 0
+ for session in list(self._webrtc._sessions.values()):
+ if session._group_id == group_id:
+ try:
+ session._send(sync_msg)
+ pushed += 1
+ except Exception:
+ pass
+ if pushed:
+ log.info("Index pushed to %d WebRTC peers", pushed)
+
+ # 11.9 — Register file hashes with hub swarm table
+ if self._hub and self._state.get("endpoint_hint"):
+ hashes = [e.id for e in idx.entries]
+ if hashes:
+ endpoint = f"webrtc:{self._config.node.quic_port}"
+ asyncio.ensure_future(self._register_swarm(hashes, endpoint))
+
+ 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"