feat(任务派发): 客户端兜底拉取页面未推送成功的任务 + 修店铺匹配定时任务误杀

问题:任务派发链路的"推送"只存在于页面里(Java 解析落库 PENDING → activate →
pywebview 桥 enqueue_json 推本机队列)。只解析没点启动、推送前关页面、在纯浏览器
打开,任务都会停在 PENDING,2 小时后被 StaleTaskRepairService 标失败
(「任务长期未被领取,已自动失败」,09-11 生产清理过 263 条同画像)。

- 服务端新增 GET /api/tasks/pull-pending(TaskClientPullController,身份从 JWT 取,
  不接受 user_id 参数):只挑创建超 5 分钟仍 PENDING 的本用户任务,逐条条件更新认领
  (PENDING→RUNNING + 接管 owner_instance_id)——与页面 activate 同一谓词,天然互斥,
  不会重复执行;认领后组装不出载荷则标 FAILED,不留 RUNNING 孤儿
- payload 由各业务模块实现 ClientTaskPullSpi 组装(task 侧不 import 业务模块,同 G5):
  首批 SIMILAR_ASIN / COLLECT_DATA / APPEARANCE_PATENT——这三个 Python 消费端会自行
  回拉明细,故载荷极简、客户端零模块知识;开关 aiimage.client-task-pull.enabled 默认 false
- 防双执行:三处 activate 由「非终态即可」收紧为只认 PENDING,未命中抛
  「任务已在执行中(可能已由客户端自动接管),无需重复启动」(顺带堵住整行 updateById
  把认领写入的 owner 覆盖回去的竞态);两个前端页 activate 失败即提示并停止入队
- 客户端:amazon/main.py 新增 pending_task_pull_worker,启动点挂在 app_client/main.py
  的 start_task_monitor(独立入口的 worker 线上并不生效);开关 pending_pull_enabled 用
  getattr 读取,避免 test/ 下的旧 config 缺键导致整包导入失败
- 同批修:StaleTaskRepairService 的 SCHEDULED 分支改按 scheduled_at + 120min 判死
  (原按 updated_at 会必杀排期 >2h 的店铺匹配定时任务,而 activate 又被 scheduledAt-90s 挡住)
- 测试:TaskClientPullServiceTest / StaleTaskRepairServiceTest / CollectDataTaskPullSpiImplTest /
  CollectDataActivateGuardTest 共 20 例;客户端 pending_task_pull_worker 7 例并更新启动顺序契约测试
