feat: improve schedule cleanup and handover

This commit is contained in:
Winnie
2026-08-20 19:12:44 +08:00
parent ae3b6d63e2
commit bc62034241
7 changed files with 306 additions and 63 deletions
+6
View File
@@ -97,6 +97,12 @@ class WorkflowVersionRequest(StrictModel):
workflow_version: int = Field(ge=1)
class DeleteScheduleNodeRequest(WorkflowVersionRequest):
"""删除节点时可显式确认一并清理其已经结束的运行记录。"""
delete_execution_history: bool = False
class CreateScheduleNodeRequest(StrictModel):
workflow_version: int = Field(ge=1)
node_key: str = Field(
+145 -8
View File
@@ -17,13 +17,15 @@ from common.db.models import (
ScheduleEdges,
ScheduleNodeRuns,
ScheduleNodes,
ScheduleRuns,
Schedules,
Scripts,
StorageObjects,
Versions,
)
from common.ids import new_ulid
from croniter import CroniterBadCronError, croniter
from fastapi import APIRouter, Depends, HTTPException, Query, status
from fastapi import APIRouter, Depends, HTTPException, Query, Request, status
from sqlalchemy import delete, func, or_, select
from sqlalchemy.ext.asyncio import AsyncSession
@@ -37,14 +39,18 @@ from backend.schedule_schemas import (
CreateScheduleNodeRequest,
CreateScheduleRequest,
CronPreviewRequest,
DeleteScheduleNodeRequest,
UpdateScheduleEdgeRequest,
UpdateScheduleNodeRequest,
UpdateScheduleRequest,
WorkflowVersionRequest,
)
from backend.services.storage import soft_delete_object
router = APIRouter(tags=["schedules"])
_ACTIVE_RUN_STATUSES = ("queued", "running")
def _iso(value: datetime | None) -> str | None:
if value is None:
@@ -60,6 +66,48 @@ def _mysql_utc(value: datetime) -> datetime:
return value.astimezone(UTC).replace(tzinfo=None)
def _execution_artifact_ids(
items: list[ScheduleRuns | ScheduleNodeRuns],
) -> set[str]:
"""收集运行日志和结果产物,供删除记录时一并移入回收站。"""
return {
storage_object_id
for item in items
for storage_object_id in (item.logs_object_id, item.result_object_id)
if storage_object_id
}
async def _delete_execution_artifacts(
storage_object_ids: set[str],
request: Request,
session: AsyncSession,
) -> None:
"""软删除可删除的运行产物。
不可变对象是运行审计原件,存储层不允许移动或删除它们。调用方随后会
删除运行记录本身,因此不可变原件不会再通过该调度节点暴露;保留原件也
不应阻塞节点或调度方案的删除。
"""
if not storage_object_ids:
return
mutable_storage_object_ids = set(
(
await session.scalars(
select(StorageObjects.storage_object_id).where(
StorageObjects.storage_object_id.in_(storage_object_ids),
StorageObjects.is_immutable == 0,
)
)
).all()
)
for storage_object_id in sorted(mutable_storage_object_ids):
await soft_delete_object(storage_object_id, request, session)
def _timezone(value: str) -> ZoneInfo:
try:
return ZoneInfo(value)
@@ -740,6 +788,7 @@ async def update_schedule(
async def delete_schedule(
schedule_id: str,
payload: WorkflowVersionRequest,
request: Request,
context: RequestContext = Depends(request_context),
session: AsyncSession = Depends(database_session),
) -> dict[str, Any]:
@@ -750,6 +799,64 @@ async def delete_schedule(
for_update=True,
)
require_workflow_version(item, payload.workflow_version)
# 正在执行的 DAG 依赖运行快照和日志对象。此时删除会让执行器无法安全收尾,
# 因此要求先等待运行结束,避免影响现有运行功能。
active_run_id = await session.scalar(
select(ScheduleRuns.run_id)
.where(
ScheduleRuns.schedule_id == schedule_id,
ScheduleRuns.run_status.in_(_ACTIVE_RUN_STATUSES),
)
.limit(1)
)
if active_run_id is not None:
raise HTTPException(
status.HTTP_409_CONFLICT,
"schedule has active runs; wait for completion before deletion",
)
schedule_runs = list(
(
await session.scalars(
select(ScheduleRuns).where(ScheduleRuns.schedule_id == schedule_id)
)
).all()
)
run_ids = [run.run_id for run in schedule_runs]
node_runs = (
list(
(
await session.scalars(
select(ScheduleNodeRuns).where(
ScheduleNodeRuns.run_id.in_(run_ids)
)
)
).all()
)
if run_ids
else []
)
# 调度删除会清理节点、边和所有运行记录;运行日志/结果文件同时移入回收站。
await _delete_execution_artifacts(
_execution_artifact_ids([*schedule_runs, *node_runs]),
request,
session,
)
if run_ids:
await session.execute(
delete(ScheduleNodeRuns).where(ScheduleNodeRuns.run_id.in_(run_ids))
)
await session.execute(
delete(ScheduleRuns).where(ScheduleRuns.run_id.in_(run_ids))
)
await session.execute(
delete(ScheduleEdges).where(ScheduleEdges.schedule_id == schedule_id)
)
await session.execute(
delete(ScheduleNodes).where(ScheduleNodes.schedule_id == schedule_id)
)
item.enabled = 0
item.next_run_at = None
item.deleted_at = _mysql_utc(datetime.now(UTC))
@@ -884,7 +991,8 @@ async def update_schedule_node(
async def delete_schedule_node(
schedule_id: str,
node_id: str,
payload: WorkflowVersionRequest,
payload: DeleteScheduleNodeRequest,
request: Request,
context: RequestContext = Depends(request_context),
session: AsyncSession = Depends(database_session),
) -> dict[str, Any]:
@@ -903,15 +1011,44 @@ async def delete_schedule_node(
)
if node is None:
raise HTTPException(status.HTTP_404_NOT_FOUND, "schedule node not found")
has_runs = await session.scalar(
select(func.count())
.select_from(ScheduleNodeRuns)
.where(ScheduleNodeRuns.node_id == node_id)
active_run_id = await session.scalar(
select(ScheduleRuns.run_id)
.join(ScheduleNodeRuns, ScheduleNodeRuns.run_id == ScheduleRuns.run_id)
.where(
ScheduleRuns.schedule_id == schedule_id,
ScheduleRuns.run_status.in_(_ACTIVE_RUN_STATUSES),
ScheduleNodeRuns.node_id == node_id,
)
.limit(1)
)
if int(has_runs or 0):
if active_run_id is not None:
raise HTTPException(
status.HTTP_409_CONFLICT,
"a node with execution history cannot be deleted",
"node has active execution; wait for the run to finish before deletion",
)
node_runs = list(
(
await session.scalars(
select(ScheduleNodeRuns).where(ScheduleNodeRuns.node_id == node_id)
)
).all()
)
if node_runs and not payload.delete_execution_history:
raise HTTPException(
status.HTTP_409_CONFLICT,
detail={
"code": "node_execution_history_exists",
"message": "node has execution history; confirmation required",
},
)
if node_runs:
await _delete_execution_artifacts(
_execution_artifact_ids(node_runs),
request,
session,
)
await session.execute(
delete(ScheduleNodeRuns).where(ScheduleNodeRuns.node_id == node_id)
)
await session.execute(
delete(ScheduleEdges).where(