Newer
Older
navi-1 / tests / unit / core / test_synapse_reactions.py
"""Tests for the Synapse reaction runner (dispatcher meta-pass → session run)."""

import asyncio
import json
from datetime import datetime, timezone
from types import SimpleNamespace
from unittest.mock import AsyncMock, MagicMock

import pytest

import navi.api.deps as deps
import navi.synapse.reactions as reactions
from navi.core.registry import ProfileRegistry
from navi.profiles import ALL_PROFILES
from navi.tools._internal.base import current_user_role
from tests.conftest_factory import FakeConnection, FakePool, FakeRecord


def _registry() -> ProfileRegistry:
    reg = ProfileRegistry()
    for p in ALL_PROFILES:
        reg.register(p)
    return reg


# ── _parse_decision ─────────────────────────────────────────────────────────

def test_parse_decision_plain_and_fenced():
    plain = json.dumps({"profile_id": "developer", "task": "t", "understood": "u"})
    assert reactions._parse_decision(plain) == {
        "profile_id": "developer", "task": "t", "understood": "u"
    }

    fenced = "Sure:\n```json\n" + plain + "\n```"
    assert reactions._parse_decision(fenced) == {
        "profile_id": "developer", "task": "t", "understood": "u"
    }


def test_parse_decision_garbage_and_skip():
    assert reactions._parse_decision("no json here") is None
    assert reactions._parse_decision("") is None
    decision = reactions._parse_decision('{"skip": true, "reason": "junk delivery"}')
    assert decision is not None and decision["skip"] is True


# ── _normalise_session_key ───────────────────────────────────────────────────

def test_one_conversation_gets_one_spelling():
    """The spelling comes from a model, and a key that differs by case or by a stray
    space is a different conversation to the session lookup: the thread forks."""
    assert reactions._normalise_session_key("TG:42") == "tg:42"
    assert reactions._normalise_session_key("  tg:42  ") == "tg:42"
    assert reactions._normalise_session_key("tg :  42") == "tg : 42"
    assert reactions._normalise_session_key("tg:42\n") == "tg:42"


def test_a_key_that_is_not_a_string_is_no_key():
    assert reactions._normalise_session_key(None) is None
    assert reactions._normalise_session_key(42) is None
    assert reactions._normalise_session_key(["tg:42"]) is None
    assert reactions._normalise_session_key("   ") is None


# ── _resolve_profile ────────────────────────────────────────────────────────

def test_a_user_visible_profile_resolves():
    reg = _registry()
    assert reactions._resolve_profile(reg, "assistant").id == "assistant"


def test_an_admin_only_profile_is_refused_for_a_regular_user():
    """The role is the deciding factor, exactly as everywhere else: a run carrying
    role "user" must not land on a profile that user's own profile list withholds."""
    reg = _registry()
    assert reg.get("secretary").is_admin_only

    resolved = reactions._resolve_profile(reg, "secretary", role="user")

    assert resolved.id != "secretary"
    assert not resolved.is_admin_only


def test_an_admin_only_profile_resolves_for_an_admin():
    """Regression: with the role hardcoded to "user" the owner's own events could not
    reach a single profile they know by name — every one of them is admin-only."""
    reg = _registry()
    assert reg.get("server_admin").is_admin_only

    assert reactions._resolve_profile(reg, "server_admin", role="admin").id == "server_admin"


def test_the_fallback_depends_on_the_role_not_on_the_order_of_all():
    """`secretary` is admin-only, so a regular user was refused it and the event slid
    onto whichever profile `all()` yielded first — an accident of ordering."""
    reg = _registry()
    assert reactions._resolve_profile(reg, "astronaut", role="admin").id == "secretary"
    assert reactions._resolve_profile(reg, "astronaut", role="user").id == "assistant"


