From 8dc03df95b19e69aa6e95f550c94d557447f18e0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=BB=84=E8=87=AA=E8=BE=BE?= <980324341@qq.com> Date: Fri, 11 Sep 2026 16:09:42 +0800 Subject: [PATCH] =?UTF-8?q?feat(task):=20=E5=AE=A2=E6=88=B7=E7=AB=AF?= =?UTF-8?q?=E5=B4=A9=E6=BA=83=E6=81=A2=E5=A4=8D=E6=8E=A5=E5=8F=A3=20POST?= =?UTF-8?q?=20/heartbeat/{taskId}/interrupted=E2=80=94=E2=80=94=E5=AE=A2?= =?UTF-8?q?=E6=88=B7=E7=AB=AF=E9=87=8D=E5=90=AF=E5=8F=91=E7=8E=B0=E4=B8=8A?= =?UTF-8?q?=E6=AC=A1=E8=BF=9B=E7=A8=8B=E5=B4=A9=E6=BA=83=E6=97=B6=E7=AB=8B?= =?UTF-8?q?=E5=8D=B3=E6=8A=8A=E6=AE=8B=E7=95=99=20RUNNING=20=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E6=A0=87=E7=BB=88=E6=80=81=EF=BC=8C=E6=9B=BF=E4=BB=A3?= =?UTF-8?q?=E6=9C=80=E9=95=BF=E7=AD=89=2030=20=E5=88=86=E9=92=9F=E7=9A=84?= =?UTF-8?q?=E5=BF=83=E8=B7=B3=E8=B6=85=E6=97=B6=E5=85=9C=E5=BA=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../controller/TaskHeartbeatController.java | 13 +++++ .../task/service/TaskHeartbeatService.java | 50 +++++++++++++++++++ 2 files changed, 63 insertions(+) diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/controller/TaskHeartbeatController.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/controller/TaskHeartbeatController.java index 94c70050..3ad6d6f3 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/controller/TaskHeartbeatController.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/controller/TaskHeartbeatController.java @@ -33,4 +33,17 @@ public class TaskHeartbeatController { @Valid @RequestBody(required = false) TaskHeartbeatRequest request) { return ApiResponse.success(taskHeartbeatService.heartbeat(taskId, request)); } + + @PostMapping("/{taskId}/interrupted") + @Operation( + summary = "上报客户端异常中断", + description = "客户端重启后发现自己上次进程崩溃时调用:将仍在 RUNNING 的任务立即标为终态" + + "(file_task→FAILED,brand_crawl_tasks→cancelled),替代最长 30 分钟的 stale 兜底。幂等,非 RUNNING 状态不修改。") + public ApiResponse interrupted( + @Parameter(description = "任务 ID", required = true, example = "200") + @PathVariable Long taskId, + @RequestBody(required = false) TaskHeartbeatRequest request) { + String reason = request == null ? null : request.getPhase(); + return ApiResponse.success(taskHeartbeatService.markInterrupted(taskId, reason)); + } } diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskHeartbeatService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskHeartbeatService.java index dfb2973f..fdaf48d1 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskHeartbeatService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskHeartbeatService.java @@ -117,6 +117,56 @@ public class TaskHeartbeatService { return fileTaskMapper.selectOne(fileQuery); } + /** + * 客户端崩溃恢复:客户端重启后发现自己上次进程级崩溃(journal 残留), + * 调此接口把上次还在 RUNNING 的任务立即标终态,替代最长等 30 分钟的 stale 兜底。 + * file_task RUNNING → FAILED;brand_crawl_tasks running → cancelled;其余状态不动(幂等)。 + */ + public TaskHeartbeatVo markInterrupted(Long taskId, String reason) { + if (taskId == null || taskId <= 0) { + log.warn("[task-interrupted] ignored invalid taskId={}", taskId); + return TaskHeartbeatVo.notAlive(null, null, "invalid taskId"); + } + String safeReason = (reason == null || reason.isBlank()) + ? "客户端异常中断,任务已自动失败" + : "客户端异常中断: " + reason.trim(); + FileTaskEntity fileTask = selectFileTask(taskId); + if (fileTask != null) { + String status = fileTask.getStatus(); + int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper() + .eq(FileTaskEntity::getId, taskId) + .eq(FileTaskEntity::getStatus, STATUS_RUNNING) + .set(FileTaskEntity::getStatus, "FAILED") + .set(FileTaskEntity::getErrorMessage, safeReason) + .set(FileTaskEntity::getFinishedAt, LocalDateTime.now())); + if (updated > 0) { + log.warn("[task-interrupted] file task marked failed by client restart taskId={} moduleType={} reason={}", + taskId, fileTask.getModuleType(), safeReason); + return TaskHeartbeatVo.notAlive(fileTask.getModuleType(), "FAILED", "marked failed"); + } + log.info("[task-interrupted] file task not in RUNNING, skipped taskId={} status={}", taskId, status); + return TaskHeartbeatVo.notAlive(fileTask.getModuleType(), status, "task is not running"); + } + BrandCrawlTaskEntity brandTask = selectBrandTask(taskId); + if (brandTask != null) { + String status = brandTask.getStatus(); + if ("running".equalsIgnoreCase(status) || "pending".equalsIgnoreCase(status)) { + brandTask.setStatus("cancelled"); + brandTask.setErrorMessage(safeReason); + int updated = brandCrawlTaskMapper.updateById(brandTask); + if (updated > 0) { + log.warn("[task-interrupted] brand task marked cancelled by client restart taskId={} reason={}", + taskId, safeReason); + return TaskHeartbeatVo.notAlive(MODULE_BRAND, "cancelled", "marked cancelled"); + } + } + log.info("[task-interrupted] brand task not in running/pending, skipped taskId={} status={}", taskId, status); + return TaskHeartbeatVo.notAlive(MODULE_BRAND, status, "task is not running"); + } + log.warn("[task-interrupted] task not found taskId={}", taskId); + return TaskHeartbeatVo.notAlive(null, null, "task not found"); + } + private BrandCrawlTaskEntity selectBrandTask(Long taskId) { LambdaQueryWrapper brandQuery = new LambdaQueryWrapper() .eq(BrandCrawlTaskEntity::getId, taskId)