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