def test_a_refused_profile_is_logged(monkeypatch):
    """A substitution nobody can see is indistinguishable from a rule that does not
    work — which is how this bug survived a live test."""
    from unittest.mock import MagicMock

    reg = _registry()
    log = MagicMock()
    monkeypatch.setattr(reactions, "log", log)

    reactions._resolve_profile(reg, "server_admin", role="user")

    event, kwargs = log.warning.call_args.args[0], log.warning.call_args.kwargs
    assert event == "synapse.reaction_profile_downgraded"
    assert kwargs["requested"] == "server_admin"
    assert kwargs["resolved"] == "assistant"
    assert kwargs["role"] == "user"
    assert kwargs["admin_only"] is True


def test_a_profile_that_resolves_is_not_logged(monkeypatch):
    from unittest.mock import MagicMock

    log = MagicMock()
    monkeypatch.setattr(reactions, "log", log)

    reactions._resolve_profile(_registry(), "server_admin", role="admin")

    assert log.warning.call_count == 0


def test_hidden_profile_falls_back_to_a_user_visible_one():
    reg = _registry()
    # dispatcher exists but is hidden → must never host a session
    assert reg.get("dispatcher").is_hidden
    resolved = reactions._resolve_profile(reg, "dispatcher")
    assert not resolved.is_hidden
    assert not resolved.is_admin_only


def test_unknown_profile_falls_back_to_a_user_visible_one():
    reg = _registry()
    resolved = reactions._resolve_profile(reg, "astronaut")
    assert not resolved.is_hidden
    assert not resolved.is_admin_only


# ── the records the event leaves behind ──────────────────────────────────────

def test_event_card_is_a_system_record_carrying_the_request():
    card = reactions._event_card(
        {"event_id": "e1"}, "gntodo.task.created", "tg:42", "Напиши в чат",
    )
    assert card.role == "system"
    assert card.content == "Напиши в чат"
    assert card.is_display is True   # the reader sees it
    assert card.is_context is True   # and so does the model
    assert card.metadata["source"] == "synapse_event"
    assert card.metadata["event_type"] == "gntodo.task.created"
    assert card.metadata["session_key"] == "tg:42"


def test_context_message_frames_a_new_thread_and_not_a_continued_one():
    envelope = {"event_id": "e1", "payload": {"chat_id": 42}}

    fresh = reactions._reaction_context_message(envelope, fresh=True)
    assert "background reaction session" in fresh
    assert '"chat_id": 42' in fresh
    # the task travels in the visible card, never twice
    assert "## Task" not in fresh

    continued = reactions._reaction_context_message(envelope, fresh=False)
    assert "background reaction session" not in continued
    assert '"chat_id": 42' in continued


def test_record_event_keeps_the_latest_event_on_top_and_a_short_trail():
    class _S:
        session_metadata: dict = {}

    s = _S()
    for i in range(reactions._EVENTS_KEPT + 5):
        reactions._record_event(s, {"event_id": f"e{i}"}, "tg.message", "u", "tg:1")

    meta = s.session_metadata["synapse"]
    assert meta["event_id"] == f"e{reactions._EVENTS_KEPT + 4}"
    assert meta["session_key"] == "tg:1"
    assert len(meta["events"]) == reactions._EVENTS_KEPT
    assert meta["events"][0]["event_id"] == "e5"


# ── _dispatch ───────────────────────────────────────────────────────────────

async def test_dispatch_returns_decision(monkeypatch):
    decision = json.dumps({
        "profile_id": "developer", "task": "fix it", "understood": "ok",
        "session_key": "tg:42",
    })
    backend = MagicMock()
    backend.complete = AsyncMock(
        return_value=SimpleNamespace(content=f"```json\n{decision}\n```")
    )
    monkeypatch.setattr(
        deps, "get_profile_registry", lambda: _registry()
    )
    monkeypatch.setattr(
        deps, "get_backend_registry", lambda: SimpleNamespace(get=lambda key: backend)
    )

    profile_id, task, understood, session_key = await reactions._dispatch(
        {"event_id": "e1"}, "reaction rules", "route tg.* to assistant"
    )
    assert (profile_id, task, understood, session_key) == ("developer", "fix it", "ok", "tg:42")
    # the dispatcher pass sees both documents and the envelope — nothing else
    user_msg = backend.complete.await_args.args[0][1].content
    assert "reaction rules" in user_msg
    assert "route tg.* to assistant" in user_msg
    assert "event_id" in user_msg


