From 967eedcad7855b5e45fee28523c7e1cbdf7d512d Mon Sep 17 00:00:00 2001 From: super <2903208875@qq.com> Date: Mon, 8 Jun 2026 18:31:55 +0800 Subject: [PATCH] =?UTF-8?q?=E6=8F=90=E4=BA=A4=E6=9B=B4=E6=96=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../docs/stale-task-finalize-guidelines.md | 79 +++++++++++ .../service/AppearancePatentTaskService.java | 78 ++++++----- .../model/entity/PriceTrackLoopRunEntity.java | 3 + .../service/PriceTrackLoopRunService.java | 108 +++++++++++---- .../service/SimilarAsinTaskService.java | 79 ++++++----- .../service/PriceTrackLoopRunServiceTest.java | 129 ++++++++++++++++++ .../brand/components/BrandPriceTrackTab.vue | 16 ++- 7 files changed, 394 insertions(+), 98 deletions(-) create mode 100644 backend-java/docs/stale-task-finalize-guidelines.md create mode 100644 backend-java/src/test/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunServiceTest.java diff --git a/backend-java/docs/stale-task-finalize-guidelines.md b/backend-java/docs/stale-task-finalize-guidelines.md new file mode 100644 index 0000000..9f2af1b --- /dev/null +++ b/backend-java/docs/stale-task-finalize-guidelines.md @@ -0,0 +1,79 @@ +# Stale 任务收尾约束(公共心跳模块) + +> 适用范围:所有使用 [TaskHeartbeatService](../src/main/java/com/nanri/aiimage/modules/task/service/TaskHeartbeatService.java) 的 file_task 模块 +> (PRODUCT_RISK_RESOLVE / PRICE_TRACK / SHOP_MATCH / PATROL_DELETE / QUERY_ASIN / +> APPEARANCE_PATENT / SIMILAR_ASIN / DELETE_BRAND)。 + +## 历史故障(2026-06-08) + +SIMILAR_ASIN taskId=13102 的 file_job 长期持有 RUNNING 等待 Python 凑齐 +batchSize 行再发 Coze。Python 端 13:20 异常掉线,但任务到下午 16 点仍卡在 12% 不前进, +30 分钟兜底失效。 + +链路: + +1. [resetStuckJobs](../src/main/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorker.java) + 每 60 秒扫一次,把 RUNNING 超 30 分钟的 file_job 改回 PENDING。 +2. worker 立刻把它捞起来重跑 `processResultFileJob`。 +3. 重跑分支命中 "等 Python 上传更多 row" → 调用 `touchJavaSideTaskActivity`, + 把 `file_task.updated_at` 刷成 now。 +4. `finalizeStaleTasks` 用 `lt(updated_at, threshold)` 预过滤,于是这个任务 + 永远进不到候选集,Redis 心跳判定那一步根本走不到。 + +→ Python 已断线 N 小时,任务仍在循环,永不收尾。 + +## 约束 + +### 1. stale-finalize 必须以 Redis 心跳为主信号 + +`TaskCacheService.touchTaskHeartbeat` 写的 Redis key 仅在 **Python 真实活动** +(心跳接口、上传 result chunk、提交 row)时刷新。Java 内部循环不应触碰它。 + +stale 判定时应: + +- 以 Redis heartbeat 为主信号; +- 仅在 Redis 不可用 / 缺失时回落到 `file_task.updated_at`; +- 不要再用 `lt(file_task.updated_at, threshold)` 做预过滤 + (除非该模块 100% 不存在长期 RUNNING 等待循环)。 + +参考实现:[SimilarAsinTaskService.finalizeStaleTasks / isHeartbeatStale](../src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java) +和 [AppearancePatentTaskService](../src/main/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTaskService.java) 中 P1-7 注释段。 + +### 2. 不要在"等 Python"分支刷 file_task.updated_at + +如果模块的 file_job 会持有 RUNNING 状态等待外部上传更多数据, +**这条等待路径上不允许调 `setUpdatedAt(now)` 类方法**。 + +否则会触发 [resetStuckJobs](../src/main/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorker.java) +→ 重跑 → 续命 → 再 reset 的死循环。 + +如必须刷 `updated_at` 推动其他逻辑,应迁移到独立字段(例如 +`file_job.updated_at`、Redis pendingScope 字段),与 stale 判定信号解耦。 + +### 3. 新增/修改 stale-finalize 流程的 PR 检查清单 + +- [ ] stale 判定的主信号是 Redis 心跳(不是 `file_task.updated_at`) +- [ ] 列出该模块所有刷 `file_task.updated_at` 的代码位置, + 其中没有任何一处会被 `processResultFileJob` 在"等 Python"分支重复触发 +- [ ] 如果 file_job 会持有长 RUNNING(pull-style 工作流),明确写明 stuck-reset + 重跑后的幂等行为(不会续命 updated_at) +- [ ] 给 finalize 路径补一个 debug 接口(参考 + [DebugTaskRecoveryController](../src/main/java/com/nanri/aiimage/modules/debug/controller/DebugTaskRecoveryController.java)), + 绕过过滤直接对单 task 强制收尾,便于线上应急 + +## 当前各模块状态(2026-06-08 复核) + +| 模块 | stale 入口 | 是否曾出现自我续命 | 是否需要改 | +|---|---|---|---| +| SIMILAR_ASIN | `SimilarAsinTaskService.finalizeStaleTasks` | 是(已修) | 已切 Redis-first | +| APPEARANCE_PATENT | `AppearancePatentTaskService.finalizeStaleTasks` | 是(已修) | 已切 Redis-first | +| PRODUCT_RISK_RESOLVE | `DeleteBrandStaleTaskService.failStaleProductRiskResolveTasks` | 否 | 当前不需要 | +| PRICE_TRACK | `DeleteBrandStaleTaskService.failStalePriceTrackTasks` | 否 | 当前不需要 | +| SHOP_MATCH | `DeleteBrandStaleTaskService.failStaleShopMatchTasks` | 否 | 当前不需要 | +| PATROL_DELETE | `DeleteBrandStaleTaskService.failStalePatrolDeleteTasks` | 否 | 当前不需要 | +| QUERY_ASIN | `DeleteBrandStaleTaskService.failStaleQueryAsinTasks` | 否 | 当前不需要 | +| DELETE_BRAND | `DeleteBrandStaleTaskService.failStaleDeleteBrandTasks` | 否 | 当前不需要 | +| BRAND | `BrandTaskService.failStaleRunningTasks`(独立心跳源) | 否 | 当前不需要 | + +下次新增"等 row 凑批" / "pull-style 长 RUNNING file_job"型功能时, +请先读本约束。 diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTaskService.java index 10c8cd9..431b1b3 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTaskService.java @@ -569,42 +569,20 @@ public class AppearancePatentTaskService { } public void finalizeStaleTasks() { - if (transactionManager != null) { - LocalDateTime now = LocalDateTime.now(); - LocalDateTime threshold = now.minusMinutes(Math.max(1, properties.getStaleTimeoutMinutes())); - long thresholdMillis = threshold.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli(); - List tasks = listStaleFinalizeCandidates(threshold); - for (FileTaskEntity task : tasks) { - if (!isOwnerCurrent(ownerFromTask(task))) { - continue; - } - long heartbeatMillis = taskCacheService.getTaskHeartbeatMillis(task.getId()); - if (heartbeatMillis > thresholdMillis) { - continue; - } - TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(task.getId(), 0L); - if (lockHandle == null) { - continue; - } - try (lockHandle) { - inNewTransaction(() -> { - finalizeStaleTask(task.getId(), "Python interrupted before uploading final appearance patent result"); - return null; - }); - } - } - return; - } - LocalDateTime now = LocalDateTime.now(); - LocalDateTime threshold = now.minusMinutes(Math.max(1, properties.getStaleTimeoutMinutes())); - long thresholdMillis = threshold.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli(); - List tasks = listStaleFinalizeCandidates(threshold); + // P1-7:stale 判定改用 Redis heartbeat 作为主信号。 + // 历史实现先按 file_task.updated_at 过滤再二次校验 Redis 心跳, + // 但 file_task.updated_at 会被 Java 端在 "等 Python 上传" 分支自我续命 + // (resetStuckJobs 把 file_job 改回 PENDING → 重跑 processResultFileJob → touchJavaSideTaskActivity)。 + // Redis heartbeat 仅在 Python 心跳/上传 result 时写入,不会被 Java 内部循环刷新。 + long thresholdMillis = LocalDateTime.now() + .minusMinutes(Math.max(1, properties.getStaleTimeoutMinutes())) + .atZone(ZoneId.systemDefault()).toInstant().toEpochMilli(); + List tasks = listStaleFinalizeCandidates(); for (FileTaskEntity task : tasks) { if (!isOwnerCurrent(ownerFromTask(task))) { continue; } - long heartbeatMillis = taskCacheService.getTaskHeartbeatMillis(task.getId()); - if (heartbeatMillis > thresholdMillis) { + if (!isHeartbeatStale(task, thresholdMillis)) { continue; } TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(task.getId(), 0L); @@ -612,7 +590,14 @@ public class AppearancePatentTaskService { continue; } try (lockHandle) { - finalizeStaleTask(task.getId(), "Python interrupted before uploading final appearance patent result"); + if (transactionManager != null) { + inNewTransaction(() -> { + finalizeStaleTask(task.getId(), "Python interrupted before uploading final appearance patent result"); + return null; + }); + } else { + finalizeStaleTask(task.getId(), "Python interrupted before uploading final appearance patent result"); + } } } } @@ -643,15 +628,38 @@ public class AppearancePatentTaskService { } } - private List listStaleFinalizeCandidates(LocalDateTime threshold) { + private List listStaleFinalizeCandidates() { + // P1-7:不再用 file_task.updated_at 预过滤,由调用方按 Redis heartbeat 判 stale。 + // APPEARANCE_PATENT 在线 RUNNING 任务量级有限(限 200),全扫成本可忽略。 return fileTaskMapper.selectList(new LambdaQueryWrapper() .eq(FileTaskEntity::getModuleType, MODULE_TYPE) .eq(FileTaskEntity::getStatus, STATUS_RUNNING) - .lt(FileTaskEntity::getUpdatedAt, threshold) .orderByAsc(FileTaskEntity::getUpdatedAt) .last("limit 200")); } + /** + * P1-7:判断 APPEARANCE_PATENT 任务是否已 stale。 + * 主信号:Redis heartbeat(仅 Python 心跳/上传 result 时刷新,不会被 Java 内部循环续命)。 + * Redis 命中:以 heartbeatMillis 与 thresholdMillis 比较。 + * Redis 缺失/降级(返回 0):回落到 file_task.updated_at,避免 Redis 故障时全量误杀。 + */ + private boolean isHeartbeatStale(FileTaskEntity task, long thresholdMillis) { + if (task == null || task.getId() == null) { + return false; + } + long heartbeatMillis = taskCacheService.getTaskHeartbeatMillis(task.getId()); + if (heartbeatMillis > 0L) { + return heartbeatMillis <= thresholdMillis; + } + LocalDateTime updatedAt = task.getUpdatedAt(); + if (updatedAt == null) { + return false; + } + long updatedMillis = updatedAt.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli(); + return updatedMillis <= thresholdMillis; + } + private SubmitContext persistSubmittedChunk(Long taskId, AppearancePatentSubmitResultRequest request) { FileTaskEntity task = fileTaskMapper.selectById(taskId); if (task == null || !MODULE_TYPE.equals(task.getModuleType())) { diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/model/entity/PriceTrackLoopRunEntity.java b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/model/entity/PriceTrackLoopRunEntity.java index 4a55547..f25352e 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/model/entity/PriceTrackLoopRunEntity.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/model/entity/PriceTrackLoopRunEntity.java @@ -1,6 +1,8 @@ package com.nanri.aiimage.modules.pricetrack.model.entity; +import com.baomidou.mybatisplus.annotation.FieldStrategy; import com.baomidou.mybatisplus.annotation.IdType; +import com.baomidou.mybatisplus.annotation.TableField; import com.baomidou.mybatisplus.annotation.TableId; import com.baomidou.mybatisplus.annotation.TableName; import lombok.Data; @@ -24,6 +26,7 @@ public class PriceTrackLoopRunEntity { private String asinFilesJson; private String countryCodesJson; private String shopsJson; + @TableField(value = "active_task_id", updateStrategy = FieldStrategy.ALWAYS) private Long activeTaskId; private Boolean stopRequested; private String errorMessage; diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunService.java index 5ce71eb..1696524 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunService.java @@ -108,7 +108,7 @@ public class PriceTrackLoopRunService { } int shopIndex = Math.max(0, entity.getCurrentShopIndex() == null ? 0 : entity.getCurrentShopIndex()); if (shopIndex >= items.size()) { - advanceAfterSuccessfulRound(entity, items.size()); + advanceAfterSuccessfulRound(entity); if (isTerminal(entity.getStatus())) { PriceTrackLoopRunDispatchVo vo = new PriceTrackLoopRunDispatchVo(); vo.setLoopRun(toVo(entity)); @@ -175,6 +175,8 @@ public class PriceTrackLoopRunService { entity.setActiveTaskId(childTaskId); entity.setUpdatedAt(LocalDateTime.now()); loopRunMapper.updateById(entity); + log.info("[price-track-loop] child bound loopRunId={} childTaskId={} roundIndex={} shopIndex={}", + loopRunId, childTaskId, roundIndex, shopIndex); } @Transactional @@ -182,12 +184,7 @@ public class PriceTrackLoopRunService { if (childTaskId == null || childTaskId <= 0) { return; } - List loops = loopRunMapper.selectList(new LambdaQueryWrapper() - .eq(PriceTrackLoopRunEntity::getActiveTaskId, childTaskId) - .last("limit 5")); - for (PriceTrackLoopRunEntity entity : loops) { - handleChildFinished(entity, childTaskId, false); - } + forEachLoopForChildTask(childTaskId, entity -> handleChildFinished(entity, childTaskId, false)); } @Transactional @@ -229,16 +226,22 @@ public class PriceTrackLoopRunService { if (!belongsToLoop(task, entity.getId())) { throw new BusinessException("子任务不属于当前循环"); } + ChildTaskContext ctx = parseChildTaskContext(task); + if (ctx == null) { + log.warn("[price-track-loop] child context missing loopRunId={} childTaskId={}", entity.getId(), childTaskId); + return; + } + boolean stateMatchesChild = ctx.roundIndex().equals(entity.getCurrentRound()) + && ctx.shopIndex().equals(entity.getCurrentShopIndex()); + if (!stateMatchesChild && entity.getActiveTaskId() != null && entity.getActiveTaskId().equals(childTaskId)) { + log.warn("[price-track-loop] child context drift detected loopRunId={} childTaskId={} loopRound={} loopShopIndex={} childRound={} childShopIndex={}", + entity.getId(), childTaskId, entity.getCurrentRound(), entity.getCurrentShopIndex(), ctx.roundIndex(), ctx.shopIndex()); + } if (entity.getActiveTaskId() != null && !entity.getActiveTaskId().equals(childTaskId)) { return; } - if (entity.getActiveTaskId() == null) { - ChildTaskContext ctx = parseChildTaskContext(task); - if (ctx == null - || !ctx.roundIndex().equals(entity.getCurrentRound()) - || !ctx.shopIndex().equals(entity.getCurrentShopIndex())) { - return; - } + if (entity.getActiveTaskId() == null && !stateMatchesChild) { + return; } entity.setActiveTaskId(null); entity.setUpdatedAt(LocalDateTime.now()); @@ -249,6 +252,8 @@ public class PriceTrackLoopRunService { : task.getErrorMessage()); entity.setFinishedAt(LocalDateTime.now()); loopRunMapper.updateById(entity); + log.warn("[price-track-loop] child finished as failed loopRunId={} childTaskId={} roundIndex={} shopIndex={} error={}", + entity.getId(), childTaskId, entity.getCurrentRound(), entity.getCurrentShopIndex(), entity.getErrorMessage()); return; } List items = parseShops(entity); @@ -263,13 +268,23 @@ public class PriceTrackLoopRunService { markStopped(entity, null); return; } - int nextShopIndex = (entity.getCurrentShopIndex() == null ? 0 : entity.getCurrentShopIndex()) + 1; + int currentRound = Math.max(1, ctx.roundIndex()); + int currentShopIndex = Math.max(0, ctx.shopIndex()); + int nextShopIndex = currentShopIndex + 1; if (nextShopIndex < items.size()) { + entity.setStatus(STATUS_RUNNING); + entity.setErrorMessage(null); + entity.setFinishedAt(null); + entity.setCurrentRound(currentRound); entity.setCurrentShopIndex(nextShopIndex); loopRunMapper.updateById(entity); + log.info("[price-track-loop] child finished advanced shop loopRunId={} childTaskId={} nextRound={} nextShopIndex={}", + entity.getId(), childTaskId, currentRound, nextShopIndex); return; } - advanceAfterSuccessfulRound(entity, items.size()); + advanceAfterSuccessfulRound(entity, currentRound); + log.info("[price-track-loop] child finished advanced round loopRunId={} childTaskId={} status={} currentRound={} currentShopIndex={} activeTaskId={}", + entity.getId(), childTaskId, entity.getStatus(), entity.getCurrentRound(), entity.getCurrentShopIndex(), entity.getActiveTaskId()); } private void reconcileWithTerminalChild(PriceTrackLoopRunEntity entity) { @@ -281,26 +296,37 @@ public class PriceTrackLoopRunService { } } - private void advanceAfterSuccessfulRound(PriceTrackLoopRunEntity entity, int shopCount) { + private void advanceAfterSuccessfulRound(PriceTrackLoopRunEntity entity) { + int currentRound = entity.getCurrentRound() == null ? 1 : Math.max(1, entity.getCurrentRound()); + advanceAfterSuccessfulRound(entity, currentRound); + } + + private void advanceAfterSuccessfulRound(PriceTrackLoopRunEntity entity, int completedRound) { if (Boolean.TRUE.equals(entity.getStopRequested())) { markStopped(entity, null); return; } + int normalizedCompletedRound = Math.max(1, completedRound); if (EXECUTION_MODE_FINITE.equals(entity.getExecutionMode())) { int targetRounds = entity.getTargetRounds() == null ? 1 : Math.max(1, entity.getTargetRounds()); - int currentRound = entity.getCurrentRound() == null ? 1 : Math.max(1, entity.getCurrentRound()); - if (currentRound >= targetRounds) { + if (normalizedCompletedRound >= targetRounds) { entity.setStatus(STATUS_SUCCESS); + entity.setCurrentRound(normalizedCompletedRound); + entity.setCurrentShopIndex(0); entity.setFinishedAt(LocalDateTime.now()); } else { - entity.setCurrentRound(currentRound + 1); + entity.setStatus(STATUS_RUNNING); + entity.setCurrentRound(normalizedCompletedRound + 1); entity.setCurrentShopIndex(0); + entity.setFinishedAt(null); } } else { - int currentRound = entity.getCurrentRound() == null ? 1 : Math.max(1, entity.getCurrentRound()); - entity.setCurrentRound(currentRound + 1); + entity.setStatus(STATUS_RUNNING); + entity.setCurrentRound(normalizedCompletedRound + 1); entity.setCurrentShopIndex(0); + entity.setFinishedAt(null); } + entity.setErrorMessage(null); entity.setActiveTaskId(null); entity.setUpdatedAt(LocalDateTime.now()); loopRunMapper.updateById(entity); @@ -320,23 +346,51 @@ public class PriceTrackLoopRunService { return task != null && isTerminal(task.getStatus()); } + private void forEachLoopForChildTask(Long childTaskId, java.util.function.Consumer consumer) { + Map loopMap = new LinkedHashMap<>(); + List activeLoops = loopRunMapper.selectList(new LambdaQueryWrapper() + .eq(PriceTrackLoopRunEntity::getActiveTaskId, childTaskId) + .last("limit 5")); + for (PriceTrackLoopRunEntity entity : activeLoops) { + if (entity != null && entity.getId() != null) { + loopMap.put(entity.getId(), entity); + } + } + FileTaskEntity task = fileTaskMapper.selectById(childTaskId); + Long loopRunId = parseLoopRunId(task); + if (loopRunId != null && !loopMap.containsKey(loopRunId)) { + PriceTrackLoopRunEntity loop = loopRunMapper.selectById(loopRunId); + if (loop != null) { + loopMap.put(loopRunId, loop); + } + } + for (PriceTrackLoopRunEntity entity : loopMap.values()) { + consumer.accept(entity); + } + } + private boolean belongsToLoop(FileTaskEntity task, Long loopRunId) { + Long childLoopRunId = parseLoopRunId(task); + return childLoopRunId != null && loopRunId.equals(childLoopRunId); + } + + private Long parseLoopRunId(FileTaskEntity task) { if (task == null || task.getRequestJson() == null || task.getRequestJson().isBlank()) { - return false; + return null; } try { Map payload = objectMapper.readValue(task.getRequestJson(), new TypeReference>() {}); Object raw = payload.get("loopRunId"); if (raw instanceof Number number) { - return loopRunId.equals(number.longValue()); + return number.longValue(); } if (raw instanceof String text && !text.isBlank()) { - return loopRunId.equals(Long.parseLong(text)); + return Long.parseLong(text); } - return false; + return null; } catch (Exception ex) { log.warn("[price-track-loop] parse child loopRunId failed taskId={} msg={}", task.getId(), ex.getMessage()); - return false; + return null; } } diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java index a1c5a80..aa391ca 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskService.java @@ -802,42 +802,21 @@ public class SimilarAsinTaskService { } public void finalizeStaleTasks() { - if (transactionManager != null) { - LocalDateTime now = LocalDateTime.now(); - LocalDateTime threshold = now.minusMinutes(Math.max(1, properties.getStaleTimeoutMinutes())); - long thresholdMillis = threshold.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli(); - List tasks = listStaleFinalizeCandidates(threshold); - for (FileTaskEntity task : tasks) { - if (!isOwnerCurrent(ownerFromTask(task))) { - continue; - } - long heartbeatMillis = taskCacheService.getTaskHeartbeatMillis(task.getId()); - if (heartbeatMillis > thresholdMillis) { - continue; - } - TaskDistributedLockService.LockHandle taskLockHandle = acquireTaskLock(task.getId(), 0L); - if (taskLockHandle == null) { - continue; - } - try (taskLockHandle) { - inNewTransaction(() -> { - finalizeStaleTask(task.getId(), "Python interrupted before uploading final similar ASIN result"); - return null; - }); - } - } - return; - } - LocalDateTime now = LocalDateTime.now(); - LocalDateTime threshold = now.minusMinutes(Math.max(1, properties.getStaleTimeoutMinutes())); - long thresholdMillis = threshold.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli(); - List tasks = listStaleFinalizeCandidates(threshold); + // P1-7:stale 判定改用 Redis heartbeat 作为主信号。 + // 历史实现先用 file_task.updated_at 过滤,再二次校验 Redis 心跳。 + // 但 file_task.updated_at 会被 Java 端 touchJavaSideTaskActivity 在 + // "等 Python 上传" 分支每 30 分钟自我续命(resetStuckJobs → 重跑 processResultFileJob), + // 导致 Python 已断线的 SIMILAR_ASIN 任务永远进不到候选集,30 分钟兜底失效。 + // Redis heartbeat 仅在 Python 心跳/上传 result 时写入,不会被 Java 内部循环刷新。 + long thresholdMillis = LocalDateTime.now() + .minusMinutes(Math.max(1, properties.getStaleTimeoutMinutes())) + .atZone(ZoneId.systemDefault()).toInstant().toEpochMilli(); + List tasks = listStaleFinalizeCandidates(); for (FileTaskEntity task : tasks) { if (!isOwnerCurrent(ownerFromTask(task))) { continue; } - long heartbeatMillis = taskCacheService.getTaskHeartbeatMillis(task.getId()); - if (heartbeatMillis > thresholdMillis) { + if (!isHeartbeatStale(task, thresholdMillis)) { continue; } TaskDistributedLockService.LockHandle taskLockHandle = acquireTaskLock(task.getId(), 0L); @@ -845,7 +824,14 @@ public class SimilarAsinTaskService { continue; } try (taskLockHandle) { - finalizeStaleTask(task.getId(), "Python interrupted before uploading final similar ASIN result"); + if (transactionManager != null) { + inNewTransaction(() -> { + finalizeStaleTask(task.getId(), "Python interrupted before uploading final similar ASIN result"); + return null; + }); + } else { + finalizeStaleTask(task.getId(), "Python interrupted before uploading final similar ASIN result"); + } } } } @@ -876,15 +862,38 @@ public class SimilarAsinTaskService { } } - private List listStaleFinalizeCandidates(LocalDateTime threshold) { + private List listStaleFinalizeCandidates() { + // P1-7:不再用 file_task.updated_at 预过滤,由调用方按 Redis heartbeat 判 stale。 + // SIMILAR_ASIN 在线 RUNNING 任务量级有限(限 200),全扫成本可忽略。 return fileTaskMapper.selectList(new LambdaQueryWrapper() .eq(FileTaskEntity::getModuleType, MODULE_TYPE) .eq(FileTaskEntity::getStatus, STATUS_RUNNING) - .lt(FileTaskEntity::getUpdatedAt, threshold) .orderByAsc(FileTaskEntity::getUpdatedAt) .last("limit 200")); } + /** + * P1-7:判断 SIMILAR_ASIN 任务是否已 stale。 + * 主信号:Redis heartbeat(仅 Python 心跳/上传 result 时刷新,不会被 Java 内部循环续命)。 + * Redis 命中:以 heartbeatMillis 与 thresholdMillis 比较。 + * Redis 缺失/降级(返回 0):回落到 file_task.updated_at,避免 Redis 故障时全量误杀。 + */ + private boolean isHeartbeatStale(FileTaskEntity task, long thresholdMillis) { + if (task == null || task.getId() == null) { + return false; + } + long heartbeatMillis = taskCacheService.getTaskHeartbeatMillis(task.getId()); + if (heartbeatMillis > 0L) { + return heartbeatMillis <= thresholdMillis; + } + LocalDateTime updatedAt = task.getUpdatedAt(); + if (updatedAt == null) { + return false; + } + long updatedMillis = updatedAt.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli(); + return updatedMillis <= thresholdMillis; + } + private SubmitContext persistSubmittedChunk(Long taskId, SimilarAsinSubmitResultRequest request) { FileTaskEntity task = fileTaskMapper.selectById(taskId); if (task == null || !MODULE_TYPE.equals(task.getModuleType())) { diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunServiceTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunServiceTest.java new file mode 100644 index 0000000..7ef9d32 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/pricetrack/service/PriceTrackLoopRunServiceTest.java @@ -0,0 +1,129 @@ +package com.nanri.aiimage.modules.pricetrack.service; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.nanri.aiimage.modules.pricetrack.mapper.PriceTrackLoopRunMapper; +import com.nanri.aiimage.modules.pricetrack.model.entity.PriceTrackLoopRunEntity; +import com.nanri.aiimage.modules.pricetrack.model.vo.PriceTrackMatchShopsVo; +import com.nanri.aiimage.modules.task.mapper.FileTaskMapper; +import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity; +import com.nanri.aiimage.modules.ziniao.service.ZiniaoShopSwitchService; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class PriceTrackLoopRunServiceTest { + + @Test + void syncLoopRunAfterChildTerminalUsesChildContextWhenLoopRowAlreadyDrifted() throws Exception { + ObjectMapper objectMapper = new ObjectMapper(); + PriceTrackLoopRunMapper loopRunMapper = mock(PriceTrackLoopRunMapper.class); + FileTaskMapper fileTaskMapper = mock(FileTaskMapper.class); + ZiniaoShopSwitchService ziniaoShopSwitchService = mock(ZiniaoShopSwitchService.class); + PriceTrackLoopRunService service = new PriceTrackLoopRunService( + loopRunMapper, + fileTaskMapper, + objectMapper, + ziniaoShopSwitchService); + + PriceTrackLoopRunEntity loop = new PriceTrackLoopRunEntity(); + loop.setId(743L); + loop.setUserId(1L); + loop.setStatus("RUNNING"); + loop.setExecutionMode("FINITE"); + loop.setTargetRounds(2); + loop.setCurrentRound(2); + loop.setCurrentShopIndex(0); + loop.setActiveTaskId(13184L); + loop.setStopRequested(false); + loop.setShopsJson(objectMapper.writeValueAsString(List.of(shop("魏振峰")))); + + FileTaskEntity task = new FileTaskEntity(); + task.setId(13184L); + task.setUserId(1L); + task.setModuleType("PRICE_TRACK"); + task.setStatus("SUCCESS"); + task.setRequestJson(objectMapper.writeValueAsString(Map.of( + "loopRunId", 743, + "roundIndex", 1, + "shopIndex", 0 + ))); + + when(loopRunMapper.selectList(any())).thenReturn(List.of(loop)); + when(fileTaskMapper.selectById(13184L)).thenReturn(task); + + service.syncLoopRunAfterChildTerminal(13184L); + + assertEquals("RUNNING", loop.getStatus()); + assertEquals(2, loop.getCurrentRound()); + assertEquals(0, loop.getCurrentShopIndex()); + assertNull(loop.getActiveTaskId()); + assertNull(loop.getFinishedAt()); + assertNull(loop.getErrorMessage()); + verify(loopRunMapper).updateById(loop); + } + + @Test + void syncLoopRunAfterChildTerminalIgnoresAlreadyAdvancedChildWithoutActiveTask() throws Exception { + ObjectMapper objectMapper = new ObjectMapper(); + PriceTrackLoopRunMapper loopRunMapper = mock(PriceTrackLoopRunMapper.class); + FileTaskMapper fileTaskMapper = mock(FileTaskMapper.class); + ZiniaoShopSwitchService ziniaoShopSwitchService = mock(ZiniaoShopSwitchService.class); + PriceTrackLoopRunService service = new PriceTrackLoopRunService( + loopRunMapper, + fileTaskMapper, + objectMapper, + ziniaoShopSwitchService); + + PriceTrackLoopRunEntity loop = new PriceTrackLoopRunEntity(); + loop.setId(743L); + loop.setUserId(1L); + loop.setStatus("RUNNING"); + loop.setExecutionMode("FINITE"); + loop.setTargetRounds(2); + loop.setCurrentRound(2); + loop.setCurrentShopIndex(0); + loop.setActiveTaskId(null); + loop.setStopRequested(false); + loop.setShopsJson(objectMapper.writeValueAsString(List.of(shop("魏振峰")))); + + FileTaskEntity task = new FileTaskEntity(); + task.setId(13184L); + task.setUserId(1L); + task.setModuleType("PRICE_TRACK"); + task.setStatus("SUCCESS"); + task.setRequestJson(objectMapper.writeValueAsString(Map.of( + "loopRunId", 743, + "roundIndex", 1, + "shopIndex", 0 + ))); + + when(loopRunMapper.selectList(any())).thenReturn(List.of(loop)); + when(fileTaskMapper.selectById(13184L)).thenReturn(task); + + service.syncLoopRunAfterChildTerminal(13184L); + + assertEquals("RUNNING", loop.getStatus()); + assertEquals(2, loop.getCurrentRound()); + assertEquals(0, loop.getCurrentShopIndex()); + assertNull(loop.getActiveTaskId()); + assertFalse(Boolean.TRUE.equals(loop.getStopRequested())); + verify(loopRunMapper, never()).updateById(loop); + } + + private PriceTrackMatchShopsVo.PriceTrackShopQueueItem shop(String name) { + PriceTrackMatchShopsVo.PriceTrackShopQueueItem item = new PriceTrackMatchShopsVo.PriceTrackShopQueueItem(); + item.setShopName(name); + item.setMatched(true); + return item; + } +} diff --git a/frontend-vue/src/pages/brand/components/BrandPriceTrackTab.vue b/frontend-vue/src/pages/brand/components/BrandPriceTrackTab.vue index b6cae8e..57fbf3e 100644 --- a/frontend-vue/src/pages/brand/components/BrandPriceTrackTab.vue +++ b/frontend-vue/src/pages/brand/components/BrandPriceTrackTab.vue @@ -1498,12 +1498,26 @@ async function refreshTaskBatch() { saveTaskDetailsToStorage() await loadHistory() syncPollingIdsWithHistory() - if (changed) await loadDashboard() + if (changed) { + await loadDashboard() + void resumeLoopExecutionIfNeeded() + } } catch { /* polling noise */ } } +async function resumeLoopExecutionIfNeeded() { + if (disposed || loopDispatching.value || !activeLoopRunId.value) return + try { + const loop = await syncActiveLoopRun() + if (!loop || loop.status === 'SUCCESS' || loop.status === 'FAILED' || loop.status === 'STOPPED') return + await runLoopExecution(loop.id) + } catch { + /* loop recovery noise */ + } +} + function scheduleNextPoll(immediate = false) { if (disposed) return if (pollTimer.value) {