diff --git a/app/api/admin_routes.py b/app/api/admin_routes.py new file mode 100644 index 0000000..9f32125 --- /dev/null +++ b/app/api/admin_routes.py @@ -0,0 +1,485 @@ +"""Admin API: CRUD конфигурации + просмотр потока. Всё под require_admin. + +Это бэкенд SPA админки (#32+): списки и создание/удаление — то, без чего +CLI не нужен в проде. Полное редактирование правил (PATCH) — в этой же +схеме: только enabled/name/template/throttle/weight/conditions/actions. +""" + +import uuid +from datetime import UTC, datetime + +from fastapi import APIRouter, Depends, HTTPException, Query, status +from pydantic import ValidationError +from sqlalchemy import delete, func, select +from sqlalchemy.orm import Session + +from app.auth.apikeys import generate_token, token_hash_of +from app.auth.deps import AuthenticatedUser, require_admin +from app.database import get_db +from app.models import ( + ApiKey, + ChannelTarget, + Delivery, + Event, + NotificationType, + RoutingRule, + RoutingRuleAction, + Source, +) +from app.api.admin_schemas import ( + AdminEventOut, + DeliveryOut, + KeyCreated, + KeyIn, + KeyOut, + RuleActionIn, + RuleActionOut, + RuleIn, + RuleOut, + RulePatch, + SourceIn, + SourceOut, + TargetIn, + TargetOut, + TypeIn, + TypeOut, +) + +router = APIRouter(prefix="/api/v1/admin", tags=["admin"]) + + +def _conflict(detail: str) -> HTTPException: + return HTTPException(status_code=status.HTTP_409_CONFLICT, detail=detail) + + +# --- sources --- + +@router.get("/sources", response_model=list[SourceOut]) +def list_sources( + user: AuthenticatedUser = Depends(require_admin), db: Session = Depends(get_db) +) -> list[SourceOut]: + sources = db.execute(select(Source).order_by(Source.name)).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)) + .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), + types_count=type_counts.get(s.id, 0), + ) + for s in sources + ] + + +@router.post("/sources", response_model=SourceOut, status_code=201) +def create_source( + payload: SourceIn, + 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(): + 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()) + + +@router.delete("/sources/{source_id}", status_code=204) +def delete_source( + source_id: int, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> None: + source = db.get(Source, source_id) + if source is 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) + db.commit() + + +# --- keys --- + +@router.get("/sources/{source_id}/keys", response_model=list[KeyOut]) +def list_keys( + source_id: int, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> list[KeyOut]: + if db.get(Source, source_id) is None: + raise HTTPException(status_code=404, detail="Источник не найден") + keys = db.execute( + select(ApiKey).where(ApiKey.source_id == source_id).order_by(ApiKey.created_at) + ).scalars().all() + return [ + KeyOut( + id=k.id, name=k.name, token_hint=k.token_hint, created_at=k.created_at, + last_used_at=k.last_used_at, revoked_at=k.revoked_at, + ) + for k in keys + ] + + +@router.post("/sources/{source_id}/keys", response_model=KeyCreated, status_code=201) +def create_key( + source_id: int, + payload: KeyIn, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> KeyCreated: + source = db.get(Source, source_id) + if source is None: + raise HTTPException(status_code=404, detail="Источник не найден") + token = generate_token() + key = ApiKey( + source_id=source_id, + name=payload.name, + token_hash=token_hash_of(token), + token_hint=token[-4:], + ) + db.add(key) + db.commit() + db.refresh(key) + return KeyCreated( + id=key.id, name=key.name, token_hint=key.token_hint, created_at=key.created_at, + last_used_at=None, revoked_at=None, token=token, + ) + + +@router.post("/sources/{source_id}/keys/{key_id}/revoke", response_model=KeyOut) +def revoke_key( + source_id: int, + key_id: int, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> KeyOut: + key = db.get(ApiKey, key_id) + if key is None or key.source_id != source_id: + raise HTTPException(status_code=404, detail="Ключ не найден") + if key.revoked_at is not None: + raise _conflict("Ключ уже отозван") + key.revoked_at = datetime.now(UTC) + db.commit() + db.refresh(key) + return KeyOut( + id=key.id, name=key.name, token_hint=key.token_hint, created_at=key.created_at, + last_used_at=key.last_used_at, revoked_at=key.revoked_at, + ) + + +# --- types --- + +@router.get("/types", response_model=list[TypeOut]) +def list_types( + 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() + ) + 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, + ) + for nt, source_name in rows + ] + + +@router.post("/types", response_model=TypeOut, status_code=201) +def create_type( + payload: TypeIn, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> TypeOut: + source = db.get(Source, payload.source_id) + if source is None: + raise HTTPException(status_code=422, detail="Источник не найден") + dup = db.execute( + select(NotificationType.id).where( + NotificationType.source_id == payload.source_id, + NotificationType.subject == payload.subject, + NotificationType.action == payload.action, + ) + ).scalar_one_or_none() + if dup: + 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()) + + +@router.delete("/types/{type_id}", status_code=204) +def delete_type( + type_id: int, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> None: + if db.get(NotificationType, type_id) is None: + raise HTTPException(status_code=404, detail="Тип не найден") + db.execute(delete(NotificationType).where(NotificationType.id == type_id)) + db.commit() + + +# --- targets --- + +@router.get("/targets", response_model=list[TargetOut]) +def list_targets( + 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() + 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, + ) + for t in targets + ] + + +@router.post("/targets", response_model=TargetOut, status_code=201) +def create_target( + payload: TargetIn, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> TargetOut: + dup = db.execute( + select(ChannelTarget.id).where( + ChannelTarget.channel == payload.channel, ChannelTarget.name == payload.name + ) + ).scalar_one_or_none() + if dup: + 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()) + + +@router.delete("/targets/{target_id}", status_code=204) +def delete_target( + target_id: int, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> None: + if db.get(ChannelTarget, target_id) is None: + raise HTTPException(status_code=404, detail="Цель не найдена") + db.execute(delete(ChannelTarget).where(ChannelTarget.id == target_id)) + # ссылки из routing_rule_actions/deliveries обнуляет FK SET NULL + db.commit() + + +# --- rules --- + +def _rule_to_out(db: Session, rule: RoutingRule) -> RuleOut: + targets = { + t.id: t.name + for t in db.execute(select(ChannelTarget)).scalars().all() + } + actions = db.execute( + select(RoutingRuleAction).where(RoutingRuleAction.rule_id == rule.id).order_by(RoutingRuleAction.id) + ).scalars().all() + return RuleOut( + 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, + actions=[ + RuleActionOut( + id=act.id, channel=act.channel, target_id=act.target_id, + template=act.template, target_name=targets.get(act.target_id), + ) + for act in actions + ], + ) + + +@router.get("/rules", response_model=list[RuleOut]) +def list_rules( + 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() + return [_rule_to_out(db, r) for r in rules] + + +@router.post("/rules", response_model=RuleOut, status_code=201) +def create_rule( + payload: RuleIn, + user: AuthenticatedUser = Depends(require_admin), + 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} не найдена") + data = payload.model_dump() + actions = data.pop("actions") + rule = RoutingRule(**data) + db.add(rule) + db.flush() + for act in actions: + db.add(RoutingRuleAction(rule_id=rule.id, **act)) + db.commit() + db.refresh(rule) + return _rule_to_out(db, rule) + + +@router.patch("/rules/{rule_id}", response_model=RuleOut) +def patch_rule( + rule_id: int, + payload: RulePatch, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> RuleOut: + rule = db.get(RoutingRule, rule_id) + if rule is 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']} не найдена") + db.execute(delete(RoutingRuleAction).where(RoutingRuleAction.rule_id == rule.id)) + for act in actions: + db.add(RoutingRuleAction(rule_id=rule.id, **act)) + for field, value in data.items(): + setattr(rule, field, value) + db.commit() + db.refresh(rule) + return _rule_to_out(db, rule) + + +@router.delete("/rules/{rule_id}", status_code=204) +def delete_rule( + rule_id: int, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> None: + if db.get(RoutingRule, rule_id) is 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)) + db.commit() + + +# --- events / deliveries (чтение) --- + +def _target_names(db: Session) -> dict[int, str]: + return {t.id: t.name for t in db.execute(select(ChannelTarget)).scalars().all()} + + +def _deliveries_out(db: Session, event_ids) -> dict[uuid.UUID, list[DeliveryOut]]: + if not event_ids: + return {} + names = _target_names(db) + rows = ( + db.execute(select(Delivery).where(Delivery.event_id.in_(event_ids)).order_by(Delivery.id)) + .scalars() + .all() + ) + out: dict[uuid.UUID, list[DeliveryOut]] = {} + for d in rows: + out.setdefault(d.event_id, []).append( + DeliveryOut( + id=d.id, event_id=d.event_id, rule_id=d.rule_id, channel=d.channel, + target_id=d.channel_target_id, target_name=names.get(d.channel_target_id), + status=d.status, attempts=d.attempts, last_error=d.last_error, + next_retry_at=d.next_retry_at, rendered_message=d.rendered_message, + created_at=d.created_at, delivered_at=d.delivered_at, + ) + ) + return out + + +@router.get("/events", response_model=list[AdminEventOut]) +def list_events( + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), + limit: int = Query(50, ge=1, le=200), + status_filter: str | None = Query(None, alias="status"), + source_id: int | None = None, +) -> list[AdminEventOut]: + q = select(Event).order_by(Event.created_at.desc()).limit(limit) + if status_filter: + q = q.where(Event.status == status_filter) + if source_id: + q = q.where(Event.source_id == source_id) + events = db.execute(q).scalars().all() + names = {s.id: s.name for s in db.execute(select(Source)).scalars().all()} + group = _deliveries_out(db, [e.id for e in events]) + return [ + AdminEventOut( + id=e.id, source_id=e.source_id, source_name=names.get(e.source_id, "?"), + subject=e.subject, action=e.action, priority=e.priority, status=e.status, + payload=dict(e.payload or {}), dedup_key=e.dedup_key, created_at=e.created_at, + expires_at=e.expires_at, deliveries=group.get(e.id, []), + ) + for e in events + ] + + +@router.get("/events/{event_id}", response_model=AdminEventOut) +def get_event( + event_id: uuid.UUID, + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), +) -> AdminEventOut: + event = db.get(Event, event_id) + if event is None: + raise HTTPException(status_code=404, detail="Событие не найдено") + source = db.get(Source, event.source_id) + group = _deliveries_out(db, [event.id]) + return AdminEventOut( + id=event.id, source_id=event.source_id, source_name=source.name if source else "?", + subject=event.subject, action=event.action, priority=event.priority, status=event.status, + payload=dict(event.payload or {}), dedup_key=event.dedup_key, created_at=event.created_at, + expires_at=event.expires_at, deliveries=group.get(event.id, []), + ) + + +@router.get("/deliveries", response_model=list[DeliveryOut]) +def list_deliveries( + user: AuthenticatedUser = Depends(require_admin), + db: Session = Depends(get_db), + limit: int = Query(100, ge=1, le=500), + status_filter: str | None = Query(None, alias="status"), + channel: str | None = None, +) -> list[DeliveryOut]: + q = select(Delivery).order_by(Delivery.created_at.desc(), Delivery.id.desc()).limit(limit) + if status_filter: + q = q.where(Delivery.status == status_filter) + if channel: + q = q.where(Delivery.channel == channel) + rows = db.execute(q).scalars().all() + names = _target_names(db) + return [ + DeliveryOut( + id=d.id, event_id=d.event_id, rule_id=d.rule_id, channel=d.channel, + target_id=d.channel_target_id, target_name=names.get(d.channel_target_id), + status=d.status, attempts=d.attempts, last_error=d.last_error, + next_retry_at=d.next_retry_at, rendered_message=d.rendered_message, + created_at=d.created_at, delivered_at=d.delivered_at, + ) + for d in rows + ] \ No newline at end of file diff --git a/app/api/admin_schemas.py b/app/api/admin_schemas.py new file mode 100644 index 0000000..76e16f2 --- /dev/null +++ b/app/api/admin_schemas.py @@ -0,0 +1,148 @@ +"""Pydantic-схемы админ-API (/api/v1/admin, гейт require_admin).""" + +import uuid +from datetime import datetime + +from pydantic import BaseModel, Field + +SlugPattern = r"^[a-z0-9]([a-z0-9._-]*[a-z0-9])?$" + + +# --- sources / keys --- + +class SourceIn(BaseModel): + name: str = Field(min_length=2, max_length=64, pattern=SlugPattern) + label: str | None = Field(None, max_length=120) + description: str | None = None + + +class SourceOut(BaseModel): + id: int + name: str + label: str | None + description: str | None + created_at: datetime + keys_count: int = 0 + types_count: int = 0 + + +class KeyIn(BaseModel): + name: str = Field(min_length=1, max_length=120) + + +class KeyOut(BaseModel): + id: int + name: str + token_hint: str + created_at: datetime + last_used_at: datetime | None + revoked_at: datetime | None + + +class KeyCreated(KeyOut): + """Токен печатается/показывается МИНУС один раз — только в этом ответе.""" + + token: str + + +# --- types --- + +class TypeIn(BaseModel): + source_id: int + subject: str = Field(min_length=1, max_length=64, pattern=SlugPattern) + action: str = Field(min_length=1, max_length=64, pattern=SlugPattern) + payload_schema: dict | None = None + description: str | None = None + + +class TypeOut(TypeIn): + id: int + source_name: str + created_at: datetime + + +# --- targets --- + +class TargetIn(BaseModel): + channel: str = Field(pattern=r"^(telegram|email|s2s|internal_log)$") + name: str = Field(min_length=1, max_length=120) + description: str | None = None + config: dict = {} + + +class TargetOut(TargetIn): + id: int + enabled: bool + created_at: datetime + + +# --- rules --- + +class RuleActionIn(BaseModel): + channel: str = Field(pattern=r"^(telegram|email|s2s|internal_log)$") + target_id: int | None = None + template: str | None = None + + +class RuleActionOut(RuleActionIn): + id: int + target_name: str | None = None + + +class RuleIn(BaseModel): + name: str = Field(min_length=1, max_length=120) + conditions: dict = {} + template: str | None = None + throttle_seconds: int | None = Field(None, ge=1) + weight: int = 0 + actions: list[RuleActionIn] = Field(min_length=1) + + +class RuleOut(RuleIn): + id: int + enabled: bool + created_at: datetime + actions: list[RuleActionOut] + + +class RulePatch(BaseModel): + name: str | None = Field(None, max_length=120) + enabled: bool | None = None + template: str | None = None + throttle_seconds: int | None = Field(None, ge=1) + weight: int | None = None + conditions: dict | None = None + actions: list[RuleActionIn] | None = None + + +# --- events / deliveries --- + +class DeliveryOut(BaseModel): + id: int + event_id: uuid.UUID + rule_id: int | None + channel: str + target_id: int | None + target_name: str | None = None + status: str + attempts: int + last_error: str | None + next_retry_at: datetime | None + rendered_message: str | None + created_at: datetime + delivered_at: datetime | None + + +class AdminEventOut(BaseModel): + id: uuid.UUID + source_id: int + source_name: str + subject: str + action: str + priority: str + status: str + payload: dict + dedup_key: str | None + created_at: datetime + expires_at: datetime | None + deliveries: list[DeliveryOut] \ No newline at end of file diff --git a/app/main.py b/app/main.py index 2280e35..26c84a9 100644 --- a/app/main.py +++ b/app/main.py @@ -6,6 +6,7 @@ from fastapi.responses import FileResponse from app.api import events_routes +from app.api import admin_routes from app.api.auth_routes import router as auth_router from app.api.routes import router as api_router from app.config import get_settings @@ -28,6 +29,7 @@ app.include_router(api_router) app.include_router(auth_router) app.include_router(events_routes.router) + app.include_router(admin_routes.router) # Собранная SPA: в Docker-образе кладётся в spa_static/. Если собранной # статики нет (локальный запуск api без фронта) — просто пропускаем.