"""Reaction runner — inbound Synapse event → hidden special session.

Scheduling contract: the webhook gateway (`navi.api.routes.synapse`) verifies
the delivery and acks it fast; everything heavier — the dispatcher meta-pass,
session creation, the agent run — happens here, fire-and-forget.

Flow:
1) per-user settings gate (reactions_enabled);
2) dispatcher meta-pass: envelope + user instructions →
   JSON {profile_id, task} or {skip: true};
3) a fresh `special=True` session is created for the target profile;
4) the agent runs the task headlessly (no WS subscriber — but the orchestrator
   run registry lets the user watch the stream live by just opening the
   session);
5) completion: session name generation, an app push per the user's
   completion_notify setting and a low-level "reaction.finished" event to
   Synapse (pure log entry there, best effort).
"""

from __future__ import annotations

import asyncio
import json

import structlog

from navi.config import settings
from navi.llm.base import Message

log = structlog.get_logger()

_DISPATCHER_PROFILE_ID = "dispatcher"
_FALLBACK_PROFILE_ID = "secretary"


def synapse_source_ready() -> bool:
    """True when navi can emit events to Synapse (source key configured)."""
    return bool(settings.synapse_source_url and settings.synapse_source_api_key)


# Strong references so the event loop does not garbage-collect running tasks.
_TASKS: set[asyncio.Task] = set()


def schedule_reaction(envelope: dict, user_id: str, event_type: str) -> None:
    """Fire-and-forget entry — never touches the ack path with heavy work."""
    task = asyncio.create_task(_run_reaction_guarded(envelope, user_id, event_type))
    _TASKS.add(task)
    task.add_done_callback(_TASKS.discard)


async def _run_reaction_guarded(envelope: dict, user_id: str, event_type: str) -> None:
    try:
        await run_reaction(envelope, user_id=user_id, event_type=event_type)
    except Exception:
        log.exception("synapse.reaction_failed", event_id=envelope.get("event_id"), user_id=user_id)


async def run_reaction(envelope: dict, *, user_id: str, event_type: str) -> None:
    from navi.api.deps import (
        get_backend_registry,
        get_orchestrator,
        get_profile_registry,
        get_session_store,
    )
    from navi.synapse.settings_store import SynapseSettingsStore

    session_store = get_session_store()
    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

    profile_id, task_text, understood = await _dispatch(envelope, settings_row.instructions)
    if profile_id is None:
        log.info("synapse.reaction_skipped", event_id=envelope.get("event_id"), user_id=user_id)
        return

    profiles = get_profile_registry()
    profile = _resolve_profile(profiles, profile_id)
    task_text = (task_text or "").strip()
    if not task_text:
        log.warning("synapse.reaction_empty_task", event_id=envelope.get("event_id"))
        return

    session = await session_store.create(profile.id, user_id=user_id)
    session.special = True
    session.name = f"Synapse: {event_type}"
    session.session_metadata["synapse"] = {
        "event_id": envelope.get("event_id"),
        "event_type": event_type,
        "understood": understood,
    }
    user_message = _reaction_message(envelope, event_type, task_text, settings_row.instructions)
    # Persist the session + the opening message before the run: a crash
    # mid-run still leaves a visible, deletable trail in the special list.
    await session_store.save(session)

    orchestrator = get_orchestrator()
    async with orchestrator.session_lock(session.id):
        run = orchestrator.create_run(session.id)
    queue = run.subscribe()

    run.task = _spawn_run_task(orchestrator, session.id, user_message, session_store, user_id)

    errors: list[str] = []
    final_text = ""
    try:
        while True:
            kind, payload = await queue.get()
            if kind == "event":
                full = getattr(payload, "full_content", None)
                if full:
                    final_text = full
            elif kind == "error":
                errors.append(str(payload))
            elif kind == "done":
                break
            elif kind == "stopped":
                errors.append("Run stopped by the user.")
    finally:
        run.unsubscribe(queue)

    await _finalise(
        session, session_store, settings_row, user_id,
        errors=errors, final_text=final_text, event_type=event_type,
    )


def _spawn_run_task(orchestrator, session_id: str, user_message: str, session_store, user_id: str):
    """Start the agent run with the target user's tool-sandbox context."""
    from navi.tools._internal.base import (
        current_user_id,
        current_user_info,
        current_user_role,
    )

    current_user_info.set(None)
    role_token = current_user_role.set("user")
    uid_token = current_user_id.set(user_id)
    try:
        return asyncio.create_task(
            orchestrator.run_agent(
                session_id, user_message, None, None, None, session_store,
            )
        )
    finally:
        current_user_id.reset(uid_token)
        current_user_role.reset(role_token)


