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