async def test_dispatch_without_a_session_key_yields_none(monkeypatch):
    """A missing or non-string session_key must never become a thread key."""
    backend = MagicMock()
    backend.complete = AsyncMock(return_value=SimpleNamespace(content=json.dumps({
        "profile_id": "developer", "task": "t", "understood": "u", "session_key": None,
    })))
    monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry())
    monkeypatch.setattr(
        deps, "get_backend_registry", lambda: SimpleNamespace(get=lambda key: backend)
    )

    assert (await reactions._dispatch({}, "i", ""))[3] is None

    backend.complete = AsyncMock(return_value=SimpleNamespace(content=json.dumps({
        "profile_id": "developer", "task": "t", "understood": "u", "session_key": "   ",
    })))
    assert (await reactions._dispatch({}, "i", ""))[3] is None


async def test_the_offered_profiles_depend_on_the_role_and_carry_their_names(monkeypatch):
    """Two halves of one bug: a routing document names profiles in words, so ids alone
    give the model nothing to match on; and an id the run cannot reach must not be
    offered at all, or the rule looks accepted and is discarded one step later."""
    backend = MagicMock()
    backend.complete = AsyncMock(return_value=SimpleNamespace(content=json.dumps(
        {"profile_id": "assistant", "task": "t", "understood": "u", "session_key": None}
    )))
    monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry())
    monkeypatch.setattr(
        deps, "get_backend_registry", lambda: SimpleNamespace(get=lambda key: backend)
    )

    async def offered(role: str) -> str:
        await reactions._dispatch({}, "i", "", role=role)
        return backend.complete.await_args.args[0][0].content

    as_admin = await offered("admin")
    as_user = await offered("user")

    assert "- server_admin (Server Administrator)" in as_admin
    assert "- server_admin" not in as_user
    assert "- dispatcher" not in as_admin          # hidden for everyone
    # the description is what a phrase like "the sysadmin profile" can match against
    assert "Server administration, monitoring" in as_admin
    assert "- assistant (Assistant)" in as_user


async def test_dispatch_normalises_the_key_it_was_given(monkeypatch):
    backend = MagicMock()
    backend.complete = AsyncMock(return_value=SimpleNamespace(content=json.dumps({
        "profile_id": "assistant", "task": "t", "understood": "u", "session_key": " TG:42 ",
    })))
    monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry())
    monkeypatch.setattr(
        deps, "get_backend_registry", lambda: SimpleNamespace(get=lambda key: backend)
    )

    assert (await reactions._dispatch({}, "i", ""))[3] == "tg:42"


async def test_dispatch_backend_failure_yields_skip(monkeypatch):
    backend = MagicMock()
    backend.complete = AsyncMock(side_effect=RuntimeError("boom"))
    monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry())
    monkeypatch.setattr(
        deps, "get_backend_registry", lambda: SimpleNamespace(get=lambda key: backend)
    )
    assert await reactions._dispatch({}, "", "") == (None, "", "", None)


# ── run_reaction ────────────────────────────────────────────────────────────

class _FakeSession:
    def __init__(self, session_id, profile_id="assistant"):
        self.id = session_id
        self.profile_id = profile_id
        self.special = False
        self.name = ""
        self.messages: list = []
        self.context: list = []
        self.session_metadata: dict = {}
        self.last_active = datetime.now(timezone.utc)


def _fake_store(saved, *, found=None, role="user"):
    store = MagicMock()

    async def fake_create(profile_id, user_id=None):
        session = _FakeSession(f"sess-{len(saved) + 1}", profile_id)
        saved.append((profile_id, user_id, session))
        return session

    store.create = fake_create
    store.save = AsyncMock()
    store.set_profile = AsyncMock(return_value=True)
    # AsyncMock(return_value=None), never a bare MagicMock: a MagicMock returns a
    # truthy mock and every event would look like a continuation.
    store.find_reaction_session = AsyncMock(return_value=found)
    # The runner reads the triggering account's role through this pool before it can
    # ask `admin_only_blocked` anything; the settings store is mocked out in these
    # tests, so the role row is the first thing off the queue.
    pool = FakeConnection()
    pool.enqueue(FakeRecord(role=role))
    store._get_pool = AsyncMock(return_value=FakePool(pool))
    return store


