Newer
Older
navi-1 / navi / synapse / reactions.py
"""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
from navi.profiles.base import admin_only_blocked

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

    # Admin-only profiles count as hidden here: the reaction session inherits
    # the triggering user's role ("user"), so routing an ordinary user's event
    # to an admin profile would hand out a wider tool surface than the profile
    # list offers that user.
    visible = [
        p.id
        for p in profiles.all()
        if not getattr(p, "is_hidden", False) and not admin_only_blocked(p, "user")
    ]
    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.

    Admin-only profiles count as hidden: the run inherits the triggering user's
    role, so landing on one would hand that user an admin tool surface. The
    configured fallback is subject to the same rule, otherwise an unroutable
    event would reach exactly the profile the dispatcher was not offered.
    """

    def _allowed(profile) -> bool:
        return not getattr(profile, "is_hidden", False) and not admin_only_blocked(profile, "user")

    try:
        profile = profiles.get(profile_id)
    except Exception:
        profile = None
    if profile is not None and _allowed(profile):
        return profile

    fallback = profiles.get(_FALLBACK_PROFILE_ID)
    if _allowed(fallback):
        return fallback
    return next((p for p in profiles.all() if _allowed(p)), fallback)


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