Newer
Older
navi-1 / navi / api / routes / synapse.py
"""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]
    }