def _fake_orchestrator(queue, *, running=lambda sid: False):
    orchestrator = MagicMock()

    session_locks: dict = {}

    def fake_lock(session_id):
        # mirror the real orchestrator: session_lock returns an asyncio.Lock
        return session_locks.setdefault(session_id, asyncio.Lock())

    orchestrator.session_lock = fake_lock
    orchestrator.is_running = running

    run = MagicMock()
    run.subscribe = lambda: queue
    run.unsubscribe = lambda q: None
    orchestrator.create_run = MagicMock(side_effect=lambda sid: run)
    orchestrator.run_agent = AsyncMock()
    return orchestrator, run


def _settings(monkeypatch, **overrides):
    """Patch SynapseSettingsStore.get to return one settings row."""
    from navi.synapse.settings_store import SynapseSettings, SynapseSettingsStore

    current = SynapseSettings(user_id="u1", reactions_enabled=True, **overrides)

    async def fake_get(self, user_id):
        return current

    monkeypatch.setattr(SynapseSettingsStore, "get", fake_get)
    return current


async def test_run_reaction_disabled_gate(monkeypatch):
    """reactions_enabled=False → no session is created, nothing is dispatched."""
    from navi.synapse.settings_store import SynapseSettings

    saved: list = []
    store = _fake_store(saved)
    monkeypatch.setattr(deps, "get_session_store", lambda: store)
    async def fake_get(self, user_id):
        return SynapseSettings(user_id=user_id)
    monkeypatch.setattr(
        "navi.synapse.settings_store.SynapseSettingsStore.get", fake_get
    )
    dispatched: list = []

    async def fake_dispatch(*args, **kwargs):
        dispatched.append(args)
        return "developer", "t", "u", None

    monkeypatch.setattr(reactions, "_dispatch", fake_dispatch)

    await reactions.run_reaction({"event_id": "e1"}, user_id="u1", event_type="todo.created")

    assert dispatched == []
    assert saved == []


async def test_run_reaction_happy_path(monkeypatch):
    """Gate passes → dispatcher picks a profile → special session is created and run.

    The account behind the event is an admin here, and the profile it is routed to is
    admin-only: that is the whole point of the role lookup. Before it, the id would have
    been refused and swapped for a user-visible profile (see _resolve_profile).
    """
    saved: list = []
    store = _fake_store(saved, role="admin")
    monkeypatch.setattr(deps, "get_session_store", lambda: store)
    _settings(monkeypatch, instructions="watch gntodo")

    monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry())
    dispatch = AsyncMock(return_value=("server_admin", "fix it", "ok", "tg:42"))
    monkeypatch.setattr(reactions, "_dispatch", dispatch)

    queue: asyncio.Queue = asyncio.Queue()
    orchestrator, _run = _fake_orchestrator(queue)
    monkeypatch.setattr(deps, "get_orchestrator", lambda: orchestrator)

    # run_agent completes right away: drain must see "done" from queue.
    ran: list = []
    roles: list = []

    async def fake_run_agent(*args, **kwargs):
        ran.append((args, kwargs))
        roles.append(current_user_role.get())
        await queue.put(("done", None))
    orchestrator.run_agent = fake_run_agent

    # No push service in deps, and no synapse source: finalise must be silent.
    monkeypatch.setattr(deps, "get_push_service", lambda: None)
    monkeypatch.setattr(reactions, "synapse_source_ready", lambda: False)

    await reactions.run_reaction(
        {"event_id": "e1", "subject": "gntodo", "action": "created"},
        user_id="u1", event_type="gntodo.task.created",
    )

    # The role the pool handed over reaches the dispatcher and the run itself.
    assert dispatch.await_args.kwargs["role"] == "admin"
    assert roles == ["admin"]

    profile_id, user_id, session = saved[0]
    assert (profile_id, user_id) == ("server_admin", "u1")
    assert session.special is True
    assert session.name == "Synapse: gntodo.task.created"
    assert session.session_metadata["synapse"]["event_id"] == "e1"
    assert session.session_metadata["synapse"]["session_key"] == "tg:42"
    assert len(session.session_metadata["synapse"]["events"]) == 1
    store.save.assert_awaited()

    # The visible record: the request as a system entry, and nothing else.
    card = session.messages[0]
    assert (card.role, card.content) == ("system", "fix it")
    assert card.metadata["source"] == "synapse_event"

    assert len(ran) == 1
    args, kwargs = ran[0]
    assert args[0] == session.id
    assert kwargs["hidden"] is True
    assert args[1].startswith("This is a background reaction session")
    orchestrator.broadcast_session_sync.assert_called_once_with(session.id)


