Newer
Older
navi-1 / navi / synapse / reactions.py
"""Reaction runner — inbound Synapse event → special session (fresh or continued).

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

Flow:
1) per-user settings gate (reactions_enabled);
2) dispatcher meta-pass: envelope + the user's routing and reaction documents →
   JSON {profile_id, task, understood, session_key} or {skip: true};
3) the event's thread is resolved. `session_key` is the conversation the
   dispatcher derived from the envelope; when an idle special session of the same
   user already carries that key and is younger than the user's TTL, the event
   joins it — that is what keeps a Telegram conversation in one Navi session.
   Otherwise a fresh `special=True` session is created for the target profile,
   and a continued session whose profile no longer matches is repointed;
4) the visible record of the event is a system-role card holding only the task, so
   the transcript shows the request and nothing else. The envelope — and, on a
   fresh session, the framing — travel as a hidden user message only the model
   sees;
5) 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). The
   user's reaction instructions reach the run through
   `current_reaction_instructions`, so they shape every request of the run without
   being stored as a message;
6) completion: 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
from datetime import datetime, timedelta, timezone

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"

# How many past events a continued thread remembers in its metadata. The trail is
# for the reader of the session row, not for the model, so a short tail is enough.
_EVENTS_KEPT = 20

_FRAMING = (
    "This is a background reaction session: a Synapse platform event was delivered "
    "and the reaction system formed the task above. 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."
)


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

# Serialises "find the event's thread, then claim it". Without it two events of one
# conversation arriving together both find the same idle session and both call
# create_run() on it, and the first run's subscribers are orphaned. The critical
# section is a handful of lookups — the run itself starts after it is released.
_thread_claim_lock = asyncio.Lock()


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_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, session_key = await _dispatch(
        envelope, settings_row.instructions, settings_row.dispatcher_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

    orchestrator = get_orchestrator()
    session, fresh, run = await _open_thread(
        session_store,
        orchestrator,
        envelope,
        user_id=user_id,
        event_type=event_type,
        understood=understood,
        profile=profile,
        session_key=session_key,
        ttl_minutes=settings_row.reaction_session_ttl_minutes,
    )
    # The visible record of the event — the request itself, as a system entry.
    session.messages.append(_event_card(envelope, event_type, session_key, task_text))
    # Persist the session + the opening records before the run: a crash mid-run
    # still leaves a visible, deletable trail in the special list.
    await session_store.save(session)
    # Someone may be watching this session live: the card is on disk now, so tell
    # them to reload rather than let them wait for the run to end.
    orchestrator.broadcast_session_sync(session.id)

    run.task = _spawn_run_task(
        orchestrator,
        session.id,
        _reaction_context_message(envelope, fresh=fresh),
        session_store,
        user_id,
        instructions=settings_row.instructions,
    )

    errors: list[str] = []
    final_text = ""
    queue = run.subscribe()
    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,
    )


async def _open_thread(
    session_store,
    orchestrator,
    envelope: dict,
    *,
    user_id: str,
    event_type: str,
    understood: str,
    profile,
    session_key: str | None,
    ttl_minutes: int,
):
    """Resolve the session this event belongs to and register its run.

    Returns (session, fresh, run); `fresh` tells the caller whether the framing
    still has to introduce the thread. The claim (session_lock + is_running +
    create_run) is the same one the websocket path uses, so a reaction never
    starts a second run over a live turn — a busy thread simply gets a new session
    instead, and the event is never dropped.
    """
    async with _thread_claim_lock:
        session = await _locate_thread(
            session_store, user_id=user_id, session_key=session_key, ttl_minutes=ttl_minutes,
        )
        run = await _claim_run(orchestrator, session.id) if session is not None else None
        if session is not None and run is None:
            log.info(
                "synapse.reaction_thread_busy",
                session_id=session.id,
                session_key=session_key,
                user_id=user_id,
            )
            session = None

        if session is None:
            session = await session_store.create(profile.id, user_id=user_id)
            session.special = True
            session.name = f"Synapse: {event_type}"
            session.session_metadata["synapse"] = {}
            # Nobody else can hold this id, so the claim cannot fail here.
            run = orchestrator.create_run(session.id)
            fresh = True
        else:
            fresh = False
            if session.profile_id != profile.id:
                # Switch, don't fork: the conversation keeps its history and the
                # new event is answered by the profile the dispatcher picked now.
                # set_profile() writes the one column; save() would never touch it.
                await session_store.set_profile(session.id, profile.id)
                session.profile_id = profile.id

        _record_event(session, envelope, event_type, understood, session_key)
        return session, fresh, run


async def _locate_thread(
    session_store, *, user_id: str, session_key: str | None, ttl_minutes: int,
):
    """The user's live-enough session for this conversation, or None."""
    if not session_key or ttl_minutes <= 0:
        return None
    cutoff = datetime.now(timezone.utc) - timedelta(minutes=ttl_minutes)
    return await session_store.find_reaction_session(
        user_id=user_id, session_key=session_key, not_active_before=cutoff,
    )


