"""Tests for the Synapse reaction runner (dispatcher meta-pass → session run)."""

import asyncio
import json
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


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


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

def test_visible_profile_resolves():
    reg = _registry()
    assert reactions._resolve_profile(reg, "developer").id == "developer"


def test_hidden_profile_falls_back_to_secretary():
    reg = _registry()
    # dispatcher exists but is hidden → must never host a session
    assert reg.get("dispatcher").is_hidden
    assert reactions._resolve_profile(reg, "dispatcher").id == "secretary"


def test_unknown_profile_falls_back_to_secretary():
    reg = _registry()
    assert reactions._resolve_profile(reg, "astronaut").id == "secretary"


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

async def test_dispatch_returns_decision(monkeypatch):
    decision = json.dumps({"profile_id": "developer", "task": "fix it", "understood": "ok"})
    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 = await reactions._dispatch({"event_id": "e1"}, "instructions")
    assert (profile_id, task, understood) == ("developer", "fix it", "ok")
    # the dispatcher pass must only see the user's instructions + the envelope
    user_msg = backend.complete.await_args.args[0][1].content
    assert "instructions" in user_msg and "event_id" in user_msg


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, "", "")


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

class _FakeSession:
    def __init__(self, session_id):
        self.id = session_id
        self.special = False
        self.name = ""
        self.session_metadata = {}


def _fake_store(saved):
    store = MagicMock()

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

    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


def _fake_orchestrator(queue):
    import contextlib

    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

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


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"

    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."""
    from navi.synapse.settings_store import SynapseSettings, SynapseSettingsStore

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

    current = SynapseSettings(user_id="u1", reactions_enabled=True, instructions="watch gntodo")
    async def fake_get(self, user_id):
        return current
    monkeypatch.setattr(SynapseSettingsStore, "get", fake_get)

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

    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 = []

    async def fake_run_agent(*args, **kwargs):
        ran.append(args)
        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",
    )

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


async def test_run_reaction_finalise_pushes_on_importance(monkeypatch):
    """completion_notify='important' + errors → app push is sent."""
    from navi.synapse.settings_store import SynapseSettings, SynapseSettingsStore

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

    current = SynapseSettings(user_id="u1", reactions_enabled=True, push_target="app")
    async def fake_get(self, user_id):
        return current
    monkeypatch.setattr(SynapseSettingsStore, "get", fake_get)

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

    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_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"]