diff --git a/API.md b/API.md index 2e20559..8902d0b 100644 --- a/API.md +++ b/API.md @@ -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` | diff --git a/backend/src/backend/auth.py b/backend/src/backend/auth.py index f302db3..a49d789 100644 --- a/backend/src/backend/auth.py +++ b/backend/src/backend/auth.py @@ -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, }, diff --git a/backend/src/backend/main.py b/backend/src/backend/main.py index 6752caf..c26fde8 100644 --- a/backend/src/backend/main.py +++ b/backend/src/backend/main.py @@ -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: diff --git a/backend/src/backend/platform.py b/backend/src/backend/platform.py new file mode 100644 index 0000000..5bc48cb --- /dev/null +++ b/backend/src/backend/platform.py @@ -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", +] \ No newline at end of file diff --git a/common/src/common/db/models/__init__.py b/common/src/common/db/models/__init__.py index 15ad04a..b7fca0d 100644 --- a/common/src/common/db/models/__init__.py +++ b/common/src/common/db/models/__init__.py @@ -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", diff --git a/common/src/common/db/models/audit.py b/common/src/common/db/models/audit.py deleted file mode 100644 index ed9edba..0000000 --- a/common/src/common/db/models/audit.py +++ /dev/null @@ -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)) diff --git a/migrations/data/migrate_legacy_workspaces.py b/migrations/data/migrate_legacy_workspaces.py index 6953d78..da081c7 100644 --- a/migrations/data/migrate_legacy_workspaces.py +++ b/migrations/data/migrate_legacy_workspaces.py @@ -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() diff --git a/migrations/data/migrate_system_json.py b/migrations/data/migrate_system_json.py index 8a52322..a608ddc 100644 --- a/migrations/data/migrate_system_json.py +++ b/migrations/data/migrate_system_json.py @@ -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() \ No newline at end of file diff --git a/migrations/versions/8d86e2f82860_initial_baseline_20_active_tables.py b/migrations/versions/8d86e2f82860_initial_baseline_20_active_tables.py index c9f4673..b9a6129 100644 --- a/migrations/versions/8d86e2f82860_initial_baseline_20_active_tables.py +++ b/migrations/versions/8d86e2f82860_initial_baseline_20_active_tables.py @@ -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 ###