refactor: delete table

This commit is contained in:
tao.chen
2026-07-31 13:21:15 +08:00
parent f288c90d37
commit 28069f516b
9 changed files with 510 additions and 1074 deletions
+60 -60
View File
@@ -21,7 +21,6 @@ from sqlalchemy import func, select
from sqlalchemy.ext.asyncio import AsyncSession
from common.db.models import (
EditSessions,
Scripts,
StorageObjects,
Versions,
@@ -227,24 +226,31 @@ async def get_script_row(
return script, storage_object
async def require_no_active_edit_session(
session: AsyncSession,
storage_object_ids: list[str],
def require_script_modify_access(
script: Scripts,
*,
user_id: str,
is_admin: bool,
) -> None:
if not storage_object_ids:
"""Enforce the V3.1 §4 access rules for write operations on a script.
Rules:
* admin: always allowed
* owner: always allowed
* non-owner: allowed iff ``is_locked`` is False
Reads (``list`` / ``get``) intentionally do not call this helper — the
design contract is "everyone in the workspace can see the script
list, but only the owner (or admin) can mutate when locked".
"""
if is_admin or script.owner_user_id == user_id:
return
active_session = await session.scalar(
select(EditSessions).where(
EditSessions.storage_object_id.in_(storage_object_ids),
EditSessions.session_status == "active",
EditSessions.expires_at > datetime.now(UTC).replace(tzinfo=None),
)
if not script.is_locked:
return
raise HTTPException(
status.HTTP_403_FORBIDDEN,
"script is locked; only the owner (or an administrator) may modify it",
)
if active_session is not None:
raise HTTPException(
status.HTTP_409_CONFLICT,
"file is being edited; end the editing session before deletion",
)
async def create_script_record(
@@ -542,10 +548,6 @@ async def delete_workspace_directory(
)
)
).all()
await require_no_active_edit_session(
session,
[script.current_object_id for script, _storage in rows],
)
for script, _storage in rows:
await request.app.state.storage_client.delete_object(
script.current_object_id
@@ -624,14 +626,11 @@ async def update_script(
session,
for_update=True,
)
if (
script.owner_user_id != context.user.user_id
and not context.is_admin
):
raise HTTPException(
status.HTTP_403_FORBIDDEN,
"script can only be changed by its owner or an administrator",
)
require_script_modify_access(
script,
user_id=context.user.user_id,
is_admin=context.is_admin,
)
content = validate_script_content(payload.content, script.script_type)
if not storage_object.relative_path:
raise HTTPException(
@@ -677,17 +676,10 @@ async def delete_script(
session,
for_update=True,
)
if (
script.owner_user_id != context.user.user_id
and not context.is_admin
):
raise HTTPException(
status.HTTP_403_FORBIDDEN,
"script can only be deleted by its owner or an administrator",
)
await require_no_active_edit_session(
session,
[script.current_object_id],
require_script_modify_access(
script,
user_id=context.user.user_id,
is_admin=context.is_admin,
)
await request.app.state.storage_client.delete_object(
script.current_object_id
@@ -722,14 +714,11 @@ async def publish_version(
session,
for_update=True,
)
if (
script.owner_user_id != context.user.user_id
and not context.is_admin
):
raise HTTPException(
status.HTTP_403_FORBIDDEN,
"script can only be published by its owner or an administrator",
)
require_script_modify_access(
script,
user_id=context.user.user_id,
is_admin=context.is_admin,
)
if (
payload.source_object_id
and payload.source_object_id != script.current_object_id
@@ -867,20 +856,31 @@ async def delete_version(
context: RequestContext = Depends(request_context),
session: AsyncSession = Depends(database_session),
) -> dict[str, Any]:
version = await session.get(Versions, versions_id)
if (
version is None
or version.workspace_id != context.workspace.workspace_id
):
raise HTTPException(status.HTTP_404_NOT_FOUND, "稳定版本不存在")
if (
version.created_by != context.user.user_id
and not context.is_admin
):
raise HTTPException(
status.HTTP_403_FORBIDDEN,
"仅稳定版本发布者或管理员可以删除",
# Join to Scripts so the lock + owner check rides on the script record,
# not on whoever happened to publish this specific version. Lock state
# is a property of the script as a whole (architecture V3.1 §4), not of
# any one version of it.
row = (
await session.execute(
select(Versions, Scripts)
.join(
Scripts,
Scripts.script_id == Versions.script_id,
)
.where(
Versions.versions_id == versions_id,
Versions.workspace_id == context.workspace.workspace_id,
)
)
).one_or_none()
if row is None:
raise HTTPException(status.HTTP_404_NOT_FOUND, "稳定版本不存在")
version, script = row
require_script_modify_access(
script,
user_id=context.user.user_id,
is_admin=context.is_admin,
)
version.schedule_hidden_at = datetime.now(UTC).replace(tzinfo=None)
await session.flush()
+2 -13
View File
@@ -1,13 +1,8 @@
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.experiments import (
ExperimentMetrics,
ExperimentResources,
Experiments,
)
from common.db.models.identity import Permissions, RolePermissions, Roles, Users
from common.db.models.runtime import EditSessions, RuntimeInstances, WorkspaceOperations
from common.db.models.runtime import WorkspaceOperations
from common.db.models.schedules import (
ScheduleEdges,
ScheduleNodeRuns,
@@ -15,7 +10,7 @@ from common.db.models.schedules import (
ScheduleRuns,
Schedules,
)
from common.db.models.scripts import NotebookSnapshots, Scripts, Versions
from common.db.models.scripts import Scripts, Versions
from common.db.models.storage import DataResources, StorageObjects, UploadSessions
from common.db.models.workspaces import WorkspaceMembers, Workspaces
@@ -24,16 +19,10 @@ __all__ = [
"AuditLogs",
"ConsumerInbox",
"DataResources",
"EditSessions",
"ExperimentMetrics",
"ExperimentResources",
"Experiments",
"NotebookSnapshots",
"OutboxEvents",
"Permissions",
"RolePermissions",
"Roles",
"RuntimeInstances",
"ScheduleEdges",
"ScheduleNodeRuns",
"ScheduleNodes",
-112
View File
@@ -1,112 +0,0 @@
import datetime
import decimal
from typing import Optional
from sqlalchemy import BigInteger, Double, Index, JSON, String, text
from sqlalchemy.dialects.mysql import BIGINT, CHAR, DATETIME, INTEGER, TINYINT
from sqlalchemy.orm import Mapped, mapped_column
from common.db.base import Base
class Experiments(Base):
__tablename__ = "experiments"
__table_args__ = (
Index("fk_experiments_logs", "logs_object_id"),
Index("fk_experiments_parent", "parent_experiment_id"),
Index("fk_experiments_result", "result_object_id"),
Index("fk_experiments_script", "script_id"),
Index("idx_experiments_owner", "owner_user_id", "created_at"),
Index("idx_experiments_schedule_run", "schedule_run_id"),
Index("idx_experiments_version", "versions_id"),
Index(
"idx_experiments_workspace",
"workspace_id",
"experiment_status",
"created_at",
),
{"comment": "实验记录"},
)
experiment_id: Mapped[str] = mapped_column(CHAR(26), primary_key=True)
workspace_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
owner_user_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
experiment_name: Mapped[str] = mapped_column(String(255), nullable=False)
source_type: Mapped[str] = mapped_column(
String(24), nullable=False, comment="python/notebook/schedule/rerun"
)
experiment_status: Mapped[str] = mapped_column(
String(24), nullable=False, server_default=text("'queued'")
)
state_version: Mapped[int] = mapped_column(
INTEGER, nullable=False, server_default=text("0"), comment="乐观锁版本"
)
created_at: Mapped[datetime.datetime] = mapped_column(
DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)")
)
script_id: Mapped[Optional[str]] = mapped_column(CHAR(26))
versions_id: Mapped[Optional[str]] = mapped_column(
CHAR(26), comment="工作副本运行时可为空"
)
schedule_run_id: Mapped[Optional[str]] = mapped_column(CHAR(26))
parent_experiment_id: Mapped[Optional[str]] = mapped_column(CHAR(26))
parameters_json: Mapped[Optional[dict]] = mapped_column(JSON)
environment_json: Mapped[Optional[dict]] = mapped_column(JSON)
result_summary: Mapped[Optional[str]] = mapped_column(String(2000))
logs_object_id: Mapped[Optional[str]] = mapped_column(CHAR(26))
result_object_id: Mapped[Optional[str]] = mapped_column(CHAR(26))
started_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3))
finished_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3))
duration_ms: Mapped[Optional[int]] = mapped_column(BIGINT)
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))
class ExperimentMetrics(Base):
__tablename__ = "experiment_metrics"
__table_args__ = (
Index("idx_experiment_metrics_lookup", "experiment_id", "metric_name", "step_no"),
{"comment": "实验指标,支持筛选和曲线"},
)
metric_id: Mapped[int] = mapped_column(BIGINT, primary_key=True)
experiment_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
metric_name: Mapped[str] = mapped_column(String(128), nullable=False)
recorded_at: Mapped[datetime.datetime] = mapped_column(
DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)")
)
metric_value: Mapped[Optional[decimal.Decimal]] = mapped_column(
Double(asdecimal=True)
)
metric_text: Mapped[Optional[str]] = mapped_column(String(1000))
step_no: Mapped[Optional[int]] = mapped_column(BigInteger)
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))
class ExperimentResources(Base):
__tablename__ = "experiment_resources"
__table_args__ = (
Index("fk_experiment_resources_resource", "resource_id"),
{"comment": "实验与数据资源"},
)
experiment_id: Mapped[str] = mapped_column(CHAR(26), primary_key=True)
resource_id: Mapped[str] = mapped_column(CHAR(26), primary_key=True)
resource_role: Mapped[str] = mapped_column(
String(16),
primary_key=True,
server_default=text("'input'"),
comment="input/output",
)
created_at: Mapped[datetime.datetime] = mapped_column(
DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)")
)
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))
+1 -114
View File
@@ -1,126 +1,13 @@
import datetime
from typing import Optional
from sqlalchemy import BINARY, Index, String, Text, text
from sqlalchemy import Index, String, Text, text
from sqlalchemy.dialects.mysql import CHAR, DATETIME, INTEGER, TINYINT
from sqlalchemy.orm import Mapped, mapped_column
from common.db.base import Base
class RuntimeInstances(Base):
__tablename__ = "runtime_instances"
__table_args__ = (
Index("fk_runtime_started_by", "started_by"),
Index("idx_runtime_lease", "actual_state", "lease_expires_at"),
Index("idx_runtime_owner_state", "owner_user_id", "actual_state"),
Index(
"idx_runtime_workspace_state",
"workspace_id",
"runtime_type",
"actual_state",
),
{"comment": "Jupyter/未来 VS Code、OpenCode Runtime 实例"},
)
runtime_id: Mapped[str] = mapped_column(CHAR(26), primary_key=True)
workspace_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
runtime_type: Mapped[str] = mapped_column(
String(24), nullable=False, server_default=text("'jupyter'")
)
runtime_provider: Mapped[str] = mapped_column(
String(24), nullable=False, comment="process/docker/kubernetes"
)
proxy_base_path: Mapped[str] = mapped_column(String(512), nullable=False)
desired_state: Mapped[str] = mapped_column(
String(16), nullable=False, server_default=text("'running'")
)
actual_state: Mapped[str] = mapped_column(
String(24),
nullable=False,
server_default=text("'provisioning'"),
comment="provisioning/starting/running/unhealthy/stopping/stopped/failed",
)
state_version: Mapped[int] = mapped_column(
INTEGER, nullable=False, server_default=text("0"), comment="乐观锁版本"
)
started_by: Mapped[str] = mapped_column(CHAR(26), nullable=False)
created_at: Mapped[datetime.datetime] = mapped_column(
DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)")
)
updated_at: Mapped[datetime.datetime] = mapped_column(
DATETIME(fsp=3),
nullable=False,
server_default=text("CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3)"),
)
owner_user_id: Mapped[Optional[str]] = mapped_column(
CHAR(26), comment="为空表示 Workspace 级 Runtime"
)
runtime_ref: Mapped[Optional[str]] = mapped_column(
String(255), comment="PID/container ID/pod UID"
)
host_node: Mapped[Optional[str]] = mapped_column(String(255))
internal_url: Mapped[Optional[str]] = mapped_column(String(1000))
started_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3))
last_heartbeat_at: Mapped[Optional[datetime.datetime]] = mapped_column(
DATETIME(fsp=3)
)
lease_expires_at: Mapped[Optional[datetime.datetime]] = mapped_column(
DATETIME(fsp=3)
)
stopped_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3))
error_message: Mapped[Optional[str]] = mapped_column(Text)
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))
class EditSessions(Base):
__tablename__ = "edit_sessions"
__table_args__ = (
Index("fk_edit_sessions_workspace", "workspace_id"),
Index(
"idx_edit_sessions_object",
"storage_object_id",
"session_status",
"expires_at",
),
Index("idx_edit_sessions_runtime", "runtime_id", "session_status"),
Index("idx_edit_sessions_user", "user_id", "session_status"),
{"comment": "编辑会话审计"},
)
edit_session_id: Mapped[str] = mapped_column(CHAR(26), primary_key=True)
workspace_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
storage_object_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
user_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
lock_token_hash: Mapped[bytes] = mapped_column(
BINARY(32), nullable=False, comment="不保存原始 token"
)
session_status: Mapped[str] = mapped_column(
String(16),
nullable=False,
server_default=text("'active'"),
comment="active/closed/expired/failed",
)
started_at: Mapped[datetime.datetime] = mapped_column(
DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)")
)
last_heartbeat_at: Mapped[datetime.datetime] = mapped_column(
DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)")
)
expires_at: Mapped[datetime.datetime] = mapped_column(DATETIME(fsp=3), nullable=False)
runtime_id: Mapped[Optional[str]] = mapped_column(CHAR(26))
jupyter_session_id: Mapped[Optional[str]] = mapped_column(String(255))
ended_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3))
end_reason: Mapped[Optional[str]] = mapped_column(String(64))
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))
class WorkspaceOperations(Base):
__tablename__ = "workspace_operations"
__table_args__ = (
-32
View File
@@ -63,38 +63,6 @@ class Scripts(Base):
deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3))
class NotebookSnapshots(Base):
__tablename__ = "notebook_snapshots"
__table_args__ = (
Index("fk_snapshots_created_by", "created_by"),
Index("fk_snapshots_source_object", "source_object_id"),
Index("idx_snapshots_workspace_created", "workspace_id", "created_at"),
Index("uk_snapshots_artifact", "artifact_object_id", unique=True),
Index("uk_snapshots_script_hash", "script_id", "content_hash", unique=True),
{"comment": "Notebook 开发快照,append-only"},
)
snapshot_id: Mapped[str] = mapped_column(CHAR(26), primary_key=True)
workspace_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
script_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
source_object_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
artifact_object_id: Mapped[str] = mapped_column(CHAR(26), nullable=False)
snapshot_name: Mapped[str] = mapped_column(String(255), nullable=False)
content_hash: Mapped[str] = mapped_column(CHAR(64), nullable=False)
outputs_stripped: Mapped[int] = mapped_column(
TINYINT(1), nullable=False, server_default=text("1")
)
created_by: Mapped[str] = mapped_column(CHAR(26), nullable=False)
created_at: Mapped[datetime.datetime] = mapped_column(
DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)")
)
description: Mapped[Optional[str]] = mapped_column(String(1000))
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))
class Versions(Base):
__tablename__ = "versions"
__table_args__ = (
@@ -1,112 +0,0 @@
"""demo workspaces and users
Revision ID: 20260728_0002
Revises: 20260724_0001
Create Date: 2026-07-28
"""
from collections.abc import Sequence
from alembic import op
revision: str = "20260728_0002"
down_revision: str | Sequence[str] | None = "20260724_0001"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
ADMIN_ROLE = "0000000000HNQN4KM476QNKW1C"
DEVELOPER_ROLE = "00000000005RCQ4GPGBK3WZYMM"
MODEL_WORKSPACE = "00000000000BM630VT9ARVFZPC"
RISK_WORKSPACE = "0000000000AE0NC0V5T424KK86"
ZHANG = "0000000000RF6FG1SDBXG59S13"
LI = "0000000000H2QYCGPCWQM1JSGS"
WANG = "0000000000RWG40ESZPGJT629J"
ZHAO = "00000000004CQV7WASJA6N6FW4"
def upgrade() -> None:
op.execute(
f"""
INSERT INTO roles
(role_id, role_code, role_name, role_scope, is_builtin)
VALUES
('{ADMIN_ROLE}', 'admin', '管理员', 'workspace', 1),
('{DEVELOPER_ROLE}', 'developer', '开发人员', 'workspace', 1)
ON DUPLICATE KEY UPDATE
role_name = VALUES(role_name),
role_scope = VALUES(role_scope)
"""
)
op.execute(
f"""
INSERT INTO users
(user_id, username, display_name, password_hash, status, email)
VALUES
('{ZHANG}', 'admin-zhang', '张三', 'demo-login-disabled', 'active',
'zhangsan@example.local'),
('{LI}', 'admin-li', '李四', 'demo-login-disabled', 'active',
'lisi@example.local'),
('{WANG}', 'dev-wang', '王五', 'demo-login-disabled', 'active',
'wangwu@example.local'),
('{ZHAO}', 'dev-zhao', '赵六', 'demo-login-disabled', 'active',
'zhaoliu@example.local')
ON DUPLICATE KEY UPDATE
username = VALUES(username),
display_name = VALUES(display_name),
status = 'active',
email = VALUES(email)
"""
)
op.execute(
f"""
INSERT INTO workspaces
(workspace_id, workspace_code, workspace_name, active_root_uri,
quota_bytes, used_bytes, status, created_by, description,
artifact_bucket, artifact_prefix)
VALUES
('{MODEL_WORKSPACE}', 'model-dev', '模型开发 Workspace',
'file:///workspace/workspaces/model-dev', 0, 0, 'active',
'{ZHANG}', '模型开发与脚本调度', 'model-platform',
'workspaces/model-dev'),
('{RISK_WORKSPACE}', 'risk-validation', '风险验证 Workspace',
'file:///workspace/workspaces/risk-validation', 0, 0, 'active',
'{LI}', '风险模型验证与批处理', 'model-platform',
'workspaces/risk-validation')
ON DUPLICATE KEY UPDATE
workspace_name = VALUES(workspace_name),
active_root_uri = VALUES(active_root_uri),
status = 'active',
description = VALUES(description)
"""
)
values = []
for workspace_id in (MODEL_WORKSPACE, RISK_WORKSPACE):
for user_id, role_id in (
(ZHANG, ADMIN_ROLE),
(LI, ADMIN_ROLE),
(WANG, DEVELOPER_ROLE),
(ZHAO, DEVELOPER_ROLE),
):
values.append(
f"('{workspace_id}', '{user_id}', '{role_id}', 'active')"
)
op.execute(
"""
INSERT INTO workspace_members
(workspace_id, user_id, role_id, member_status)
VALUES
"""
+ ",\n".join(values)
+ """
ON DUPLICATE KEY UPDATE
role_id = VALUES(role_id),
member_status = 'active'
"""
)
def downgrade() -> None:
# 演示身份可能已产生业务数据,降级时保留,避免破坏外键引用。
pass
@@ -1,41 +0,0 @@
"""decouple schedule artifact visibility from version history
Revision ID: 20260728_0003
Revises: 20260728_0002
Create Date: 2026-07-28
"""
from collections.abc import Sequence
from alembic import op
import sqlalchemy as sa
from sqlalchemy.dialects import mysql
revision: str = "20260728_0003"
down_revision: str | Sequence[str] | None = "20260728_0002"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
op.add_column(
"versions",
sa.Column(
"schedule_hidden_at",
mysql.DATETIME(fsp=3),
nullable=True,
comment="从调度稳定版本列表移除的时间;不影响版本和运行历史",
),
)
op.create_index(
"idx_versions_schedule_visible",
"versions",
["workspace_id", "schedule_hidden_at", "created_at"],
unique=False,
)
def downgrade() -> None:
op.drop_index("idx_versions_schedule_visible", table_name="versions")
op.drop_column("versions", "schedule_hidden_at")
@@ -1,49 +0,0 @@
"""remove Redis-specific lock column naming
Revision ID: 20260730_0004
Revises: 20260728_0003
Create Date: 2026-07-30
"""
from collections.abc import Sequence
from alembic import context, op
import sqlalchemy as sa
revision: str = "20260730_0004"
down_revision: str | None = "20260728_0003"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def _column_names() -> set[str]:
inspector = sa.inspect(op.get_bind())
return {item["name"] for item in inspector.get_columns("edit_sessions")}
def upgrade() -> None:
if context.is_offline_mode():
return
columns = _column_names()
if "redis_lock_key" in columns and "lock_key" not in columns:
op.alter_column(
"edit_sessions",
"redis_lock_key",
new_column_name="lock_key",
existing_type=sa.String(length=512),
existing_nullable=False,
)
def downgrade() -> None:
if context.is_offline_mode():
return
columns = _column_names()
if "lock_key" in columns and "redis_lock_key" not in columns:
op.alter_column(
"edit_sessions",
"lock_key",
new_column_name="redis_lock_key",
existing_type=sa.String(length=512),
existing_nullable=False,
)