1. 移除多个导出接口的显式response_model声明 2. 调整access_rule和test_case的创建接口位置,修复静态路径冲突 3. 优化适配器配置校验的异常处理逻辑 4. 重构集成路由的查询逻辑,统一使用get_integration_or_raise 5. 新增channels路由组下的capability、reports、dashboard、webhook、wizard、doctor、directory、session共8个子路由模块 6. 注册channels_router到全局路由列表
718 lines
26 KiB
Python
718 lines
26 KiB
Python
"""外部系统限界上下文 - System 子域 Router。
|
||
|
||
挂载到聚合 router 的 /systems 前缀下,覆盖系统 CRUD / 克隆 / 导出 /
|
||
状态 / 健康 / 指标用例。Request Schema 与 Input DTO 不共享类,Router 内
|
||
显式构造 DTO,操作人字段由 current_user.uid 填充。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
from datetime import UTC, datetime
|
||
from typing import Any
|
||
|
||
from fastapi import APIRouter, Depends, Query
|
||
from fastapi.responses import StreamingResponse
|
||
from pydantic import BaseModel, ConfigDict, Field
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
from yuxi.external_systems.infrastructure.container import create_use_cases_from_db
|
||
from yuxi.external_systems.use_cases.dto.system import (
|
||
BatchCloneSystemsInput,
|
||
BulkEnableSystemsInput,
|
||
CloneSystemInput,
|
||
CreateSystemInput,
|
||
DeleteSystemInput,
|
||
ExportSystemsInput,
|
||
GetCircuitBreakerStatusInput,
|
||
GetGovernanceConfigInput,
|
||
GetSystemHealthInput,
|
||
GetSystemImpactAnalysisInput,
|
||
GetSystemInput,
|
||
GetSystemMetricsInput,
|
||
GetSystemResourceSummaryInput,
|
||
GetSystemStatsInput,
|
||
ListSystemsInput,
|
||
ResetCircuitBreakerInput,
|
||
SetSystemEnabledInput,
|
||
UpdateCircuitBreakerInput,
|
||
UpdateObservabilityInput,
|
||
UpdatePoolConfigInput,
|
||
UpdateRateLimitInput,
|
||
UpdateRetryPolicyInput,
|
||
UpdateSecretRefsInput,
|
||
UpdateSystemInput,
|
||
ValidateSystemConfigInput,
|
||
)
|
||
from yuxi.storage.postgres.models_business import User
|
||
|
||
from server.utils.auth_middleware import get_admin_user, get_db, get_required_user
|
||
|
||
system_router = APIRouter(prefix="/systems", tags=["external-systems-system"])
|
||
|
||
|
||
# =============================================================================
|
||
# === Request Schemas(与 Input DTO 不共享类) ===
|
||
# =============================================================================
|
||
|
||
|
||
class CreateSystemRequest(BaseModel):
|
||
"""创建系统请求体。字段对齐 ``CreateSystemInput``(不含 created_by)。
|
||
|
||
字段长度约束对齐 ``ExternalSystem`` ORM 列定义,在边界层拦截非法输入:
|
||
``category`` 对齐 String(64)、``adapter_type``/``auth_type`` 对齐 String(32)、
|
||
``source_type`` 对齐 String(64)、``timeout`` 对齐工具子域取值范围。
|
||
"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
slug: str = Field(..., min_length=1, max_length=128, pattern=r"^[a-zA-Z_][a-zA-Z0-9_-]{0,127}$")
|
||
name: str = Field(..., min_length=1, max_length=128)
|
||
description: str = ""
|
||
category: str = Field(default="default", max_length=64)
|
||
adapter_type: str = Field(default="http", max_length=32)
|
||
source_type: str | None = Field(default=None, max_length=64)
|
||
enabled: bool = True
|
||
connection_config: dict[str, Any] = Field(default_factory=dict)
|
||
auth_type: str = Field(default="none", max_length=32)
|
||
auth_config: dict[str, Any] | None = None
|
||
secret_refs: dict[str, Any] = Field(default_factory=dict)
|
||
rate_limit: dict[str, Any] = Field(default_factory=dict)
|
||
circuit_breaker: dict[str, Any] = Field(default_factory=dict)
|
||
pool_config: dict[str, Any] = Field(default_factory=dict)
|
||
observability: dict[str, Any] = Field(default_factory=dict)
|
||
timeout: int = Field(default=30, ge=1, le=300)
|
||
retry_policy: dict[str, Any] = Field(default_factory=dict)
|
||
|
||
|
||
class UpdateSystemRequest(BaseModel):
|
||
"""更新系统请求体。字段对齐 ``UpdateSystemInput``(不含 id 与 updated_by)。
|
||
|
||
字段长度约束对齐 ``ExternalSystem`` ORM 列定义,仅透传客户端显式设置的字段
|
||
(通过 ``exclude_unset=True``),未设置字段保持 ``None`` 以保留部分更新语义。
|
||
"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
name: str | None = Field(default=None, min_length=1, max_length=128)
|
||
description: str | None = None
|
||
category: str | None = Field(default=None, max_length=64)
|
||
adapter_type: str | None = Field(default=None, max_length=32)
|
||
source_type: str | None = Field(default=None, max_length=64)
|
||
enabled: bool | None = None
|
||
connection_config: dict[str, Any] | None = None
|
||
auth_type: str | None = Field(default=None, max_length=32)
|
||
auth_config: dict[str, Any] | None = None
|
||
secret_refs: dict[str, Any] | None = None
|
||
rate_limit: dict[str, Any] | None = None
|
||
circuit_breaker: dict[str, Any] | None = None
|
||
pool_config: dict[str, Any] | None = None
|
||
observability: dict[str, Any] | None = None
|
||
timeout: int | None = Field(default=None, ge=1, le=300)
|
||
retry_policy: dict[str, Any] | None = None
|
||
|
||
|
||
class CloneSystemRequest(BaseModel):
|
||
"""克隆系统请求体。字段对齐 ``CloneSystemInput``(不含 id 与 created_by)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
new_slug: str = Field(..., min_length=1, max_length=128, pattern=r"^[a-zA-Z_][a-zA-Z0-9_-]{0,127}$")
|
||
new_name: str | None = None
|
||
clone_tools: bool = True
|
||
clone_environments: bool = True
|
||
|
||
|
||
class SetSystemEnabledRequest(BaseModel):
|
||
"""设置系统启用状态请求体。字段对齐 ``SetSystemEnabledInput``(不含 id 与 updated_by)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
enabled: bool
|
||
|
||
|
||
class UpdateRateLimitRequest(BaseModel):
|
||
"""更新限流配置请求体。字段对齐 ``UpdateRateLimitInput``(不含 system_id 与 updated_by)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
rate_limit: dict[str, Any]
|
||
|
||
|
||
class UpdateCircuitBreakerRequest(BaseModel):
|
||
"""更新熔断器配置请求体。字段对齐 ``UpdateCircuitBreakerInput``(不含 system_id 与 updated_by)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
circuit_breaker: dict[str, Any]
|
||
|
||
|
||
class UpdatePoolConfigRequest(BaseModel):
|
||
"""更新连接池配置请求体。字段对齐 ``UpdatePoolConfigInput``(不含 system_id 与 updated_by)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
pool_config: dict[str, Any]
|
||
|
||
|
||
class UpdateObservabilityRequest(BaseModel):
|
||
"""更新可观测性配置请求体。字段对齐 ``UpdateObservabilityInput``(不含 system_id 与 updated_by)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
observability: dict[str, Any]
|
||
|
||
|
||
class UpdateRetryPolicyRequest(BaseModel):
|
||
"""更新重试策略请求体。字段对齐 ``UpdateRetryPolicyInput``(不含 system_id 与 updated_by)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
retry_policy: dict[str, Any]
|
||
|
||
|
||
class UpdateSecretRefsRequest(BaseModel):
|
||
"""更新密钥引用请求体。字段对齐 ``UpdateSecretRefsInput``(不含 system_id 与 updated_by)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
secret_refs: dict[str, Any]
|
||
|
||
|
||
class BatchEnabledRequest(BaseModel):
|
||
"""批量启用/禁用请求体。字段对齐 ``BulkEnableSystemsInput``(不含 updated_by)。"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
ids: list[int] = Field(..., min_length=1, max_length=100)
|
||
enabled: bool
|
||
|
||
|
||
class BatchCloneRequest(BaseModel):
|
||
"""批量克隆请求体。字段对齐 ``BatchCloneSystemsInput``(不含 created_by)。
|
||
|
||
``slug_suffix`` 长度上限 64,为 ``source.slug``(最长 128)预留拼接余量,
|
||
拼接后的 ``new_slug`` 由 service 层 ``validate_slug`` 校验总长与格式。
|
||
"""
|
||
|
||
model_config = ConfigDict(frozen=True)
|
||
|
||
source_system_ids: list[int] = Field(..., min_length=1, max_length=100)
|
||
slug_suffix: str = Field(..., min_length=1, max_length=64)
|
||
override_config: dict[str, Any] | None = None
|
||
|
||
|
||
# =============================================================================
|
||
# === 静态路径端点(必须在 /{system_id} 之前声明) ===
|
||
# =============================================================================
|
||
|
||
|
||
@system_router.get("/export")
|
||
async def export_systems(
|
||
ids: list[int] | None = Query(None),
|
||
category: str | None = Query(None),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> StreamingResponse:
|
||
"""导出系统列表为 JSON 字节流。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = ExportSystemsInput(
|
||
ids=ids,
|
||
category=category,
|
||
)
|
||
output = await use_cases.system_service.export_systems(input_dto)
|
||
|
||
payload = json.dumps(
|
||
[item.model_dump() for item in output.items],
|
||
ensure_ascii=False,
|
||
default=str,
|
||
).encode("utf-8")
|
||
timestamp = datetime.now(UTC).strftime("%Y%m%d%H%M%S")
|
||
|
||
async def _stream():
|
||
yield payload
|
||
|
||
return StreamingResponse(
|
||
_stream(),
|
||
media_type="application/octet-stream",
|
||
headers={"Content-Disposition": f"attachment; filename=systems_{timestamp}.json"},
|
||
)
|
||
|
||
|
||
@system_router.get("", response_model=dict)
|
||
async def list_systems(
|
||
limit: int = Query(100, ge=1, le=500),
|
||
offset: int = Query(0, ge=0),
|
||
category: str | None = Query(None),
|
||
adapter_type: str | None = Query(None),
|
||
source_type: str | None = Query(None),
|
||
enabled: bool | None = Query(None),
|
||
keyword: str | None = Query(None),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""分页列出外部系统。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = ListSystemsInput(
|
||
page=offset // limit + 1,
|
||
page_size=limit,
|
||
category=category,
|
||
adapter_type=adapter_type,
|
||
source_type=source_type,
|
||
enabled=enabled,
|
||
keyword=keyword,
|
||
)
|
||
output = await use_cases.system_service.list_systems(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.post("", response_model=dict)
|
||
async def create_system(
|
||
body: CreateSystemRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""创建外部系统。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = CreateSystemInput(**body.model_dump(), created_by=current_user.uid)
|
||
output = await use_cases.system_service.create_system(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/stats", response_model=dict)
|
||
async def get_system_stats(
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""获取系统聚合统计概览。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetSystemStatsInput(requested_by=current_user.uid)
|
||
output = await use_cases.system_service.get_system_stats(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.post("/batch-enabled", response_model=dict)
|
||
async def batch_set_enabled(
|
||
body: BatchEnabledRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""批量启用或禁用外部系统。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = BulkEnableSystemsInput(
|
||
ids=body.ids,
|
||
enabled=body.enabled,
|
||
updated_by=current_user.uid,
|
||
)
|
||
output = await use_cases.system_service.batch_set_enabled(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.post("/batch-clone", response_model=dict)
|
||
async def batch_clone_systems(
|
||
body: BatchCloneRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""批量克隆外部系统配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = BatchCloneSystemsInput(
|
||
source_system_ids=body.source_system_ids,
|
||
slug_suffix=body.slug_suffix,
|
||
override_config=body.override_config,
|
||
created_by=current_user.uid,
|
||
)
|
||
output = await use_cases.system_service.batch_clone_systems(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/circuit-breakers/status", response_model=dict)
|
||
async def get_circuit_breaker_status(
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""获取全局熔断器状态总览。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetCircuitBreakerStatusInput(requested_by=current_user.uid)
|
||
output = await use_cases.system_service.get_circuit_breaker_status(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
# =============================================================================
|
||
# === 动态路径端点 /{system_id} ===
|
||
# =============================================================================
|
||
|
||
|
||
@system_router.get("/{system_id}", response_model=dict)
|
||
async def get_system(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""获取系统详情。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetSystemInput(id=system_id)
|
||
output = await use_cases.system_service.get_system(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.put("/{system_id}", response_model=dict)
|
||
async def update_system(
|
||
system_id: int,
|
||
body: UpdateSystemRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""更新外部系统。仅透传客户端显式设置的字段,保留部分更新语义。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = UpdateSystemInput(
|
||
id=system_id,
|
||
updated_by=current_user.uid,
|
||
**body.model_dump(exclude_unset=True),
|
||
)
|
||
output = await use_cases.system_service.update_system(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.delete("/{system_id}", response_model=dict)
|
||
async def delete_system(
|
||
system_id: int,
|
||
cascade: bool = Query(False),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""删除外部系统。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = DeleteSystemInput(id=system_id, cascade=cascade)
|
||
output = await use_cases.system_service.delete_system(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.patch("/{system_id}/enabled", response_model=dict)
|
||
async def set_system_enabled(
|
||
system_id: int,
|
||
body: SetSystemEnabledRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""设置系统启用/禁用状态。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = SetSystemEnabledInput(
|
||
id=system_id,
|
||
enabled=body.enabled,
|
||
updated_by=current_user.uid,
|
||
)
|
||
output = await use_cases.system_service.set_system_enabled(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.post("/{system_id}/clone", response_model=dict)
|
||
async def clone_system(
|
||
system_id: int,
|
||
body: CloneSystemRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""克隆外部系统。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = CloneSystemInput(
|
||
**body.model_dump(),
|
||
id=system_id,
|
||
created_by=current_user.uid,
|
||
)
|
||
output = await use_cases.system_service.clone_system(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/{system_id}/health", response_model=dict, deprecated=True)
|
||
async def get_system_health(
|
||
system_id: int,
|
||
env_key: str | None = Query(None),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""获取系统健康状态(已废弃)。
|
||
|
||
.. deprecated::
|
||
请使用 ``GET /system/external-systems/health-checks/by-system/{system_id}/latest`` 替代。
|
||
新端点响应字段更完整(含 id 与审计字段),本端点保留用于前端迁移期间的向后兼容。
|
||
"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetSystemHealthInput(id=system_id, env_key=env_key)
|
||
output = await use_cases.system_service.get_system_health(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/{system_id}/metrics", response_model=dict, deprecated=True)
|
||
async def get_system_metrics(
|
||
system_id: int,
|
||
env_key: str | None = Query(None),
|
||
start_at: datetime | None = Query(None, description="ISO8601 开始时间"),
|
||
end_at: datetime | None = Query(None, description="ISO8601 结束时间"),
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""获取系统调用指标(已废弃)。
|
||
|
||
.. deprecated::
|
||
由 ``GET /system/external-systems/metrics/aggregate`` 替代,
|
||
新端点响应字段对齐 ORM aggregate 返回,支持 ``adapter_type`` 过滤。
|
||
前端迁移完成后将移除本端点。
|
||
"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetSystemMetricsInput(
|
||
id=system_id,
|
||
env_key=env_key,
|
||
start_at=start_at,
|
||
end_at=end_at,
|
||
)
|
||
output = await use_cases.system_service.get_system_metrics(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
# =============================================================================
|
||
# === 运维治理配置管理 ===
|
||
# =============================================================================
|
||
|
||
|
||
@system_router.get("/{system_id}/rate-limit", response_model=dict)
|
||
async def get_rate_limit(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""查询系统限流配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetGovernanceConfigInput(system_id=system_id, requested_by=current_user.uid)
|
||
output = await use_cases.system_service.get_rate_limit(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.put("/{system_id}/rate-limit", response_model=dict)
|
||
async def update_rate_limit(
|
||
system_id: int,
|
||
body: UpdateRateLimitRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""更新系统限流配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = UpdateRateLimitInput(
|
||
system_id=system_id,
|
||
rate_limit=body.rate_limit,
|
||
updated_by=current_user.uid,
|
||
)
|
||
output = await use_cases.system_service.update_rate_limit(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/{system_id}/circuit-breaker", response_model=dict)
|
||
async def get_circuit_breaker(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""查询系统熔断器配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetGovernanceConfigInput(system_id=system_id, requested_by=current_user.uid)
|
||
output = await use_cases.system_service.get_circuit_breaker(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.put("/{system_id}/circuit-breaker", response_model=dict)
|
||
async def update_circuit_breaker(
|
||
system_id: int,
|
||
body: UpdateCircuitBreakerRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""更新系统熔断器配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = UpdateCircuitBreakerInput(
|
||
system_id=system_id,
|
||
circuit_breaker=body.circuit_breaker,
|
||
updated_by=current_user.uid,
|
||
)
|
||
output = await use_cases.system_service.update_circuit_breaker(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.post("/{system_id}/circuit-breaker/reset", response_model=dict)
|
||
async def reset_circuit_breaker(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""重置系统熔断器运行时状态。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = ResetCircuitBreakerInput(system_id=system_id, reset_by=current_user.uid)
|
||
output = await use_cases.system_service.reset_circuit_breaker(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/{system_id}/connection-pool", response_model=dict)
|
||
async def get_connection_pool(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""查询系统连接池配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetGovernanceConfigInput(system_id=system_id, requested_by=current_user.uid)
|
||
output = await use_cases.system_service.get_pool_config(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.put("/{system_id}/connection-pool", response_model=dict)
|
||
async def update_connection_pool(
|
||
system_id: int,
|
||
body: UpdatePoolConfigRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""更新系统连接池配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = UpdatePoolConfigInput(
|
||
system_id=system_id,
|
||
pool_config=body.pool_config,
|
||
updated_by=current_user.uid,
|
||
)
|
||
output = await use_cases.system_service.update_pool_config(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/{system_id}/observability", response_model=dict)
|
||
async def get_observability(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""查询系统可观测性配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetGovernanceConfigInput(system_id=system_id, requested_by=current_user.uid)
|
||
output = await use_cases.system_service.get_observability(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.put("/{system_id}/observability", response_model=dict)
|
||
async def update_observability(
|
||
system_id: int,
|
||
body: UpdateObservabilityRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""更新系统可观测性配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = UpdateObservabilityInput(
|
||
system_id=system_id,
|
||
observability=body.observability,
|
||
updated_by=current_user.uid,
|
||
)
|
||
output = await use_cases.system_service.update_observability(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/{system_id}/retry-policy", response_model=dict)
|
||
async def get_retry_policy(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""查询系统重试策略配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetGovernanceConfigInput(system_id=system_id, requested_by=current_user.uid)
|
||
output = await use_cases.system_service.get_retry_policy(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.put("/{system_id}/retry-policy", response_model=dict)
|
||
async def update_retry_policy(
|
||
system_id: int,
|
||
body: UpdateRetryPolicyRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""更新系统重试策略配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = UpdateRetryPolicyInput(
|
||
system_id=system_id,
|
||
retry_policy=body.retry_policy,
|
||
updated_by=current_user.uid,
|
||
)
|
||
output = await use_cases.system_service.update_retry_policy(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/{system_id}/secret-refs", response_model=dict)
|
||
async def get_secret_refs(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""查询系统密钥引用配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetGovernanceConfigInput(system_id=system_id, requested_by=current_user.uid)
|
||
output = await use_cases.system_service.get_secret_refs(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.put("/{system_id}/secret-refs", response_model=dict)
|
||
async def update_secret_refs(
|
||
system_id: int,
|
||
body: UpdateSecretRefsRequest,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""更新系统密钥引用配置。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = UpdateSecretRefsInput(
|
||
system_id=system_id,
|
||
secret_refs=body.secret_refs,
|
||
updated_by=current_user.uid,
|
||
)
|
||
output = await use_cases.system_service.update_secret_refs(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
# =============================================================================
|
||
# === 配置验证与依赖分析 ===
|
||
# =============================================================================
|
||
|
||
|
||
@system_router.post("/{system_id}/validate-config", response_model=dict)
|
||
async def validate_system_config(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_admin_user),
|
||
) -> dict[str, Any]:
|
||
"""校验系统配置完整性。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = ValidateSystemConfigInput(system_id=system_id, validated_by=current_user.uid)
|
||
output = await use_cases.system_service.validate_system_config(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/{system_id}/resource-summary", response_model=dict)
|
||
async def get_system_resource_summary(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""获取系统资源汇总。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetSystemResourceSummaryInput(system_id=system_id, requested_by=current_user.uid)
|
||
output = await use_cases.system_service.get_system_resource_summary(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|
||
|
||
|
||
@system_router.get("/{system_id}/impact-analysis", response_model=dict)
|
||
async def get_system_impact_analysis(
|
||
system_id: int,
|
||
db: AsyncSession = Depends(get_db),
|
||
current_user: User = Depends(get_required_user),
|
||
) -> dict[str, Any]:
|
||
"""获取系统变更影响分析。"""
|
||
use_cases = create_use_cases_from_db(db)
|
||
input_dto = GetSystemImpactAnalysisInput(system_id=system_id, requested_by=current_user.uid)
|
||
output = await use_cases.system_service.get_system_impact_analysis(input_dto)
|
||
return {"success": True, "data": output.model_dump()}
|