154 lines
5.4 KiB
Python
154 lines
5.4 KiB
Python
"""services 路由:MCP 服务注册 + per-service API Key 管理"""
|
||
|
||
import hashlib
|
||
import secrets
|
||
|
||
from fastapi import APIRouter, Depends, HTTPException, Query, status
|
||
from pydantic import BaseModel
|
||
|
||
from ..core.db import get_pool
|
||
from ..core.deps import current_admin
|
||
from ..core.redis import cache_delete
|
||
|
||
router = APIRouter(prefix="/api/services", tags=["services"])
|
||
|
||
|
||
class ServiceCreate(BaseModel):
|
||
service_name: str # erp / crm / ...
|
||
service_url: str # MCP 服务地址(如 http://10.100.154.100:8001/mcp)
|
||
description: str | None = None
|
||
|
||
|
||
def _row_to_dict(row) -> dict:
|
||
return {
|
||
"service_id": row["service_id"],
|
||
"service_name": row["service_name"],
|
||
"service_url": row["service_url"],
|
||
"api_key": row["api_key"],
|
||
"description": row["description"],
|
||
"status": row["status"],
|
||
"created_at": row["created_at"].isoformat() if row["created_at"] else None,
|
||
"created_by": row["created_by"],
|
||
"revoked_at": row["revoked_at"].isoformat() if row["revoked_at"] else None,
|
||
"last_used_at": row["last_used_at"].isoformat() if row["last_used_at"] else None,
|
||
}
|
||
|
||
|
||
@router.get("")
|
||
async def list_services(
|
||
status_filter: str | None = Query(None, alias="status"),
|
||
admin: dict = Depends(current_admin),
|
||
):
|
||
pool = await get_pool()
|
||
query = "SELECT * FROM mcp_service WHERE 1=1"
|
||
params: list = []
|
||
if status_filter:
|
||
query += f" AND status = ${len(params)+1}"
|
||
params.append(status_filter)
|
||
query += " ORDER BY service_id DESC"
|
||
rows = await pool.fetch(query, *params)
|
||
return {"total": len(rows), "services": [_row_to_dict(r) for r in rows]}
|
||
|
||
|
||
@router.post("", status_code=status.HTTP_201_CREATED)
|
||
async def register_service(req: ServiceCreate, admin: dict = Depends(current_admin)):
|
||
pool = await get_pool()
|
||
|
||
# 检查是否已存在同名 active 服务
|
||
existing = await pool.fetchrow(
|
||
"SELECT service_id FROM mcp_service WHERE service_name = $1 AND status = 'active'",
|
||
req.service_name,
|
||
)
|
||
if existing:
|
||
raise HTTPException(400, f"服务 {req.service_name} 已存在且处于 active 状态")
|
||
|
||
# MCP 服务地址唯一性校验
|
||
url_existing = await pool.fetchval(
|
||
"SELECT 1 FROM mcp_service WHERE service_url = $1",
|
||
req.service_url.strip(),
|
||
)
|
||
if url_existing:
|
||
raise HTTPException(400, f"MCP 服务地址 {req.service_url} 已被其他服务使用")
|
||
|
||
# 生成 API Key
|
||
plain = secrets.token_urlsafe(32)
|
||
api_key_hash = hashlib.sha256(plain.encode()).hexdigest()
|
||
|
||
row = await pool.fetchrow(
|
||
"""INSERT INTO mcp_service (service_name, service_url, api_key, api_key_hash, description, created_by)
|
||
VALUES ($1, $2, $3, $4, $5, $6)
|
||
RETURNING service_id, service_name, service_url, api_key, description, created_at, created_by""",
|
||
req.service_name, req.service_url.strip(), plain, api_key_hash, req.description, admin.get("username", "admin"),
|
||
)
|
||
|
||
return {
|
||
"api_key": plain,
|
||
"service_id": row["service_id"],
|
||
"service_name": row["service_name"],
|
||
"message": "请保存此 API Key,配置到 MCP 服务的 MCP_AUTH_API_KEY 环境变量",
|
||
}
|
||
|
||
|
||
@router.patch("/{service_id}/revoke")
|
||
async def revoke_service(
|
||
service_id: int,
|
||
admin: dict = Depends(current_admin),
|
||
):
|
||
pool = await get_pool()
|
||
row = await pool.fetchrow(
|
||
"SELECT service_id, status, api_key_hash FROM mcp_service WHERE service_id = $1", service_id
|
||
)
|
||
if row is None:
|
||
raise HTTPException(404, "服务不存在")
|
||
if row["status"] == "revoked":
|
||
raise HTTPException(400, "服务已吊销")
|
||
|
||
await pool.execute(
|
||
"UPDATE mcp_service SET status = 'revoked', revoked_at = NOW() WHERE service_id = $1",
|
||
service_id,
|
||
)
|
||
await cache_delete(f"mcp:svc:{row['api_key_hash']}")
|
||
return {"success": True, "service_id": service_id, "status": "revoked"}
|
||
|
||
|
||
@router.patch("/{service_id}/enable")
|
||
async def enable_service(
|
||
service_id: int,
|
||
admin: dict = Depends(current_admin),
|
||
):
|
||
"""启用服务:将已吊销的服务恢复为 active。"""
|
||
pool = await get_pool()
|
||
row = await pool.fetchrow(
|
||
"SELECT service_id, status, api_key_hash FROM mcp_service WHERE service_id = $1", service_id
|
||
)
|
||
if row is None:
|
||
raise HTTPException(404, "服务不存在")
|
||
if row["status"] == "active":
|
||
raise HTTPException(400, "服务已是启用状态")
|
||
|
||
await pool.execute(
|
||
"UPDATE mcp_service SET status = 'active', revoked_at = NULL WHERE service_id = $1",
|
||
service_id,
|
||
)
|
||
await cache_delete(f"mcp:svc:{row['api_key_hash']}")
|
||
return {"success": True, "service_id": service_id, "status": "active"}
|
||
|
||
|
||
@router.delete("/{service_id}")
|
||
async def delete_service(
|
||
service_id: int,
|
||
admin: dict = Depends(current_admin),
|
||
):
|
||
pool = await get_pool()
|
||
row = await pool.fetchrow(
|
||
"SELECT service_id, status, api_key_hash FROM mcp_service WHERE service_id = $1", service_id
|
||
)
|
||
if row is None:
|
||
raise HTTPException(404, "服务不存在")
|
||
if row["status"] != "revoked":
|
||
raise HTTPException(400, "仅允许删除已吊销的服务,请先调用吊销端点")
|
||
|
||
await pool.execute("DELETE FROM mcp_service WHERE service_id = $1", service_id)
|
||
await cache_delete(f"mcp:svc:{row['api_key_hash']}")
|
||
return {"success": True, "service_id": service_id}
|