async def test_run_reaction_continues_the_thread_within_the_ttl(monkeypatch):
    """A second event of one conversation joins the session the first one opened."""
    existing = _FakeSession("sess-existing", profile_id="assistant")
    existing.special = True
    existing.name = "Synapse: tg.message"
    existing.session_metadata["synapse"] = {"events": [{"event_id": "e0"}]}

    saved: list = []
    store = _fake_store(saved, found=existing)
    monkeypatch.setattr(deps, "get_session_store", lambda: store)
    _settings(monkeypatch, reaction_session_ttl_minutes=1440)

    monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry())
    monkeypatch.setattr(
        reactions, "_dispatch",
        AsyncMock(return_value=("assistant", "answer him", "ok", "tg:42")),
    )

    queue: asyncio.Queue = asyncio.Queue()
    orchestrator, _run = _fake_orchestrator(queue)
    monkeypatch.setattr(deps, "get_orchestrator", lambda: orchestrator)
    ran: list = []

    async def fake_run_agent(*args, **kwargs):
        ran.append((args, kwargs))
        await queue.put(("done", None))
    orchestrator.run_agent = fake_run_agent
    monkeypatch.setattr(deps, "get_push_service", lambda: None)
    monkeypatch.setattr(reactions, "synapse_source_ready", lambda: False)

    await reactions.run_reaction({"event_id": "e2"}, user_id="u1", event_type="tg.message")

    assert saved == []  # no new session
    assert store.set_profile.await_count == 0  # profile already matches
    kwargs = store.find_reaction_session.await_args.kwargs
    assert kwargs["session_key"] == "tg:42"
    assert kwargs["user_id"] == "u1"
    assert isinstance(kwargs["not_active_before"], datetime)

    meta = existing.session_metadata["synapse"]
    assert meta["event_id"] == "e2"
    assert [e["event_id"] for e in meta["events"]] == ["e0", "e2"]
    # The thread was already framed when it was opened — do not say it again.
    assert "background reaction session" not in ran[0][0][1]
    assert '"event_id": "e2"' in ran[0][0][1]
    assert existing.messages[0].content == "answer him"


async def test_run_reaction_does_not_continue_when_the_ttl_is_zero(monkeypatch):
    """ttl=0 is the pre-continuation behaviour: never look for a thread."""
    saved: list = []
    store = _fake_store(saved)
    monkeypatch.setattr(deps, "get_session_store", lambda: store)
    _settings(monkeypatch, reaction_session_ttl_minutes=0)

    monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry())
    monkeypatch.setattr(
        reactions, "_dispatch",
        AsyncMock(return_value=("assistant", "t", "u", "tg:42")),
    )

    queue: asyncio.Queue = asyncio.Queue()
    orchestrator, _run = _fake_orchestrator(queue)
    monkeypatch.setattr(deps, "get_orchestrator", lambda: orchestrator)

    async def fake_run_agent(*args, **kwargs):
        await queue.put(("done", None))
    orchestrator.run_agent = fake_run_agent
    monkeypatch.setattr(deps, "get_push_service", lambda: None)
    monkeypatch.setattr(reactions, "synapse_source_ready", lambda: False)

    await reactions.run_reaction({"event_id": "e3"}, user_id="u1", event_type="tg.message")

    assert store.find_reaction_session.await_count == 0
    assert len(saved) == 1


