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