Develop #16
@@ -25,9 +25,10 @@
|
||||
4. [调度 (`/api/v1/schedules/...` + `/api/v1/schedule-runs/...`)](#四调度)
|
||||
5. [数据资源 (`/api/v1/data-resources/...`)](#五数据资源)
|
||||
6. [管理后台 (`/api/v1/admin/...`)](#六管理后台)
|
||||
7. [Jupyter 路由 (Nginx `auth_request`)](#七jupyter-路由)
|
||||
8. [对象存储控制面 (`/internal/v1/...`,同进程 RPC)](#八对象存储控制面)
|
||||
9. [健康检查](#九健康检查)
|
||||
7. [系统管理 (`/api/v1/platform/...`)](#七系统管理-apiv1platform)
|
||||
8. [Jupyter 路由 (Nginx `auth_request`)](#八jupyter-路由)
|
||||
9. [对象存储控制面 (`/internal/v1/...`,同进程 RPC)](#九对象存储控制面)
|
||||
10. [健康检查](#十健康检查)
|
||||
|
||||
---
|
||||
|
||||
@@ -44,6 +45,8 @@
|
||||
- 缺失或过期 → HTTP `401`。
|
||||
- 有效但用户不在 workspace → HTTP `403`(由 Nginx `auth_request` 透传给客户端)。
|
||||
|
||||
> `/api/v1/auth/me` 与 `/api/v1/auth/login` 响应中的 `data.user` 对象额外携带 `is_system_admin: bool` 字段,派生自 `users.platform_role_id` 指向的角色 `role_code == 'admin'` 且用户状态为 `active`。前端据此决定是否渲染"系统管理"入口。详见 §七。
|
||||
|
||||
---
|
||||
|
||||
## 二、统一约定
|
||||
@@ -484,11 +487,139 @@ Base 前缀 `/api/v1/admin`。
|
||||
|
||||
> `PATCH` / `DELETE` 员工接口**不**涉及密码字段,也不返回密码相关信息。
|
||||
|
||||
## 七、Jupyter 路由
|
||||
## 七、系统管理 (`/api/v1/platform/...`)
|
||||
|
||||
平台级(跨 workspace)管理接口,用于管理 workspace 实体与 workspace 成员。
|
||||
所有端点要求调用者是**系统管理员**——其 `users.platform_role_id` 指向
|
||||
`role_code='admin'` 的角色行,且 `users.status == 'active'`。系统管理员判定
|
||||
通过 `GET /api/v1/auth/me` 响应中的 `data.user.is_system_admin` 字段(详见 §一)。
|
||||
|
||||
| 方法 | 路径 | 说明 |
|
||||
|---|---|---|
|
||||
| `GET` | `/api/v1/platform/workspaces` | 列 workspace(`active`/`archived`);已软删的过滤掉 |
|
||||
| `POST` | `/api/v1/platform/workspaces` | 创建 workspace(返回 201);创建者自动成为 admin 成员 |
|
||||
| `GET` | `/api/v1/platform/workspaces/{workspace_id}` | 单个 workspace(含已 disabled 的,用于恢复) |
|
||||
| `PATCH` | `/api/v1/platform/workspaces/{workspace_id}` | 改 workspace 字段;`status` 仅允许 `active`/`archived` |
|
||||
| `DELETE` | `/api/v1/platform/workspaces/{workspace_id}` | 软删 workspace;级联软删其成员 |
|
||||
| `GET` | `/api/v1/platform/workspaces/{workspace_id}/members` | 列成员 |
|
||||
| `POST` | `/api/v1/platform/workspaces/{workspace_id}/members` | 添加成员(返回 201) |
|
||||
| `PATCH` | `/api/v1/platform/workspaces/{workspace_id}/members/{user_id}` | 改成员角色/状态 |
|
||||
| `DELETE` | `/api/v1/platform/workspaces/{workspace_id}/members/{user_id}` | 软删成员 |
|
||||
|
||||
> **不变量**:
|
||||
> - 每个 workspace 必须始终保留至少一个 `admin` 角色的活跃成员;对最后 admin 做降级 / 停用 / 删除 → 409。
|
||||
> - 系统管理员不能通过 `DELETE .../members/{self}` 把自己移除(403)。唯一退出方式是 `DELETE /workspaces/{id}` 软删整个 workspace,后者会级联软删所有成员。
|
||||
> - 列表类接口静默 `pageSize=100` 上限,无客户端分页参数(YAGNI)。
|
||||
> - 跨 workspace 操作**不**需要 `?workspace_id=` query 参数,与 `/api/v1/admin/...`(workspace 内成员管理)不要混淆。
|
||||
|
||||
### 7.1 `POST /api/v1/platform/workspaces`
|
||||
|
||||
创建 workspace;创建者(当前系统管理员)自动成为该 workspace 的 `admin` 成员。
|
||||
|
||||
- **请求体字段**:
|
||||
|
||||
| 字段 | 类型 | 必填 | 限制 | 说明 |
|
||||
|---|---|---|---|---|
|
||||
| `workspace_code` | string | 是 | regex `^[a-z0-9-]{3,32}$`(类似 git repo 名) | 创建后冻结,不可改 |
|
||||
| `workspace_name` | string | 是 | 1~150 字符 | 显示名称 |
|
||||
| `quota_bytes` | int | 否 | ≥0,默认 `0` | 配额字节数,`0` 表示无配额 |
|
||||
| `description` | string | 否 | ≤1000 字符 | |
|
||||
|
||||
- **服务端自动生成字段**(不接收):
|
||||
- `workspace_id`(ULID)
|
||||
- `active_root_uri`(`s3://workspaces/{workspace_id}/`)
|
||||
- `status`(`"active"`)
|
||||
- `created_by`(当前管理员 `user_id`)
|
||||
- `created_at` / `updated_at`(DB 自动)
|
||||
|
||||
- **响应 201**:见下 §7.2 `WorkspacePayload`。
|
||||
|
||||
### 7.2 `GET /api/v1/platform/workspaces/{workspace_id}` / `WorkspacePayload`
|
||||
|
||||
- **响应 200**:
|
||||
```json
|
||||
{
|
||||
"request_id": "...",
|
||||
"data": {
|
||||
"workspace_id": "01HXY...",
|
||||
"workspace_code": "model-development",
|
||||
"workspace_name": "模型开发 Workspace",
|
||||
"active_root_uri": "s3://workspaces/01HXY.../",
|
||||
"quota_bytes": 0,
|
||||
"status": "active",
|
||||
"description": null,
|
||||
"created_by": "01HXY...",
|
||||
"created_at": "2026-08-04T12:00:00.000",
|
||||
"updated_at": null
|
||||
},
|
||||
"meta": {"count": ..., "page_size": 100}
|
||||
}
|
||||
```
|
||||
|
||||
### 7.3 `PATCH /api/v1/platform/workspaces/{workspace_id}`
|
||||
|
||||
部分更新。**不可改**:`workspace_id`、`workspace_code`、`active_root_uri`、`created_by`、时间戳、软删标记。
|
||||
|
||||
- **请求体字段**(全部可选):
|
||||
|
||||
| 字段 | 类型 | 限制 | 说明 |
|
||||
|---|---|---|---|
|
||||
| `workspace_name` | string | 1~150 | |
|
||||
| `quota_bytes` | int | ≥0 | |
|
||||
| `description` | string | ≤1000 | |
|
||||
| `status` | string | `active` \| `archived` | **不允许 `disabled`**——软删须走 DELETE |
|
||||
|
||||
- 错误:`status="disabled"` → 422;已 disabled 的 workspace → 409。
|
||||
|
||||
### 7.4 `DELETE /api/v1/platform/workspaces/{workspace_id}`
|
||||
|
||||
软删除。允许从 `active` 或 `archived` 状态调用。
|
||||
|
||||
- **副作用**:
|
||||
- 该 workspace 行:`status='disabled'`、`is_deleted=1`、`deleted_at=NOW()`
|
||||
- **级联**:所有未删除的 `workspace_members` 行同步 `is_deleted=1`、`deleted_at=NOW()`
|
||||
- 已 disabled 的 workspace 再删 → 409。
|
||||
|
||||
### 7.5 `POST /api/v1/platform/workspaces/{workspace_id}/members`
|
||||
|
||||
添加成员。
|
||||
|
||||
- **请求体字段**:
|
||||
|
||||
| 字段 | 类型 | 必填 | 限制 | 说明 |
|
||||
|---|---|---|---|---|
|
||||
| `user_id` | string | 是 | 26 字符 ULID | |
|
||||
| `role_code` | string | 是 | `admin` \| `developer` | **不可填 `system_admin`**(那是用户级身份,不是 workspace 角色) |
|
||||
|
||||
- 服务端默认 `member_status='active'`。
|
||||
- 用户不存在 → 404;用户已是该 workspace 成员 → 409。
|
||||
|
||||
### 7.6 `PATCH /api/v1/platform/workspaces/{workspace_id}/members/{user_id}`
|
||||
|
||||
修改成员的角色或状态。
|
||||
|
||||
- **请求体字段**(全部可选):
|
||||
|
||||
| 字段 | 类型 | 限制 | 说明 |
|
||||
|---|---|---|---|
|
||||
| `role_code` | string | `admin` \| `developer` | 降级最后 admin → 409 |
|
||||
| `member_status` | string | `active` \| `disabled` \| `locked` | 停用 / 锁定最后 admin → 409 |
|
||||
|
||||
### 7.7 `DELETE /api/v1/platform/workspaces/{workspace_id}/members/{user_id}`
|
||||
|
||||
软删除成员。
|
||||
|
||||
- **自我移除保护**:`user_id == 当前管理员 user_id` → 403 "系统管理员不能把自己从 workspace 移除;如需退出,请删除整个 workspace"
|
||||
- **末位 admin 保护**:若删除的是最后一个 `admin` 角色活跃成员 → 409
|
||||
- 不存在的成员 → 404
|
||||
|
||||
---
|
||||
|
||||
## 八、Jupyter 路由
|
||||
|
||||
> **本节是 Nginx 行为,不是直接 HTTP 端点**。前端**不要**直接调用。
|
||||
|
||||
### 7.1 浏览器 → 用户打开 notebook
|
||||
### 8.1 浏览器 → 用户打开 notebook
|
||||
|
||||
用户在前端点击某个 notebook,前端拼出 URL:
|
||||
```
|
||||
@@ -498,7 +629,7 @@ GET /jupyter/{workspace_id}/api/contents/{相对路径}.ipynb
|
||||
WS /jupyter/{workspace_id}/api/kernels/...
|
||||
```
|
||||
|
||||
### 7.2 Nginx `auth_request` 鉴权
|
||||
### 8.2 Nginx `auth_request` 鉴权
|
||||
|
||||
Nginx 收到上述请求后,**先**发一个内部子请求:
|
||||
```
|
||||
@@ -525,7 +656,7 @@ Nginx: auth_request_set 捕获这两个变量,proxy_pass 到子进程并注入
|
||||
浏览器收到响应,**自始至终未接触 Jupyter Token**
|
||||
```
|
||||
|
||||
### 7.3 鉴权失败码
|
||||
### 8.3 鉴权失败码
|
||||
|
||||
| 状态 | 触发条件 |
|
||||
|---|---|
|
||||
@@ -538,14 +669,14 @@ Nginx 把这些状态原样透传给浏览器,前端可在 `onerror` 里判断
|
||||
|
||||
---
|
||||
|
||||
## 八、对象存储控制面
|
||||
## 九、对象存储控制面
|
||||
|
||||
> 路径前缀 `/internal/v1/...`,**前端不要直接调用**。这是 backend 内部
|
||||
> 异步消息处理(Schedule worker)用的 RPC 端点,经 `StorageClient` HTTP
|
||||
> 客户端访问。Backend 通过 ASGI `auth_request_set` 路由转发,外部无法
|
||||
> 访问。
|
||||
|
||||
### 8.1 `POST /internal/v1/uploads`
|
||||
### 9.1 `POST /internal/v1/uploads`
|
||||
|
||||
创建上传会话。`Idempotency-Key` 必填,同 key + 同元数据 → 复用;同 key + 不同元数据 → 409。
|
||||
|
||||
@@ -564,15 +695,15 @@ Nginx 把这些状态原样透传给浏览器,前端可在 `onerror` 里判断
|
||||
|
||||
返回 `{upload_id, bucket_name, object_key, presigned_url, expires_in_seconds}`。
|
||||
|
||||
### 8.2 `POST /internal/v1/uploads/{upload_id}/complete`
|
||||
### 9.2 `POST /internal/v1/uploads/{upload_id}/complete`
|
||||
|
||||
完成上传。从 RustFS 读 HEAD → 校验 hash → 写 `StorageObjects` 行。
|
||||
|
||||
### 8.3 `POST /internal/v1/uploads/{upload_id}/abort`
|
||||
### 9.3 `POST /internal/v1/uploads/{upload_id}/abort`
|
||||
|
||||
主动放弃。释放 `UploadSessions` 行,对象不入库。
|
||||
|
||||
### 8.4 `POST /internal/v1/objects`
|
||||
### 9.4 `POST /internal/v1/objects`
|
||||
|
||||
**单步创建**(不走 presigned PUT,字节随请求体直传)。适用 < 100 KiB 对象。
|
||||
|
||||
@@ -590,15 +721,15 @@ Nginx 把这些状态原样透传给浏览器,前端可在 `onerror` 里判断
|
||||
}
|
||||
```
|
||||
|
||||
### 8.5 `POST /internal/v1/objects/{storage_object_id}/download-url`
|
||||
### 9.5 `POST /internal/v1/objects/{storage_object_id}/download-url`
|
||||
|
||||
生成 presigned GET URL。
|
||||
|
||||
### 8.6 `DELETE /internal/v1/objects/{storage_object_id}`
|
||||
### 9.6 `DELETE /internal/v1/objects/{storage_object_id}`
|
||||
|
||||
软删。`is_immutable == 1` 的对象拒绝删除。
|
||||
|
||||
### 8.7 usage_type → 桶路由(自动)
|
||||
### 9.7 usage_type → 桶路由(自动)
|
||||
|
||||
| usage_type | 实际桶(env var) | 默认桶名 |
|
||||
|---|---|---|
|
||||
@@ -610,7 +741,7 @@ Nginx 把这些状态原样透传给浏览器,前端可在 `onerror` 里判断
|
||||
|
||||
---
|
||||
|
||||
## 九、健康检查
|
||||
## 十、健康检查
|
||||
|
||||
| 方法 | 路径 | 用途 |
|
||||
|---|---|---|
|
||||
@@ -630,18 +761,20 @@ Nginx 把这些状态原样透传给浏览器,前端可在 `onerror` 里判断
|
||||
|---|---|---|
|
||||
| 400 | 参数错误 | Pydantic 校验失败 |
|
||||
| 401 | 未鉴权 | JWT 缺失/无效 |
|
||||
| 403 | 鉴权失败 | 非 workspace 成员 / `is_locked` 阻写 |
|
||||
| 403 | 鉴权失败 | 非 workspace 成员 / `is_locked` 阻写 / **非系统管理员访问 `/api/v1/platform/*`** / 系统管理员自我移除 workspace 成员 |
|
||||
| 404 | 不存在 | resource_id / script_id / schedule_id 找不到 |
|
||||
| 409 | 冲突 | DAG 无效 / 同 idempotency_key 不同元数据 / 目标已存在 / `is_immutable` 阻删 |
|
||||
| 409 | 冲突 | DAG 无效 / 同 idempotency_key 不同元数据 / 目标已存在 / `is_immutable` 阻删 / **workspace 末位 admin 保护** |
|
||||
| 412 | 条件失败 | `source_object_id` 与当前工作副本不一致 |
|
||||
| 413 | 太大 | 内容超过 100 MiB / 10 MiB |
|
||||
| 422 | 语义错误 | 文件名非法 / cron 表达式非法 / 路径逃逸 |
|
||||
| 422 | 语义错误 | 文件名非法 / cron 表达式非法 / 路径逃逸 / **`workspace_code` 不匹配 `^[a-z0-9-]{3,32}$` / `status="disabled"` 走 PATCH** |
|
||||
| 500 | 内部错误 | DB / 存储不可达 |
|
||||
|
||||
## 附录 B — 状态枚举
|
||||
|
||||
| 类型 | 取值 |
|
||||
|---|---|
|
||||
| `Workspaces.status` | `active` / `archived` / `disabled`(`disabled` 由 DELETE 设置,PATCH 不允许设) |
|
||||
| `WorkspaceMembers.member_status` | `active` / `disabled` / `locked` |
|
||||
| `StorageObjects.usage_type` | `data_resource` / `version_artifact` / `snapshot` / `run_log` / `run_result` / `working_copy` / `public_script` |
|
||||
| `StorageObjects.object_status` | `available` / `deleted` |
|
||||
| `StorageObjects.visibility` | `private` / `workspace` / `public` |
|
||||
|
||||
@@ -55,7 +55,12 @@ def _clear_session_cookie(response: Response) -> None:
|
||||
response.delete_cookie(key=COOKIE_NAME, path="/")
|
||||
|
||||
|
||||
def _user_payload(user: Users, role_code: str | None = None) -> dict[str, Any]:
|
||||
def _user_payload(
|
||||
user: Users,
|
||||
role_code: str | None = None,
|
||||
*,
|
||||
is_system_admin: bool = False,
|
||||
) -> dict[str, Any]:
|
||||
return {
|
||||
"user_id": user.user_id,
|
||||
"username": user.username,
|
||||
@@ -63,9 +68,30 @@ def _user_payload(user: Users, role_code: str | None = None) -> dict[str, Any]:
|
||||
"email": user.email,
|
||||
"status": user.status,
|
||||
"role_code": role_code,
|
||||
"is_system_admin": is_system_admin,
|
||||
}
|
||||
|
||||
|
||||
async def _resolve_is_system_admin(
|
||||
session: AsyncSession,
|
||||
user: Users,
|
||||
) -> bool:
|
||||
"""Return True iff the user holds a platform-scoped admin role.
|
||||
|
||||
The check is: ``Users.status == 'active'`` AND
|
||||
``Users.platform_role_id`` points to a ``Roles`` row whose
|
||||
``role_code == 'admin'``. Any other shape (no platform_role_id,
|
||||
disabled user, wrong role code) returns False — the frontend reads
|
||||
this to decide whether to show the system-admin entry point.
|
||||
"""
|
||||
if user.status != "active" or user.platform_role_id is None:
|
||||
return False
|
||||
platform_role = await session.scalar(
|
||||
select(Roles).where(Roles.role_id == user.platform_role_id)
|
||||
)
|
||||
return platform_role is not None and platform_role.role_code == "admin"
|
||||
|
||||
|
||||
def _workspace_payload(
|
||||
workspace: Workspaces,
|
||||
role: Roles,
|
||||
@@ -161,10 +187,12 @@ async def login(
|
||||
token = issue_jwt(user.user_id, ttl_seconds=COOKIE_TTL_SECONDS)
|
||||
_set_session_cookie(request, response, token)
|
||||
|
||||
is_system_admin = await _resolve_is_system_admin(session, user)
|
||||
|
||||
return {
|
||||
"request_id": new_ulid(),
|
||||
"data": {
|
||||
"user": _user_payload(user, user_role_code),
|
||||
"user": _user_payload(user, user_role_code, is_system_admin=is_system_admin),
|
||||
"workspaces": workspaces,
|
||||
"default_workspace_id": default_workspace_id,
|
||||
},
|
||||
@@ -238,10 +266,12 @@ async def me(
|
||||
if user_role_code is None and rows:
|
||||
user_role_code = rows[0][1].role_code
|
||||
|
||||
is_system_admin = await _resolve_is_system_admin(session, user)
|
||||
|
||||
return {
|
||||
"request_id": new_ulid(),
|
||||
"data": {
|
||||
"user": _user_payload(user, user_role_code),
|
||||
"user": _user_payload(user, user_role_code, is_system_admin=is_system_admin),
|
||||
"workspaces": workspaces,
|
||||
"default_workspace_id": default_workspace_id,
|
||||
},
|
||||
|
||||
@@ -13,6 +13,7 @@ from common.service_app import create_service_app
|
||||
from common.storage import RustFSObjectStore
|
||||
from backend.admin import router as admin_router
|
||||
from backend.auth import router as auth_router
|
||||
from backend.platform import router as platform_router
|
||||
from backend.jupyter import router as jupyter_router
|
||||
from backend.resources import router as resources_router
|
||||
from backend.runtime_client import RuntimeClient
|
||||
@@ -86,6 +87,7 @@ app.include_router(schedule_runs_router)
|
||||
app.include_router(schedules_router)
|
||||
app.include_router(scripts_router)
|
||||
app.include_router(admin_router)
|
||||
app.include_router(platform_router)
|
||||
|
||||
# Reuse the proven storage endpoints without running another FastAPI service.
|
||||
for route in storage_app.routes:
|
||||
|
||||
@@ -0,0 +1,607 @@
|
||||
"""System-admin (platform-scope) endpoints for workspace & membership management.
|
||||
|
||||
All routes under ``/api/v1/platform/*`` are gated by
|
||||
:func:`system_admin_context`, which requires the requester to hold a
|
||||
``Users.platform_role_id`` pointing to a ``Roles`` row whose
|
||||
``role_code == 'admin'``. Unlike ``backend.dependencies.request_context``,
|
||||
this dependency does NOT require an active workspace membership — system
|
||||
admins can manage workspaces before/without being a member of any.
|
||||
|
||||
Endpoints
|
||||
---------
|
||||
|
||||
Workspace CRUD::
|
||||
|
||||
GET /workspaces — list non-deleted workspaces
|
||||
POST /workspaces — create a new workspace
|
||||
GET /workspaces/{workspace_id} — single workspace (incl. disabled)
|
||||
PATCH /workspaces/{workspace_id} — update editable fields
|
||||
DELETE /workspaces/{workspace_id} — soft delete (cascades memberships)
|
||||
|
||||
Workspace membership CRUD::
|
||||
|
||||
GET /workspaces/{workspace_id}/members — list active members
|
||||
POST /workspaces/{workspace_id}/members — add a member
|
||||
PATCH /workspaces/{workspace_id}/members/{user_id} — update role/status
|
||||
DELETE /workspaces/{workspace_id}/members/{user_id} — remove a member
|
||||
|
||||
Invariants
|
||||
----------
|
||||
|
||||
* Every workspace must always retain at least one active ``admin`` member.
|
||||
* A system admin cannot remove their own workspace membership via
|
||||
``DELETE .../members/{self}``; the only escape is to delete the entire
|
||||
workspace, which cascades membership soft-deletion.
|
||||
* ``DELETE /workspaces/{id}`` is allowed from any non-disabled status and
|
||||
sets ``status='disabled'`` + ``is_deleted=1`` + ``deleted_at`` on the
|
||||
workspace and every one of its active memberships.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import datetime
|
||||
import re
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, Literal
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Request, status
|
||||
from pydantic import BaseModel, ConfigDict, Field
|
||||
from sqlalchemy import func, select, update
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from backend.dependencies import current_user, database_session
|
||||
from common.db.models import Roles, Users, WorkspaceMembers, Workspaces
|
||||
from common.ids import new_ulid
|
||||
|
||||
|
||||
router = APIRouter(prefix="/api/v1/platform", tags=["platform"])
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Constants
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
WORKSPACE_CODE_PATTERN = re.compile(r"^[a-z0-9-]{3,32}$")
|
||||
LIST_PAGE_SIZE = 100
|
||||
|
||||
WORKSPACE_EDITABLE_STATUS = ("active", "archived")
|
||||
MEMBER_ROLE_CODES = ("admin", "developer")
|
||||
MEMBER_STATUS_VALUES = ("active", "disabled", "locked")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Schemas
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class WorkspaceCreate(BaseModel):
|
||||
model_config = ConfigDict(extra="forbid")
|
||||
|
||||
workspace_code: str = Field(min_length=3, max_length=32)
|
||||
workspace_name: str = Field(min_length=1, max_length=150)
|
||||
quota_bytes: int = Field(default=0, ge=0)
|
||||
description: str | None = Field(default=None, max_length=1000)
|
||||
|
||||
|
||||
class WorkspaceUpdate(BaseModel):
|
||||
model_config = ConfigDict(extra="forbid")
|
||||
|
||||
workspace_name: str | None = Field(default=None, min_length=1, max_length=150)
|
||||
quota_bytes: int | None = Field(default=None, ge=0)
|
||||
description: str | None = Field(default=None, max_length=1000)
|
||||
# 'disabled' is rejected here on purpose — soft delete must go through DELETE.
|
||||
status: Literal["active", "archived"] | None = None
|
||||
|
||||
|
||||
class MemberCreate(BaseModel):
|
||||
model_config = ConfigDict(extra="forbid")
|
||||
|
||||
user_id: str = Field(min_length=26, max_length=26)
|
||||
role_code: Literal["admin", "developer"]
|
||||
|
||||
|
||||
class MemberUpdate(BaseModel):
|
||||
model_config = ConfigDict(extra="forbid")
|
||||
|
||||
role_code: Literal["admin", "developer"] | None = None
|
||||
member_status: Literal["active", "disabled", "locked"] | None = None
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# System-admin context dependency
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class SystemAdminContext:
|
||||
"""Resolved identity for a system-admin request.
|
||||
|
||||
Carries the request id, the authenticated user row, and the resolved
|
||||
``Roles`` row the user holds via ``Users.platform_role_id``. By
|
||||
construction the role's ``role_code`` is ``"admin"``.
|
||||
"""
|
||||
|
||||
request_id: str
|
||||
user: Users
|
||||
platform_role: Roles
|
||||
|
||||
|
||||
async def system_admin_context(
|
||||
request: Request,
|
||||
session: AsyncSession = Depends(database_session),
|
||||
) -> SystemAdminContext:
|
||||
"""Resolve the requester as a system admin.
|
||||
|
||||
Steps:
|
||||
1. Reuse :func:`backend.dependencies.current_user` to validate the JWT
|
||||
cookie and fetch the active ``Users`` row (raises 401 on failure).
|
||||
2. Require ``Users.platform_role_id`` to point to a row whose
|
||||
``role_code == 'admin'`` — anything else is 403.
|
||||
"""
|
||||
user = await current_user(request, session)
|
||||
if user.platform_role_id is None:
|
||||
raise HTTPException(
|
||||
status.HTTP_403_FORBIDDEN,
|
||||
"需要系统管理员权限",
|
||||
)
|
||||
platform_role = await session.scalar(
|
||||
select(Roles).where(Roles.role_id == user.platform_role_id)
|
||||
)
|
||||
if platform_role is None or platform_role.role_code != "admin":
|
||||
raise HTTPException(
|
||||
status.HTTP_403_FORBIDDEN,
|
||||
"需要系统管理员权限",
|
||||
)
|
||||
request_id = request.headers.get("X-Request-ID") or new_ulid()
|
||||
return SystemAdminContext(
|
||||
request_id=request_id,
|
||||
user=user,
|
||||
platform_role=platform_role,
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Payload helpers
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def workspace_payload(workspace: Workspaces) -> dict[str, Any]:
|
||||
return {
|
||||
"workspace_id": workspace.workspace_id,
|
||||
"workspace_code": workspace.workspace_code,
|
||||
"workspace_name": workspace.workspace_name,
|
||||
"active_root_uri": workspace.active_root_uri,
|
||||
"quota_bytes": workspace.quota_bytes,
|
||||
"status": workspace.status,
|
||||
"description": workspace.description,
|
||||
"created_by": workspace.created_by,
|
||||
"created_at": workspace.created_at.isoformat(),
|
||||
"updated_at": (
|
||||
workspace.updated_at.isoformat() if workspace.updated_at else None
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def member_payload(
|
||||
user: Users,
|
||||
role: Roles,
|
||||
membership: WorkspaceMembers,
|
||||
) -> dict[str, Any]:
|
||||
return {
|
||||
"user_id": user.user_id,
|
||||
"username": user.username,
|
||||
"display_name": user.display_name,
|
||||
"email": user.email,
|
||||
"user_status": user.status,
|
||||
"role_code": role.role_code,
|
||||
"role_name": role.role_name,
|
||||
"member_status": membership.member_status,
|
||||
"joined_at": membership.joined_at.isoformat(),
|
||||
}
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Internal helpers
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
async def _load_workspace(session: AsyncSession, workspace_id: str) -> Workspaces:
|
||||
workspace = await session.get(Workspaces, workspace_id)
|
||||
if workspace is None:
|
||||
raise HTTPException(status.HTTP_404_NOT_FOUND, "workspace 不存在")
|
||||
return workspace
|
||||
|
||||
|
||||
async def _load_role_by_code(session: AsyncSession, role_code: str) -> Roles:
|
||||
role = await session.scalar(select(Roles).where(Roles.role_code == role_code))
|
||||
if role is None:
|
||||
raise HTTPException(
|
||||
status.HTTP_422_UNPROCESSABLE_ENTITY,
|
||||
f"角色 {role_code} 不存在",
|
||||
)
|
||||
return role
|
||||
|
||||
|
||||
async def _count_active_admins(
|
||||
session: AsyncSession,
|
||||
workspace_id: str,
|
||||
exclude_user_id: str | None = None,
|
||||
) -> int:
|
||||
"""Count active admin members of ``workspace_id``.
|
||||
|
||||
Pass ``exclude_user_id`` when checking "would X be the last admin?"
|
||||
before mutating X.
|
||||
"""
|
||||
admin_role = await _load_role_by_code(session, "admin")
|
||||
stmt = (
|
||||
select(func.count())
|
||||
.select_from(WorkspaceMembers)
|
||||
.where(
|
||||
WorkspaceMembers.workspace_id == workspace_id,
|
||||
WorkspaceMembers.role_id == admin_role.role_id,
|
||||
WorkspaceMembers.member_status == "active",
|
||||
WorkspaceMembers.is_deleted == 0,
|
||||
)
|
||||
)
|
||||
if exclude_user_id is not None:
|
||||
stmt = stmt.where(WorkspaceMembers.user_id != exclude_user_id)
|
||||
return int(await session.scalar(stmt) or 0)
|
||||
|
||||
|
||||
def _envelope(request_id: str, data: Any, meta: dict[str, Any] | None = None) -> dict[str, Any]:
|
||||
return {
|
||||
"request_id": request_id,
|
||||
"data": data,
|
||||
"meta": meta or {},
|
||||
}
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Workspace CRUD
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
@router.get("/workspaces")
|
||||
async def list_workspaces(
|
||||
context: SystemAdminContext = Depends(system_admin_context),
|
||||
session: AsyncSession = Depends(database_session),
|
||||
) -> dict[str, Any]:
|
||||
"""List active/archived workspaces. Soft-deleted rows are filtered out.
|
||||
|
||||
Silent ``pageSize=100`` cap — YAGNI on real pagination until needed.
|
||||
"""
|
||||
rows = (
|
||||
await session.execute(
|
||||
select(Workspaces)
|
||||
.where(
|
||||
Workspaces.status != "disabled",
|
||||
Workspaces.is_deleted == 0,
|
||||
)
|
||||
.order_by(Workspaces.created_at, Workspaces.workspace_id)
|
||||
.limit(LIST_PAGE_SIZE)
|
||||
)
|
||||
).scalars().all()
|
||||
return _envelope(
|
||||
context.request_id,
|
||||
[workspace_payload(w) for w in rows],
|
||||
{"count": len(rows), "page_size": LIST_PAGE_SIZE},
|
||||
)
|
||||
|
||||
|
||||
@router.post("/workspaces", status_code=status.HTTP_201_CREATED)
|
||||
async def create_workspace(
|
||||
payload: WorkspaceCreate,
|
||||
context: SystemAdminContext = Depends(system_admin_context),
|
||||
session: AsyncSession = Depends(database_session),
|
||||
) -> dict[str, Any]:
|
||||
"""Create a workspace and auto-join the creator as an admin member."""
|
||||
if not WORKSPACE_CODE_PATTERN.fullmatch(payload.workspace_code):
|
||||
raise HTTPException(
|
||||
status.HTTP_422_UNPROCESSABLE_ENTITY,
|
||||
"workspace_code 必须匹配 ^[a-z0-9-]{3,32}$",
|
||||
)
|
||||
duplicate = await session.scalar(
|
||||
select(Workspaces.workspace_id).where(
|
||||
Workspaces.workspace_code == payload.workspace_code,
|
||||
)
|
||||
)
|
||||
if duplicate is not None:
|
||||
raise HTTPException(status.HTTP_409_CONFLICT, "workspace_code 已存在")
|
||||
|
||||
admin_role = await _load_role_by_code(session, "admin")
|
||||
workspace_id = new_ulid()
|
||||
workspace = Workspaces(
|
||||
workspace_id=workspace_id,
|
||||
workspace_code=payload.workspace_code,
|
||||
workspace_name=payload.workspace_name,
|
||||
active_root_uri=f"s3://workspaces/{workspace_id}/",
|
||||
quota_bytes=payload.quota_bytes,
|
||||
status="active",
|
||||
created_by=context.user.user_id,
|
||||
description=payload.description,
|
||||
)
|
||||
session.add(workspace)
|
||||
session.add(
|
||||
WorkspaceMembers(
|
||||
workspace_id=workspace_id,
|
||||
user_id=context.user.user_id,
|
||||
role_id=admin_role.role_id,
|
||||
member_status="active",
|
||||
)
|
||||
)
|
||||
await session.flush()
|
||||
await session.refresh(workspace)
|
||||
return _envelope(context.request_id, workspace_payload(workspace))
|
||||
|
||||
|
||||
@router.get("/workspaces/{workspace_id}")
|
||||
async def get_workspace(
|
||||
workspace_id: str,
|
||||
context: SystemAdminContext = Depends(system_admin_context),
|
||||
session: AsyncSession = Depends(database_session),
|
||||
) -> dict[str, Any]:
|
||||
"""Fetch a single workspace — even soft-deleted ones are reachable."""
|
||||
workspace = await _load_workspace(session, workspace_id)
|
||||
return _envelope(context.request_id, workspace_payload(workspace))
|
||||
|
||||
|
||||
@router.patch("/workspaces/{workspace_id}")
|
||||
async def update_workspace(
|
||||
workspace_id: str,
|
||||
payload: WorkspaceUpdate,
|
||||
context: SystemAdminContext = Depends(system_admin_context),
|
||||
session: AsyncSession = Depends(database_session),
|
||||
) -> dict[str, Any]:
|
||||
"""Patch editable workspace fields. ``status='disabled'`` is rejected."""
|
||||
workspace = await _load_workspace(session, workspace_id)
|
||||
if workspace.status == "disabled":
|
||||
raise HTTPException(
|
||||
status.HTTP_409_CONFLICT,
|
||||
"workspace 已删除,无法修改",
|
||||
)
|
||||
if payload.workspace_name is not None:
|
||||
workspace.workspace_name = payload.workspace_name.strip()
|
||||
if payload.quota_bytes is not None:
|
||||
workspace.quota_bytes = payload.quota_bytes
|
||||
if payload.description is not None:
|
||||
workspace.description = payload.description
|
||||
if payload.status is not None:
|
||||
workspace.status = payload.status
|
||||
await session.flush()
|
||||
await session.refresh(workspace)
|
||||
return _envelope(context.request_id, workspace_payload(workspace))
|
||||
|
||||
|
||||
@router.delete("/workspaces/{workspace_id}")
|
||||
async def delete_workspace(
|
||||
workspace_id: str,
|
||||
context: SystemAdminContext = Depends(system_admin_context),
|
||||
session: AsyncSession = Depends(database_session),
|
||||
) -> dict[str, Any]:
|
||||
"""Soft-delete a workspace and cascade-soft-delete its memberships.
|
||||
|
||||
Allowed from any non-disabled status (active or archived). The
|
||||
membership cascade is what lets system admins leave a workspace —
|
||||
there is no per-member DELETE escape for self-removal.
|
||||
"""
|
||||
workspace = await _load_workspace(session, workspace_id)
|
||||
if workspace.status == "disabled":
|
||||
raise HTTPException(
|
||||
status.HTTP_409_CONFLICT,
|
||||
"workspace 已被删除",
|
||||
)
|
||||
now = datetime.datetime.utcnow()
|
||||
workspace.status = "disabled"
|
||||
workspace.is_deleted = 1
|
||||
workspace.deleted_at = now
|
||||
await session.execute(
|
||||
update(WorkspaceMembers)
|
||||
.where(
|
||||
WorkspaceMembers.workspace_id == workspace_id,
|
||||
WorkspaceMembers.is_deleted == 0,
|
||||
)
|
||||
.values(is_deleted=1, deleted_at=now)
|
||||
)
|
||||
await session.flush()
|
||||
await session.refresh(workspace)
|
||||
return _envelope(context.request_id, workspace_payload(workspace))
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Workspace membership CRUD
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
@router.get("/workspaces/{workspace_id}/members")
|
||||
async def list_members(
|
||||
workspace_id: str,
|
||||
context: SystemAdminContext = Depends(system_admin_context),
|
||||
session: AsyncSession = Depends(database_session),
|
||||
) -> dict[str, Any]:
|
||||
"""List active and historical (non-soft-deleted) members of a workspace."""
|
||||
await _load_workspace(session, workspace_id)
|
||||
rows = (
|
||||
await session.execute(
|
||||
select(Users, Roles, WorkspaceMembers)
|
||||
.join(
|
||||
WorkspaceMembers,
|
||||
WorkspaceMembers.user_id == Users.user_id,
|
||||
)
|
||||
.join(Roles, Roles.role_id == WorkspaceMembers.role_id)
|
||||
.where(
|
||||
WorkspaceMembers.workspace_id == workspace_id,
|
||||
WorkspaceMembers.is_deleted == 0,
|
||||
)
|
||||
.order_by(WorkspaceMembers.joined_at, Users.user_id)
|
||||
.limit(LIST_PAGE_SIZE)
|
||||
)
|
||||
).all()
|
||||
return _envelope(
|
||||
context.request_id,
|
||||
[member_payload(u, r, m) for u, r, m in rows],
|
||||
{"count": len(rows), "page_size": LIST_PAGE_SIZE},
|
||||
)
|
||||
|
||||
|
||||
@router.post(
|
||||
"/workspaces/{workspace_id}/members",
|
||||
status_code=status.HTTP_201_CREATED,
|
||||
)
|
||||
async def add_member(
|
||||
workspace_id: str,
|
||||
payload: MemberCreate,
|
||||
context: SystemAdminContext = Depends(system_admin_context),
|
||||
session: AsyncSession = Depends(database_session),
|
||||
) -> dict[str, Any]:
|
||||
"""Add a user to a workspace. The new row starts with member_status='active'."""
|
||||
await _load_workspace(session, workspace_id)
|
||||
user = await session.get(Users, payload.user_id)
|
||||
if user is None:
|
||||
raise HTTPException(status.HTTP_404_NOT_FOUND, "用户不存在")
|
||||
role = await _load_role_by_code(session, payload.role_code)
|
||||
duplicate = await session.scalar(
|
||||
select(WorkspaceMembers.user_id).where(
|
||||
WorkspaceMembers.workspace_id == workspace_id,
|
||||
WorkspaceMembers.user_id == payload.user_id,
|
||||
WorkspaceMembers.is_deleted == 0,
|
||||
)
|
||||
)
|
||||
if duplicate is not None:
|
||||
raise HTTPException(
|
||||
status.HTTP_409_CONFLICT,
|
||||
"用户已是该 workspace 成员",
|
||||
)
|
||||
membership = WorkspaceMembers(
|
||||
workspace_id=workspace_id,
|
||||
user_id=payload.user_id,
|
||||
role_id=role.role_id,
|
||||
member_status="active",
|
||||
)
|
||||
session.add(membership)
|
||||
await session.flush()
|
||||
await session.refresh(membership)
|
||||
return _envelope(context.request_id, member_payload(user, role, membership))
|
||||
|
||||
|
||||
@router.patch("/workspaces/{workspace_id}/members/{user_id}")
|
||||
async def update_member(
|
||||
workspace_id: str,
|
||||
user_id: str,
|
||||
payload: MemberUpdate,
|
||||
context: SystemAdminContext = Depends(system_admin_context),
|
||||
session: AsyncSession = Depends(database_session),
|
||||
) -> dict[str, Any]:
|
||||
"""Update a member's role and/or status. Last-admin guard applies."""
|
||||
await _load_workspace(session, workspace_id)
|
||||
row = (
|
||||
await session.execute(
|
||||
select(Users, Roles, WorkspaceMembers)
|
||||
.join(
|
||||
WorkspaceMembers,
|
||||
WorkspaceMembers.user_id == Users.user_id,
|
||||
)
|
||||
.join(Roles, Roles.role_id == WorkspaceMembers.role_id)
|
||||
.where(
|
||||
WorkspaceMembers.workspace_id == workspace_id,
|
||||
WorkspaceMembers.user_id == user_id,
|
||||
WorkspaceMembers.is_deleted == 0,
|
||||
)
|
||||
)
|
||||
).first()
|
||||
if row is None:
|
||||
raise HTTPException(status.HTTP_404_NOT_FOUND, "成员不存在")
|
||||
user, role, membership = row
|
||||
|
||||
next_role = role
|
||||
if payload.role_code is not None and payload.role_code != role.role_code:
|
||||
if (
|
||||
role.role_code == "admin"
|
||||
and payload.role_code != "admin"
|
||||
and membership.member_status == "active"
|
||||
):
|
||||
remaining = await _count_active_admins(
|
||||
session, workspace_id, exclude_user_id=user_id,
|
||||
)
|
||||
if remaining == 0:
|
||||
raise HTTPException(
|
||||
status.HTTP_409_CONFLICT,
|
||||
"workspace 必须保留至少一个 admin",
|
||||
)
|
||||
next_role = await _load_role_by_code(session, payload.role_code)
|
||||
membership.role_id = next_role.role_id
|
||||
|
||||
if payload.member_status is not None and payload.member_status != membership.member_status:
|
||||
if (
|
||||
role.role_code == "admin"
|
||||
and payload.member_status != "active"
|
||||
):
|
||||
remaining = await _count_active_admins(
|
||||
session, workspace_id, exclude_user_id=user_id,
|
||||
)
|
||||
if remaining == 0:
|
||||
raise HTTPException(
|
||||
status.HTTP_409_CONFLICT,
|
||||
"workspace 必须保留至少一个 admin",
|
||||
)
|
||||
membership.member_status = payload.member_status
|
||||
|
||||
await session.flush()
|
||||
await session.refresh(membership)
|
||||
return _envelope(context.request_id, member_payload(user, next_role, membership))
|
||||
|
||||
|
||||
@router.delete("/workspaces/{workspace_id}/members/{user_id}")
|
||||
async def remove_member(
|
||||
workspace_id: str,
|
||||
user_id: str,
|
||||
context: SystemAdminContext = Depends(system_admin_context),
|
||||
session: AsyncSession = Depends(database_session),
|
||||
) -> dict[str, Any]:
|
||||
"""Soft-delete a workspace membership.
|
||||
|
||||
System admins cannot remove themselves — the only escape is to delete
|
||||
the entire workspace, which cascades membership soft-deletion.
|
||||
"""
|
||||
await _load_workspace(session, workspace_id)
|
||||
if user_id == context.user.user_id:
|
||||
raise HTTPException(
|
||||
status.HTTP_403_FORBIDDEN,
|
||||
"系统管理员不能把自己从 workspace 移除;如需退出,请删除整个 workspace",
|
||||
)
|
||||
row = (
|
||||
await session.execute(
|
||||
select(Roles, WorkspaceMembers)
|
||||
.join(Roles, Roles.role_id == WorkspaceMembers.role_id)
|
||||
.where(
|
||||
WorkspaceMembers.workspace_id == workspace_id,
|
||||
WorkspaceMembers.user_id == user_id,
|
||||
WorkspaceMembers.is_deleted == 0,
|
||||
)
|
||||
)
|
||||
).first()
|
||||
if row is None:
|
||||
raise HTTPException(status.HTTP_404_NOT_FOUND, "成员不存在")
|
||||
role, membership = row
|
||||
if role.role_code == "admin" and membership.member_status == "active":
|
||||
remaining = await _count_active_admins(
|
||||
session, workspace_id, exclude_user_id=user_id,
|
||||
)
|
||||
if remaining == 0:
|
||||
raise HTTPException(
|
||||
status.HTTP_409_CONFLICT,
|
||||
"workspace 必须保留至少一个 admin",
|
||||
)
|
||||
membership.is_deleted = 1
|
||||
membership.deleted_at = datetime.datetime.utcnow()
|
||||
await session.flush()
|
||||
return _envelope(
|
||||
context.request_id,
|
||||
{"workspace_id": workspace_id, "user_id": user_id, "removed": True},
|
||||
)
|
||||
|
||||
|
||||
__all__ = [
|
||||
"router",
|
||||
"SystemAdminContext",
|
||||
"system_admin_context",
|
||||
]
|
||||
@@ -1,5 +1,4 @@
|
||||
from common.db.base import Base
|
||||
from common.db.models.audit import AuditLogs
|
||||
from common.db.models.events import ConsumerInbox, OutboxEvents
|
||||
from common.db.models.identity import Permissions, RolePermissions, Roles, Users
|
||||
from common.db.models.runtime import WorkspaceOperations
|
||||
@@ -16,7 +15,6 @@ from common.db.models.workspaces import WorkspaceMembers, Workspaces
|
||||
|
||||
__all__ = [
|
||||
"Base",
|
||||
"AuditLogs",
|
||||
"ConsumerInbox",
|
||||
"DataResources",
|
||||
"OutboxEvents",
|
||||
|
||||
@@ -1,38 +0,0 @@
|
||||
import datetime
|
||||
from typing import Optional
|
||||
|
||||
from sqlalchemy import Index, JSON, String, text
|
||||
from sqlalchemy.dialects.mysql import BIGINT, CHAR, DATETIME, TINYINT
|
||||
from sqlalchemy.orm import Mapped, mapped_column
|
||||
|
||||
from common.db.base import Base
|
||||
|
||||
|
||||
class AuditLogs(Base):
|
||||
__tablename__ = "audit_logs"
|
||||
__table_args__ = (
|
||||
Index("idx_audit_action_time", "action_code", "created_at"),
|
||||
Index("idx_audit_actor_time", "actor_user_id", "created_at"),
|
||||
Index("idx_audit_workspace_time", "workspace_id", "created_at"),
|
||||
{"comment": "操作审计日志"},
|
||||
)
|
||||
|
||||
audit_id: Mapped[int] = mapped_column(BIGINT, primary_key=True)
|
||||
action_code: Mapped[str] = mapped_column(String(128), nullable=False)
|
||||
target_type: Mapped[str] = mapped_column(String(64), nullable=False)
|
||||
operation_status: Mapped[str] = mapped_column(
|
||||
String(16), nullable=False, server_default=text("'success'")
|
||||
)
|
||||
created_at: Mapped[datetime.datetime] = mapped_column(
|
||||
DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)")
|
||||
)
|
||||
workspace_id: Mapped[Optional[str]] = mapped_column(CHAR(26))
|
||||
actor_user_id: Mapped[Optional[str]] = mapped_column(CHAR(26))
|
||||
target_id: Mapped[Optional[str]] = mapped_column(String(128))
|
||||
client_ip: Mapped[Optional[str]] = mapped_column(String(45))
|
||||
user_agent: Mapped[Optional[str]] = mapped_column(String(1000))
|
||||
detail_json: Mapped[Optional[dict]] = mapped_column(JSON)
|
||||
is_deleted: Mapped[int] = mapped_column(
|
||||
TINYINT(1), nullable=False, server_default=text("0")
|
||||
)
|
||||
deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3))
|
||||
@@ -6,7 +6,6 @@ import asyncio
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
from collections import defaultdict
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from urllib.parse import quote
|
||||
@@ -17,7 +16,6 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from common.db import create_database_engine, create_session_factory
|
||||
from common.db.models import (
|
||||
AuditLogs,
|
||||
Roles,
|
||||
Users,
|
||||
WorkspaceMembers,
|
||||
@@ -95,7 +93,6 @@ def new_stats() -> dict[str, int]:
|
||||
"workspaces_updated": 0,
|
||||
"workspace_members_inserted": 0,
|
||||
"workspace_members_updated": 0,
|
||||
"audit_logs_workspace_backfilled": 0,
|
||||
}
|
||||
|
||||
|
||||
@@ -266,44 +263,6 @@ async def migrate_members(
|
||||
stats["workspace_members_updated"] += 1
|
||||
|
||||
|
||||
async def backfill_audit_workspaces(
|
||||
session: AsyncSession,
|
||||
source: list[LegacyWorkspace],
|
||||
workspace_ids: dict[str, str],
|
||||
users: dict[str, Users],
|
||||
stats: dict[str, int],
|
||||
) -> None:
|
||||
workspace_codes_by_username: dict[str, set[str]] = defaultdict(set)
|
||||
for workspace in source:
|
||||
for username in workspace.userIds:
|
||||
workspace_codes_by_username[username].add(workspace.id)
|
||||
|
||||
workspace_by_user_id = {
|
||||
users[username].user_id: workspace_ids[next(iter(codes))]
|
||||
for username, codes in workspace_codes_by_username.items()
|
||||
if len(codes) == 1
|
||||
}
|
||||
|
||||
audit_logs = (
|
||||
await session.scalars(
|
||||
select(AuditLogs).where(AuditLogs.workspace_id.is_(None))
|
||||
)
|
||||
).all()
|
||||
for audit_log in audit_logs:
|
||||
if (
|
||||
not isinstance(audit_log.detail_json, dict)
|
||||
or audit_log.detail_json.get("migration_source")
|
||||
!= "platform_data/system.json"
|
||||
):
|
||||
continue
|
||||
workspace_id = workspace_by_user_id.get(
|
||||
audit_log.actor_user_id
|
||||
)
|
||||
if workspace_id is not None:
|
||||
audit_log.workspace_id = workspace_id
|
||||
stats["audit_logs_workspace_backfilled"] += 1
|
||||
|
||||
|
||||
async def run_migration(
|
||||
database_url: str,
|
||||
source: list[LegacyWorkspace],
|
||||
@@ -339,13 +298,6 @@ async def run_migration(
|
||||
users,
|
||||
stats,
|
||||
)
|
||||
await backfill_audit_workspaces(
|
||||
session,
|
||||
source,
|
||||
workspace_ids,
|
||||
users,
|
||||
stats,
|
||||
)
|
||||
|
||||
if apply_changes:
|
||||
await session.commit()
|
||||
|
||||
@@ -5,18 +5,16 @@ import asyncio
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from urllib.parse import quote
|
||||
|
||||
from pydantic import BaseModel, ConfigDict, Field
|
||||
from pydantic import BaseModel, ConfigDict
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from common.db import create_database_engine, create_session_factory
|
||||
from common.db.models import (
|
||||
AuditLogs,
|
||||
Permissions,
|
||||
RolePermissions,
|
||||
Roles,
|
||||
@@ -39,19 +37,6 @@ PERMISSION_NAMES = {
|
||||
"resource.public.manage": "管理公共资源",
|
||||
"system.view": "查看系统管理",
|
||||
"system.manage": "管理系统配置",
|
||||
"audit.view": "查看审计日志",
|
||||
}
|
||||
|
||||
ACTION_CATALOG = {
|
||||
"保存调度配置": ("schedule.save", "schedule"),
|
||||
"删除脚本对象": ("script.delete", "script"),
|
||||
"删除实验记录": ("experiment.delete", "experiment"),
|
||||
"删除数据资源": ("data_resource.delete", "data_resource"),
|
||||
"上传数据资源": ("data_resource.upload", "data_resource"),
|
||||
"新建脚本对象": ("script.create", "script"),
|
||||
"修改用户角色": ("user.role.update", "user"),
|
||||
"运行 Python": ("script.run_python", "script"),
|
||||
"运行调度": ("schedule.run", "schedule"),
|
||||
}
|
||||
|
||||
|
||||
@@ -74,26 +59,18 @@ class LegacyRole(BaseModel):
|
||||
permissions: list[str]
|
||||
|
||||
|
||||
class LegacyAuditLog(BaseModel):
|
||||
model_config = ConfigDict(extra="forbid")
|
||||
|
||||
id: str
|
||||
actorId: str
|
||||
actorName: str
|
||||
role: str
|
||||
action: str
|
||||
target: str
|
||||
detail: str
|
||||
status: str
|
||||
createdAt: str
|
||||
|
||||
|
||||
class LegacySystem(BaseModel):
|
||||
model_config = ConfigDict(extra="forbid")
|
||||
"""Legacy platform_data/system.json shape.
|
||||
|
||||
``extra='ignore'`` so legacy files that include other top-level
|
||||
keys (e.g. historical audit-log payloads) can still be parsed —
|
||||
only the fields this script actually consumes are listed below.
|
||||
"""
|
||||
|
||||
model_config = ConfigDict(extra="ignore")
|
||||
|
||||
users: list[LegacyUser]
|
||||
roles: list[LegacyRole]
|
||||
audit_logs: list[LegacyAuditLog] = Field(alias="auditLogs")
|
||||
|
||||
|
||||
def deterministic_legacy_ulid(entity_type: str, legacy_key: str) -> str:
|
||||
@@ -120,10 +97,8 @@ def require_unique(values: list[str], label: str) -> None:
|
||||
def validate_source(source: LegacySystem) -> None:
|
||||
role_codes = [role.key for role in source.roles]
|
||||
user_codes = [user.id for user in source.users]
|
||||
audit_ids = [item.id for item in source.audit_logs]
|
||||
require_unique(role_codes, "role keys")
|
||||
require_unique(user_codes, "user ids")
|
||||
require_unique(audit_ids, "audit ids")
|
||||
|
||||
role_code_set = set(role_codes)
|
||||
unknown_roles = sorted(
|
||||
@@ -134,15 +109,6 @@ def validate_source(source: LegacySystem) -> None:
|
||||
if unknown_roles:
|
||||
raise ValueError(f"users reference unknown roles: {unknown_roles}")
|
||||
|
||||
user_code_set = set(user_codes)
|
||||
unknown_actors = sorted(
|
||||
item.actorId
|
||||
for item in source.audit_logs
|
||||
if item.actorId not in user_code_set
|
||||
)
|
||||
if unknown_actors:
|
||||
raise ValueError(f"audit logs reference unknown users: {unknown_actors}")
|
||||
|
||||
for role in source.roles:
|
||||
require_unique(role.permissions, f"permissions of role {role.key}")
|
||||
for permission_code in role.permissions:
|
||||
@@ -177,8 +143,6 @@ def new_stats() -> dict[str, int]:
|
||||
"role_permissions_inserted": 0,
|
||||
"users_inserted": 0,
|
||||
"users_updated": 0,
|
||||
"audit_logs_inserted": 0,
|
||||
"audit_logs_skipped": 0,
|
||||
}
|
||||
|
||||
|
||||
@@ -343,63 +307,6 @@ async def migrate_users(
|
||||
return user_ids
|
||||
|
||||
|
||||
def audit_catalog(action: str) -> tuple[str, str]:
|
||||
known = ACTION_CATALOG.get(action)
|
||||
if known is not None:
|
||||
return known
|
||||
digest = hashlib.sha256(action.encode("utf-8")).hexdigest()[:16]
|
||||
return (f"legacy.action.{digest}", "legacy")
|
||||
|
||||
|
||||
async def migrate_audit_logs(
|
||||
session: AsyncSession,
|
||||
source: LegacySystem,
|
||||
user_ids: dict[str, str],
|
||||
stats: dict[str, int],
|
||||
) -> None:
|
||||
existing_payloads = (
|
||||
await session.scalars(select(AuditLogs.detail_json))
|
||||
).all()
|
||||
existing_legacy_ids = {
|
||||
payload.get("legacy_id")
|
||||
for payload in existing_payloads
|
||||
if isinstance(payload, dict) and payload.get("legacy_id")
|
||||
}
|
||||
|
||||
for legacy in sorted(source.audit_logs, key=lambda item: item.createdAt):
|
||||
if legacy.id in existing_legacy_ids:
|
||||
stats["audit_logs_skipped"] += 1
|
||||
continue
|
||||
|
||||
action_code, target_type = audit_catalog(legacy.action)
|
||||
session.add(
|
||||
AuditLogs(
|
||||
actor_user_id=user_ids[legacy.actorId],
|
||||
action_code=action_code,
|
||||
target_type=target_type,
|
||||
target_id=legacy.target[:128] or None,
|
||||
operation_status=(
|
||||
"success" if legacy.status == "成功" else "failed"
|
||||
),
|
||||
created_at=datetime.strptime(
|
||||
legacy.createdAt, "%Y-%m-%d %H:%M:%S"
|
||||
),
|
||||
detail_json={
|
||||
"legacy_id": legacy.id,
|
||||
"legacy_actor_name": legacy.actorName,
|
||||
"legacy_role": legacy.role,
|
||||
"legacy_action": legacy.action,
|
||||
"legacy_target": legacy.target,
|
||||
"legacy_detail": legacy.detail,
|
||||
"legacy_status": legacy.status,
|
||||
"migration_source": "platform_data/system.json",
|
||||
},
|
||||
)
|
||||
)
|
||||
existing_legacy_ids.add(legacy.id)
|
||||
stats["audit_logs_inserted"] += 1
|
||||
|
||||
|
||||
async def run_migration(
|
||||
database_url: str,
|
||||
source: LegacySystem,
|
||||
@@ -424,12 +331,9 @@ async def run_migration(
|
||||
permission_ids,
|
||||
stats,
|
||||
)
|
||||
user_ids = await migrate_users(
|
||||
await migrate_users(
|
||||
session, source, role_ids, stats
|
||||
)
|
||||
await migrate_audit_logs(
|
||||
session, source, user_ids, stats
|
||||
)
|
||||
|
||||
if apply_changes:
|
||||
await session.commit()
|
||||
@@ -489,7 +393,6 @@ async def async_main() -> None:
|
||||
len(role.permissions) for role in source.roles
|
||||
),
|
||||
"users": len(source.users),
|
||||
"audit_logs": len(source.audit_logs),
|
||||
},
|
||||
"changes": stats,
|
||||
}
|
||||
@@ -501,4 +404,4 @@ def main() -> None:
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
main()
|
||||
@@ -1,4 +1,4 @@
|
||||
"""initial baseline (20 active tables)
|
||||
"""initial baseline (19 active tables)
|
||||
|
||||
Revision ID: 8d86e2f82860
|
||||
Revises:
|
||||
@@ -21,26 +21,6 @@ depends_on: str | Sequence[str] | None = None
|
||||
def upgrade() -> None:
|
||||
"""Upgrade schema."""
|
||||
# ### commands auto generated by Alembic - please adjust! ###
|
||||
op.create_table('audit_logs',
|
||||
sa.Column('audit_id', mysql.BIGINT(), nullable=False),
|
||||
sa.Column('action_code', sa.String(length=128), nullable=False),
|
||||
sa.Column('target_type', sa.String(length=64), nullable=False),
|
||||
sa.Column('operation_status', sa.String(length=16), server_default=sa.text("'success'"), nullable=False),
|
||||
sa.Column('created_at', mysql.DATETIME(fsp=3), server_default=sa.text('CURRENT_TIMESTAMP(3)'), nullable=False),
|
||||
sa.Column('workspace_id', mysql.CHAR(length=26), nullable=True),
|
||||
sa.Column('actor_user_id', mysql.CHAR(length=26), nullable=True),
|
||||
sa.Column('target_id', sa.String(length=128), nullable=True),
|
||||
sa.Column('client_ip', sa.String(length=45), nullable=True),
|
||||
sa.Column('user_agent', sa.String(length=1000), nullable=True),
|
||||
sa.Column('detail_json', sa.JSON(), nullable=True),
|
||||
sa.Column('is_deleted', mysql.TINYINT(display_width=1), server_default=sa.text('0'), nullable=False),
|
||||
sa.Column('deleted_at', mysql.DATETIME(fsp=3), nullable=True),
|
||||
sa.PrimaryKeyConstraint('audit_id'),
|
||||
comment='操作审计日志'
|
||||
)
|
||||
op.create_index('idx_audit_action_time', 'audit_logs', ['action_code', 'created_at'], unique=False)
|
||||
op.create_index('idx_audit_actor_time', 'audit_logs', ['actor_user_id', 'created_at'], unique=False)
|
||||
op.create_index('idx_audit_workspace_time', 'audit_logs', ['workspace_id', 'created_at'], unique=False)
|
||||
op.create_table('consumer_inbox',
|
||||
sa.Column('consumer_name', sa.String(length=128), nullable=False),
|
||||
sa.Column('event_id', mysql.CHAR(length=26), nullable=False),
|
||||
@@ -551,8 +531,4 @@ def downgrade() -> None:
|
||||
op.drop_table('data_resources')
|
||||
op.drop_index('idx_consumer_inbox_status', table_name='consumer_inbox')
|
||||
op.drop_table('consumer_inbox')
|
||||
op.drop_index('idx_audit_workspace_time', table_name='audit_logs')
|
||||
op.drop_index('idx_audit_actor_time', table_name='audit_logs')
|
||||
op.drop_index('idx_audit_action_time', table_name='audit_logs')
|
||||
op.drop_table('audit_logs')
|
||||
# ### end Alembic commands ###
|
||||
|
||||
Reference in New Issue
Block a user