diff --git a/.env.example b/.env.example index f7049e7..201d450 100644 --- a/.env.example +++ b/.env.example @@ -38,6 +38,11 @@ # channel_targets.config несёт token_ref ("navi-rei"); значение секрета здесь: # S2S_SECRET_NAVI_REI= +# --- MCP (docs/07) --- +# Bearer-токен MCP-эндпоинта для ИИ-агента. Пусто/не задано — /mcp выключен. +# Генерация: openssl rand -hex 32. Токен — как API-ключ: в .env + gnexus-creds. +MCP_TOKEN= + # --- Misc --- # Порт API на хосте. SYNAPSE_PORT=8013 \ No newline at end of file diff --git a/CLAUDE.md b/CLAUDE.md index 38fae22..04207c7 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -64,6 +64,7 @@ - **Доступ к админке** — роли gnexus-auth (SSO): настраивать Synapse — только `admin` и выше (админ-API под 403). Пользователь роли `user` **входит в приложение** с ограниченным личным разделом «Мои события» (`/api/v1/me`; в админку не пускается ни сервером, ни роутером SPA). Своих паролей Synapse не хранит. - **Многопользовательность (2026-10-03)**: «о ком событие» — конвенция `payload.user_id` (= `sub` gnexus-auth; контракт v1 цел, получателей в конверте нет). Что видит пользователь — решают правила: целевые каналы `user` (личный лог) и `push` (web-push) без цели, адресат из `payload.user_id`; правило без них — событие «только для админа». Привязка «user_id ↔ адрес канала» для push — самопользовательская: браузер подписывается в разделе «Настройки» (docs/06). Остался открытый вопрос №3 по tg chat_id/email opt-in. В OAuth-колбэке токен выдаётся любому аутентифицированному (revoke не-админу больше не нужен). - **Настройки + PWA (2026-10-03)**: редактируемые параметры — реестр `app/settings_registry.py`, дефолт из `.env`, оверрайды в таблице `app_settings` через админку (применяются без перезапуска; docs/06). Секреты (VAPID private, тг-токен, SMTP-пароль) — в БД write-only: записываются через UI, никогда не возвращаются API, читает только воркер; исключение из правила «секреты только в .env», согласовано владельцем. Канал `push` реализован: pywebpush (DER base64url ключи, 404/410 → подписка удаляется), таблица `push_subscriptions`, PWA (manifest + sw.js) в дистрибутиве SPA. +- **MCP + архив вместо удаления (2026-10-03)**: embedded MCP-сервер (FastMCP, streamable-http) на `/mcp` того же контейнера api; доступ — статический Bearer `MCP_TOKEN` из `.env` (пусто → /mcp не монтируется). Тулы (31, docs/07) переиспользуют admin-route функции напрямую (та же валидация, автор — синтетический superadmin): источники + **выдача API-ключей клиентам** (plaintext один раз), типы, цели, правила, поток, настройки; `send_test_event` — контрольный прогон правила. Удаление справочников (sources/types/targets/rules) = архив: `deleted_at`, restore-эндпоинты, списки с `include_archived`, create поверх архива → 409 с подсказкой; архивный источник/тип не принимают события (401/422), архивная цель — доставки `skipped` с аудитом; push-подписки — исключение, удаляются физически. docs/07-mcp.md. ## Инфраструктура diff --git a/README.md b/README.md index 95a868f..753a4c2 100644 --- a/README.md +++ b/README.md @@ -77,6 +77,18 @@ Ретраи (#30): провал доставки → пауза 30 с → 2 м → 10 м → 30 м, 5 попыток → `failed`; скан due-доставок — beat воркера раз в минуту (`synapse.retry_due`). Шумодав — `--throttle N` у правила (не чаще одной доставки в цель за N с, critical проходит всегда); все подошедшие правила применяются (режим `all`, `ROUTING_MATCH_MODE=first` — только первое). Подробности — docs/05, раздел «Типы и правила». +## MCP: управление ИИ-агентом + +Synapse встраивает MCP-сервер (`/mcp`, streamable-http) — ИИ-агент (Claude Code) управляет всем: регистрирует источники и **выдаёт API-ключи клиентам** (plaintext токен показывается один раз), ведёт типы/цели/правила, смотрит поток событий и доставок, меняет настройки. Включается токеном: + +```bash +# .env: MCP_TOKEN=$(openssl rand -hex 32) и перезапуск api +claude mcp add --transport http synapse http://localhost:8013/mcp \ + --header "Authorization: Bearer " +``` + +Удаление = архив: DELETE ставит метку (`deleted_at`), `restore` возвращает; архивный источник перестаёт принимать события, цель — «сюда больше не ходим» (доставки skipped с аудитом). Каталог тулов и семантика — docs/07. + ## Вход в админку SSO-поток целиком на сервере (gnexus-auth-client-py, OAuth2 Authorization Code + PKCE): diff --git a/alembic/versions/20261003_0007_soft_delete.py b/alembic/versions/20261003_0007_soft_delete.py new file mode 100644 index 0000000..4b98026 --- /dev/null +++ b/alembic/versions/20261003_0007_soft_delete.py @@ -0,0 +1,51 @@ +"""Архив вместо удаления: deleted_at на источниках/типах/целях/правилах + +Revision ID: b7d3e8f04a92 +Revises: e2f6a9b3c4d7 +Create Date: 2026-10-03 +""" + +from alembic import op +import sqlalchemy as sa + + +revision = "b7d3e8f04a92" +down_revision = "e2f6a9b3c4d7" +branch_labels = None +depends_on = None + +TABLES = ("sources", "notification_types", "channel_targets", "routing_rules") + +# Частичные индексы: живые строки (архив исключён) — уникальность/подчистка +# имен живых записей делает код, индекс ускоряет списки и дабл-чеки. +# Для notification_types — составной (source_id, subject, action) — дабл-чек +# create_type ищет по тройке; subject/action в индексе. +PARTIAL = { + "sources": ("ix_sources_live", ["name"]), + "notification_types": ("ix_ntypes_live", ["source_id", "subject", "action"]), + "channel_targets": ("ix_targets_live", ["channel", "name"]), + "routing_rules": ("ix_rules_live", ["enabled"]), +} + + +def upgrade() -> None: + for table in TABLES: + op.add_column( + table, + sa.Column("deleted_at", sa.DateTime(timezone=True), nullable=True), + ) + for table, (iname, cols) in PARTIAL.items(): + op.create_index( + iname, + table, + cols, + unique=False, + postgresql_where=sa.text("deleted_at IS NULL"), + ) + + +def downgrade() -> None: + for table, (iname, _cols) in PARTIAL.items(): + op.drop_index(iname, table) + for table in TABLES: + op.drop_column(table, "deleted_at") \ No newline at end of file diff --git a/app/api/admin_routes.py b/app/api/admin_routes.py index 590c2c9..bb9bddc 100644 --- a/app/api/admin_routes.py +++ b/app/api/admin_routes.py @@ -62,22 +62,29 @@ @router.get("/sources", response_model=list[SourceOut]) def list_sources( - user: AuthenticatedUser = Depends(require_admin), db: Session = Depends(get_db) + include_archived: bool = False, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), ) -> list[SourceOut]: - sources = db.execute(select(Source).order_by(Source.name)).scalars().all() + q = select(Source).order_by(Source.name) + if not include_archived: + q = q.where(Source.deleted_at.is_(None)) + sources = db.execute(q).scalars().all() key_counts = dict( db.execute(select(ApiKey.source_id, func.count(ApiKey.id)).group_by(ApiKey.source_id)).all() ) type_counts = dict( db.execute( select(NotificationType.source_id, func.count(NotificationType.id)) + .where(NotificationType.deleted_at.is_(None)) .group_by(NotificationType.source_id) ).all() ) return [ SourceOut( id=s.id, name=s.name, label=s.label, description=s.description, - created_at=s.created_at, keys_count=key_counts.get(s.id, 0), + created_at=s.created_at, deleted_at=s.deleted_at, + keys_count=key_counts.get(s.id, 0), types_count=type_counts.get(s.id, 0), ) for s in sources @@ -90,13 +97,65 @@ user: AuthenticatedUser = Depends(require_admin), db: Session = Depends(get_db), ) -> SourceOut: - if db.execute(select(Source.id).where(Source.name == payload.name)).scalar_one_or_none(): + dup = db.execute(select(Source).where(Source.name == payload.name)).scalar_one_or_none() + if dup is not None: + if dup.deleted_at is not None: + raise _conflict( + f"source '{payload.name}' есть в архиве (id={dup.id}) — " + "восстанови (restore) или выбери другое имя" + ) raise _conflict(f"source '{payload.name}' уже существует") source = Source(**payload.model_dump()) db.add(source) db.commit() db.refresh(source) - return SourceOut(id=source.id, created_at=source.created_at, **payload.model_dump()) + return _source_out_full(db, source) + + +@router.post("/sources/{source_id}/restore", response_model=SourceOut) +def restore_source( + source_id: int, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> SourceOut: + source = db.get(Source, source_id) + if source is None or source.deleted_at is None: + raise HTTPException(status_code=404, detail="В архиве нет источника с таким id") + _raise_restore_conflict( + db, + f"source '{source.name}'", + select(Source.id).where( + Source.name == source.name, + Source.deleted_at.is_(None), + Source.id != source.id, + ), + ) + source.deleted_at = None + db.commit() + db.refresh(source) + return _source_out_full(db, source) + + +def _raise_restore_conflict(db: Session, what: str, live_query) -> None: + """Имя восстанавливаемой записи занято живой записью → 409 с подсказкой.""" + live_id = db.execute(live_query).scalar_one_or_none() + if live_id is not None: + raise _conflict(f"{what}: имя занято живой записью id={live_id} — restore невозможен") + + +def _source_out_full(db: Session, source: Source) -> SourceOut: + """SourceOut одной строки (create/restore): счётчики врозь, не группой.""" + keys_count = db.execute( + select(func.count(ApiKey.id)).where(ApiKey.source_id == source.id) + ).scalar_one() + types_count = db.execute( + select(func.count(NotificationType.id)).where(NotificationType.source_id == source.id) + ).scalar_one() + return SourceOut( + id=source.id, name=source.name, label=source.label, description=source.description, + created_at=source.created_at, deleted_at=source.deleted_at, + keys_count=keys_count, types_count=types_count, + ) @router.delete("/sources/{source_id}", status_code=204) @@ -105,15 +164,11 @@ user: AuthenticatedUser = Depends(require_admin), db: Session = Depends(get_db), ) -> None: + """Архив, не удаление: события и ключи остаются в БД, restore снимает метку.""" source = db.get(Source, source_id) - if source is None: + if source is None or source.deleted_at is not None: raise HTTPException(status_code=404, detail="Источник не найден") - events = db.execute(select(func.count(Event.id)).where(Event.source_id == source_id)).scalar_one() - if events: - raise _conflict( - f"Источник имеет {events} событий — удалить нельзя (история); отзови ключи и выключи типы" - ) - db.delete(source) + source.deleted_at = datetime.now(UTC) db.commit() @@ -149,6 +204,8 @@ source = db.get(Source, source_id) if source is None: raise HTTPException(status_code=404, detail="Источник не найден") + if source.deleted_at is not None: + raise HTTPException(status_code=404, detail="Источник в архиве — сначала restore") token = generate_token() key = ApiKey( source_id=source_id, @@ -190,20 +247,22 @@ @router.get("/types", response_model=list[TypeOut]) def list_types( - user: AuthenticatedUser = Depends(require_admin), db: Session = Depends(get_db) + include_archived: bool = False, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), ) -> list[TypeOut]: - rows = ( - db.execute( - select(NotificationType, Source.name) - .join(Source, NotificationType.source_id == Source.id) - .order_by(Source.name, NotificationType.subject, NotificationType.action) - ).all() - ) + q = select(NotificationType, Source.name).join( + Source, NotificationType.source_id == Source.id + ).order_by(Source.name, NotificationType.subject, NotificationType.action) + if not include_archived: + # архив скрыт двойным фильтром: архивный тип или тип архивного источника + q = q.where(NotificationType.deleted_at.is_(None), Source.deleted_at.is_(None)) + rows = db.execute(q).all() return [ TypeOut( id=nt.id, source_id=nt.source_id, source_name=source_name, subject=nt.subject, action=nt.action, payload_schema=nt.payload_schema, description=nt.description, - created_at=nt.created_at, + created_at=nt.created_at, deleted_at=nt.deleted_at, ) for nt, source_name in rows ] @@ -218,20 +277,31 @@ source = db.get(Source, payload.source_id) if source is None: raise HTTPException(status_code=422, detail="Источник не найден") + if source.deleted_at is not None: + raise HTTPException(status_code=422, detail="Источник в архиве — сначала restore") dup = db.execute( - select(NotificationType.id).where( + select(NotificationType) + .where( NotificationType.source_id == payload.source_id, NotificationType.subject == payload.subject, NotificationType.action == payload.action, ) ).scalar_one_or_none() - if dup: + if dup is not None: + if dup.deleted_at is not None: + raise _conflict( + f"Тип ({source.name}, {payload.subject}, {payload.action}) есть в архиве " + f"(id={dup.id}) — восстанови (restore) или зарегистрируй заново с другим action" + ) raise _conflict(f"Тип ({source.name}, {payload.subject}, {payload.action}) уже есть") nt = NotificationType(**payload.model_dump()) db.add(nt) db.commit() db.refresh(nt) - return TypeOut(id=nt.id, source_name=source.name, created_at=nt.created_at, **payload.model_dump()) + return TypeOut( + id=nt.id, source_name=source.name, created_at=nt.created_at, deleted_at=None, + **payload.model_dump(), + ) @router.delete("/types/{type_id}", status_code=204) @@ -240,23 +310,63 @@ user: AuthenticatedUser = Depends(require_admin), db: Session = Depends(get_db), ) -> None: - if db.get(NotificationType, type_id) is None: + nt = db.get(NotificationType, type_id) + if nt is None or nt.deleted_at is not None: raise HTTPException(status_code=404, detail="Тип не найден") - db.execute(delete(NotificationType).where(NotificationType.id == type_id)) + nt.deleted_at = datetime.now(UTC) db.commit() +@router.post("/types/{type_id}/restore", response_model=TypeOut) +def restore_type( + type_id: int, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> TypeOut: + nt = db.get(NotificationType, type_id) + if nt is None or nt.deleted_at is None: + raise HTTPException(status_code=404, detail="В архиве нет типа с таким id") + source = db.get(Source, nt.source_id) + if source is None or source.deleted_at is not None: + raise HTTPException(status_code=409, detail="Источник типа в архиве — restore невозможен") + _raise_restore_conflict( + db, + f"тип ({source.name}, {nt.subject}, {nt.action})", + select(NotificationType.id).where( + NotificationType.source_id == nt.source_id, + NotificationType.subject == nt.subject, + NotificationType.action == nt.action, + NotificationType.deleted_at.is_(None), + NotificationType.id != nt.id, + ), + ) + nt.deleted_at = None + db.commit() + db.refresh(nt) + return TypeOut( + id=nt.id, source_id=nt.source_id, source_name=source.name, subject=nt.subject, + action=nt.action, payload_schema=nt.payload_schema, description=nt.description, + created_at=nt.created_at, deleted_at=None, + ) + + # --- targets --- @router.get("/targets", response_model=list[TargetOut]) def list_targets( - user: AuthenticatedUser = Depends(require_admin), db: Session = Depends(get_db) + include_archived: bool = False, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), ) -> list[TargetOut]: - targets = db.execute(select(ChannelTarget).order_by(ChannelTarget.channel, ChannelTarget.name)).scalars().all() + q = select(ChannelTarget).order_by(ChannelTarget.channel, ChannelTarget.name) + if not include_archived: + q = q.where(ChannelTarget.deleted_at.is_(None)) + targets = db.execute(q).scalars().all() return [ TargetOut( id=t.id, channel=t.channel, name=t.name, description=t.description, config=t.config, enabled=t.enabled, created_at=t.created_at, + deleted_at=t.deleted_at, ) for t in targets ] @@ -269,17 +379,24 @@ db: Session = Depends(get_db), ) -> TargetOut: dup = db.execute( - select(ChannelTarget.id).where( - ChannelTarget.channel == payload.channel, ChannelTarget.name == payload.name - ) + select(ChannelTarget) + .where(ChannelTarget.channel == payload.channel, ChannelTarget.name == payload.name) ).scalar_one_or_none() - if dup: + if dup is not None: + if dup.deleted_at is not None: + raise _conflict( + f"Цель {payload.channel}/{payload.name} есть в архиве (id={dup.id}) — " + "восстанови (restore) или выбери другое имя" + ) raise _conflict(f"Цель {payload.channel}/{payload.name} уже есть") target = ChannelTarget(**payload.model_dump()) db.add(target) db.commit() db.refresh(target) - return TargetOut(id=target.id, enabled=target.enabled, created_at=target.created_at, **payload.model_dump()) + return TargetOut( + id=target.id, enabled=target.enabled, created_at=target.created_at, + deleted_at=None, **payload.model_dump(), + ) @router.patch("/targets/{target_id}", response_model=TargetOut) @@ -290,16 +407,23 @@ db: Session = Depends(get_db), ) -> TargetOut: target = db.get(ChannelTarget, target_id) - if target is None: + if target is None or target.deleted_at is not None: raise HTTPException(status_code=404, detail="Цель не найдена") data = payload.model_dump(exclude_unset=True) if "name" in data and (data["name"] != target.name): dup = db.execute( - select(ChannelTarget.id).where( - ChannelTarget.channel == target.channel, ChannelTarget.name == data["name"] + select(ChannelTarget) + .where( + ChannelTarget.channel == target.channel, + ChannelTarget.name == data["name"], + ChannelTarget.id != target.id, ) ).scalar_one_or_none() - if dup: + if dup is not None: + if dup.deleted_at is not None: + raise _conflict( + f"Цель {target.channel}/{data['name']} есть в архиве (id={dup.id})" + ) raise _conflict(f"Цель {target.channel}/{data['name']} уже есть") for field, value in data.items(): setattr(target, field, value) @@ -308,6 +432,7 @@ return TargetOut( id=target.id, channel=target.channel, name=target.name, description=target.description, config=target.config, enabled=target.enabled, created_at=target.created_at, + deleted_at=target.deleted_at, ) @@ -317,13 +442,43 @@ user: AuthenticatedUser = Depends(require_admin), db: Session = Depends(get_db), ) -> None: - if db.get(ChannelTarget, target_id) is None: + """Архив: ссылки из правил НЕ обнуляются (restore возвращает цель в строй).""" + target = db.get(ChannelTarget, target_id) + if target is None or target.deleted_at is not None: raise HTTPException(status_code=404, detail="Цель не найдена") - db.execute(delete(ChannelTarget).where(ChannelTarget.id == target_id)) - # ссылки из routing_rule_actions/deliveries обнуляет FK SET NULL + target.deleted_at = datetime.now(UTC) db.commit() +@router.post("/targets/{target_id}/restore", response_model=TargetOut) +def restore_target( + target_id: int, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> TargetOut: + target = db.get(ChannelTarget, target_id) + if target is None or target.deleted_at is None: + raise HTTPException(status_code=404, detail="В архиве нет цели с таким id") + _raise_restore_conflict( + db, + f"цель {target.channel}/{target.name}", + select(ChannelTarget.id).where( + ChannelTarget.channel == target.channel, + ChannelTarget.name == target.name, + ChannelTarget.deleted_at.is_(None), + ChannelTarget.id != target.id, + ), + ) + target.deleted_at = None + db.commit() + db.refresh(target) + return TargetOut( + id=target.id, channel=target.channel, name=target.name, description=target.description, + config=target.config, enabled=target.enabled, created_at=target.created_at, + deleted_at=None, + ) + + # --- rules --- def _rule_to_out(db: Session, rule: RoutingRule) -> RuleOut: @@ -338,6 +493,7 @@ id=rule.id, name=rule.name, enabled=rule.enabled, weight=rule.weight, conditions=rule.conditions, template=rule.template, throttle_seconds=rule.throttle_seconds, created_at=rule.created_at, + deleted_at=rule.deleted_at, actions=[ RuleActionOut( id=act.id, channel=act.channel, target_id=act.target_id, @@ -350,9 +506,14 @@ @router.get("/rules", response_model=list[RuleOut]) def list_rules( - user: AuthenticatedUser = Depends(require_admin), db: Session = Depends(get_db) + include_archived: bool = False, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), ) -> list[RuleOut]: - rules = db.execute(select(RoutingRule).order_by(RoutingRule.weight, RoutingRule.id)).scalars().all() + q = select(RoutingRule).order_by(RoutingRule.weight, RoutingRule.id) + if not include_archived: + q = q.where(RoutingRule.deleted_at.is_(None)) + rules = db.execute(q).scalars().all() return [_rule_to_out(db, r) for r in rules] @@ -363,9 +524,7 @@ db: Session = Depends(get_db), ) -> RuleOut: for act in payload.actions: - target = db.get(ChannelTarget, act.target_id) if act.target_id else "ok" - if not target: - raise HTTPException(status_code=422, detail=f"Цель id={act.target_id} не найдена") + _raise_action_target_error(db, act.channel, act.target_id) if act.channel in ("user", "push") and act.target_id: raise HTTPException( status_code=422, @@ -383,6 +542,21 @@ return _rule_to_out(db, rule) +def _raise_action_target_error(db: Session, channel: str, target_id: int | None) -> None: + """Цель действия: живая — ok; в архиве — 422 «цель в архиве»; нет — 422.""" + if not target_id: + return + target = db.get(ChannelTarget, target_id) + if target is None: + raise HTTPException(status_code=422, detail=f"Цель id={target_id} не найдена") + if target.deleted_at is not None: + raise HTTPException( + status_code=422, + detail=f"Цель id={target_id} ({target.channel}/{target.name}) в архиве — restore, " + "другая цель или канал без цели", + ) + + @router.patch("/rules/{rule_id}", response_model=RuleOut) def patch_rule( rule_id: int, @@ -391,15 +565,13 @@ db: Session = Depends(get_db), ) -> RuleOut: rule = db.get(RoutingRule, rule_id) - if rule is None: + if rule is None or rule.deleted_at is not None: raise HTTPException(status_code=404, detail="Правило не найдено") data = payload.model_dump(exclude_unset=True) if "actions" in data: actions = data.pop("actions") for act in actions: - target = db.get(ChannelTarget, act.get("target_id")) if act.get("target_id") else "ok" - if not target: - raise HTTPException(status_code=422, detail=f"Цель id={act['target_id']} не найдена") + _raise_action_target_error(db, act["channel"], act.get("target_id")) if act.get("channel") in ("user", "push") and act.get("target_id"): raise HTTPException( status_code=422, @@ -421,13 +593,29 @@ user: AuthenticatedUser = Depends(require_admin), db: Session = Depends(get_db), ) -> None: - if db.get(RoutingRule, rule_id) is None: + """Архив: действия правила не удаляются — restore вернёт их целиком.""" + rule = db.get(RoutingRule, rule_id) + if rule is None or rule.deleted_at is not None: raise HTTPException(status_code=404, detail="Правило не найдено") - db.execute(delete(RoutingRuleAction).where(RoutingRuleAction.rule_id == rule_id)) - db.execute(delete(RoutingRule).where(RoutingRule.id == rule_id)) + rule.deleted_at = datetime.now(UTC) db.commit() +@router.post("/rules/{rule_id}/restore", response_model=RuleOut) +def restore_rule( + rule_id: int, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> RuleOut: + rule = db.get(RoutingRule, rule_id) + if rule is None or rule.deleted_at is None: + raise HTTPException(status_code=404, detail="В архиве нет правила с таким id") + rule.deleted_at = None + db.commit() + db.refresh(rule) + return _rule_to_out(db, rule) + + # --- events / deliveries (чтение) --- def _target_names(db: Session) -> dict[int, str]: @@ -551,8 +739,10 @@ kind=defn.kind, help=defn.help, options=list(defn.options), - # secret — write-only: факт наличия в set, значение не возвращается - value=None if (is_secret(defn) or raw is None or default is None) else raw, + # secret — write-only: факт наличия в set, значение не возвращается. + # Оверрайд показывается для любого ключа реестра (у push/smtp + # .env-дефолта нет — `default None` не значит «не возвращать»). + value=None if (is_secret(defn) or raw is None) else raw, set=raw is not None or default is not None, ) ) diff --git a/app/api/admin_schemas.py b/app/api/admin_schemas.py index 021d4b1..2698250 100644 --- a/app/api/admin_schemas.py +++ b/app/api/admin_schemas.py @@ -22,6 +22,8 @@ label: str | None description: str | None created_at: datetime + # архив: дата метки; None у живых записей (UI-списки архив не показывают) + deleted_at: datetime | None = None keys_count: int = 0 types_count: int = 0 @@ -59,6 +61,7 @@ id: int source_name: str created_at: datetime + deleted_at: datetime | None = None # --- targets --- @@ -74,6 +77,7 @@ id: int enabled: bool created_at: datetime + deleted_at: datetime | None = None class TargetPatch(BaseModel): @@ -112,6 +116,7 @@ enabled: bool created_at: datetime actions: list[RuleActionOut] + deleted_at: datetime | None = None class RulePatch(BaseModel): diff --git a/app/api/events_routes.py b/app/api/events_routes.py index 0711e7e..ac8a1bd 100644 --- a/app/api/events_routes.py +++ b/app/api/events_routes.py @@ -22,7 +22,7 @@ ) from app.auth.apikeys import resolve_event_source from app.database import get_db -from app.models import ApiKey, ChannelTarget, Delivery, Event, NotificationType +from app.models import ChannelTarget, Delivery, Event, NotificationType, Source from app.settings_store import get_setting from app.worker.celery_app import celery_app @@ -37,25 +37,28 @@ self.event_status = event_status -def accept_event(db: Session, api_key: ApiKey, source_name: str, envelope: EventEnvelope) -> Event: - """Общая логика приёма (одиночный и batch): проверки + вставка. +def accept_event(db: Session, source: Source, envelope: EventEnvelope) -> Event: + """Общая логика приёма (одиночный и batch, и send_test_event из MCP): + проверки + вставка. Принимает источник (не ключ) — MCP вызывает его + для тест-события без выдачи API-ключа себе. Диспетчеризацию в Celery выполняет вызывающий код после коммита. """ - if envelope.source != source_name: + if envelope.source != source.name: raise HTTPException( status_code=status.HTTP_403_FORBIDDEN, detail=( f"Поле source='{envelope.source}' не совпадает с источником " - f"ключа '{source_name}'" + f"ключа '{source.name}'" ), ) known_type = db.execute( select(NotificationType.id).where( - NotificationType.source_id == api_key.source_id, + NotificationType.source_id == source.id, NotificationType.subject == envelope.subject, NotificationType.action == envelope.action, + NotificationType.deleted_at.is_(None), ) ).scalar_one_or_none() if known_type is None: @@ -75,7 +78,7 @@ prior = db.execute( select(Event) .where( - Event.source_id == api_key.source_id, + Event.source_id == source.id, Event.dedup_key == envelope.dedup_key, Event.created_at > datetime.now(UTC) - timedelta(seconds=window), ) @@ -85,7 +88,7 @@ raise Deduplicated(prior.id, prior.status) event = Event( - source_id=api_key.source_id, + source_id=source.id, subject=envelope.subject, action=envelope.action, priority=envelope.priority, @@ -117,7 +120,7 @@ ) -> AcceptedEvent: api_key, source = resolve_event_source(db, request) try: - event = accept_event(db, api_key, source.name, envelope) + event = accept_event(db, source, envelope) except Deduplicated as dup: return AcceptedEvent(id=str(dup.event_id), status=dup.event_status, deduplicated=True) api_key.last_used_at = datetime.now(UTC) @@ -148,7 +151,7 @@ results.append(BatchRejected(index=i, detail=err.errors()[0].get("msg", "невалидно"))) continue try: - event = accept_event(db, api_key, source.name, envelope) + event = accept_event(db, source, envelope) results.append(AcceptedEvent(id=str(event.id), status=event.status)) accepted.append(str(event.id)) except Deduplicated as dup: diff --git a/app/auth/apikeys.py b/app/auth/apikeys.py index 411961e..3c756a6 100644 --- a/app/auth/apikeys.py +++ b/app/auth/apikeys.py @@ -57,4 +57,11 @@ raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, detail="Источник ключа удалён" ) + if source.deleted_at is not None: + # архив, не удаление: ключ жив, но события не принимает — админ + # восстановит источник (restore) или выдаст ключ новому + raise HTTPException( + status_code=status.HTTP_401_UNAUTHORIZED, + detail="Источник ключа в архиве — события не принимаются", + ) return row, source \ No newline at end of file diff --git a/app/config.py b/app/config.py index 9c4a4b1..0687151 100644 --- a/app/config.py +++ b/app/config.py @@ -51,6 +51,10 @@ # SPA лежит в образе рядом с приложением; в dev используется vite-сервер. spa_dist_dir: Path = BASE_DIR / "spa_static" + # MCP-сервер (docs/07): статический Bearer-токен для ИИ-агента. + # Пусто — эндпоинт /mcp не монтируется. Генерация: openssl rand -hex 32. + mcp_token: str = "" + # CORS для dev-режима (vite :5173); в проде SPA раздаётся с того же origin. cors_origins: list[str] = ["http://localhost:5173"] diff --git a/app/main.py b/app/main.py index 45f302b..a08813f 100644 --- a/app/main.py +++ b/app/main.py @@ -1,5 +1,7 @@ """Точка входа FastAPI: API + раздача собранной SPA-статики с того же контейнера.""" +from contextlib import AsyncExitStack, asynccontextmanager + from fastapi import FastAPI from fastapi.middleware.cors import CORSMiddleware from fastapi.staticfiles import StaticFiles @@ -16,7 +18,33 @@ def create_app() -> FastAPI: settings = get_settings() - app = FastAPI(title=settings.app_name, docs_url="/api/docs", openapi_url="/api/openapi.json") + + # MCP (docs/07): embedded FastMCP на /mcp streamable-http, гейт по Bearer. + # Пустой MCP_TOKEN — сервера нет (GET /mcp уйдёт в SPA, POST — 405). + # Session manager FastMCP требует явного запуска: enter_async_context + # в lifespan держит его открытым всё время жизни приложения. + mcp_mount = None + lifespan_kwargs: dict = {} + if settings.mcp_token: + from app.mcp.server import McpTokenGuard, build_mcp_server + + raw_mcp_app = build_mcp_server().streamable_http_app() + mcp_mount = McpTokenGuard(raw_mcp_app, settings.mcp_token) + + @asynccontextmanager + async def _lifespan(app: FastAPI): + async with AsyncExitStack() as stack: + await stack.enter_async_context( + raw_mcp_app.router.lifespan_context(raw_mcp_app) + ) + yield + + lifespan_kwargs["lifespan"] = _lifespan + + app = FastAPI( + title=settings.app_name, docs_url="/api/docs", openapi_url="/api/openapi.json", + **lifespan_kwargs, + ) # В проде SPA с того же origin, CORS нужен только dev-серверу vite. app.add_middleware( @@ -33,6 +61,20 @@ app.include_router(admin_routes.router) app.include_router(me_routes.router) + # MCP монтируется до SPA catch-all (иначе GET /mcp съел бы роутер SPA). + # Именно Route, не Mount: Mount('/mcp') матчит только '/mcp/...' (regex + # с trailing slash), точный POST /mcp провалился бы в SPA (GET index / + # POST 405); Route отдаёт в саб-апп путь без root_path-трюков. + if mcp_mount is not None: + from starlette.routing import Route + + app.router.routes.append( + Route( + "/mcp", mcp_mount, methods=["GET", "POST", "DELETE", "OPTIONS"], + include_in_schema=False, + ) + ) + # Собранная SPA: в Docker-образе кладётся в spa_static/. Если собранной # статики нет (локальный запуск api без фронта) — просто пропускаем. spa = settings.spa_dist_dir diff --git a/app/mcp/__init__.py b/app/mcp/__init__.py new file mode 100644 index 0000000..c57ae3e --- /dev/null +++ b/app/mcp/__init__.py @@ -0,0 +1,6 @@ +"""MCP-сервер Synapse (docs/07): embedded FastMCP над admin-API. + +Тулы переиспользуют admin-route функции (та же валидация, те же ответы); +транспорт — streamable-http на /mcp, доступ по статическому Bearer-токену +MCP_TOKEN из .env. Пустой токен — эндпоинт не монтируется. +""" \ No newline at end of file diff --git a/app/mcp/server.py b/app/mcp/server.py new file mode 100644 index 0000000..10e0dd9 --- /dev/null +++ b/app/mcp/server.py @@ -0,0 +1,57 @@ +"""Сборка MCP-сервера Synapse (docs/07). + +FastMCP (stateless_http) — Starlette-приложение streamable-http; поверх — +McpTokenGuard: статический Bearer-токен из MCP_TOKEN (.env). Токен один, +выдаётся администратором (генерация `openssl rand -hex 32`) — полномочия +тулов = superadmin внутри API, так что токен держать как ключ прод-сервера. +""" + +from mcp.server.fastmcp import FastMCP + +from app.mcp.tools import register_tools + + +def build_mcp_server() -> FastMCP: + mcp = FastMCP( + "synapse", + stateless_http=True, + instructions=( + "Управление хабом уведомлений Gnexus Synapse: источники и выдача " + "API-ключей, типы уведомлений, цели каналов, правила маршрутизации, " + "поток событий и доставок, настройки. Удаление = архив (restore " + "возвращает); секреты настроек write-only." + ), + ) + register_tools(mcp) + return mcp + + +class McpTokenGuard: + """Чистый ASGI-гейт (не BaseHTTPMiddleware): проверяет Bearer-токен на + каждом http-запросе. Не BaseHTTPMiddleware — чтобы не ломать SSE-стримы.""" + + def __init__(self, app, token: str) -> None: + self.app = app + self.expected = f"Bearer {token}".encode() + + async def __call__(self, scope, receive, send): + if scope["type"] != "http": + await self.app(scope, receive, send) + return + headers = {k.lower(): v for k, v in scope.get("headers") or []} + if headers.get(b"authorization") == self.expected: + await self.app(scope, receive, send) + return + body = b'{"detail":"Invalid or missing MCP token"}' + await send( + { + "type": "http.response.start", + "status": 401, + "headers": [ + (b"content-type", b"application/json"), + (b"www-authenticate", b"Bearer"), + (b"content-length", str(len(body)).encode()), + ], + } + ) + await send({"type": "http.response.body", "body": body}) \ No newline at end of file diff --git a/app/mcp/tools.py b/app/mcp/tools.py new file mode 100644 index 0000000..ab4043e --- /dev/null +++ b/app/mcp/tools.py @@ -0,0 +1,541 @@ +"""Инструменты MCP: управление Synapse ИИ-агентом. + +Каждый тул — синхронная функция (FastMCP гоняет её в threadpool — совместимо +с синхронной SQLAlchemy-сессией). Логика записи НЕ дублируется: тулы +переиспользуют admin-route функции (admin_routes.py) напрямую с +user=MCP_USER, db=<сессия> — Depends не срабатывают, та же валидация и те же +коды ошибок. HTTPException -> {"error": status, "detail"}. +""" + +import uuid +from typing import Any + +from fastapi import HTTPException +from pydantic import BaseModel, ValidationError +from sqlalchemy import func, select, text + +from app.api import admin_routes +from app.api.admin_schemas import ( + KeyIn, + RuleIn, + RulePatch, + SettingsPut, + SourceIn, + TargetIn, + TargetPatch, + TypeIn, +) +from app.auth.deps import AuthenticatedUser +from app.config import get_settings +from app.database import SessionLocal +from app.models import ( + ChannelTarget, + Delivery, + Event, + NotificationType, + PushSubscription, + RoutingRule, + Source, +) +from app.worker.celery_app import celery_app + +# Синтетический superadmin: require_admin пропускает, автор — MCP-токен. +MCP_USER = AuthenticatedUser( + user_id="mcp", email="mcp@synapse.local", email_verified=True, system_role="superadmin" +) + +LIMIT_CAP = 500 + + +def _run(fn) -> Any: + """Сессия на вызов тулза; HTTPException/pydantic ValidationError -> dict ошибки.""" + try: + with SessionLocal() as db: + return fn(db) + except HTTPException as err: + return {"error": err.status_code, "detail": str(err.detail)} + except ValidationError as err: + return {"error": 422, "detail": err.errors()[0].get("msg", "невалидно")} + + +def _patch(model_cls: type[BaseModel], **kwargs) -> BaseModel: + """PATCH-схема только из переданных полей: None = «не менять» + (exclude_unset в route-функции видит незаданные поля незаданными).""" + return model_cls.model_validate({k: v for k, v in kwargs.items() if v is not None}) + + +def register_tools(mcp) -> None: + + @mcp.tool() + def system_status() -> dict: + """Здоровье Synapse: БД, Celery-воркер, конфигурация (SSO).""" + def fn(db): + db.execute(text("SELECT 1")) + return {"database": "ok"} + out = _run(fn) + if "error" in out: + return {"database": "unavailable", **out} + try: + pong = celery_app.send_task("synapse.ping").get(timeout=5) + out.update(worker=pong, sso_configured=get_settings().gauth_configured()) + except Exception: # noqa: BLE001 — недоступный воркер не ломает статус + out.update(worker="unavailable", sso_configured=get_settings().gauth_configured()) + return out + + @mcp.tool() + def stats_get() -> dict: + """Сводка: события и доставки по статусам, число живых + sources/types/targets/rules, число deliveries в ожидании.""" + def fn(db): + ev = dict( + db.execute(select(Event.status, func.count(Event.id)).group_by(Event.status)).all() + ) + dv = dict( + db.execute( + select(Delivery.status, func.count(Delivery.id)).group_by(Delivery.status) + ).all() + ) + return { + "events": ev, + "deliveries": dv, + "pending_deliveries": dv.get("pending", 0), + "sources": db.scalar( + select(func.count(Source.id)).where(Source.deleted_at.is_(None)) + ), + "types": db.scalar( + select(func.count(NotificationType.id)).where( + NotificationType.deleted_at.is_(None) + ) + ), + "targets": db.scalar( + select(func.count(ChannelTarget.id)).where( + ChannelTarget.deleted_at.is_(None) + ) + ), + "rules": db.scalar( + select(func.count(RoutingRule.id)).where(RoutingRule.deleted_at.is_(None)) + ), + } + return _run(fn) + + # --- источники и ключи (выдача ключей клиентам, docs/07) --- + + @mcp.tool() + def sources_list( + query: str | None = None, include_archived: bool = False, limit: int = 100 + ) -> list[dict]: + """Источники событий. query — подстрока в name/label/description; + include_archived=true — показать и заархивированные (deleted_at != null).""" + def fn(db): + rows = admin_routes.list_sources( + include_archived=include_archived, user=MCP_USER, db=db + ) + out = [r.model_dump(mode="json") for r in rows] + if query: + ql = query.lower() + out = [ + r for r in out + if ql in r["name"].lower() + or ql in (r.get("label") or "").lower() + or ql in (r.get("description") or "").lower() + ] + return out[: max(1, min(limit, LIMIT_CAP))] + return _run(fn) + + @mcp.tool() + def source_create(name: str, label: str | None = None, description: str | None = None) -> dict: + """Зарегистрировать источник событий (сервис). name — slug [a-z0-9._-], + label/description — человекочитаемое «что это за сервис».""" + def fn(db): + src = admin_routes.create_source( + SourceIn(name=name, label=label, description=description), + user=MCP_USER, db=db, + ) + return src.model_dump(mode="json") + return _run(fn) + + @mcp.tool() + def source_archive(source_id: int) -> dict: + """Заархивировать источник (это не удаление): ключ перестаёт принимать + события (401), source_restore возвращает в строй.""" + def fn(db): + admin_routes.delete_source(source_id, user=MCP_USER, db=db) + return {"ok": True, "archived": source_id} + return _run(fn) + + @mcp.tool() + def source_restore(source_id: int) -> dict: + """Вернуть заархивированный источник в строй. Конфликт имени → 409.""" + def fn(db): + src = admin_routes.restore_source(source_id, user=MCP_USER, db=db) + return src.model_dump(mode="json") + return _run(fn) + + @mcp.tool() + def keys_list(source_id: int) -> list[dict]: + """API-ключи источника: хэш и хвост — plaintext никогда не возвращается, + полный токен виден один раз в key_issue.""" + def fn(db): + keys = admin_routes.list_keys(source_id, user=MCP_USER, db=db) + return [k.model_dump(mode="json") for k in keys] + return _run(fn) + + @mcp.tool() + def key_issue(source_id: int, name: str) -> dict: + """Выдать API-ключ клиенту (сервису-источнику). Полный токен `syn_...` + в поле token возвращается ОДИН РАЗ — передать клиенту (в его .env); + в Synapse остаётся только хэш. Повторный вызов = новый токен.""" + def fn(db): + created = admin_routes.create_key( + source_id, KeyIn(name=name), user=MCP_USER, db=db + ) + return created.model_dump(mode="json") + return _run(fn) + + @mcp.tool() + def key_revoke(source_id: int, key_id: int) -> dict: + """Отозвать API-ключ источника (события с ним начнут получать 401).""" + def fn(db): + key = admin_routes.revoke_key(source_id, key_id, user=MCP_USER, db=db) + return key.model_dump(mode="json") + return _run(fn) + + # --- типы уведомлений (тройка source/subject/action) --- + + @mcp.tool() + def types_list( + source_id: int | None = None, + query: str | None = None, + include_archived: bool = False, + limit: int = 200, + ) -> list[dict]: + """Зарегистрированные типы уведомлений (тройка source/subject/action). + Событие с незарегистрированной тройкой приём отклоняет (422).""" + def fn(db): + rows = admin_routes.list_types( + include_archived=include_archived, user=MCP_USER, db=db + ) + out = [r.model_dump(mode="json") for r in rows] + if source_id is not None: + out = [r for r in out if r["source_id"] == source_id] + if query: + ql = query.lower() + out = [ + r for r in out + if ql in r["subject"].lower() + or ql in r["action"].lower() + or ql in (r.get("description") or "").lower() + ] + return out[: max(1, min(limit, LIMIT_CAP))] + return _run(fn) + + @mcp.tool() + def type_register( + source_id: int, + subject: str, + action: str, + description: str | None = None, + payload_schema: dict | None = None, + ) -> dict: + """Зарегистрировать тип уведомления: тройка (source_id, subject, action). + payload_schema — опциональная JSON Schema payload'а (валидация мягкая, в воркере).""" + def fn(db): + nt = admin_routes.create_type( + TypeIn( + source_id=source_id, subject=subject, action=action, + payload_schema=payload_schema, description=description, + ), + user=MCP_USER, db=db, + ) + return nt.model_dump(mode="json") + return _run(fn) + + @mcp.tool() + def type_archive(type_id: int) -> dict: + """Заархивировать тип: приём его тройки начнёт отвечать + 422 «тип не зарегистрирован»; type_restore вернёт в строй.""" + def fn(db): + admin_routes.delete_type(type_id, user=MCP_USER, db=db) + return {"ok": True, "archived": type_id} + return _run(fn) + + @mcp.tool() + def type_restore(type_id: int) -> dict: + """Вернуть заархивированный тип в строй (источник должен быть жив).""" + def fn(db): + nt = admin_routes.restore_type(type_id, user=MCP_USER, db=db) + return nt.model_dump(mode="json") + return _run(fn) + + # --- цели каналов --- + + @mcp.tool() + def targets_list( + channel: str | None = None, + enabled: bool | None = None, + include_archived: bool = False, + limit: int = 200, + ) -> list[dict]: + """Цели каналов (TG-чат, email-адрес, s2s-точка). channel — фильтр.""" + def fn(db): + rows = admin_routes.list_targets( + include_archived=include_archived, user=MCP_USER, db=db + ) + out = [r.model_dump(mode="json") for r in rows] + if channel: + out = [r for r in out if r["channel"] == channel] + if enabled is not None: + out = [r for r in out if r["enabled"] == enabled] + return out[: max(1, min(limit, LIMIT_CAP))] + return _run(fn) + + @mcp.tool() + def target_create( + channel: str, name: str, config: dict, description: str | None = None + ) -> dict: + """Создать цель канала. channel ∈ telegram|email|s2s|internal_log; + config — идентификатор цели (s2s: {"endpoint": …, "token_ref": "navi-rei"}). + Секреты целей — в .env/gnexus-creds, в config только token_ref.""" + def fn(db): + t = admin_routes.create_target( + TargetIn(channel=channel, name=name, config=config, description=description), + user=MCP_USER, db=db, + ) + return t.model_dump(mode="json") + return _run(fn) + + @mcp.tool() + def target_patch( + target_id: int, name: str | None = None, description: str | None = None, + config: dict | None = None, enabled: bool | None = None, + ) -> dict: + """Частичное обновление цели (переданное — меняется, None — «не менять»). + Канал не меняется (ссылки из правил).""" + def fn(db): + t = admin_routes.patch_target( + target_id, + _patch(TargetPatch, name=name, description=description, + config=config, enabled=enabled), + user=MCP_USER, db=db, + ) + return t.model_dump(mode="json") + return _run(fn) + + @mcp.tool() + def target_archive(target_id: int) -> dict: + """Заархивировать цель: правила продолжат матчить, но доставки в неё + запишутся skipped «цель в архиве»; restore вернёт в строй.""" + def fn(db): + admin_routes.delete_target(target_id, user=MCP_USER, db=db) + return {"ok": True, "archived": target_id} + return _run(fn) + + @mcp.tool() + def target_restore(target_id: int) -> dict: + """Вернуть заархивированную цель в строй. Конфликт имени → 409.""" + def fn(db): + t = admin_routes.restore_target(target_id, user=MCP_USER, db=db) + return t.model_dump(mode="json") + return _run(fn) + + # --- правила маршрутизации --- + + @mcp.tool() + def rules_list( + enabled: bool | None = None, include_archived: bool = False, limit: int = 200 + ) -> list[dict]: + """Правила маршрутизации (условия -> действия). weight — порядок в режиме first.""" + def fn(db): + rows = admin_routes.list_rules(include_archived=include_archived, user=MCP_USER, db=db) + out = [r.model_dump(mode="json") for r in rows] + if enabled is not None: + out = [r for r in out if r["enabled"] == enabled] + return out[: max(1, min(limit, LIMIT_CAP))] + return _run(fn) + + @mcp.tool() + def rule_create( + name: str, + actions: list[dict], + conditions: dict | None = None, + template: str | None = None, + throttle_seconds: int | None = None, + weight: int = 0, + ) -> dict: + """Создать правило: действия [{"channel": "telegram|email|s2s|internal_log|user|push", + "target_id": id|null, "template": "{{ payload.x }}"}]; условия + {"source": "имя"|null, "subjects": [], "actions": [], "priority_min": "normal", + "payload": {...}}. Каналы user/push целей не имеют (адресат payload.user_id).""" + def fn(db): + rule = admin_routes.create_rule( + RuleIn.model_validate({ + "name": name, "actions": actions, + "conditions": conditions or {}, "template": template, + "throttle_seconds": throttle_seconds, "weight": weight, + }), + user=MCP_USER, db=db, + ) + return rule.model_dump(mode="json") + return _run(fn) + + @mcp.tool() + def rule_patch( + rule_id: int, name: str | None = None, enabled: bool | None = None, + conditions: dict | None = None, template: str | None = None, + throttle_seconds: int | None = None, weight: int | None = None, + actions: list[dict] | None = None, + ) -> dict: + """Частично обновить правило (вкл/выкл, условия, действия, шаблон, вес).""" + def fn(db): + rule = admin_routes.patch_rule( + rule_id, + _patch(RulePatch, name=name, enabled=enabled, conditions=conditions, + template=template, throttle_seconds=throttle_seconds, + weight=weight, actions=actions), + user=MCP_USER, db=db, + ) + return rule.model_dump(mode="json") + return _run(fn) + + @mcp.tool() + def rule_archive(rule_id: int) -> dict: + """Заархивировать правило (действия сохраняются — restore вернёт целиком).""" + def fn(db): + admin_routes.delete_rule(rule_id, user=MCP_USER, db=db) + return {"ok": True, "archived": rule_id} + return _run(fn) + + @mcp.tool() + def rule_restore(rule_id: int) -> dict: + """Вернуть заархивированное правило в строй.""" + def fn(db): + rule = admin_routes.restore_rule(rule_id, user=MCP_USER, db=db) + return rule.model_dump(mode="json") + return _run(fn) + + # --- поток: события и доставки --- + + @mcp.tool() + def events_list( + limit: int = 50, status: str | None = None, source_name: str | None = None + ) -> list[dict]: + """Последние события (конверт + payload + доставки). status ∈ + queued|processing|done|failed|expired; source_name — фильтр по имени источника.""" + def fn(db): + source_id = None + if source_name: + src = db.execute( + select(Source.id).where(Source.name == source_name) + ).scalar_one_or_none() + source_id = src if src is not None else -1 # нет источника — пусто + rows = admin_routes.list_events( + user=MCP_USER, db=db, limit=limit, status_filter=status, source_id=source_id + ) + return [r.model_dump(mode="json") for r in rows] + return _run(fn) + + @mcp.tool() + def event_get(event_id: str) -> dict: + """Событие по id: конверт, статусы и доставки в каналы.""" + def fn(db): + ev = admin_routes.get_event(uuid.UUID(event_id), user=MCP_USER, db=db) + return ev.model_dump(mode="json") + return _run(fn) + + @mcp.tool() + def deliveries_list( + limit: int = 100, status: str | None = None, channel: str | None = None + ) -> list[dict]: + """Журнал доставок: канал, цель, статус, попытки, ошибки ретраев. + status ∈ pending|delivered|skipped|failed.""" + def fn(db): + rows = admin_routes.list_deliveries( + user=MCP_USER, db=db, limit=limit, status_filter=status, channel=channel + ) + return [r.model_dump(mode="json") for r in rows] + return _run(fn) + + @mcp.tool() + def push_subscriptions_list(user_id: str | None = None, limit: int = 100) -> list[dict]: + """Web-push подписки пользователей (адресат — payload.user_id). + user_id — фильтр по sub gnexus-auth.""" + def fn(db): + q = select(PushSubscription).order_by(PushSubscription.id.desc()) + if user_id: + q = q.where(PushSubscription.user_id == user_id) + rows = db.execute(q.limit(max(1, min(limit, LIMIT_CAP)))).scalars().all() + return [ + { + "id": s.id, "user_id": s.user_id, "endpoint": s.endpoint, + "ua": s.ua, "created_at": s.created_at.isoformat(), + } + for s in rows + ] + return _run(fn) + + @mcp.tool() + def push_subscription_delete(subscription_id: int) -> dict: + """Удалить web-push подписку: единственное ЖЁСТКОЕ удаление в Synapse — + устройство пользователя больше не подписано.""" + def fn(db): + sub = db.get(PushSubscription, subscription_id) + if sub is None: + raise HTTPException(status_code=404, detail="Подписка не найдена") + db.delete(sub) + db.commit() + return {"ok": True, "deleted": subscription_id} + return _run(fn) + + # --- настройки --- + + @mcp.tool() + def settings_get() -> list[dict]: + """Реестр настроек Synapse: дефолт .env + оверрайды админки. + Секреты write-only — value всегда null, факт наличия в set.""" + def fn(db): + rows = admin_routes.list_settings(user=MCP_USER, db=db) + return [r.model_dump(mode="json") for r in rows] + return _run(fn) + + @mcp.tool() + def settings_put(values: dict) -> dict: + """Установить оверрайды настроек (поверх .env). Правила: строка — set, + null — вернуть дефолт .env, "" у секрета — «не менять». Нарушение — 422.""" + def fn(db): + return admin_routes.put_settings( + SettingsPut.model_validate({"values": values}), user=MCP_USER, db=db + ) + return _run(fn) + + # --- приём контрольного события --- + + @mcp.tool() + def send_test_event( + source_name: str, subject: str, action: str, payload: dict | None = None, + priority: str = "normal", ttl_seconds: int | None = None, + ) -> dict: + """Контрольное событие «как бы от источника» (source_name): полный + маршрут в воркере -> доставки по правилам. Для проверки нового правила. + Тройка должна быть зарегистрирована (type_register).""" + from app.api.events_routes import accept_event + from app.api.schemas import EventEnvelope + + def fn(db): + src = db.execute( + select(Source).where(Source.name == source_name) + ).scalar_one_or_none() + if src is None: + raise HTTPException(status_code=404, detail=f"Источник '{source_name}' не найден") + if src.deleted_at is not None: + raise HTTPException( + status_code=422, detail=f"Источник '{source_name}' в архиве — сначала restore" + ) + envelope = EventEnvelope( + source=source_name, subject=subject, action=action, priority=priority, + payload=payload or {}, ttl_seconds=ttl_seconds, + ) + event = accept_event(db, src, envelope) + db.commit() + celery_app.send_task("synapse.ingest", args=[str(event.id)]) + return {"id": str(event.id), "status": event.status, "source": source_name} + return _run(fn) \ No newline at end of file diff --git a/app/models/dicts.py b/app/models/dicts.py index 2b39a9f..055ffad 100644 --- a/app/models/dicts.py +++ b/app/models/dicts.py @@ -51,6 +51,9 @@ created_at: Mapped[datetime] = mapped_column( DateTime(timezone=True), server_default=func.now() ) + # Архив вместо удаления (MCP): DELETE ставит метку, restore снимает. + # Живые пути фильтруют deleted_at.is_(None); история — нет. + deleted_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) class ApiKey(Base): @@ -98,6 +101,7 @@ created_at: Mapped[datetime] = mapped_column( DateTime(timezone=True), server_default=func.now() ) + deleted_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) class ChannelTarget(Base): @@ -121,4 +125,5 @@ enabled: Mapped[bool] = mapped_column(default=True) created_at: Mapped[datetime] = mapped_column( DateTime(timezone=True), server_default=func.now() - ) \ No newline at end of file + ) + deleted_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) \ No newline at end of file diff --git a/app/models/rules.py b/app/models/rules.py index 11e73f7..8632dea 100644 --- a/app/models/rules.py +++ b/app/models/rules.py @@ -57,6 +57,7 @@ created_at: Mapped[datetime] = mapped_column( DateTime(timezone=True), server_default=func.now() ) + deleted_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True)) updated_at: Mapped[datetime | None] = mapped_column( DateTime(timezone=True), onupdate=func.now() ) diff --git a/app/worker/tasks.py b/app/worker/tasks.py index 75709ac..3efd660 100644 --- a/app/worker/tasks.py +++ b/app/worker/tasks.py @@ -61,6 +61,9 @@ ) if target is None: return False, "s2s: у доставки нет цели (target_id пуст)" + if target.deleted_at is not None: + # ретрай по цели, заархивированной после маршрутизации + return False, "цель в архиве — restore в админке вернёт её в строй" return send_s2s(target, event, source_name) if delivery.channel == "push": # Web-push: подписки получателя хранятся в push_subscriptions (docs/06). @@ -203,7 +206,10 @@ event.status = "processing" rules = db.execute( - select(RoutingRule).where(RoutingRule.enabled.is_(True)).order_by(RoutingRule.weight) + select(RoutingRule).where( + RoutingRule.enabled.is_(True), + RoutingRule.deleted_at.is_(None), + ).order_by(RoutingRule.weight) ).scalars().all() matched = [ rule for rule in rules @@ -233,6 +239,29 @@ if act.channel in _USER_CHANNELS: recipient, user_reason = _user_recipient(event) + # Архивная цель — «сюда больше не ходим»: правило живо (или + # цель заархивировали позже правила), доставка уходит в + # skipped с аудитом (аналог user без адресата). + if act.channel not in _USER_CHANNELS and act.target_id is not None: + t_row = db.get(ChannelTarget, act.target_id) + if t_row is None or t_row.deleted_at is not None: + db.add( + Delivery( + event_id=event.id, + rule_id=rule.id, + channel=act.channel, + channel_target_id=act.target_id, + status="skipped", + attempts=0, + rendered_message=_render( + act.template or rule.template or "", event, source.name + ), + last_error="цель в архиве — restore в админке вернёт её в строй", + ) + ) + n_skipped += 1 + continue + # шумодав: срезаем (с аудит-записью), critical проходит if event.priority != "critical" and _throttled( db, rule, act.channel, act.target_id, now, recipient=recipient diff --git a/docs/07-mcp.md b/docs/07-mcp.md new file mode 100644 index 0000000..0029478 --- /dev/null +++ b/docs/07-mcp.md @@ -0,0 +1,116 @@ +# 07 — MCP-сервер: управление Synapse ИИ-агентом + +Synapse встраивает **FastMCP** (пакет `mcp`) — тот же контейнер API, тот же +процесс uvicorn. Эндпоинт `/mcp`, транспорт **streamable-http** (stateless), +доступ — статический Bearer-токен `MCP_TOKEN` из `.env`. + +Зачем: ИИ-агент (Claude Code и совместимые) управляет Synapse целиком — +выдаёт API-ключи клиентам (сервисам-источникам), регистрирует типы, ведёт +правила маршрутизации, смотрит поток событий и доставок, меняет настройки — +без ручной работы в админке. Правка 2026-10-03. + +## Включение + +1. `MCP_TOKEN=` в `.env` (и gnexus-creds). Пусто/не задано — `/mcp` + не существует: GET уходит в SPA, POST отвечает 405. +2. Перезапуск контейнера api. +3. Регистрация в Claude Code: + +```bash +claude mcp add --transport http synapse http://localhost:8013/mcp \ + --header "Authorization: Bearer " +``` + +Генерация токена: `openssl rand -hex 32`. Токен — как ключ прод-сервера +(внутри API полномочия MCD = superadmin), хранить в `.env` + gnexus-creds, +никогда не в репозитории. + +## Архитектура + +- **Тулы не дублируют логику**: сервер переиспользует admin-route функции + (`app/api/admin_routes.py`) напрямую — та же pydantic-валидация, те же коды + ошибок и тексты. Автор запросов — синтетический `superadmin` (токен и есть + доверие). +- Ошибка тулза возвращается как данные: `{"error": 409, "detail": "…"}` — + агент читает причину и может её исправить (например restore архивной записи). +- Синхронная SQLAlchemy-сессия на вызов; FastMCP исполняет синхронные тулы + в threadpool. +- Гейт — чистый ASGI-класс (`McpTokenGuard`, app/mcp/server.py): проверяет + `Authorization: Bearer` на каждом запросе, 401 (`www-authenticate: Bearer`) + при несовпадении; не BaseHTTPMiddleware — SSE-стримы не ломает. +- Выключение: убрать `MCP_TOKEN` → сервер не монтируется. + +## Каталог тулов (31) + +| Группа | Тулы | +|---|---| +| Состояние | `system_status` (БД + Celery ping + SSO) · `stats_get` (счётчики по статусам, живые справочники) | +| Источники | `sources_list` (query, include_archived) · `source_create` · `source_archive` · `source_restore` | +| Ключи | `keys_list` · `key_issue` (**plaintext один раз** в `token`) · `key_revoke` | +| Типы | `types_list` (source_id, query, include_archived) · `type_register` · `type_archive` · `type_restore` | +| Цели | `targets_list` (channel, enabled, include_archived) · `target_create` · `target_patch` · `target_archive` · `target_restore` | +| Правила | `rules_list` (enabled, include_archived) · `rule_create` · `rule_patch` · `rule_archive` · `rule_restore` | +| Поток | `events_list` (status, source_name) · `event_get` · `deliveries_list` (status, channel) | +| Push | `push_subscriptions_list` (user_id) · `push_subscription_delete` (жёстко — устройства пользователей) | +| Настройки | `settings_get` (секреты write-only: value=null, факт в set) · `settings_put` (правила docs/06: "" — не менять, null — сброс) | +| Контроль | `send_test_event` (source_name, subject, action, payload, priority) — полный маршрут в воркере | + +## Архив вместо удаления (семантика DELETE) + +DELETE не удаляет — ставит `deleted_at` (архив). Действует на **sources, +notification_types, channel_targets, routing_rules**. Пуш-подписки — +исключение: устройства пользователей удаляются физически (push-нотификации +не переживают несуществующую подписку). + +Что это значит: + +- DELETE отвечает 204; вторая архивация того же id — 404 («не найден»). +- Списки по умолчанию живое; `include_archived=true` покажет архив + (`deleted_at` в объектах; в UI-списках поле всегда null). +- Архивный источник: ключ получает 401 «Источник ключа в архиве» — события + не принимаются; ключи отозвать нельзя, но и выдать новый нельзя (404 + «сначала restore»). +- Архивный тип: тройка приёма → 422 «не зарегистрирован». +- Архивная цель в живом правиле: матчится, но доставка уходит в `skipped` + с причиной «цель в архиве» (аудит, видно в deliveries). +- Create поверх архива: имя занято архивной записью → **409** с подсказкой + «есть в архиве (id=N) — восстанови (restore) или выбери другое имя». +- Restore: `POST /api/v1/admin/{sources|types|targets|rules}/{id}/restore` + → 200 с объектом; конфликт имени с живой записью → 409. Restore типа + требует живого источника (409). +- История не режется: старые события/доставки показывают имена архивных + источников и целей. + +REST-эквивалент доступен и из админ-API (те же эндпоинты), UI-клиент SPA +ничего не заметил (204-ответы прежние). + +## Пример сценариев + +Выдать ключ клиенту: + +``` +1. source_create(name="bugtrail", description="Трекер ошибок…") +2. key_issue(source_id=, name="основной") + → token: "syn_…" — передать клиенту один раз +3. type_register(source_id, subject="issue", action="created") +4. Третьесторонний сервис шлёт POST /api/v1/events с Bearer syn_… +``` + +Проверить правило: + +``` +1. rule_create(name="critical → Navi rei", actions=[{"channel":"s2s","target_id":3}], + conditions={"priority_min":"critical"}) +2. send_test_event(source_name="monitoring", subject="container", action="down", + payload={"container":"api"}, priority="critical") +3. events_list(source_name="monitoring") → доставed в цель +``` + +## Ограничения + +- Stateless транспорт: без сессий между вызовами (каждый вызов независим) — + агент просто вызывает тула последовательно. +- Секреты настроек остаются write-only (`settings_get` их не отдаёт). +- `send_test_event` и просмотр потока требуют живого воркера (Celery) — + `system_status` покажет. +- Транспорт — внутренняя сеть/локальный порт API: наружу `/mcp` не публикуем. \ No newline at end of file diff --git a/requirements.txt b/requirements.txt index 31750f4..f43526c 100644 --- a/requirements.txt +++ b/requirements.txt @@ -12,3 +12,4 @@ gnexus-gauth @ git+https://git.gnexus.space/git/root/gnexus-auth-client-py.git pywebpush>=2.0 jinja2>=3.1 +mcp>=1.16,<2