This commit is contained in:
2026-09-14 16:23:34 +08:00
parent 8cab9d4bad
commit 24c5a09c7f
18 changed files with 1254 additions and 16 deletions
@@ -0,0 +1,60 @@
package com.nanri.aiimage.modules.appearancepatent.service;
import com.nanri.aiimage.modules.appearancepatent.model.dto.AppearancePatentParsedPayloadDto;
import com.nanri.aiimage.modules.appearancepatent.model.vo.AppearancePatentParsedGroupVo;
import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
import com.nanri.aiimage.modules.task.spi.ClientTaskPullSpi;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
/**
* 外观专利的客户端兜底拉取实现。
*
* <p>Python 消费端需要 groups(解析分组,页面上是现拉 /queue-payload 再入队),
* 这里直接复用同一个 service 方法;payload 与页面保持一致(含 prompt / api_key)。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class AppearancePatentTaskPullSpiImpl implements ClientTaskPullSpi {
private static final String QUEUE_TYPE = "appearance-patent-run";
private final AppearancePatentTaskService taskService;
private final AppearancePatentTaskCacheService taskCacheService;
@Override
public String moduleType() {
return AppearancePatentTaskService.MODULE_TYPE;
}
@Override
public Map<String, Object> buildQueuePayload(FileTaskEntity task) {
AppearancePatentParsedPayloadDto payload = taskService.queuePayload(task.getId(), task.getUserId());
List<AppearancePatentParsedGroupVo> groups = payload.getGroups() == null ? List.of() : payload.getGroups();
if (groups.isEmpty()) {
// 空 groups 在 Python 侧会被静默跳过("groups/rows is empty, skip"),宁可在服务端直接判失败
log.warn("[appearance-patent] 兜底拉取失败:解析分组为空 taskId={}", task.getId());
return null;
}
Map<String, Object> data = new LinkedHashMap<>();
data.put("taskId", task.getId());
data.put("user_id", task.getUserId());
data.put("prompt", payload.getAiPrompt());
data.put("api_key", payload.getApiKey());
data.put("groups", groups);
log.info("[appearance-patent] 兜底载荷已组装 taskId={} groups={}", task.getId(), groups.size());
return Map.of("type", QUEUE_TYPE, "data", data);
}
@Override
public void onClaimed(FileTaskEntity task) {
// 对齐 activate:刷新模块缓存心跳,让页面立刻看到 RUNNING
taskCacheService.touchTaskHeartbeat(task.getId());
}
}
@@ -280,12 +280,20 @@ public class AppearancePatentTaskService {
throw new BusinessException("任务不存在");
}
ensureTaskOwnedByCurrentInstance(task, "activate");
if (STATUS_SUCCESS.equals(task.getStatus()) || STATUS_FAILED.equals(task.getStatus())) {
// 只允许 PENDING→RUNNING(条件更新):与客户端「兜底拉取」的原子认领互斥,
// 谁先翻转谁执行,避免页面与客户端重复执行同一任务
int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getId, taskId)
.eq(FileTaskEntity::getStatus, STATUS_PENDING)
.set(FileTaskEntity::getStatus, STATUS_RUNNING)
.set(FileTaskEntity::getUpdatedAt, LocalDateTime.now()));
if (updated == 0) {
FileTaskEntity latest = fileTaskMapper.selectById(taskId);
if (latest != null && STATUS_RUNNING.equals(latest.getStatus())) {
throw new BusinessException("任务已在执行中(可能已由客户端自动接管),无需重复启动");
}
throw new BusinessException("任务已结束");
}
task.setStatus(STATUS_RUNNING);
task.setUpdatedAt(LocalDateTime.now());
fileTaskMapper.updateById(task);
taskCacheService.touchTaskHeartbeat(taskId);
}
@@ -381,14 +381,18 @@ public class CollectDataService {
@Transactional
public void activateTask(Long taskId, Long userId) {
FileTaskEntity task = requireTask(taskId, userId);
// 条件更新:上面的「已结束」判断与写入之间存在窗口(TOCTOU),期间 /fail 可能已把任务
// 标为 FAILED —— 整行 updateById 会把它复活成 RUNNING(前端显示"执行中"但无人推进)
// 只允许 PENDING→RUNNING(条件更新):既堵住 TOCTOU/fail 抢先标 FAILED 后被整行
// updateById 复活成 RUNNING),又与客户端「兜底拉取」的原子认领互斥,谁先翻转谁执行
int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getId, task.getId())
.notIn(FileTaskEntity::getStatus, STATUS_SUCCESS, STATUS_FAILED)
.eq(FileTaskEntity::getStatus, STATUS_PENDING)
.set(FileTaskEntity::getStatus, STATUS_RUNNING)
.set(FileTaskEntity::getUpdatedAt, LocalDateTime.now()));
if (updated == 0) {
FileTaskEntity latest = fileTaskMapper.selectById(task.getId());
if (latest != null && STATUS_RUNNING.equals(latest.getStatus())) {
throw new BusinessException("任务已在执行中(可能已由客户端自动接管),无需重复启动");
}
throw new BusinessException("任务已结束");
}
}
@@ -0,0 +1,71 @@
package com.nanri.aiimage.modules.collectdata.service;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
import com.nanri.aiimage.modules.task.spi.ClientTaskPullSpi;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.LinkedHashMap;
import java.util.Map;
/**
* 集采(collect-data)的客户端兜底拉取实现。
*
* <p>Python 消费端需要 taskId / totalRows / pageSize / filters(明细行自行按 /items 分页拉取)。
* filters 取自任务行 request_json 里 parse 时落库的那份(CollectDataService#persistParsedTask),
* 序列化后就是 Python 读取的 camelCase 键(countryCode/minAmount/...)。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class CollectDataTaskPullSpiImpl implements ClientTaskPullSpi {
private static final String QUEUE_TYPE = "collect-data-run";
private static final String TASK_TYPE = "collect-data";
private final ObjectMapper objectMapper;
@Override
public String moduleType() {
return CollectDataService.MODULE_TYPE;
}
@Override
public Map<String, Object> buildQueuePayload(FileTaskEntity task) {
JsonNode request = parseJson(task.getRequestJson());
if (request == null) {
log.warn("[collect-data] 兜底拉取失败:任务请求参数缺失或不可解析 taskId={}", task.getId());
return null;
}
JsonNode filtersNode = request.get("filters");
Map<String, Object> filters = filtersNode == null || filtersNode.isNull()
? Map.of()
: objectMapper.convertValue(filtersNode, Map.class);
JsonNode stats = parseJson(task.getResultJson());
int totalRows = stats == null ? 0 : stats.path("totalRows").asInt(0);
Map<String, Object> data = new LinkedHashMap<>();
data.put("taskId", task.getId());
data.put("taskNo", task.getTaskNo());
data.put("taskType", TASK_TYPE);
data.put("totalRows", totalRows);
data.put("pageSize", CollectDataService.DEFAULT_PAGE_SIZE);
data.put("filters", filters);
log.info("[collect-data] 兜底载荷已组装 taskId={} totalRows={} filters={}", task.getId(), totalRows, filters);
return Map.of("type", QUEUE_TYPE, "data", data);
}
private JsonNode parseJson(String json) {
if (json == null || json.isBlank()) {
return null;
}
try {
return objectMapper.readTree(json);
} catch (Exception ex) {
log.warn("[collect-data] 兜底拉取解析任务 JSON 失败 err={}", ex.getMessage());
return null;
}
}
}
@@ -0,0 +1,47 @@
package com.nanri.aiimage.modules.similarasin.service;
import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
import com.nanri.aiimage.modules.task.spi.ClientTaskPullSpi;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.LinkedHashMap;
import java.util.Map;
/**
* 相似 ASIN 的客户端兜底拉取实现。
*
* <p>Python 消费端只需 taskId 即可自行回拉解析载荷(parsed-payload)与本地 aliprice 账号/代理,
* 因此这里只补 user_id(价格回调需要)—— 与页面 payload 相比刻意不带 counts/api_key
* 避免与 Python 侧已有的拉取、服务端密钥兜底逻辑重复。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class SimilarAsinTaskPullSpiImpl implements ClientTaskPullSpi {
private static final String QUEUE_TYPE = "similar-asin-run";
private final SimilarAsinTaskCacheService taskCacheService;
@Override
public String moduleType() {
return SimilarAsinTaskService.MODULE_TYPE;
}
@Override
public Map<String, Object> buildQueuePayload(FileTaskEntity task) {
Map<String, Object> data = new LinkedHashMap<>();
data.put("taskId", task.getId());
data.put("user_id", task.getUserId());
log.info("[similar-asin] 兜底载荷已组装 taskId={} userId={}", task.getId(), task.getUserId());
return Map.of("type", QUEUE_TYPE, "data", data);
}
@Override
public void onClaimed(FileTaskEntity task) {
// 对齐 activate:刷新模块缓存心跳,让页面立刻看到 RUNNING
taskCacheService.touchTaskHeartbeat(task.getId());
}
}
@@ -532,12 +532,20 @@ public class SimilarAsinTaskService implements SimilarAsinPipelineHost {
throw new BusinessException("任务不存在");
}
ownershipSupport().ensureTaskOwnedByCurrentInstance(task, "activate");
if (STATUS_SUCCESS.equals(task.getStatus()) || STATUS_FAILED.equals(task.getStatus())) {
// 只允许 PENDING→RUNNING(条件更新):与客户端「兜底拉取」的原子认领互斥,
// 谁先翻转谁执行,避免页面与客户端重复执行同一任务
int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getId, taskId)
.eq(FileTaskEntity::getStatus, STATUS_PENDING)
.set(FileTaskEntity::getStatus, STATUS_RUNNING)
.set(FileTaskEntity::getUpdatedAt, LocalDateTime.now()));
if (updated == 0) {
FileTaskEntity latest = fileTaskMapper.selectById(taskId);
if (latest != null && STATUS_RUNNING.equals(latest.getStatus())) {
throw new BusinessException("任务已在执行中(可能已由客户端自动接管),无需重复启动");
}
throw new BusinessException("任务已结束");
}
task.setStatus(STATUS_RUNNING);
task.setUpdatedAt(LocalDateTime.now());
fileTaskMapper.updateById(task);
taskCacheService.touchTaskHeartbeat(taskId);
}
}
@@ -0,0 +1,42 @@
package com.nanri.aiimage.modules.task.controller;
import com.nanri.aiimage.common.api.ApiResponse;
import com.nanri.aiimage.common.model.entity.AdminUserEntity;
import com.nanri.aiimage.common.security.AdminAuthSupport;
import com.nanri.aiimage.modules.task.model.vo.TaskClientPullVo;
import com.nanri.aiimage.modules.task.service.TaskClientPullService;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.tags.Tag;
import jakarta.servlet.http.HttpServletRequest;
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
/**
* 客户端兜底拉取:桌面客户端定时调用,把"页面没推送成功、长期停在 PENDING"的任务领走执行。
*
* <p>用户身份一律从 JWT 解析(客户端 4.0.14+ 的请求钩子自动带 Bearer),未登录直接 401 ——
* 不允许匿名按 user_id 参数领别人的任务。
*/
@RestController
@RequiredArgsConstructor
@RequestMapping("/api/tasks")
@Tag(name = "任务兜底拉取", description = "供桌面客户端领取页面未推送成功的待执行任务")
public class TaskClientPullController {
private final TaskClientPullService taskClientPullService;
private final AdminAuthSupport adminAuthSupport;
@GetMapping("/pull-pending")
@Operation(
summary = "拉取并认领长期未推送的任务",
description = "只返回当前登录用户、创建超过 N 分钟仍为 PENDING 的任务;服务端原子认领(PENDING→RUNNING),"
+ "与页面推送互斥,不会重复执行。开关关闭时返回空列表。")
public ApiResponse<List<TaskClientPullVo>> pullPending(HttpServletRequest request) {
AdminUserEntity me = adminAuthSupport.requireUser(request);
return ApiResponse.success(taskClientPullService.pullPendingTasks(me.getId()));
}
}
@@ -0,0 +1,27 @@
package com.nanri.aiimage.modules.task.model.vo;
import io.swagger.v3.oas.annotations.media.Schema;
import lombok.Data;
import java.util.Map;
/** 客户端兜底拉取返回的任务项:任务元信息 + 客户端可直接入队的载荷。 */
@Data
@Schema(description = "客户端兜底拉取的任务项")
public class TaskClientPullVo {
@Schema(description = "任务 ID", example = "7004")
private Long taskId;
@Schema(description = "模块类型", example = "SIMILAR_ASIN")
private String moduleType;
@Schema(description = "任务编号", example = "SIMILAR_ASIN-2063533925058785280")
private String taskNo;
@Schema(description = "任务创建时间", example = "2026-09-14 10:00:00")
private String createdAt;
@Schema(description = "客户端可直接入队的载荷:{type, data}")
private Map<String, Object> queuePayload;
}
@@ -55,6 +55,7 @@ public class StaleTaskRepairService {
}
try (lock) {
repairFileTaskStaleIdle();
repairFileTaskStaleScheduled();
repairFileTaskStaleRunning();
repairBrandStale();
} catch (Exception ex) {
@@ -62,12 +63,12 @@ public class StaleTaskRepairService {
}
}
/** 中间态 PENDING/SCHEDULED 超时未接单 → FAILED。 */
/** PENDING 长时间未被领取 → FAILED(页面没推进客户端队列,客户端兜底拉取也没领到)。 */
private void repairFileTaskStaleIdle() {
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(STALE_IDLE_MINUTES);
List<FileTaskEntity> stale = fileTaskMapper.selectList(new LambdaQueryWrapper<FileTaskEntity>()
.select(FileTaskEntity::getId, FileTaskEntity::getModuleType)
.in(FileTaskEntity::getStatus, STATUS_PENDING, STATUS_SCHEDULED)
.eq(FileTaskEntity::getStatus, STATUS_PENDING)
.lt(FileTaskEntity::getUpdatedAt, cutoff)
.last("limit 500"));
if (stale.isEmpty()) {
@@ -77,7 +78,7 @@ public class StaleTaskRepairService {
LocalDateTime now = LocalDateTime.now();
int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper<FileTaskEntity>()
.in(FileTaskEntity::getId, stale.stream().map(FileTaskEntity::getId).toList())
.in(FileTaskEntity::getStatus, STATUS_PENDING, STATUS_SCHEDULED)
.eq(FileTaskEntity::getStatus, STATUS_PENDING)
.set(FileTaskEntity::getStatus, STATUS_FAILED)
.set(FileTaskEntity::getErrorMessage, reason)
.set(FileTaskEntity::getFinishedAt, now));
@@ -87,6 +88,42 @@ public class StaleTaskRepairService {
}
}
/**
* SCHEDULED 到时后长时间未启动 → FAILED。
*
* <p>只按 updated_at 判死会误杀:定时任务(店铺匹配)到点前本来就不会有人动它,
* 排期在 STALE_IDLE_MINUTES 之后的任务必被判死,而 activate 又被「未到定时执行时间」挡住。
* 因此有 scheduled_at 时以「到点时间」为判死基准(到点后再宽限 STALE_IDLE_MINUTES),
* scheduled_at 为空时退回 updated_at 逻辑。completeStage 会把 scheduled_at 刷成下一轮时间,
* 多轮定时任务不会被上一轮的时间误判。
*/
private void repairFileTaskStaleScheduled() {
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(STALE_IDLE_MINUTES);
List<FileTaskEntity> stale = fileTaskMapper.selectList(new LambdaQueryWrapper<FileTaskEntity>()
.select(FileTaskEntity::getId)
.eq(FileTaskEntity::getStatus, STATUS_SCHEDULED)
.and(w -> w.isNull(FileTaskEntity::getScheduledAt)
.lt(FileTaskEntity::getUpdatedAt, cutoff)
.or()
.lt(FileTaskEntity::getScheduledAt, cutoff))
.last("limit 500"));
if (stale.isEmpty()) {
return;
}
String reason = "定时任务到点后长时间未启动,已自动失败";
LocalDateTime now = LocalDateTime.now();
int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper<FileTaskEntity>()
.in(FileTaskEntity::getId, stale.stream().map(FileTaskEntity::getId).toList())
.eq(FileTaskEntity::getStatus, STATUS_SCHEDULED)
.set(FileTaskEntity::getStatus, STATUS_FAILED)
.set(FileTaskEntity::getErrorMessage, reason)
.set(FileTaskEntity::getFinishedAt, now));
if (updated > 0) {
log.warn("[stale-task-repair] 陈旧定时任务已标失败 count={} ids={}",
updated, stale.stream().map(FileTaskEntity::getId).toList());
}
}
/** RUNNING 但心跳超过 2h 未刷新 → FAILED(心跳残留续命即僵尸)。 */
private void repairFileTaskStaleRunning() {
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(STALE_IDLE_MINUTES);
@@ -0,0 +1,202 @@
package com.nanri.aiimage.modules.task.service;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.nanri.aiimage.config.InstanceMetadata;
import com.nanri.aiimage.modules.task.mapper.FileTaskMapper;
import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
import com.nanri.aiimage.modules.task.model.vo.TaskClientPullVo;
import com.nanri.aiimage.modules.task.spi.ClientTaskPullSpi;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
/**
* 客户端兜底拉取:把"页面没推送成功、长期停在 PENDING"的任务交给在线客户端执行。
*
* <p>PENDING 表示任务已解析落库、等页面把它推进本机 Python 队列;页面这一环缺失时任务会一直
* 停在这里,2 小时后被 {@code StaleTaskRepairService} 标失败。本服务让客户端主动来领:
* 只挑创建超过 N 分钟仍是 PENDING 的本用户任务,逐条用条件更新认领(PENDING→RUNNING),
* 只有把状态翻过来的调用方算领取成功 —— 与页面 activate 天然互斥,不会重复执行。
*
* <p>开关 {@code aiimage.client-task-pull.enabled} 默认关闭;模块白名单默认只放三个
* "执行参数全在服务端"的模块,其余模块实现 {@link ClientTaskPullSpi} 后可逐步加入。
*/
@Slf4j
@Service
public class TaskClientPullService {
private static final String STATUS_PENDING = "PENDING";
private static final String STATUS_RUNNING = "RUNNING";
private static final String STATUS_FAILED = "FAILED";
private static final DateTimeFormatter TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
/** 单次拉取上限,防止配置写大后一次认领过多把客户端队列压满。 */
private static final int MAX_LIMIT = 20;
private final FileTaskMapper fileTaskMapper;
private final InstanceMetadata instanceMetadata;
/** moduleType → 模块兜底实现(启动时建索引并校验重复注册)。 */
private final Map<String, ClientTaskPullSpi> pullHandlers;
@Value("${aiimage.client-task-pull.enabled:false}")
private boolean enabled;
@Value("${aiimage.client-task-pull.module-types:SIMILAR_ASIN,COLLECT_DATA,APPEARANCE_PATENT}")
private String moduleTypes;
@Value("${aiimage.client-task-pull.min-pending-minutes:5}")
private long minPendingMinutes;
@Value("${aiimage.client-task-pull.limit:5}")
private int limit;
public TaskClientPullService(FileTaskMapper fileTaskMapper,
InstanceMetadata instanceMetadata,
List<ClientTaskPullSpi> handlers) {
this.fileTaskMapper = fileTaskMapper;
this.instanceMetadata = instanceMetadata;
Map<String, ClientTaskPullSpi> index = new LinkedHashMap<>();
for (ClientTaskPullSpi handler : handlers == null ? List.<ClientTaskPullSpi>of() : handlers) {
String moduleType = handler.moduleType();
if (moduleType == null || moduleType.isBlank()) {
throw new IllegalStateException("ClientTaskPullSpi 未声明 moduleType: "
+ handler.getClass().getName());
}
ClientTaskPullSpi exists = index.put(moduleType.trim().toUpperCase(Locale.ROOT), handler);
if (exists != null) {
throw new IllegalStateException("moduleType=" + moduleType + " 注册了多个兜底拉取实现: "
+ exists.getClass().getName() + " / " + handler.getClass().getName());
}
}
this.pullHandlers = Map.copyOf(index);
log.info("[client-task-pull] 模块兜底实现注册完成 count={} modules={}", index.size(), index.keySet());
}
/** 拉取并认领该用户"长期未被领取"的任务,按创建时间升序返回,最多 limit 条。 */
public List<TaskClientPullVo> pullPendingTasks(Long userId) {
if (!enabled) {
return List.of();
}
if (userId == null || userId <= 0) {
return List.of();
}
List<String> types = resolveModuleTypes();
if (types.isEmpty()) {
log.warn("[client-task-pull] 模块白名单为空,本次不拉取 userId={}", userId);
return List.of();
}
long startedAt = System.currentTimeMillis();
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(Math.max(1L, minPendingMinutes));
int safeLimit = Math.max(1, Math.min(limit, MAX_LIMIT));
List<FileTaskEntity> candidates = fileTaskMapper.selectList(new LambdaQueryWrapper<FileTaskEntity>()
.select(FileTaskEntity::getId, FileTaskEntity::getModuleType, FileTaskEntity::getCreatedAt)
.eq(FileTaskEntity::getUserId, userId)
.eq(FileTaskEntity::getStatus, STATUS_PENDING)
.in(FileTaskEntity::getModuleType, types)
.lt(FileTaskEntity::getCreatedAt, cutoff)
.orderByAsc(FileTaskEntity::getCreatedAt)
.last("limit " + safeLimit));
if (candidates.isEmpty()) {
return List.of();
}
List<TaskClientPullVo> claimed = new ArrayList<>();
for (FileTaskEntity candidate : candidates) {
ClientTaskPullSpi handler = pullHandlers.get(candidate.getModuleType());
if (handler == null) {
// 白名单配了未实现 SPI 的模块:只告警、不动状态(页面仍可正常推送)
log.warn("[client-task-pull] 模块未实现兜底拉取,跳过 taskId={} moduleType={}",
candidate.getId(), candidate.getModuleType());
continue;
}
if (!claim(candidate.getId())) {
log.info("[client-task-pull] 任务已被页面领取或状态已变化,跳过 taskId={}", candidate.getId());
continue;
}
FileTaskEntity task = fileTaskMapper.selectById(candidate.getId());
if (task == null) {
continue;
}
Map<String, Object> queuePayload = null;
try {
queuePayload = handler.buildQueuePayload(task);
} catch (Exception ex) {
log.error("[client-task-pull] 组装兜底载荷异常 taskId={} moduleType={} err={}",
task.getId(), task.getModuleType(), ex.getMessage(), ex);
}
if (queuePayload == null) {
failClaimedTask(task, "任务数据不完整,无法自动执行,请重新提交");
continue;
}
try {
handler.onClaimed(task);
} catch (Exception ex) {
// 缓存/心跳刷新失败不影响执行:客户端已拿到载荷,心跳会由执行侧补上
log.warn("[client-task-pull] 认领后刷新模块缓存失败 taskId={} err={}", task.getId(), ex.getMessage());
}
TaskClientPullVo vo = new TaskClientPullVo();
vo.setTaskId(task.getId());
vo.setModuleType(task.getModuleType());
vo.setTaskNo(task.getTaskNo());
vo.setCreatedAt(task.getCreatedAt() == null ? null : TIME_FORMATTER.format(task.getCreatedAt()));
vo.setQueuePayload(queuePayload);
claimed.add(vo);
log.warn("[client-task-pull] 兜底认领成功 taskId={} moduleType={} userId={} taskNo={} 创建于={}",
task.getId(), task.getModuleType(), task.getUserId(), task.getTaskNo(), vo.getCreatedAt());
}
log.info("[client-task-pull] 拉取完成 userId={} 候选={} 认领={} 耗时={}ms",
userId, candidates.size(), claimed.size(), System.currentTimeMillis() - startedAt);
return claimed;
}
/**
* 条件更新认领:只有把 PENDING 翻成 RUNNING 的一方算领取成功。
* PENDING 无活跃 owner(页面尚未执行),认领即接管,因此不做归属转发。
*/
private boolean claim(Long taskId) {
return fileTaskMapper.update(null, new LambdaUpdateWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getId, taskId)
.eq(FileTaskEntity::getStatus, STATUS_PENDING)
.set(FileTaskEntity::getStatus, STATUS_RUNNING)
.set(FileTaskEntity::getUpdatedAt, LocalDateTime.now())
.set(FileTaskEntity::getOwnerInstanceId, currentInstanceId())) > 0;
}
/** 认领后组装不出载荷:直接标失败,避免留成 RUNNING 孤儿等 2 小时心跳线。 */
private void failClaimedTask(FileTaskEntity task, String reason) {
int updated = fileTaskMapper.update(null, new LambdaUpdateWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getId, task.getId())
.eq(FileTaskEntity::getStatus, STATUS_RUNNING)
.set(FileTaskEntity::getStatus, STATUS_FAILED)
.set(FileTaskEntity::getErrorMessage, reason)
.set(FileTaskEntity::getFinishedAt, LocalDateTime.now()));
log.warn("[client-task-pull] 认领后组装载荷失败,已标失败 taskId={} moduleType={} 命中={} 原因={}",
task.getId(), task.getModuleType(), updated, reason);
}
private List<String> resolveModuleTypes() {
if (moduleTypes == null || moduleTypes.isBlank()) {
return List.of();
}
List<String> types = new ArrayList<>();
for (String raw : moduleTypes.split(",")) {
String value = raw == null ? "" : raw.trim().toUpperCase(Locale.ROOT);
if (!value.isEmpty() && !types.contains(value)) {
types.add(value);
}
}
return types;
}
private String currentInstanceId() {
String instanceId = instanceMetadata == null ? null : instanceMetadata.getInstanceId();
return instanceId == null || instanceId.isBlank() ? "unknown-instance" : instanceId;
}
}
@@ -0,0 +1,36 @@
package com.nanri.aiimage.modules.task.spi;
import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
import java.util.Map;
/**
* 客户端兜底拉取的模块侧扩展点。
*
* <p>背景:任务由页面"解析落库 PENDING → activate → 经 pywebview 桥推给本机 Python 队列"派发,
* 页面这一段缺失(只解析没点启动 / 推送前关页面 / 在纯浏览器打开)时任务永远停在 PENDING,
* 2 小时后被 StaleTaskRepairService 标失败(「任务长期未被领取,已自动失败」)。
* 本接口让业务模块自行组装"客户端可直接入队"的 payload,由 {@code TaskClientPullService}
* 认领(PENDING→RUNNING 条件更新)后下发给客户端执行 —— 等于替页面补上"推送"这一步。
*
* <p>与 {@link TaskModuleHeartbeatSpi} 同理(G5):task 侧不 import 业务模块,
* 各模块实现本接口,Spring 注入 List 后由 task 侧建索引。
*/
public interface ClientTaskPullSpi {
/** 本实现负责的 moduleType(与 biz_file_task.module_type 一致)。 */
String moduleType();
/**
* 组装客户端可直接入队的载荷:{@code {"type": "...", "data": {...}}}。
* 返回 null 表示该任务当前不具备兜底执行条件(调用方会把任务标失败并回写原因)。
*/
Map<String, Object> buildQueuePayload(FileTaskEntity task);
/**
* 认领成功后的模块侧动作:对齐各自的 activate(多数模块实现为
* {@code cacheService.touchTaskHeartbeat(taskId)},让页面立刻看到 RUNNING)。
*/
default void onClaimed(FileTaskEntity task) {
}
}
@@ -234,6 +234,14 @@ aiimage:
stuck-timeout-minutes: ${AIIMAGE_RESULT_FILE_JOB_STUCK_TIMEOUT_MINUTES:30}
heartbeat-interval-ms: ${AIIMAGE_RESULT_FILE_JOB_HEARTBEAT_INTERVAL_MS:60000}
batch-size: ${AIIMAGE_RESULT_FILE_JOB_BATCH_SIZE:20}
# 客户端兜底拉取:页面没推送成功(只解析没点启动 / 推送前关页面 / 纯浏览器打开)、
# 长期停在 PENDING 的任务由在线客户端领走执行,替代 2 小时后的"任务长期未被领取"标失败。
# 默认关闭;打开前需先发布带兜底拉取线程的客户端版本,旧客户端不受影响。
client-task-pull:
enabled: ${AIIMAGE_CLIENT_TASK_PULL_ENABLED:false}
module-types: ${AIIMAGE_CLIENT_TASK_PULL_MODULE_TYPES:SIMILAR_ASIN,COLLECT_DATA,APPEARANCE_PATENT}
min-pending-minutes: ${AIIMAGE_CLIENT_TASK_PULL_MIN_PENDING_MINUTES:5}
limit: ${AIIMAGE_CLIENT_TASK_PULL_LIMIT:5}
coze-task:
max-concurrent: ${AIIMAGE_COZE_TASK_MAX_CONCURRENT:12}
brand-check: