feat: 增加采集任务重新执行入口 (#63)

This commit is contained in:
chengma
2026-08-09 21:44:29 +08:00
parent 72821eaaef
commit eda3661604
12 changed files with 509 additions and 27 deletions
+84 -3
View File
@@ -34,6 +34,10 @@ class DuplicateTaskError(ValueError):
"""相同远程任务编号已经存在,不能覆盖。"""
class CollectRerunError(ValueError):
"""当前任务不满足重新采集条件。"""
def utc_now_iso() -> str:
"""返回精确到秒的 UTC ISO 8601 时间。"""
@@ -206,6 +210,81 @@ class TaskRepository:
connection.close()
return self._to_detail(row) if row is not None else None
def validate_collect_rerun(self, remote_task_id: str) -> TaskDetail:
"""校验任务能否重新采集,成功时返回任务详情。"""
connection = open_database(self._db_path)
try:
row = connection.execute(
"SELECT * FROM pdd_tasks WHERE remote_task_id = ?",
(remote_task_id,),
).fetchone()
if row is None:
raise CollectRerunError(f"任务 {remote_task_id} 不存在")
self._check_collect_rerun(connection, row)
finally:
connection.close()
return self._to_detail(row)
def prepare_collect_rerun(self, remote_task_id: str) -> TaskDetail:
"""事务内把已结束的采集任务恢复为待执行,保留旧结果。"""
now = utc_now_iso()
connection = open_database(self._db_path)
try:
with connection:
row = connection.execute(
"SELECT * FROM pdd_tasks WHERE remote_task_id = ?",
(remote_task_id,),
).fetchone()
if row is None:
raise CollectRerunError(f"任务 {remote_task_id} 不存在")
self._check_collect_rerun(connection, row)
connection.execute(
"UPDATE pdd_tasks SET status = 'claimed',"
" current_step = 'rerun_requested', finished_at = NULL,"
" last_error_code = NULL, last_error_message = NULL,"
" updated_at = ? WHERE id = ?",
(now, row["id"]),
)
finally:
connection.close()
task = self.get_task(remote_task_id)
assert task is not None
return task
@staticmethod
def _check_collect_rerun(
connection: sqlite3.Connection, row: sqlite3.Row
) -> None:
"""检查重新采集的类型、终态和 Outbox 约束。"""
if row["task_type"] != TaskType.COLLECT.value:
raise CollectRerunError("采购任务不能重新执行,以免重复下单")
allowed_statuses = {
TaskStatus.SUCCEEDED.value,
TaskStatus.FAILED.value,
TaskStatus.CANCELLED.value,
TaskStatus.MANUAL_REVIEW.value,
}
if row["status"] not in allowed_statuses:
status_name = {
TaskStatus.CLAIMED.value: "待执行",
TaskStatus.RUNNING.value: "执行中",
TaskStatus.RESULT_PENDING.value: "结果待提交",
TaskStatus.RETRY_WAIT.value: "等待重试",
}.get(row["status"], row["status"])
raise CollectRerunError(f"任务当前为“{status_name}”,不能重新采集")
unsent_count = int(
connection.execute(
"SELECT COUNT(*) FROM outbox_events"
" WHERE task_id = ? AND status != 'sent'",
(row["id"],),
).fetchone()[0]
)
if unsent_count:
raise CollectRerunError("任务仍有未发送的结果,请先完成提交")
def start_collect_run(
self, remote_task_id: str, device_address: str
) -> StartedTaskRun:
@@ -284,22 +363,24 @@ class TaskRepository:
"pdd_data": pdd_data,
}
idempotency_key = f"{remote_task_id}:{attempt_id}:result-v1"
result_json = json.dumps(pdd_data, ensure_ascii=False)
connection.execute(
"UPDATE pdd_tasks SET status = 'result_pending',"
" current_step = 'submit_result', pdd_data = ?, goods_id = ?,"
" title = ?, price_cent = ?, finished_at = ?, updated_at = ?"
" WHERE id = ?",
(
json.dumps(pdd_data, ensure_ascii=False),
result_json,
pdd_data.get("goods_id"), pdd_data.get("title"),
self._summary_price(pdd_data), now, now, task["id"],
),
)
connection.execute(
"UPDATE task_runs SET run_status = 'succeeded',"
" current_step = 'submit_result', finished_at = ?, updated_at = ?"
" current_step = 'submit_result', result_data = ?,"
" finished_at = ?, updated_at = ?"
" WHERE attempt_id = ?",
(now, now, attempt_id),
(result_json, now, now, attempt_id),
)
cursor = connection.execute(
"INSERT INTO outbox_events (task_id, event_type, idempotency_key,"