diff options
Diffstat (limited to 'packages/meshbay-node')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/daemon.py | 28 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/hub_client.py | 22 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_daemon.py | 74 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_security_regressions.py | 21 |
4 files changed, 14 insertions, 131 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index 220a908..a29aa9f 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -667,12 +667,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 @@ -1468,21 +1462,6 @@ 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)) - def _drop_group_sessions(self, group_id: str) -> None: """Close live sessions for a revoked group (H4).""" if not self._webrtc or not group_id: @@ -1492,13 +1471,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..5f47090 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. @@ -461,27 +460,6 @@ class HubClient: log.warning("Could not deliver WebRTC answer to %s: %s", str(msg.get("peer_id"))[:8], e) - # ── Swarm registration ───────────────────────────────────────────────── - - async def register_swarm(self, content_hashes: list[str], endpoint: str) -> int: - """Register file hashes in the hub swarm table. Returns count registered.""" - 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 - # ── Convenience: full startup sequence ─────────────────────────────────── async def startup(self, endpoint_hint: str | None = None) -> HubSession: diff --git a/packages/meshbay-node/tests/test_daemon.py b/packages/meshbay-node/tests/test_daemon.py index 8c4da2d..acaafac 100644 --- a/packages/meshbay-node/tests/test_daemon.py +++ b/packages/meshbay-node/tests/test_daemon.py @@ -2,7 +2,7 @@ Integration test: Node daemon wires all components correctly. Phase 11 — verifies that NodeDaemon creates chat stores, WebRTC transport, -index push on change, swarm registration, and shuts down cleanly. +index push on change, and shuts down cleanly. Hub interaction is mocked. """ @@ -241,7 +241,6 @@ async def test_daemon_index_change_pushes_to_peers(tmp_path, shared_dir, gek, hu daemon = NodeDaemon(config) daemon._broadcast_coalesce_secs = 0.01 # real value would make this test wait 0.5s daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=2) daemon._state["endpoint_hint"] = "node123" sk_node = Ed25519PrivateKey.generate() @@ -268,47 +267,10 @@ async def test_daemon_index_change_pushes_to_peers(tmp_path, shared_dir, gek, hu payload = unseal(gek, PURPOSE_INDEX, "index_sync", "a" * 32, msg) assert len(payload["entries"]) == indexer.index.count - # Finding H7: this group is private, so its content hashes must NOT be - # registered with the hub. The test previously asserted the opposite — - # publishing a fingerprint of every private file was treated as expected - # behaviour. Index push to members is unaffected (asserted above). + # Finding H7: a change to the index tells the hub nothing — no content hash + # of any group reaches it. Index push to members is unaffected (above). await asyncio.sleep(0.1) - daemon._hub.register_swarm.assert_not_called() - -@pytest.mark.asyncio -async def test_daemon_index_change_registers_swarm_for_public_group( - tmp_path, shared_dir, gek, hub_pk_pem): - """Public groups still register content hashes with the hub swarm (H7).""" - config = Config( - hub=HubConfig(url="http://localhost:9999", username="testuser"), - node=NodeConfig(quic_port=_free_port(), ui_port=_free_port()), - groups=[GroupConfig( - id="a" * 32, - name="public-group", - shared_dir=str(shared_dir), - visibility="public", - quic_port=29010, - )], - keystore=KeystoreConfig(path=tmp_path / "keystore.enc"), - data_dir=tmp_path / "data", - ) - daemon = NodeDaemon(config) - daemon._broadcast_coalesce_secs = 0.01 - daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=2) - daemon._state["endpoint_hint"] = "node123" - - indexer = DirectoryIndexer( - roots=one_root(shared_dir), group_id="a" * 32, - sk_node=Ed25519PrivateKey.generate(), gek=gek) - await indexer.initial_scan() - - await daemon._on_index_change(indexer) - - await asyncio.sleep(0.1) - daemon._hub.register_swarm.assert_called_once() - assert len(daemon._hub.register_swarm.call_args[0][0]) == indexer.index.count - + assert daemon._hub.mock_calls == [] @pytest.mark.asyncio async def test_daemon_index_change_skips_other_group_peers( @@ -325,7 +287,6 @@ async def test_daemon_index_change_skips_other_group_peers( daemon = NodeDaemon(config) daemon._broadcast_coalesce_secs = 0.01 daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=0) daemon._state["endpoint_hint"] = "node123" sk_node = Ed25519PrivateKey.generate() @@ -370,7 +331,6 @@ def _new_daemon_for_group(tmp_path, shared_dir, gek, group_id="a" * 32, daemon = NodeDaemon(config) daemon._broadcast_coalesce_secs = 0.01 daemon._hub = AsyncMock() - daemon._hub.register_swarm = AsyncMock(return_value=0) daemon._state["endpoint_hint"] = "node123" return daemon @@ -532,29 +492,3 @@ async def test_a_burst_of_changes_produces_one_broadcast(tmp_path, shared_dir, g await asyncio.sleep(0.05) session._send.assert_called_once() - - -@pytest.mark.asyncio -async def test_swarm_registration_only_sends_new_hashes_after_the_first( - tmp_path, shared_dir, gek): - daemon = _new_daemon_for_group(tmp_path, shared_dir, gek, visibility="public") - indexer = DirectoryIndexer( - roots=one_root(shared_dir), group_id="a" * 32, - sk_node=Ed25519PrivateKey.generate(), gek=gek) - await indexer.initial_scan() - total_files = indexer.index.count - - await daemon._on_index_change(indexer) - await asyncio.sleep(0.05) - assert len(daemon._hub.register_swarm.call_args_list[0].args[0]) == total_files - - from meshbay_common.protocol import IndexEntry - indexer.index.add_entry(IndexEntry(id="new-file-id", name="new.mp4", - path="shared", size=10, type="video", - added_at=0)) - await daemon._on_index_change(indexer) - await asyncio.sleep(0.05) - - assert daemon._hub.register_swarm.call_count == 2 - assert daemon._hub.register_swarm.call_args_list[1].args[0] == ["new-file-id"], \ - "only the newly added hash must be (re-)registered, not the whole library" diff --git a/packages/meshbay-node/tests/test_security_regressions.py b/packages/meshbay-node/tests/test_security_regressions.py index c0c2c2c..43e88c2 100644 --- a/packages/meshbay-node/tests/test_security_regressions.py +++ b/packages/meshbay-node/tests/test_security_regressions.py @@ -602,19 +602,18 @@ def test_denylist_persists_and_honours_groups(tmp_path): assert not reloaded.is_denied("someone", "other", "g-allowed") -def test_swarm_registration_skips_private_groups(): +def test_the_node_registers_no_content_hash_with_the_hub(): """ - H7: the daemon registered content hashes for every group, private included, - handing the hub a fingerprint of every private file. The bug was masked by a - mis-mounted route, so fixing the route without this filter would have turned a - dormant leak into a live one. + H7: the daemon once registered content hashes for every group, private + included, handing the hub a fingerprint of every private file. The swarm that + received them is gone, so no group's hashes are sent to the hub at all. """ - source = daemon_source() - assert 'visibility' in source and '_register_swarm' in source - # Both registration sites must gate on public visibility. - for marker in ['gctx.get("visibility") == "public"', - 'group_cfg.visibility == "public"']: - assert marker in source, f"swarm registration not gated: {marker}" + from pathlib import Path + + import meshbay_node.hub_client as hub_client + for source in (daemon_source(), Path(hub_client.__file__).read_text(encoding="utf-8")): + assert "/v1/swarm" not in source + assert "register_swarm" not in source def test_keystore_argon2_is_production_strength(): |