"""Synapse webhook routes — s2s delivery gateway + per-user target secrets.
Receiving half of the Synapse integration (handbook 10-platform/notifications.md,
§ «Интеграция получателя»): Synapse delivers events to POST /webhooks/synapse.
The route is registered with both spellings («.../synapse» and «.../synapse/») —
protection against an SPA catch-all swallowing one of them, mirroring the
gnexus-auth webhook contract.
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
from typing import Annotated, Literal
import structlog
from fastapi import APIRouter, Depends, HTTPException, Query, Request
from gnexus_synapse import ack, verify_webhook
from gnexus_synapse.exceptions import SynapseWebhookError
from pydantic import BaseModel
from navi.api.deps import require_user
from navi.auth import User
from navi.synapse.outbound import synapse_source_ready
from navi.synapse.settings_store import (
CompletionNotify,
PushTarget,
SynapseSettingsStore,
)
from navi.synapse.store import SynapseTargetStore
log = structlog.get_logger()
webhook_router = APIRouter(prefix="/webhooks", tags=["synapse"])
targets_router = APIRouter(prefix="/synapse-targets", tags=["synapse"])
settings_router = APIRouter(prefix="/synapse-settings", tags=["synapse"])
# Idempotency: Synapse retries until a 2xx, so the same delivery may arrive
# more than once. An insertion-ordered map of event ids, hard-capped, is
# enough — retries always concern minutes-old deliveries.
_seen_events: OrderedDict[str, bool] = OrderedDict()
_SEEN_EVENTS_MAX = 1024
def _mark_event(event_id: str) -> bool:
"""True if this event id was not seen before."""
if event_id in _seen_events:
_seen_events.move_to_end(event_id)
return False
_seen_events[event_id] = True
while len(_seen_events) > _SEEN_EVENTS_MAX:
_seen_events.popitem(last=False)
return True
def _pool():
from navi.api.deps import get_session_store
return get_session_store()._get_pool()
@webhook_router.post("/synapse")
@webhook_router.post("/synapse/")
async def synapse_delivery(request: Request) -> dict:
"""Verify and accept one Synapse s2s delivery."""
raw = await request.body()
store = SynapseTargetStore(await _pool())
targets = await store.active_secrets()
if not targets:
raise HTTPException(status_code=503, detail="No navi-side Synapse targets configured")
# Try every active per-user secret until one verifies; a mismatch tells
# nothing about which target failed — never expose that.
envelope: dict | None = None
matched: dict | None = None
for target in targets:
try:
envelope = verify_webhook(raw, dict(request.headers), target["secret"])
matched = target
break
except SynapseWebhookError:
continue
if envelope is None:
log.warning("synapse.deliver_rejected", detail="no target secret matched")
raise HTTPException(status_code=401, detail="Invalid webhook signature")
event_id = envelope.get("event_id") or ""
payload = envelope.get("payload") if isinstance(envelope.get("payload"), dict) else {}
if event_id and not _mark_event(event_id):
# Already accepted (Synapse retry) — ack again, don't act twice.
return ack(event_id)
log.info(
"synapse.delivery_accepted",
event_id=event_id,
event_type=request.headers.get("x-gnexus-event-type") or envelope.get("type"),
source=request.headers.get("x-synapse-source") or envelope.get("source"),
subject=envelope.get("subject"),
action=envelope.get("action"),
payload_user_id=payload.get("user_id"),
target_user_id=matched["user_id"],
target_ref=matched["token_ref"],
)
# 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)
class CreateSynapseTargetRequest(BaseModel):
token_ref: str
secret: str
class SynapseTargetItem(BaseModel):
id: int
token_ref: str
created_at: str
@targets_router.post("")
async def create_target(
payload: CreateSynapseTargetRequest,
user: Annotated[User, Depends(require_user)],
) -> dict:
ref = payload.token_ref.strip()
if not ref or len(ref) > 64:
raise HTTPException(status_code=422, detail="token_ref must be 1..64 characters")
if len(payload.secret) < 16:
raise HTTPException(
status_code=422, detail="Secret must be at least 16 characters"
)
store = SynapseTargetStore(await _pool())
target = await store.create(user.id, ref, payload.secret)
log.info("synapse.target_created", user_id=user.id, token_ref=ref, target_id=target.id)
return {
"id": target.id,
"token_ref": target.token_ref,
"created_at": target.created_at.isoformat(),
}
@targets_router.get("")
async def list_targets(user: Annotated[User, Depends(require_user)]) -> dict:
store = SynapseTargetStore(await _pool())
items = await store.list_for_user(user.id)
return {
"items": [
SynapseTargetItem(
id=t.id,
token_ref=t.token_ref,
created_at=t.created_at.isoformat(),
).model_dump()
for t in items
]
}
@targets_router.delete("/{target_id}", status_code=204)
async def revoke_target(
target_id: int,
user: Annotated[User, Depends(require_user)],
) -> None:
store = SynapseTargetStore(await _pool())
if not await store.revoke(target_id, user.id):
raise HTTPException(status_code=404, detail="Target not found")
log.info("synapse.target_revoked", user_id=user.id, target_id=target_id)
class UpdateSynapseSettingsRequest(BaseModel):
reactions_enabled: bool
push_target: PushTarget
completion_notify: CompletionNotify
instructions: str
@settings_router.get("")
async def get_synapse_settings(
user: Annotated[User, Depends(require_user)],
) -> dict:
store = SynapseSettingsStore(await _pool())
settings_row = await store.get(user.id)
data = settings_row.model_dump(mode="json", exclude={"user_id"})
# Not persisted — tells the UI whether Synapse-linked options are usable.
data["source_ready"] = synapse_source_ready()
return data
@settings_router.put("")
async def update_synapse_settings(
payload: UpdateSynapseSettingsRequest,
user: Annotated[User, Depends(require_user)],
) -> dict:
store = SynapseSettingsStore(await _pool())
current = await store.get(user.id)
current.reactions_enabled = payload.reactions_enabled
current.push_target = payload.push_target
current.completion_notify = payload.completion_notify
if payload.instructions != current.instructions:
await store.save_instructions(
user.id, payload.instructions, edited_by="user"
)
current.instructions = payload.instructions
await store.save(current)
data = current.model_dump(mode="json", exclude={"user_id"})
data["source_ready"] = synapse_source_ready()
return data
@settings_router.get("/versions")
async def list_instruction_versions(
user: Annotated[User, Depends(require_user)],
limit: Annotated[int, Query(ge=1, le=100)] = 20,
) -> dict:
store = SynapseSettingsStore(await _pool())
versions = await store.list_versions(user.id, limit=limit)
return {
"items": [v.model_dump(mode="json") for v in versions]
}