async def test_run_reaction_starts_a_new_session_when_the_thread_is_busy(monkeypatch):
    """A thread with a live turn must not get a second run over it — the event
    still runs, in a new session."""
    existing = _FakeSession("sess-busy", profile_id="assistant")
    existing.special = True

    saved: list = []
    store = _fake_store(saved, found=existing)
    monkeypatch.setattr(deps, "get_session_store", lambda: store)
    _settings(monkeypatch)

    monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry())
    monkeypatch.setattr(
        reactions, "_dispatch",
        AsyncMock(return_value=("assistant", "t", "u", "tg:42")),
    )

    queue: asyncio.Queue = asyncio.Queue()
    orchestrator, _run = _fake_orchestrator(queue, running=lambda sid: sid == "sess-busy")
    monkeypatch.setattr(deps, "get_orchestrator", lambda: orchestrator)

    async def fake_run_agent(*args, **kwargs):
        await queue.put(("done", None))
    orchestrator.run_agent = fake_run_agent
    monkeypatch.setattr(deps, "get_push_service", lambda: None)
    monkeypatch.setattr(reactions, "synapse_source_ready", lambda: False)

    await reactions.run_reaction({"event_id": "e4"}, user_id="u1", event_type="tg.message")

    assert len(saved) == 1  # a fresh session was created
    assert saved[0][2].id != "sess-busy"
    assert orchestrator.create_run.call_args_list[-1].args[0] == saved[0][2].id


async def test_run_reaction_switches_the_profile_of_a_continued_session(monkeypatch):
    """The dispatcher may pick another profile for a later event of the same
    thread; the session follows it instead of forking."""
    existing = _FakeSession("sess-existing", profile_id="secretary")
    existing.special = True

    saved: list = []
    store = _fake_store(saved, found=existing)
    monkeypatch.setattr(deps, "get_session_store", lambda: store)
    _settings(monkeypatch)

    monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry())
    monkeypatch.setattr(
        reactions, "_dispatch",
        AsyncMock(return_value=("assistant", "t", "u", "tg:42")),
    )

    queue: asyncio.Queue = asyncio.Queue()
    orchestrator, _run = _fake_orchestrator(queue)
    monkeypatch.setattr(deps, "get_orchestrator", lambda: orchestrator)

    async def fake_run_agent(*args, **kwargs):
        await queue.put(("done", None))
    orchestrator.run_agent = fake_run_agent
    monkeypatch.setattr(deps, "get_push_service", lambda: None)
    monkeypatch.setattr(reactions, "synapse_source_ready", lambda: False)

    await reactions.run_reaction({"event_id": "e5"}, user_id="u1", event_type="tg.message")

    store.set_profile.assert_awaited_once_with("sess-existing", "assistant")
    assert saved == []
    assert existing.profile_id == "assistant"


