diff options
Diffstat (limited to 'packages/meshbay-node/tests/test_scan_settings_policy.py')
| -rw-r--r-- | packages/meshbay-node/tests/test_scan_settings_policy.py | 205 |
1 files changed, 205 insertions, 0 deletions
diff --git a/packages/meshbay-node/tests/test_scan_settings_policy.py b/packages/meshbay-node/tests/test_scan_settings_policy.py new file mode 100644 index 0000000..719b988 --- /dev/null +++ b/packages/meshbay-node/tests/test_scan_settings_policy.py @@ -0,0 +1,205 @@ +""" +The operator can tune how often the indexer's reconciliation backstop runs, +and how long it waits after a file's last write before hashing it. + +Same shape as test_apps_enabled_policy.py / test_member_upload_policy.py: +changed by a signed operator instruction, stored on the node rather than the +hub. Unlike those two, there is also a *live* DirectoryIndexer object to +update — see test_set_scan_settings_updates_the_live_indexer below. +""" + +import os +from pathlib import Path + +import pytest +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey + +from meshbay_common.adminop import OP_SET_SCAN_SETTINGS +from meshbay_common.crypto import generate_gek +from meshbay_node import ops +from meshbay_node.indexer.group_index import GroupIndex +from meshbay_node.indexer.indexer import DirectoryIndexer +from meshbay_node.roster import Roster +from meshbay_node.transport.webrtc_server import WebRTCPeerSession + +from conftest import one_root + +pytestmark = pytest.mark.asyncio + + +@pytest.fixture +def gek(): + return generate_gek() + + +@pytest.fixture +def shared_dir(tmp_path): + d = tmp_path / "shared" + d.mkdir() + (d / "video.mkv").write_bytes(os.urandom(256)) + return d + + +def _session(tmp_path: Path, user_id: str, *, operator: str | None = None) -> WebRTCPeerSession: + shared_root = tmp_path / "shared" + shared_root.mkdir(exist_ok=True) + index = GroupIndex(group_id="g" * 32, sk_node=Ed25519PrivateKey.generate()) + ctx = { + "roots": one_root(shared_root), + "index": index, + "sk_node": index.sk_node, + "node_user_id": operator, + } + session = WebRTCPeerSession.__new__(WebRTCPeerSession) + session._ctx = ctx + session._group_id = None + session._user_id = user_id + session._pk_user = "" + session.sent = [] + session._send = session.sent.append + session._audit = lambda *a, **k: None + return session + + +# ── Refused before a challenge is even issued ─────────────────────────────── + +async def test_out_of_range_reconcile_interval_is_refused(tmp_path): + session = _session(tmp_path, "op", operator="op") + session._has_admin_authority = lambda: True + issued = [] + session._issue_admin_challenge = lambda op, subject: issued.append((op, subject)) + + session._do_set_scan_settings( + {"reconcile_interval_secs": 1.0, "debounce_secs": 2.0}) + + assert not issued + assert [m for m in session.sent if m.get("type") == "error"] + + +async def test_out_of_range_debounce_is_refused(tmp_path): + session = _session(tmp_path, "op", operator="op") + session._has_admin_authority = lambda: True + issued = [] + session._issue_admin_challenge = lambda op, subject: issued.append((op, subject)) + + session._do_set_scan_settings( + {"reconcile_interval_secs": 600.0, "debounce_secs": 99999.0}) + + assert not issued + assert [m for m in session.sent if m.get("type") == "error"] + + +async def test_non_numeric_values_are_refused(tmp_path): + session = _session(tmp_path, "op", operator="op") + session._has_admin_authority = lambda: True + issued = [] + session._issue_admin_challenge = lambda op, subject: issued.append((op, subject)) + + session._do_set_scan_settings( + {"reconcile_interval_secs": "not-a-number", "debounce_secs": 2.0}) + + assert not issued + assert [m for m in session.sent if m.get("type") == "error"] + + +async def test_a_request_with_nobody_to_authorize_it_is_refused(tmp_path): + session = _session(tmp_path, "member-1", operator="the-operator") + session._has_admin_authority = lambda: False + + session._do_set_scan_settings( + {"reconcile_interval_secs": 600.0, "debounce_secs": 2.0}) + + assert [m for m in session.sent if m.get("type") == "error"] + + +# ── Who may change it ─────────────────────────────────────────────────────── + +async def test_changing_it_needs_a_signature(tmp_path): + """The request only ever produces a challenge — nothing is applied + until a signature over the transcript verifies.""" + session = _session(tmp_path, "op", operator="op") + session._has_admin_authority = lambda: True + issued = [] + session._issue_admin_challenge = lambda op, subject: issued.append((op, subject)) + + session._do_set_scan_settings( + {"reconcile_interval_secs": 600.0, "debounce_secs": 2.0}) + + assert issued == [(OP_SET_SCAN_SETTINGS, "600,2")] + + +# ── Where it is stored ────────────────────────────────────────────────────── + +async def test_the_setting_lives_on_the_node_and_survives_a_restart(tmp_path): + roster = Roster(db_path=tmp_path / "roster.db") + await roster.open() + try: + defaults = await roster.scan_settings("g1") + assert defaults == { + "reconcile_interval_secs": DirectoryIndexer.DEFAULT_RECONCILE_SECS, + "debounce_secs": DirectoryIndexer.DEFAULT_DEBOUNCE_SECS, + }, "unset must mean the indexer's own defaults, or an upgrade " \ + "changes behaviour for every existing group" + + await roster.set_scan_settings("g1", 1200.0, 5.0, set_by="op") + assert await roster.scan_settings("g1") == { + "reconcile_interval_secs": 1200.0, "debounce_secs": 5.0} + finally: + await roster.close() + + reopened = Roster(db_path=tmp_path / "roster.db") + await reopened.open() + try: + assert await reopened.scan_settings("g1") == { + "reconcile_interval_secs": 1200.0, "debounce_secs": 5.0} + assert await reopened.scan_settings("g2") == { + "reconcile_interval_secs": DirectoryIndexer.DEFAULT_RECONCILE_SECS, + "debounce_secs": DirectoryIndexer.DEFAULT_DEBOUNCE_SECS, + }, "one group's setting must not answer for another" + finally: + await reopened.close() + + +# ── Applying it to the live indexer ───────────────────────────────────────── + +async def test_set_scan_settings_updates_the_live_indexer(tmp_path, shared_dir, gek): + roster = Roster(db_path=tmp_path / "roster.db") + await roster.open() + indexer = DirectoryIndexer( + roots=one_root(shared_dir), group_id="g1", + sk_node=Ed25519PrivateKey.generate(), gek=gek) + await indexer.initial_scan() + indexer._reconcile_delay = 5000.0 # simulate a long-idle backoff + state = {"roster": roster, "indexers": {"g1": indexer}} + + try: + result = await ops.set_scan_settings(state, "g1", 1800.0, 3.0) + + assert result == {"reconcile_interval_secs": 1800.0, "debounce_secs": 3.0, + "group_id": "g1"} + assert indexer.reconcile_secs == 1800.0 + assert indexer.debounce_secs == 3.0 + assert indexer._reconcile_delay == 1800.0, \ + "the new interval must apply right away, not after whatever " \ + "backoff had already stretched the wait to" + assert await roster.scan_settings("g1") == { + "reconcile_interval_secs": 1800.0, "debounce_secs": 3.0} + finally: + await roster.close() + + +async def test_set_scan_settings_without_a_live_indexer_still_persists(tmp_path): + """A group hosted on the node but with no running indexer in this + process (e.g. a test, or a group not yet hot-loaded) must not crash — + the setting still lands in roster.db for whenever it is.""" + roster = Roster(db_path=tmp_path / "roster.db") + await roster.open() + state = {"roster": roster, "indexers": {}} + + try: + result = await ops.set_scan_settings(state, "g1", 1800.0, 3.0) + assert result["reconcile_interval_secs"] == 1800.0 + assert await roster.scan_settings("g1") == { + "reconcile_interval_secs": 1800.0, "debounce_secs": 3.0} + finally: + await roster.close() |