""" 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_root_writable_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()