async def test_run_reaction_finalise_pushes_on_importance(monkeypatch):
    """completion_notify='important' + errors → app push is sent."""
    saved: list = []
    store = _fake_store(saved)
    monkeypatch.setattr(deps, "get_session_store", lambda: store)
    _settings(monkeypatch, push_target="app")

    monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry())
    monkeypatch.setattr(
        reactions, "_dispatch", AsyncMock(return_value=("secretary", "t", "u", None))
    )

    queue: asyncio.Queue = asyncio.Queue()
    orchestrator, _run = _fake_orchestrator(queue)
    monkeypatch.setattr(deps, "get_orchestrator", lambda: orchestrator)

    async def failing_run_agent(*args, **kwargs):
        await queue.put(("error", "model exploded"))
        await queue.put(("done", None))
    orchestrator.run_agent = failing_run_agent

    pushes: list = []
    class _Push:
        async def notify_custom(self, session_id, user_id, title, body):
            pushes.append((session_id, user_id, title, body))
    monkeypatch.setattr(deps, "get_push_service", lambda: _Push())
    monkeypatch.setattr(reactions, "synapse_source_ready", lambda: True)
    emitted: list = []

    async def fake_emit_low_level(subject, action, payload, priority="low", dedup_key=None):
        emitted.append({"subject": subject, "action": action, "payload": payload})
    import navi.synapse.outbound as outbound_mod
    monkeypatch.setattr(outbound_mod, "emit_low_level", fake_emit_low_level)

    await reactions.run_reaction({"event_id": "e7"}, user_id="u1", event_type="gntodo.x")

    assert len(pushes) == 1
    title = pushes[0][2]
    assert "упала" in title
    assert "model exploded" in pushes[0][3]
    assert emitted[0]["action"] == "failed"
    assert emitted[0]["payload"]["event_id"] == "e7"


# ── real pool wiring ─────────────────────────────────────────────────────────

async def test_pool_of_awaits_the_async_session_pool():
    """Regression: `_pool_of` returned `_get_pool()`'s coroutine, so every
    reaction run died on `.fetchrow` before it could read the settings."""
    import inspect

    pool = object()

    class _SessionStore:
        async def _get_pool(self):
            return pool

    assert inspect.iscoroutinefunction(reactions._pool_of)
    assert await reactions._pool_of(_SessionStore()) is pool


async def test_user_role_reads_the_account_s_role_from_navi_users():
    conn = FakeConnection()
    conn.enqueue(FakeRecord(role="admin"))
    assert await reactions._user_role(FakePool(conn), "u1") == "admin"
    assert "navi_users" in conn.calls[0][1]


async def test_user_role_fails_closed():
    """An unreadable role must never *widen* a run, so "user" is the answer and the
    failure is logged rather than guessed at."""
    assert await reactions._user_role(FakePool(FakeConnection()), "u1") == "user"

    conn = FakeConnection()
    conn.enqueue(FakeRecord(no_role_column="here"))
    assert await reactions._user_role(FakePool(conn), "u1") == "user"

    class _Pool:
        async def fetchrow(self, *args):
            raise RuntimeError("db down")

    assert await reactions._user_role(_Pool(), "u1") == "user"


async def test_spawn_run_task_carries_the_role_into_the_run():
    """The ContextVar is set only long enough to snapshot the new task's context, so
    the role is observable inside the run — which is where the tools read it."""
    seen: list[str] = []

    class _Orchestrator:
        async def run_agent(self, *args, **kwargs):
            seen.append(current_user_role.get())

    await reactions._spawn_run_task(_Orchestrator(), "s1", "hi", None, "u1", role="admin")

    assert seen == ["admin"]
    assert current_user_role.get() == "user"  # the caller's own context is untouched


async def test_run_reaction_reads_settings_through_the_real_pool(monkeypatch):
    """The disabled gate is decided by the row that comes out of the pool —
    SynapseSettingsStore and the pool are both real here."""
    from tests.conftest_factory import FakeConnection, FakePool, FakeRecord

    conn = FakeConnection()
    conn.enqueue(FakeRecord(
        user_id="u1", reactions_enabled=False, push_target="app",
        completion_notify="important", instructions="", dispatcher_instructions="",
        reaction_session_ttl_minutes=1440, updated_at=None,
    ))
    pool = FakePool(conn)

    class _SessionStore:
        async def _get_pool(self):
            return pool

    monkeypatch.setattr(deps, "get_session_store", lambda: _SessionStore())

    dispatched: list = []

    async def fake_dispatch(*args, **kwargs):
        dispatched.append(args)
        return "developer", "t", "u", None

    monkeypatch.setattr(reactions, "_dispatch", fake_dispatch)

    await reactions.run_reaction({"event_id": "e9"}, user_id="u1", event_type="todo.created")

    assert dispatched == []
    assert [c[0] for c in conn.calls] == ["fetchrow"]