async def _dispatch(envelope: dict, instructions: str) -> tuple[str | None, str, str]:
    """One-shot dispatcher call. Returns (profile_id, task, understood) or (None, '', '')."""
    from navi.api.deps import get_backend_registry, get_profile_registry

    try:
        profiles = get_profile_registry()
        dispatcher = profiles.get(_DISPATCHER_PROFILE_ID)
        backend = get_backend_registry().get(dispatcher.llm_backend)
    except Exception:
        log.exception("synapse.reaction_dispatcher_unavailable")
        return None, "", ""

    visible = [p.id for p in profiles.all() if not getattr(p, "is_hidden", False)]
    system = dispatcher.system_prompt + (
        "\n\n## Available profile ids\n" + "\n".join(f"- {pid}" for pid in visible)
    )
    user = (
        "## The user's reaction instructions\n"
        + (instructions.strip() or "(empty — pick the safest sensible route)")
        + "\n\n## Event envelope\n"
        + json.dumps(envelope, ensure_ascii=False, indent=2)
    )
    try:
        resp = await backend.complete(
            [Message(role="system", content=system), Message(role="user", content=user)],
            tools=[],
            model=dispatcher.model[0],
            temperature=dispatcher.temperature,
        )
    except Exception:
        log.exception("synapse.reaction_dispatch_call_failed")
        return None, "", ""

    parsed = _parse_decision(resp.content or "")
    if parsed is None or parsed.get("skip"):
        return None, "", ""
    return parsed.get("profile_id") or "", parsed.get("task") or "", parsed.get("understood") or ""


def _parse_decision(text: str) -> dict | None:
    text = text.strip()
    if not text:
        return None
    if text.startswith("```"):
        text = text.strip("` \n\t")
        if text.startswith("json"):
            text = text[4:].lstrip()
    try:
        start, end = text.find("{"), text.rfind("}")
        if start == -1 or end <= start:
            return None
        return json.loads(text[start:end + 1])
    except ValueError:
        return None


def _resolve_profile(profiles, profile_id: str):
    """Only user-visible profiles may host a reaction session."""
    from navi.exceptions import ProfileNotFound

    try:
        profile = profiles.get(profile_id)
        if getattr(profile, "is_hidden", False):
            raise ProfileNotFound(profile_id)
        return profile
    except Exception:
        return profiles.get(_FALLBACK_PROFILE_ID)


def _reaction_message(envelope: dict, event_type: str, task_text: str, instructions: str) -> str:
    parts = [
        "This is a background reaction session: a Synapse platform event was delivered "
        "and the reaction system formed the task below. React on it as instructed — "
        "this session is a service one, the user reads it later. Work autonomously: "
        "there is no live dialog partner here.",
        "",
        "## Task (from the reaction dispatcher)",
        task_text,
        "",
        "## Event envelope",
        "```json",
        json.dumps(envelope, ensure_ascii=False, indent=2),
        "```",
    ]
    if instructions.strip():
        parts += [
            "",
            "## The user's reaction instructions",
            instructions.strip(),
        ]
    return "\n".join(parts)


async def _finalise(
    session,
    session_store,
    settings_row,
    user_id: str,
    *,
    errors: list[str],
    final_text: str,
    event_type: str,
) -> None:
    from navi.api.deps import get_push_service
    from navi.push.service import _preview
    from navi.synapse.outbound import emit_low_level

    # App push, per the user's completion_notify preference.
    important_case = bool(errors)
    want_push = (
        settings_row.completion_notify == "always"
        or (settings_row.completion_notify == "important" and important_case)
    )
    if want_push and settings_row.push_target in ("app", "app_synapse"):
        push_service = get_push_service()
        if push_service is not None:
            try:
                if errors:
                    title = f"Реакция {event_type} упала"
                    body = "(ошибка) " + _preview(errors[-1])
                else:
                    title = f"Реакция {event_type} завершена"
                    body = _preview(final_text)
                await push_service.notify_custom(session.id, user_id, title, body)
            except Exception:
                log.exception("synapse.reaction_push_failed", session_id=session.id)

    # Low-level lifecycle event to Synapse — a pure log entry on their side,
    # independent of the push settings; best effort, silent on failure.
    if synapse_source_ready():
        try:
            await emit_low_level(
                subject="reaction",
                action="finished" if not errors else "failed",
                payload={
                    "session_id": session.id,
                    "event_id": session.session_metadata.get("synapse", {}).get("event_id"),
                    "event_type": event_type,
                    "errors": errors,
                },
            )
        except Exception:
            log.debug("synapse.reaction_lifecycle_event_skipped", session_id=session.id)


async def _pool_of(session_store):
    return await session_store._get_pool()