diff --git a/navi/api/routes/synapse.py b/navi/api/routes/synapse.py index 0c388c8..90b978a 100644 --- a/navi/api/routes/synapse.py +++ b/navi/api/routes/synapse.py @@ -6,9 +6,9 @@ protection against an SPA catch-all swallowing one of them, mirroring the gnexus-auth webhook contract. -Delivery handling today is verify-and-log: navi generates no outbound events -yet and no concrete inbound event types are wired to an action — the dispatch -layer arrives together with the first real event type. +Delivery handling: verify-and-ack, then a fire-and-forget reaction run +(`navi.synapse.reactions`) for the matched user — gated by their per-user +reaction settings. Outbound events are emitted from `navi.synapse.outbound`. """ from collections import OrderedDict @@ -104,8 +104,14 @@ target_user_id=matched["user_id"], target_ref=matched["token_ref"], ) - # Dispatch to the resolved user's UI arrives with concrete event types; - # today the gateway only verifies and acknowledges. + # Reaction layer: a fire-and-forget agent run under the matched user's + # reaction settings (gated there). The ack above must not wait for it. + from navi.synapse.reactions import schedule_reaction + + event_type = ( + request.headers.get("x-gnexus-event-type") or envelope.get("type") or "" + ) + schedule_reaction(envelope, user_id=matched["user_id"], event_type=event_type) return ack(event_id) diff --git a/navi/core/registry.py b/navi/core/registry.py index 9301837..681b905 100644 --- a/navi/core/registry.py +++ b/navi/core/registry.py @@ -16,6 +16,7 @@ ListProfilesTool, ManageRecallTool, MemoryTool, + NotifyTool, PeerTool, ReflectTool, PlanTool, @@ -240,6 +241,7 @@ synapse_instructions_tool = SynapseInstructionsTool( pool_provider=getattr(session_store, "_get_pool", None), ) + notify_tool = NotifyTool() builtins = [FilesystemTool(ai_helper=ai_helper), CodeExecTool(), TerminalTool(terminal_manager=terminal_manager), SshExecTool(), ImageViewTool(), @@ -252,7 +254,7 @@ mcp_status_tool, create_mcp_server_tool, test_mcp_tool_tool, schedule_recall_tool, manage_recall_tool, spawn_tool, switch_tool, list_profiles_tool, - synapse_instructions_tool, + synapse_instructions_tool, notify_tool, TasksTool()] if memory_tool: builtins.append(memory_tool) diff --git a/navi/profiles/developer/config.json b/navi/profiles/developer/config.json index e5bba67..26812e4 100644 --- a/navi/profiles/developer/config.json +++ b/navi/profiles/developer/config.json @@ -58,7 +58,8 @@ "schedule_recall", "manage_recall", "peer", - "synapse_instructions" + "synapse_instructions", + "notify" ], "mcp": { "navi-web": [ diff --git a/navi/profiles/discuss/config.json b/navi/profiles/discuss/config.json index d2c0bef..fc7530b 100644 --- a/navi/profiles/discuss/config.json +++ b/navi/profiles/discuss/config.json @@ -43,7 +43,8 @@ "filesystem", "schedule_recall", "manage_recall", - "synapse_instructions" + "synapse_instructions", + "notify" ], "mcp": { "gnexus-book": [ diff --git a/navi/profiles/modeler_3d/config.json b/navi/profiles/modeler_3d/config.json index 8c7c769..9269749 100644 --- a/navi/profiles/modeler_3d/config.json +++ b/navi/profiles/modeler_3d/config.json @@ -56,7 +56,8 @@ "content_publish", "schedule_recall", "manage_recall", - "synapse_instructions" + "synapse_instructions", + "notify" ], "mcp": { "navi-3d": [ diff --git a/navi/profiles/navi_code/config.json b/navi/profiles/navi_code/config.json index 8aa566d..d556dd2 100644 --- a/navi/profiles/navi_code/config.json +++ b/navi/profiles/navi_code/config.json @@ -59,7 +59,8 @@ "schedule_recall", "manage_recall", "peer", - "synapse_instructions" + "synapse_instructions", + "notify" ], "mcp": { "navi-web": [ diff --git a/navi/profiles/secretary/config.json b/navi/profiles/secretary/config.json index 563d803..25b9ff5 100644 --- a/navi/profiles/secretary/config.json +++ b/navi/profiles/secretary/config.json @@ -55,7 +55,8 @@ "gmail", "schedule_recall", "manage_recall", - "synapse_instructions" + "synapse_instructions", + "notify" ], "mcp": { "navi-web": [ diff --git a/navi/profiles/server_admin/config.json b/navi/profiles/server_admin/config.json index e300aef..44622d3 100644 --- a/navi/profiles/server_admin/config.json +++ b/navi/profiles/server_admin/config.json @@ -57,7 +57,8 @@ "schedule_recall", "manage_recall", "peer", - "synapse_instructions" + "synapse_instructions", + "notify" ], "mcp": { "gnexus-book": [ diff --git a/navi/profiles/tool_developer/config.json b/navi/profiles/tool_developer/config.json index 330d155..de5abc2 100644 --- a/navi/profiles/tool_developer/config.json +++ b/navi/profiles/tool_developer/config.json @@ -59,7 +59,8 @@ "mcp_status", "schedule_recall", "manage_recall", - "synapse_instructions" + "synapse_instructions", + "notify" ], "mcp": { "navi-web": [ diff --git a/navi/push/service.py b/navi/push/service.py index 2fbf5b3..4aea7e7 100644 --- a/navi/push/service.py +++ b/navi/push/service.py @@ -78,6 +78,21 @@ } asyncio.create_task(self._send_all(user_id, payload)) + async def notify_custom( + self, session_id: str, user_id: str | None, title: str, body: str + ) -> None: + """Arbitrary notification push (reaction completions, intervention + signals). No cooldown — an explicit signal must not be muffled.""" + if not self.enabled or not user_id: + return + payload = { + "title": title, + "body": body, + "url": f"/#{session_id}", + "session_id": session_id, + } + asyncio.create_task(self._send_all(user_id, payload)) + async def _send_all(self, user_id: str, payload: dict) -> None: try: subs = await self._store.list_for_user(user_id) diff --git a/navi/synapse/outbound.py b/navi/synapse/outbound.py new file mode 100644 index 0000000..3811a00 --- /dev/null +++ b/navi/synapse/outbound.py @@ -0,0 +1,49 @@ +"""Outbound navi→Synapse events — the source side of the integration. + +Navi emits low-level lifecycle facts (reaction finished / failed) through the +Synapse ingest API, where they are simply logged. Requires the source key +(`syn_*`) registered in the Synapse admin panel; without it nothing is sent. +""" + +import structlog + +from navi.config import settings + +log = structlog.get_logger() + + +def synapse_source_ready() -> bool: + return bool(settings.synapse_source_url and settings.synapse_source_api_key) + + +async def emit_low_level( + subject: str, + action: str, + payload: dict, + priority: str = "low", + dedup_key: str | None = None, +) -> None: + """Fire one ingest event; silently skip when source is not configured.""" + if not synapse_source_ready(): + return + from gnexus_synapse import AsyncSynapseClient + + client = AsyncSynapseClient( + url=settings.synapse_source_url, + api_key=settings.synapse_source_api_key, + default_source=settings.synapse_source_name, + ) + try: + event = await client.emit( + subject=subject, + action=action, + priority=priority, + payload=payload, + dedup_key=dedup_key, + ) + if event is not None: + log.info( + "synapse.outbound_emitted", subject=subject, action=action, event_id=event.id, + ) + finally: + await client.aclose() \ No newline at end of file diff --git a/navi/synapse/reactions.py b/navi/synapse/reactions.py new file mode 100644 index 0000000..2eb62c3 --- /dev/null +++ b/navi/synapse/reactions.py @@ -0,0 +1,298 @@ +"""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(_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) + + +def _pool_of(session_store): + return session_store._get_pool() \ No newline at end of file diff --git a/navi/tools/__init__.py b/navi/tools/__init__.py index 95bdca9..485f807 100644 --- a/navi/tools/__init__.py +++ b/navi/tools/__init__.py @@ -15,6 +15,7 @@ from .switch_profile import SwitchProfileTool from .list_profiles import ListProfilesTool from .synapse_instructions import SynapseInstructionsTool +from .notify import NotifyTool from .reflect import ReflectTool from .plan import PlanRunner, PlanTool from .peer import PeerTool @@ -40,6 +41,7 @@ "SwitchProfileTool", "ListProfilesTool", "SynapseInstructionsTool", + "NotifyTool", "ReflectTool", "PlanRunner", "PlanTool", diff --git a/navi/tools/notify.py b/navi/tools/notify.py new file mode 100644 index 0000000..afb578f --- /dev/null +++ b/navi/tools/notify.py @@ -0,0 +1,132 @@ +"""notify — the agent's user-facing signal channel (application + Synapse). + +Respects the user's push_target setting: app-only, app+Synapse, or Synapse-only. +The Synapse leg needs the source key; without it that leg is skipped and the +result says so explicitly, so the agent never wastes retries. +""" + +from navi.synapse.settings_store import SynapseSettingsStore +from navi.tools._internal.base import ( + Tool, + ToolContext, + ToolResult, + current_session_id, + current_user_id, +) + +_LEVEL_TO_PRIORITY = {"info": "normal", "warning": "high", "intervention": "critical"} + + +class NotifyTool(Tool): + name = "notify" + description = ( + "Push a notification to the user (application web push and/or a Synapse " + "event, per the user's notification settings). Use it to signal that " + "something needs the user's attention or that a background task finished. " + "For 'intervention' level state clearly WHAT decision you need from the " + "user and by when — the push opens this session." + ) + parameters = { + "type": "object", + "properties": { + "message": {"type": "string", "description": "Notification text, one plain sentence or two."}, + "level": { + "type": "string", + "enum": ["info", "warning", "intervention"], + "description": "info: FYI. warning: something went wrong / degraded. intervention: the user's decision is needed now.", + }, + }, + "required": ["message"], + } + + async def execute(self, params: dict, ctx: ToolContext | None = None) -> ToolResult: + user_id = ctx.user_id if ctx is not None else current_user_id.get(None) + if user_id is None: + return ToolResult( + success=False, output="", + error="No user context — cannot address a notification.", + ) + message = (params.get("message") or "").strip() + if not message: + return ToolResult(success=False, output="", error="Empty notification text.") + level = params.get("level") or "info" + if level not in _LEVEL_TO_PRIORITY: + return ToolResult( + success=False, output="", + error=f"Unknown level '{level}': info, warning, intervention.", + ) + session_id = ctx.session_id if ctx is not None else current_session_id.get() + + delivered: list[str] = [] + skipped: list[str] = [] + + settings_row = await SynapseSettingsStore(self._store_pool()).get(user_id) + push_target = settings_row.push_target + + if push_target in ("app", "app_synapse"): + push_service = self._push_service() + if push_service is not None and push_service.enabled: + try: + await push_service.notify_custom( + session_id or "", user_id, + title=_title_for(level), + body=message, + ) + delivered.append("app") + except Exception as e: + skipped.append(f"app ({e})") + else: + skipped.append("app (web push not configured)") + + if push_target in ("app_synapse", "synapse"): + from navi.synapse.outbound import emit_low_level, synapse_source_ready + + if synapse_source_ready(): + try: + from navi.synapse.outbound import emit_low_level + + await emit_low_level( + subject="navi-notification", + action=level, + payload={ + "session_id": session_id, + "user_id": user_id, + "message": message, + }, + priority=_LEVEL_TO_PRIORITY[level], + dedup_key=f"notify:{session_id}:{abs(hash(message)) % 10**10}", + ) + delivered.append("synapse") + except Exception as e: + skipped.append(f"synapse ({e})") + else: + skipped.append("synapse (source key not configured)") + + parts = [] + if delivered: + parts.append("Delivered: " + ", ".join(delivered) + ".") + if skipped: + parts.append("Skipped: " + "; ".join(skipped) + ".") + return ToolResult(success=True, output=" ".join(parts)) + + def _store_pool(self): + from navi.api.deps import get_session_store + + return get_session_store()._get_pool() + + @staticmethod + def _push_service(): + try: + from navi.api.deps import get_push_service + + return get_push_service() + except Exception: + return None + + +def _title_for(level: str) -> str: + return { + "info": "Navi: уведомление", + "warning": "Navi: предупреждение", + "intervention": "Navi: нужно ваше решение", + }[level] \ No newline at end of file diff --git a/tests/unit/api/test_synapse.py b/tests/unit/api/test_synapse.py index edcc558..9673f83 100644 --- a/tests/unit/api/test_synapse.py +++ b/tests/unit/api/test_synapse.py @@ -15,6 +15,7 @@ def synapse_client(monkeypatch): import navi.api.routes.synapse as synapse_mod import navi.auth.deps as auth_deps + import navi.synapse.reactions as reactions # Anonymous admin (auth disabled) — CRUD routes resolve without a cookie. monkeypatch.setattr( @@ -23,9 +24,20 @@ ) monkeypatch.setattr(synapse_mod, "_pool", AsyncMock(return_value=MagicMock())) + # Don't spawn real reaction runs from tests; collect the calls instead. + reactions_scheduled: list = [] + monkeypatch.setattr( + reactions, "schedule_reaction", + lambda envelope, user_id, event_type: reactions_scheduled.append( + (envelope, user_id, event_type) + ), + ) + from navi.main import app - return TestClient(app) + client = TestClient(app) + client.reactions_scheduled = reactions_scheduled + return client def _target(**overrides): @@ -124,6 +136,45 @@ assert synapse_mod._seen_events["e4"] is True +def test_synapse_delivery_schedules_reaction(synapse_client, monkeypatch): + """A fresh, verified delivery hands the envelope to the reaction layer.""" + import navi.api.routes.synapse as synapse_mod + + monkeypatch.setattr( + synapse_mod.SynapseTargetStore, "active_secrets", + AsyncMock(return_value=[_target()]), + ) + synapse_mod._seen_events.clear() + + envelope = { + "event_id": "e9", "subject": "task", "action": "created", + "payload": {"user_id": "u1"}, "source": "gntodo", + } + resp = _post(synapse_client, "/webhooks/synapse", envelope) + assert resp.status_code == 200 + + assert synapse_client.reactions_scheduled == [( + envelope, "u1", "gntodo.task.created", + )] + + +def test_synapse_replayed_delivery_skips_reaction(synapse_client, monkeypatch): + """A replayed event id is acked again but does not react twice.""" + import navi.api.routes.synapse as synapse_mod + + monkeypatch.setattr( + synapse_mod.SynapseTargetStore, "active_secrets", + AsyncMock(return_value=[_target()]), + ) + synapse_mod._seen_events.clear() + + envelope = {"event_id": "e10", "subject": "t", "action": "a"} + first = _post(synapse_client, "/webhooks/synapse", envelope) + second = _post(synapse_client, "/webhooks/synapse", envelope) + assert first.status_code == second.status_code == 200 + assert len(synapse_client.reactions_scheduled) == 1 + + def test_synapse_target_crud(synapse_client, monkeypatch): import navi.api.routes.synapse as synapse_mod import navi.synapse.store as store_mod diff --git a/tests/unit/core/test_synapse_reactions.py b/tests/unit/core/test_synapse_reactions.py new file mode 100644 index 0000000..11805e6 --- /dev/null +++ b/tests/unit/core/test_synapse_reactions.py @@ -0,0 +1,257 @@ +"""Tests for the Synapse reaction runner (dispatcher meta-pass → session run).""" + +import asyncio +import json +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock + +import pytest + +import navi.api.deps as deps +import navi.synapse.reactions as reactions +from navi.core.registry import ProfileRegistry +from navi.profiles import ALL_PROFILES + + +def _registry() -> ProfileRegistry: + reg = ProfileRegistry() + for p in ALL_PROFILES: + reg.register(p) + return reg + + +# ── _parse_decision ───────────────────────────────────────────────────────── + +def test_parse_decision_plain_and_fenced(): + plain = json.dumps({"profile_id": "developer", "task": "t", "understood": "u"}) + assert reactions._parse_decision(plain) == { + "profile_id": "developer", "task": "t", "understood": "u" + } + + fenced = "Sure:\n```json\n" + plain + "\n```" + assert reactions._parse_decision(fenced) == { + "profile_id": "developer", "task": "t", "understood": "u" + } + + +def test_parse_decision_garbage_and_skip(): + assert reactions._parse_decision("no json here") is None + assert reactions._parse_decision("") is None + decision = reactions._parse_decision('{"skip": true, "reason": "junk delivery"}') + assert decision is not None and decision["skip"] is True + + +# ── _resolve_profile ──────────────────────────────────────────────────────── + +def test_visible_profile_resolves(): + reg = _registry() + assert reactions._resolve_profile(reg, "developer").id == "developer" + + +def test_hidden_profile_falls_back_to_secretary(): + reg = _registry() + # dispatcher exists but is hidden → must never host a session + assert reg.get("dispatcher").is_hidden + assert reactions._resolve_profile(reg, "dispatcher").id == "secretary" + + +def test_unknown_profile_falls_back_to_secretary(): + reg = _registry() + assert reactions._resolve_profile(reg, "astronaut").id == "secretary" + + +# ── _dispatch ─────────────────────────────────────────────────────────────── + +async def test_dispatch_returns_decision(monkeypatch): + decision = json.dumps({"profile_id": "developer", "task": "fix it", "understood": "ok"}) + backend = MagicMock() + backend.complete = AsyncMock( + return_value=SimpleNamespace(content=f"```json\n{decision}\n```") + ) + monkeypatch.setattr( + deps, "get_profile_registry", lambda: _registry() + ) + monkeypatch.setattr( + deps, "get_backend_registry", lambda: SimpleNamespace(get=lambda key: backend) + ) + + profile_id, task, understood = await reactions._dispatch({"event_id": "e1"}, "instructions") + assert (profile_id, task, understood) == ("developer", "fix it", "ok") + # the dispatcher pass must only see the user's instructions + the envelope + user_msg = backend.complete.await_args.args[0][1].content + assert "instructions" in user_msg and "event_id" in user_msg + + +async def test_dispatch_backend_failure_yields_skip(monkeypatch): + backend = MagicMock() + backend.complete = AsyncMock(side_effect=RuntimeError("boom")) + monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry()) + monkeypatch.setattr( + deps, "get_backend_registry", lambda: SimpleNamespace(get=lambda key: backend) + ) + assert await reactions._dispatch({}, "") == (None, "", "") + + +# ── run_reaction ──────────────────────────────────────────────────────────── + +class _FakeSession: + def __init__(self, session_id): + self.id = session_id + self.special = False + self.name = "" + self.session_metadata = {} + + +def _fake_store(saved): + store = MagicMock() + + async def fake_create(profile_id, user_id=None): + session = _FakeSession("sess-1") + saved.append((profile_id, user_id, session)) + return session + + store.create = fake_create + store.save = AsyncMock() + return store + + +def _fake_orchestrator(queue): + import contextlib + + orchestrator = MagicMock() + + session_locks: dict = {} + + def fake_lock(session_id): + # mirror the real orchestrator: session_lock returns an asyncio.Lock + return session_locks.setdefault(session_id, asyncio.Lock()) + + orchestrator.session_lock = fake_lock + + run = MagicMock() + run.subscribe = lambda: queue + run.unsubscribe = lambda q: None + orchestrator.create_run = lambda sid: run + orchestrator.run_agent = AsyncMock() + return orchestrator, run + + +async def test_run_reaction_disabled_gate(monkeypatch): + """reactions_enabled=False → no session is created, nothing is dispatched.""" + from navi.synapse.settings_store import SynapseSettings + + saved: list = [] + store = _fake_store(saved) + monkeypatch.setattr(deps, "get_session_store", lambda: store) + async def fake_get(self, user_id): + return SynapseSettings(user_id=user_id) + monkeypatch.setattr( + "navi.synapse.settings_store.SynapseSettingsStore.get", fake_get + ) + dispatched: list = [] + + async def fake_dispatch(*args, **kwargs): + dispatched.append(args) + return "developer", "t", "u" + + monkeypatch.setattr(reactions, "_dispatch", fake_dispatch) + + await reactions.run_reaction({"event_id": "e1"}, user_id="u1", event_type="todo.created") + + assert dispatched == [] + assert saved == [] + + +async def test_run_reaction_happy_path(monkeypatch): + """Gate passes → dispatcher picks a profile → special session is created and run.""" + from navi.synapse.settings_store import SynapseSettings, SynapseSettingsStore + + saved: list = [] + store = _fake_store(saved) + monkeypatch.setattr(deps, "get_session_store", lambda: store) + + current = SynapseSettings(user_id="u1", reactions_enabled=True, instructions="watch gntodo") + async def fake_get(self, user_id): + return current + monkeypatch.setattr(SynapseSettingsStore, "get", fake_get) + + monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry()) + monkeypatch.setattr(reactions, "_dispatch", AsyncMock(return_value=("developer", "fix it", "ok"))) + + queue: asyncio.Queue = asyncio.Queue() + orchestrator, run = _fake_orchestrator(queue) + monkeypatch.setattr(deps, "get_orchestrator", lambda: orchestrator) + + # run_agent completes right away: drain must see "done" from queue. + ran: list = [] + + async def fake_run_agent(*args, **kwargs): + ran.append(args) + await queue.put(("done", None)) + orchestrator.run_agent = fake_run_agent + + # No push service in deps, and no synapse source: finalise must be silent. + monkeypatch.setattr(deps, "get_push_service", lambda: None) + monkeypatch.setattr(reactions, "synapse_source_ready", lambda: False) + + await reactions.run_reaction( + {"event_id": "e1", "subject": "gntodo", "action": "created"}, + user_id="u1", event_type="gntodo.task.created", + ) + + profile_id, user_id, session = saved[0] + assert (profile_id, user_id) == ("developer", "u1") + assert session.special is True + assert session.name == "Synapse: gntodo.task.created" + assert session.session_metadata["synapse"]["event_id"] == "e1" + store.save.assert_awaited() + assert len(ran) == 1 + assert ran[0][0] == session.id + + +async def test_run_reaction_finalise_pushes_on_importance(monkeypatch): + """completion_notify='important' + errors → app push is sent.""" + from navi.synapse.settings_store import SynapseSettings, SynapseSettingsStore + + saved: list = [] + store = _fake_store(saved) + monkeypatch.setattr(deps, "get_session_store", lambda: store) + + current = SynapseSettings(user_id="u1", reactions_enabled=True, push_target="app") + async def fake_get(self, user_id): + return current + monkeypatch.setattr(SynapseSettingsStore, "get", fake_get) + + monkeypatch.setattr(deps, "get_profile_registry", lambda: _registry()) + monkeypatch.setattr(reactions, "_dispatch", AsyncMock(return_value=("secretary", "t", "u"))) + + queue: asyncio.Queue = asyncio.Queue() + orchestrator, _run = _fake_orchestrator(queue) + monkeypatch.setattr(deps, "get_orchestrator", lambda: orchestrator) + + async def failing_run_agent(*args, **kwargs): + await queue.put(("error", "model exploded")) + await queue.put(("done", None)) + orchestrator.run_agent = failing_run_agent + + pushes: list = [] + class _Push: + async def notify_custom(self, session_id, user_id, title, body): + pushes.append((session_id, user_id, title, body)) + monkeypatch.setattr(deps, "get_push_service", lambda: _Push()) + monkeypatch.setattr(reactions, "synapse_source_ready", lambda: True) + emitted: list = [] + + async def fake_emit_low_level(subject, action, payload, priority="low", dedup_key=None): + emitted.append({"subject": subject, "action": action, "payload": payload}) + import navi.synapse.outbound as outbound_mod + monkeypatch.setattr(outbound_mod, "emit_low_level", fake_emit_low_level) + + await reactions.run_reaction({"event_id": "e7"}, user_id="u1", event_type="gntodo.x") + + assert len(pushes) == 1 + title = pushes[0][2] + assert "упала" in title + assert "model exploded" in pushes[0][3] + assert emitted[0]["action"] == "failed" + assert emitted[0]["payload"]["event_id"] == "e7" \ No newline at end of file diff --git a/tests/unit/tools/test_notify.py b/tests/unit/tools/test_notify.py new file mode 100644 index 0000000..70cf7cb --- /dev/null +++ b/tests/unit/tools/test_notify.py @@ -0,0 +1,139 @@ +"""Tests for the notify tool — app push / Synapse legs per push_target.""" + +import pytest +from unittest.mock import AsyncMock, MagicMock + +from navi.synapse.settings_store import SynapseSettings +from navi.tools import NotifyTool +from navi.tools._internal.base import ToolContext + + +def _tool(monkeypatch, settings_row, *, push_service=None, source_ready=True): + async def fake_get(self, user_id): + return settings_row + + monkeypatch.setattr("navi.synapse.settings_store.SynapseSettingsStore.get", fake_get) + monkeypatch.setattr(NotifyTool, "_store_pool", lambda self: None) + + monkeypatch.setattr(NotifyTool, "_push_service", staticmethod(lambda: push_service)) + + import navi.synapse.outbound as outbound + + monkeypatch.setattr(outbound, "synapse_source_ready", lambda: source_ready) + emitted = [] + + async def fake_emit(**kwargs): + emitted.append(kwargs) + return MagicMock(id="ev-1") + + async def fake_emit_low_level(subject, action, payload, priority="low", dedup_key=None): + emitted.append( + {"subject": subject, "action": action, "payload": payload, + "priority": priority, "dedup_key": dedup_key} + ) + + monkeypatch.setattr(outbound, "emit_low_level", fake_emit_low_level) + return NotifyTool(), emitted + + +async def test_no_user_context_fails(monkeypatch): + tool, _ = _tool(monkeypatch, SynapseSettings(user_id="u1")) + result = await tool.execute({"message": "hi"}) + assert not result.success + assert "user" in result.error.lower() + + +async def test_empty_message_fails(monkeypatch): + tool, _ = _tool(monkeypatch, SynapseSettings(user_id="u1")) + result = await tool.execute({"message": " "}, ctx=ToolContext(user_id="u1")) + assert not result.success + + +async def test_unknown_level_fails(monkeypatch): + tool, _ = _tool(monkeypatch, SynapseSettings(user_id="u1")) + result = await tool.execute( + {"message": "hi", "level": "panic"}, ctx=ToolContext(user_id="u1") + ) + assert not result.success + assert "panic" in result.error + + +async def test_app_target_delivers_app_only(monkeypatch): + push = MagicMock(enabled=True, notify_custom=AsyncMock()) + tool, emitted = _tool(monkeypatch, SynapseSettings(user_id="u1", push_target="app"), + push_service=push) + result = await tool.execute({"message": "done"}, ctx=ToolContext(user_id="u1", session_id="s1")) + assert result.success + assert "Delivered: app" in result.output + assert "synapse" not in result.output + push.notify_custom.assert_awaited_once() + session_arg, user_arg = push.notify_custom.await_args.args + title, body = push.notify_custom.await_args.kwargs["title"], push.notify_custom.await_args.kwargs["body"] + assert (session_arg, user_arg) == ("s1", "u1") + assert body == "done" + assert emitted == [] + + +async def test_app_synapse_delivers_both(monkeypatch): + push = MagicMock(enabled=True, notify_custom=AsyncMock()) + tool, emitted = _tool(monkeypatch, SynapseSettings(user_id="u1", push_target="app_synapse"), + push_service=push) + result = await tool.execute( + {"message": "need your call", "level": "intervention"}, + ctx=ToolContext(user_id="u1"), + ) + assert result.success + assert "Delivered: app, synapse" in result.output + assert len(emitted) == 1 + assert emitted[0]["subject"] == "navi-notification" + assert emitted[0]["action"] == "intervention" + assert emitted[0]["priority"] == "critical" + + +async def test_synapse_only_skips_app(monkeypatch): + push = MagicMock(enabled=True, notify_custom=AsyncMock()) + tool, emitted = _tool(monkeypatch, SynapseSettings(user_id="u1", push_target="synapse"), + push_service=push) + result = await tool.execute({"message": "done"}, ctx=ToolContext(user_id="u1")) + assert result.success + assert "Delivered: synapse" in result.output + # the app leg is not requested at all, not "skipped" + assert "Skipped" not in result.output + push.notify_custom.assert_not_awaited() + assert len(emitted) == 1 + + +async def test_synapse_leg_skipped_when_not_configured(monkeypatch): + push = MagicMock(enabled=True, notify_custom=AsyncMock()) + tool, _ = _tool(monkeypatch, SynapseSettings(user_id="u1", push_target="synapse"), + push_service=push, source_ready=False) + result = await tool.execute({"message": "done"}, ctx=ToolContext(user_id="u1")) + assert result.success + assert "Skipped: synapse (source key not configured)" in result.output + assert "Delivered" not in result.output + + +async def test_app_leg_skipped_when_push_service_missing(monkeypatch): + tool, emitted = _tool( + monkeypatch, SynapseSettings(user_id="u1", push_target="app"), + push_service=None, source_ready=False, + ) + result = await tool.execute({"message": "done"}, ctx=ToolContext(user_id="u1")) + assert result.success + assert "Skipped: app (web push not configured)" in result.output + assert emitted == [] + + +async def test_warning_maps_to_high_priority(monkeypatch): + tool, emitted = _tool(monkeypatch, SynapseSettings(user_id="u1", push_target="synapse")) + await tool.execute({"message": "hm", "level": "warning"}, ctx=ToolContext(user_id="u1")) + assert emitted[0]["priority"] == "high" + + +def test_registered_as_builtin(): + """The tool is wired into the default registry by name.""" + from navi.core.registry import ToolRegistry + + tools = ToolRegistry() + tools.register(NotifyTool(), builtin=True) + assert tools.get("notify") is not None \ No newline at end of file