async def _claim_run(orchestrator, session_id: str):
    """Register a run for `session_id`, or None when a turn is already running."""
    async with orchestrator.session_lock(session_id):
        if orchestrator.is_running(session_id):
            return None
        return orchestrator.create_run(session_id)


def _record_event(
    session, envelope: dict, event_type: str, understood: str, session_key: str | None,
) -> None:
    """Update the session's synapse metadata: latest event on top, short trail below.

    The top-level fields are what `_finalise` and the UI read; `events` is the trail
    of a thread that may answer many events, trimmed so a long-lived conversation
    does not grow the row without bound.
    """
    meta = session.session_metadata.setdefault("synapse", {})
    meta["event_id"] = envelope.get("event_id")
    meta["event_type"] = event_type
    meta["understood"] = understood
    if session_key:
        meta["session_key"] = session_key
    events = meta.setdefault("events", [])
    events.append({
        "event_id": envelope.get("event_id"),
        "event_type": event_type,
        "at": datetime.now(timezone.utc).isoformat(),
    })
    del events[:-_EVENTS_KEPT]


def _event_card(
    envelope: dict, event_type: str, session_key: str | None, task_text: str,
) -> Message:
    """The transcript record of one Synapse event: the request, and nothing else.

    A system-role message, not a user one — the user did not write this. It is both
    displayed and sent to the model (`ContextBuilder._LLM_SYSTEM_SOURCES`), so the
    reader and the agent see the same request.
    """
    return Message(
        role="system",
        content=task_text,
        metadata={
            "source": "synapse_event",
            "event_id": envelope.get("event_id"),
            "event_type": event_type,
            "session_key": session_key,
        },
        is_display=True,
        is_context=True,
        created_at=datetime.now(timezone.utc),
    )


def _spawn_run_task(
    orchestrator,
    session_id: str,
    user_message: str,
    session_store,
    user_id: str,
    *,
    instructions: str = "",
):
    """Start the agent run with the target user's tool-sandbox context.

    `instructions` — the user's reaction document — is handed to the run through a
    ContextVar: ContextBuilder folds it into the system prompt of every request the
    run makes, so the model always has it while the transcript and session_messages
    never do. create_task snapshots the context, so sub-agents inherit it too.
    """
    from navi.tools._internal.base import (
        current_reaction_instructions,
        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)
    instr_token = current_reaction_instructions.set(instructions or None)
    try:
        return asyncio.create_task(
            orchestrator.run_agent(
                session_id, user_message, None, None, None, session_store,
                hidden=True,
            )
        )
    finally:
        current_user_id.reset(uid_token)
        current_user_role.reset(role_token)
        current_reaction_instructions.reset(instr_token)


async def _dispatch(
    envelope: dict, instructions: str, routing_instructions: str = "",
) -> tuple[str | None, str, str, str | None]:
    """One-shot dispatcher call.

    Returns (profile_id, task, understood, session_key), or (None, '', '', None)
    when the dispatcher skips the event or cannot be reached.
    """
    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, "", "", 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 dispatcher instructions (routing)\n"
        + (routing_instructions.strip() or "(empty — route by the event type alone)")
        + "\n\n## 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, "", "", None

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


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_context_message(envelope: dict, *, fresh: bool) -> str:
    """The hidden user message the model reads: framing + envelope.

    The framing introduces the session, so it belongs to a new thread only — an
    event that joins an existing conversation already carries that context in its
    history and does not need to be told again that it is a background session.

    The visible card carries the task, so it is not repeated here.
    """
    parts: list[str] = []
    if fresh:
        parts += [_FRAMING, ""]
    parts += [
        "## Event envelope",
        "```json",
        json.dumps(envelope, ensure_ascii=False, indent=2),
        "```",
    ]
    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()