diff --git a/navi/mcp/keystore.py b/navi/mcp/keystore.py index addf076..71b7eb8 100644 --- a/navi/mcp/keystore.py +++ b/navi/mcp/keystore.py @@ -7,6 +7,7 @@ from __future__ import annotations +import inspect import logging import time from datetime import datetime, timezone @@ -96,7 +97,10 @@ return key del self._cache[(user_id, server_name)] try: - key = await self._store_factory().get(user_id, server_name) + store = self._store_factory() + if inspect.isawaitable(store): # the pool behind the store is async + store = await store + key = await store.get(user_id, server_name) except Exception: log.warning( "mcp user key resolve failed: user=%s server=%s — falling back to default credential", @@ -138,7 +142,8 @@ if _resolver is None: from navi.api.deps import get_session_store - _resolver = KeyResolver( - lambda: McpKeyStore(get_session_store()._get_pool()), - ) + async def _factory() -> McpKeyStore: + return McpKeyStore(await get_session_store()._get_pool()) + + _resolver = KeyResolver(_factory) return _resolver \ No newline at end of file diff --git a/navi/synapse/reactions.py b/navi/synapse/reactions.py index 2eb62c3..76b1045 100644 --- a/navi/synapse/reactions.py +++ b/navi/synapse/reactions.py @@ -66,7 +66,7 @@ from navi.synapse.settings_store import SynapseSettingsStore session_store = get_session_store() - settings_row = await SynapseSettingsStore(_pool_of(session_store)).get(user_id) + settings_row = await SynapseSettingsStore(await _pool_of(session_store)).get(user_id) if not settings_row.reactions_enabled: log.info("synapse.reaction_disabled", event_id=envelope.get("event_id"), user_id=user_id) return @@ -294,5 +294,5 @@ log.debug("synapse.reaction_lifecycle_event_skipped", session_id=session.id) -def _pool_of(session_store): - return session_store._get_pool() \ No newline at end of file +async def _pool_of(session_store): + return await session_store._get_pool() \ No newline at end of file diff --git a/navi/tools/notify.py b/navi/tools/notify.py index afb578f..aec1d0c 100644 --- a/navi/tools/notify.py +++ b/navi/tools/notify.py @@ -60,7 +60,7 @@ delivered: list[str] = [] skipped: list[str] = [] - settings_row = await SynapseSettingsStore(self._store_pool()).get(user_id) + settings_row = await SynapseSettingsStore(await self._store_pool()).get(user_id) push_target = settings_row.push_target if push_target in ("app", "app_synapse"): @@ -83,8 +83,6 @@ if synapse_source_ready(): try: - from navi.synapse.outbound import emit_low_level - await emit_low_level( subject="navi-notification", action=level, @@ -109,10 +107,11 @@ parts.append("Skipped: " + "; ".join(skipped) + ".") return ToolResult(success=True, output=" ".join(parts)) - def _store_pool(self): + @staticmethod + async def _store_pool(): from navi.api.deps import get_session_store - return get_session_store()._get_pool() + return await get_session_store()._get_pool() @staticmethod def _push_service(): diff --git a/navi/tools/synapse_instructions.py b/navi/tools/synapse_instructions.py index 7d5c9d0..b6b9949 100644 --- a/navi/tools/synapse_instructions.py +++ b/navi/tools/synapse_instructions.py @@ -6,6 +6,8 @@ is recorded with edited_by='navi' in the version history. """ +import inspect + from navi.synapse.settings_store import SynapseSettingsStore from navi.tools._internal.base import Tool, ToolContext, ToolResult, current_user_id @@ -54,7 +56,10 @@ success=False, output="", error="Reaction settings storage is not available (no pool).", ) - store = SynapseSettingsStore(self._pool_provider()) + pool = self._pool_provider() + if inspect.isawaitable(pool): # PgSessionStore._get_pool is async + pool = await pool + store = SynapseSettingsStore(pool) settings_row = await store.get(user_id) op = params.get("op", "read") diff --git a/tests/conftest_factory.py b/tests/conftest_factory.py index fd6aa7b..397a01e 100644 --- a/tests/conftest_factory.py +++ b/tests/conftest_factory.py @@ -253,6 +253,20 @@ async def close(self): pass + # Stores that query the pool directly (`self._pool.fetchrow(...)`) instead of + # acquiring a connection hit the same in-memory connection. + async def execute(self, query: str, *args): + return await self._conn.execute(query, *args) + + async def fetch(self, query: str, *args): + return await self._conn.fetch(query, *args) + + async def fetchrow(self, query: str, *args): + return await self._conn.fetchrow(query, *args) + + async def fetchval(self, query: str, *args): + return await self._conn.fetchval(query, *args) + async def __aenter__(self): return self diff --git a/tests/unit/core/test_synapse_reactions.py b/tests/unit/core/test_synapse_reactions.py index 11805e6..05d22df 100644 --- a/tests/unit/core/test_synapse_reactions.py +++ b/tests/unit/core/test_synapse_reactions.py @@ -112,6 +112,9 @@ store.create = fake_create store.save = AsyncMock() + # PgSessionStore._get_pool is async; the settings store is mocked out in + # these tests, so the pool itself is never used. + store._get_pool = AsyncMock(return_value=None) return store @@ -254,4 +257,53 @@ assert "упала" in title assert "model exploded" in pushes[0][3] assert emitted[0]["action"] == "failed" - assert emitted[0]["payload"]["event_id"] == "e7" \ No newline at end of file + 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_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="", 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" + + 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"] \ No newline at end of file diff --git a/tests/unit/mcp/test_keystore.py b/tests/unit/mcp/test_keystore.py index 3058f2c..2de2dac 100644 --- a/tests/unit/mcp/test_keystore.py +++ b/tests/unit/mcp/test_keystore.py @@ -10,7 +10,7 @@ from navi.mcp.config import McpServerConfig, McpUserKey from navi.mcp.keystore import KeyResolver, McpKeyStore - +import navi.api.deps as deps from tests.conftest_factory import FakeRecord @@ -237,4 +237,51 @@ store.get = AsyncMock(side_effect=[RuntimeError("down"), "sk-recovered"]) resolver = KeyResolver(lambda: store) assert await resolver.resolve("u1", "srv") is None - assert await resolver.resolve("u1", "srv") == "sk-recovered" \ No newline at end of file + assert await resolver.resolve("u1", "srv") == "sk-recovered" + + +class TestResolverWiring: + """The lazily built process-wide resolver, on a real (fake) asyncpg pool.""" + + async def test_default_resolver_awaits_the_session_pool(self, encryptor, monkeypatch): + """Regression: the factory used to wrap `_get_pool()`'s coroutine in + McpKeyStore, so `resolve` always raised AttributeError, swallowed it and + silently fell back to the default plaintext credential. + """ + import navi.mcp.keystore as ks + from tests.conftest_factory import FakeConnection, FakePool + + conn = FakeConnection() + conn.enqueue(encryptor.encrypt("sk-user")) + pool = FakePool(conn) + + class _SessionStore: + async def _get_pool(self): + return pool + + monkeypatch.setattr(deps, "get_session_store", lambda: _SessionStore()) + monkeypatch.setattr(ks, "_resolver", None) + + resolver = ks.get_key_resolver() + assert await resolver.resolve("u1", "srv") == "sk-user" + assert [c[0] for c in conn.calls] == ["fetchval"] + + async def test_missing_key_resolves_none(self, encryptor, monkeypatch): + """The store really answers `None` (no row) — not an exception that the + resolver turned into a fallback: the query must have been issued. + """ + import navi.mcp.keystore as ks + from tests.conftest_factory import FakeConnection, FakePool + + conn = FakeConnection() # nothing enqueued → no row + pool = FakePool(conn) + + class _SessionStore: + async def _get_pool(self): + return pool + + monkeypatch.setattr(deps, "get_session_store", lambda: _SessionStore()) + monkeypatch.setattr(ks, "_resolver", None) + + assert await ks.get_key_resolver().resolve("u1", "srv") is None + assert [c[0] for c in conn.calls] == ["fetchval"] \ No newline at end of file diff --git a/tests/unit/tools/test_notify.py b/tests/unit/tools/test_notify.py index 70cf7cb..0e8ba7c 100644 --- a/tests/unit/tools/test_notify.py +++ b/tests/unit/tools/test_notify.py @@ -13,7 +13,11 @@ return settings_row monkeypatch.setattr("navi.synapse.settings_store.SynapseSettingsStore.get", fake_get) - monkeypatch.setattr(NotifyTool, "_store_pool", lambda self: None) + + async def fake_pool(self): + return None + + monkeypatch.setattr(NotifyTool, "_store_pool", fake_pool) monkeypatch.setattr(NotifyTool, "_push_service", staticmethod(lambda: push_service)) @@ -136,4 +140,55 @@ tools = ToolRegistry() tools.register(NotifyTool(), builtin=True) - assert tools.get("notify") is not None \ No newline at end of file + assert tools.get("notify") is not None + + +# ── real pool wiring ───────────────────────────────────────────────────────── + +async def test_reads_settings_through_the_real_pool(monkeypatch): + """Regression: `PgSessionStore._get_pool` is async, and the tool used to hand + its coroutine straight to SynapseSettingsStore — so every notify call died + with "'coroutine' object has no attribute 'fetchrow'". + + Nothing is faked below the tool: the settings row comes out of the fake + asyncpg connection, and push_target is taken from it (the row says "synapse", + the default would be "app"). + """ + import inspect + + from tests.conftest_factory import FakeConnection, FakePool, FakeRecord + + assert inspect.iscoroutinefunction(NotifyTool._store_pool) + + conn = FakeConnection() + conn.enqueue(FakeRecord( + user_id="u1", reactions_enabled=False, push_target="synapse", + completion_notify="important", instructions="", updated_at=None, + )) + pool = FakePool(conn) + + class _SessionStore: + async def _get_pool(self): + return pool + + monkeypatch.setattr("navi.api.deps.get_session_store", lambda: _SessionStore()) + + import navi.synapse.outbound as outbound + + emitted = [] + + async def fake_emit_low_level(subject, action, payload, priority="low", dedup_key=None): + emitted.append({"subject": subject, "action": action, "priority": priority}) + return MagicMock(id="ev-1") + + monkeypatch.setattr(outbound, "synapse_source_ready", lambda: True) + monkeypatch.setattr(outbound, "emit_low_level", fake_emit_low_level) + + result = await NotifyTool().execute( + {"message": "ping"}, ctx=ToolContext(user_id="u1", session_id="s1") + ) + + assert result.success, result.error + assert "Delivered: synapse" in result.output + assert emitted[0]["priority"] == "normal" + assert [c[0] for c in conn.calls] == ["fetchrow"] \ No newline at end of file diff --git a/tests/unit/tools/test_synapse_instructions.py b/tests/unit/tools/test_synapse_instructions.py index d7c386f..3b7fa10 100644 --- a/tests/unit/tools/test_synapse_instructions.py +++ b/tests/unit/tools/test_synapse_instructions.py @@ -84,4 +84,66 @@ assert dispatcher.is_hidden is True listed = [p.id for p in reg.all() if not getattr(p, "is_hidden", False)] - assert "dispatcher" not in listed \ No newline at end of file + assert "dispatcher" not in listed + + +# ── real pool wiring ───────────────────────────────────────────────────────── + +def _settings_row(conn, **overrides): + from tests.conftest_factory import FakeRecord + + row = dict( + user_id="u1", reactions_enabled=True, push_target="app", + completion_notify="important", instructions="react to mentions", + updated_at=None, + ) + row.update(overrides) + conn.enqueue(FakeRecord(**row)) + return conn + + +async def test_read_goes_through_the_real_pool(monkeypatch): + """Regression: the registry wires `pool_provider` to PgSessionStore._get_pool, + which is async — the tool used to pass that coroutine into + SynapseSettingsStore and die with + "'coroutine' object has no attribute 'fetchrow'". + + SynapseSettingsStore is not mocked here: the document comes from the fake + asyncpg connection. + """ + import inspect + + from tests.conftest_factory import FakeConnection, FakePool + + from navi.tools import SynapseInstructionsTool + + conn = _settings_row(FakeConnection()) + pool = FakePool(conn) + + async def provider(): # what navi/core/registry.py passes: a bound _get_pool + return pool + + result = await SynapseInstructionsTool(pool_provider=provider).execute( + {"op": "read"}, ctx=ToolContext(session_id="s1", user_id="u1") + ) + + assert result.success, result.error + assert result.output == "react to mentions" + assert [c[0] for c in conn.calls] == ["fetchrow"] + + +async def test_sync_pool_provider_still_accepted(monkeypatch): + """A provider that hands back an already-resolved pool keeps working.""" + from tests.conftest_factory import FakeConnection, FakePool + + from navi.tools import SynapseInstructionsTool + + conn = _settings_row(FakeConnection(), instructions="watch gntodo") + pool = FakePool(conn) + + result = await SynapseInstructionsTool(pool_provider=lambda: pool).execute( + {"op": "read"}, ctx=ToolContext(session_id="s1", user_id="u1") + ) + + assert result.success, result.error + assert result.output == "watch gntodo" \ No newline at end of file