feat: 安全移除本地任务记录 (#108)

This commit is contained in:
chengma
2026-08-10 17:23:24 +08:00
parent 9d227b431a
commit b2721071c9
10 changed files with 619 additions and 30 deletions
+98 -2
View File
@@ -7,7 +7,7 @@ import json
import sqlite3
from datetime import datetime, timezone
from pathlib import Path
from typing import Dict, List, Optional, Tuple, Union
from typing import Dict, Iterable, List, Optional, Tuple, Union
from uuid import uuid4
from .db import initialize_database, open_database
@@ -39,6 +39,10 @@ class CollectRerunError(ValueError):
"""当前任务不满足重新采集条件。"""
class TaskRemovalError(ValueError):
"""勾选任务不满足从普通列表移除的安全条件。"""
def utc_now_iso() -> str:
"""返回精确到秒的 UTC ISO 8601 时间。"""
@@ -145,6 +149,98 @@ class TaskRepository:
connection.close()
return int(row[0])
def remove_tasks_from_list(self, remote_task_ids: Iterable[str]) -> int:
"""把符合安全条件的终态任务从普通列表中软移除。
全部任务会在同一个事务中完成校验和更新;任一任务不安全时,
所有任务都保持原样。执行记录和 Outbox 不会被删除或修改。
"""
task_ids = tuple(
dict.fromkeys(
str(value).strip()
for value in remote_task_ids
if value is not None and str(value).strip()
)
)
if not task_ids:
raise TaskRemovalError("没有可删除的任务,请重新勾选。")
placeholders = ", ".join("?" for _ in task_ids)
terminal_statuses = {
TaskStatus.SUCCEEDED.value,
TaskStatus.FAILED.value,
TaskStatus.CANCELLED.value,
}
now = utc_now_iso()
connection = open_database(self._db_path)
try:
with connection:
rows = connection.execute(
"SELECT id, remote_task_id, status, removed_at FROM pdd_tasks"
f" WHERE remote_task_id IN ({placeholders})",
task_ids,
).fetchall()
found = {row["remote_task_id"]: row for row in rows}
missing = [task_id for task_id in task_ids if task_id not in found]
if missing:
raise TaskRemovalError(
f"本地找不到任务 {missing[0]},请刷新列表后重试。"
)
removed = next(
(row for row in rows if row["removed_at"] is not None), None
)
if removed is not None:
raise TaskRemovalError(
f"任务 {removed['remote_task_id']} 已不在普通列表,请刷新后重试。"
)
non_terminal = next(
(row for row in rows if row["status"] not in terminal_statuses),
None,
)
if non_terminal is not None:
raise TaskRemovalError(
f"任务 {non_terminal['remote_task_id']} 尚未结束,不能删除。"
)
irreversible = connection.execute(
"SELECT t.remote_task_id FROM task_runs r"
" JOIN pdd_tasks t ON t.id = r.task_id"
f" WHERE t.remote_task_id IN ({placeholders})"
" AND r.irreversible_action_at IS NOT NULL LIMIT 1",
task_ids,
).fetchone()
if irreversible is not None:
raise TaskRemovalError(
f"任务 {irreversible['remote_task_id']} 已进入不可逆阶段,"
"必须保留在列表中核对订单。"
)
unsent = connection.execute(
"SELECT t.remote_task_id, o.status FROM outbox_events o"
" JOIN pdd_tasks t ON t.id = o.task_id"
f" WHERE t.remote_task_id IN ({placeholders})"
" AND o.status <> 'sent' LIMIT 1",
task_ids,
).fetchone()
if unsent is not None:
raise TaskRemovalError(
f"任务 {unsent['remote_task_id']} 还有未发送完成的数据,"
"请先重新上报。"
)
cursor = connection.execute(
"UPDATE pdd_tasks SET removed_at = ?, updated_at = ?"
f" WHERE remote_task_id IN ({placeholders})"
" AND removed_at IS NULL",
(now, now, *task_ids),
)
if cursor.rowcount != len(task_ids):
raise TaskRemovalError("任务列表已发生变化,请刷新后重试。")
finally:
connection.close()
return len(task_ids)
def get_task(self, remote_task_id: str) -> Optional[TaskDetail]:
"""按稳定远程编号读取完整任务;不存在时返回 None。"""
@@ -1594,7 +1690,7 @@ class TaskRepository:
@staticmethod
def _build_where(filters: TaskFilters) -> Tuple[str, List[object]]:
clauses = []
clauses = ["removed_at IS NULL"]
parameters: List[object] = []
if filters.task_type is not None:
clauses.append("task_type = ?")