diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/permission/service/PermissionMenuSchemaInitializer.java b/backend-java/src/main/java/com/nanri/aiimage/modules/permission/service/PermissionMenuSchemaInitializer.java index 1ae6048f..310df905 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/permission/service/PermissionMenuSchemaInitializer.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/permission/service/PermissionMenuSchemaInitializer.java @@ -47,20 +47,34 @@ public class PermissionMenuSchemaInitializer { new DefaultAppChildMenu("ERP", "erp", "erp", 132, "brand_logistics_tools") ); + /** + * 后台一级分组:只做权限树层级与「上级菜单」候选,不映射真实页面。 + * route_path 使用 group- 前缀虚拟值以满足 (menu_type, route_path) 唯一约束; + * 前端按 column_key 前缀 admin_group_ 识别分组行,不会出现在页面选择器中。 + */ + private static final List DEFAULT_ADMIN_GROUPS = List.of( + new DefaultAdminGroup("账号与权限", "admin_group_account", "group-account", 5), + new DefaultAdminGroup("数据管理", "admin_group_data", "group-data", 15), + new DefaultAdminGroup("店铺管理", "admin_group_shop", "group-shop", 45), + new DefaultAdminGroup("记录与版本", "admin_group_record", "group-record", 55) + ); + private static final List DEFAULT_ADMIN_MENUS = List.of( - new DefaultAdminMenu("用户管理", "admin_users", "users", 10), - new DefaultAdminMenu("菜单权限配置", "admin_columns", "columns", 20), - new DefaultAdminMenu("分组管理", "admin_group_manage", "group-manage", 25), - new DefaultAdminMenu("去重数据汇总", "admin_dedupe_total_data", "dedupe-total-data", 30), - new DefaultAdminMenu("无效ASIN数据", "admin_invalid_asin_data", "invalid-asin-data", 35), - new DefaultAdminMenu("店铺密钥管理", "admin_shop_keys", "shop-keys", 40), - new DefaultAdminMenu("店铺管理", "admin_shop_manage", "shop-manage", 50), - new DefaultAdminMenu("跳过跟价ASIN", "admin_skip_price_asin", "skip-price-asin", 60), - new DefaultAdminMenu("查询ASIN", "admin_query_asin", "query-asin", 65), - new DefaultAdminMenu("商品类目", "admin_product_categories", "product-categories", 66), - new DefaultAdminMenu("生成记录", "admin_history", "history", 70), - new DefaultAdminMenu("软件版本管理", "admin_version", "version", 80), - new DefaultAdminMenu("数字人版本管理", "digital_human_version", "digital-human-version", 81) + new DefaultAdminMenu("用户管理", "admin_users", "users", 10, "admin_group_account"), + new DefaultAdminMenu("菜单权限配置", "admin_columns", "columns", 20, "admin_group_account"), + new DefaultAdminMenu("分组管理", "admin_group_manage", "group-manage", 25, "admin_group_account"), + new DefaultAdminMenu("去重数据汇总", "admin_dedupe_total_data", "dedupe-total-data", 30, "admin_group_data"), + new DefaultAdminMenu("品牌数据库", "admin_invalid_asin_data", "invalid-asin-data", 35, "admin_group_data"), + new DefaultAdminMenu("查询ASIN", "admin_query_asin", "query-asin", 65, "admin_group_data"), + new DefaultAdminMenu("商品类目", "admin_product_categories", "product-categories", 66, "admin_group_data"), + new DefaultAdminMenu("店铺密钥管理", "admin_shop_keys", "shop-keys", 40, "admin_group_shop"), + new DefaultAdminMenu("店铺管理", "admin_shop_manage", "shop-manage", 50, "admin_group_shop"), + new DefaultAdminMenu("最低价ASIN设置", "admin_skip_price_asin", "skip-price-asin", 60, "admin_group_shop"), + new DefaultAdminMenu("店铺数据记录", "admin_shop_data_crawl_tasks", "shop-data-crawl-tasks", 82, "admin_group_shop"), + new DefaultAdminMenu("生成记录", "admin_history", "history", 70, "admin_group_record"), + new DefaultAdminMenu("视频任务记录", "admin_image_video_tasks", "image-video-tasks", 75, "admin_group_record"), + new DefaultAdminMenu("软件版本管理", "admin_version", "version", 80, "admin_group_record"), + new DefaultAdminMenu("数字人版本管理", "digital_human_version", "digital-human-version", 81, "admin_group_record") ); @EventListener(ApplicationReadyEvent.class) @@ -86,6 +100,7 @@ public class PermissionMenuSchemaInitializer { executeQuietly("UPDATE columns SET menu_type = 'app' WHERE menu_type IS NULL OR menu_type = ''"); executeQuietly("UPDATE columns SET sort_order = id WHERE sort_order IS NULL OR sort_order = 0"); executeQuietly("ALTER TABLE columns ADD UNIQUE KEY uk_menu_type_route_path (menu_type, route_path)"); + ensureDefaultAdminGroups(); ensureDefaultAdminMenus(); boolean shopDataCrawlDataPermissionCreated = ensureShopDataCrawlAdminDefaults(); ensureDefaultAppMenus(); @@ -105,20 +120,37 @@ public class PermissionMenuSchemaInitializer { } } - private void ensureDefaultAdminMenus() { - for (DefaultAdminMenu menu : DEFAULT_ADMIN_MENUS) { + private void ensureDefaultAdminGroups() { + for (DefaultAdminGroup group : DEFAULT_ADMIN_GROUPS) { executeQuietly(""" - INSERT INTO columns (name, column_key, menu_type, route_path, sort_order) - SELECT '%s', '%s', 'admin', '%s', %d + INSERT INTO columns (name, column_key, menu_type, route_path, parent_id, sort_order) + SELECT '%s', '%s', 'admin', '%s', NULL, %d WHERE NOT EXISTS ( SELECT 1 FROM columns WHERE column_key = '%s' ) - """.formatted(menu.name(), menu.columnKey(), menu.routePath(), menu.sortOrder(), menu.columnKey())); + """.formatted(group.name(), group.columnKey(), group.routePath(), group.sortOrder(), group.columnKey())); + } + } + + private void ensureDefaultAdminMenus() { + for (DefaultAdminMenu menu : DEFAULT_ADMIN_MENUS) { executeQuietly(""" - UPDATE columns - SET sort_order = %d, parent_id = NULL - WHERE column_key = '%s' AND (sort_order IS NULL OR sort_order = 0) - """.formatted(menu.sortOrder(), menu.columnKey())); + INSERT INTO columns (name, column_key, menu_type, route_path, sort_order, parent_id) + SELECT '%s', '%s', 'admin', '%s', %d, parent.id + FROM columns parent + WHERE parent.column_key = '%s' + AND NOT EXISTS ( + SELECT 1 FROM columns WHERE column_key = '%s' + ) + """.formatted(menu.name(), menu.columnKey(), menu.routePath(), menu.sortOrder(), + menu.parentKey(), menu.columnKey())); + executeQuietly(""" + UPDATE columns child + JOIN columns parent ON parent.column_key = '%s' + SET child.name = '%s', child.sort_order = %d, + child.parent_id = COALESCE(child.parent_id, parent.id) + WHERE child.column_key = '%s' + """.formatted(menu.parentKey(), menu.name(), menu.sortOrder(), menu.columnKey())); } } @@ -179,8 +211,9 @@ public class PermissionMenuSchemaInitializer { private boolean ensureShopDataCrawlAdminDefaults() { executeQuietly(""" - INSERT INTO columns (name, column_key, menu_type, route_path, sort_order) - SELECT '店铺数据记录', 'admin_shop_data_crawl_tasks', 'admin', 'shop-data-crawl-tasks', 82 + INSERT INTO columns (name, column_key, menu_type, route_path, sort_order, parent_id) + SELECT '店铺数据记录', 'admin_shop_data_crawl_tasks', 'admin', 'shop-data-crawl-tasks', 82, + (SELECT id FROM columns WHERE column_key = 'admin_group_shop' LIMIT 1) WHERE NOT EXISTS ( SELECT 1 FROM columns WHERE column_key = 'admin_shop_data_crawl_tasks' ) @@ -226,7 +259,11 @@ public class PermissionMenuSchemaInitializer { } } - private record DefaultAdminMenu(String name, String columnKey, String routePath, int sortOrder) { + private record DefaultAdminMenu(String name, String columnKey, String routePath, int sortOrder, + String parentKey) { + } + + private record DefaultAdminGroup(String name, String columnKey, String routePath, int sortOrder) { } private record DefaultAppMenu(String name, String columnKey, String routePath, int sortOrder) { 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 b77848fa..80f0aae0 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 @@ -5127,10 +5127,47 @@ public class SimilarAsinTaskService { payload.setSourceFiles(sourceFiles == null ? List.of() : sourceFiles); payload.setHeaders(headers == null ? List.of() : headers); payload.setItems(allRows == null ? List.of() : new ArrayList<>(allRows)); - payload.setGroups(groups == null ? List.of() : groups); + // 兼容 Python 旧链路:groups 需内嵌完整行(Python 端 normalize_groups + // 直接返回 groups 并按 gp["items"] 取行;仅有 startIndex/endIndex 引用时 + // items 为空、任务空转)。startIndex/endIndex 保留,Java 内部 expandGroupRefs 继续使用。 + payload.setGroups(fillGroupItems(groups, allRows)); return writeJson(payload, "保存解析结果失败"); } + /** + * 按 group.startIndex/endIndex 半开区间从全量行切片,内嵌回 group.items; + * 无索引或切片越界时退回按 group.items 原样(兼容旧 payload 反序列化场景)。 + */ + static List fillGroupItems(List groups, List allRows) { + if (groups == null || groups.isEmpty()) { + return groups == null ? List.of() : groups; + } + if (allRows == null || allRows.isEmpty()) { + return groups; + } + List rows = new ArrayList<>(allRows); + for (SimilarAsinParsedGroupVo group : groups) { + if (group == null) { + continue; + } + Integer start = group.getStartIndex(); + Integer end = group.getEndIndex(); + if (start != null && end != null) { + int from = Math.max(0, Math.min(start, rows.size())); + int to = Math.max(from, Math.min(end, rows.size())); + if (to > from) { + group.setItems(new ArrayList<>(rows.subList(from, to))); + continue; + } + } + // 索引缺失/无效:保留已有 items(兼容旧 payload),不覆盖 + if (group.getItems() == null) { + group.setItems(new ArrayList<>()); + } + } + return groups; + } + private String buildTaskResultJson(String aiPrompt, String apiKey, Boolean imgSwitch, Boolean categorySwitch, List sourceFiles, String parsedPayloadPointer) { Map payload = new LinkedHashMap<>(); payload.put("aiPrompt", normalize(aiPrompt)); @@ -5172,6 +5209,11 @@ public class SimilarAsinTaskService { } private SimilarAsinParsedPayloadDto hydrateParsedPayloadRows(SimilarAsinParsedPayloadDto payload) { + return hydrateForPythonCompat(payload); + } + + /** 解析载荷 hydrate:回填 allItems/items/groups 内嵌 items(Python 旧链路兼容)。 */ + static SimilarAsinParsedPayloadDto hydrateForPythonCompat(SimilarAsinParsedPayloadDto payload) { if (payload == null) { return new SimilarAsinParsedPayloadDto(); } @@ -5182,6 +5224,13 @@ public class SimilarAsinTaskService { if (payload.getItems() == null || payload.getItems().isEmpty()) { payload.setItems(rows); } + // Python 旧链路兼容:读取时也回填 groups 内嵌 items(存量任务 payload + // 只有引用),否则 Python normalize_groups 拿到空 items 任务空转。 + if (payload.getGroups() != null && !payload.getGroups().isEmpty()) { + List filledRows = payload.getItems() == null || payload.getItems().isEmpty() + ? rows : payload.getItems(); + payload.setGroups(fillGroupItems(payload.getGroups(), filledRows)); + } return payload; } @@ -5205,6 +5254,16 @@ public class SimilarAsinTaskService { // P1-4:识别跨实例 local 指针读不到的场景,把"chunk 在另一实例"的元信息 // 通过 task_scope_state.last_error 留痕,便于排查"为何 owner 切换后 chunk 读不到"。 boolean crossInstance = msg.contains("only exists on instance="); + // 类型不兼容(如历史版本序列化数据含当前类加载器不认识的类型标记): + // 该 chunk 永久不可读,跳过而不是把组装/收尾永久拖死,留痕便于排查。 + boolean typeMismatch = ex instanceof TypeNotPresentException + || msg.contains("not present"); + if (typeMismatch) { + log.warn("[similar-asin] chunk payload type mismatch, skip chunk taskId={} chunk={} err={}", + chunk.getTaskId(), chunk.getChunkIndex(), msg); + recordChunkReadFailure(chunk, false, msg); + return rows; + } log.warn("[similar-asin] read chunk payload failed taskId={} chunk={} crossInstance={} err={}", chunk.getTaskId(), chunk.getChunkIndex(), crossInstance, msg); recordChunkReadFailure(chunk, crossInstance, msg); diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/dto/TaskFileJobDispatchEvent.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/dto/TaskFileJobDispatchEvent.java index 2145d47c..fc6ec49e 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/dto/TaskFileJobDispatchEvent.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/model/dto/TaskFileJobDispatchEvent.java @@ -9,6 +9,7 @@ public record TaskFileJobDispatchEvent( Long resultId, String scopeKey, String jobType, + String owner, LocalDateTime createdAt ) { } diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobDispatchCoordinator.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobDispatchCoordinator.java index f5b08a23..4ed0072f 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobDispatchCoordinator.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobDispatchCoordinator.java @@ -1,5 +1,6 @@ package com.nanri.aiimage.modules.task.service; +import com.nanri.aiimage.config.InstanceMetadata; import com.nanri.aiimage.modules.task.model.dto.TaskFileJobDispatchEvent; import com.nanri.aiimage.modules.task.model.dto.TaskFileJobMessage; import lombok.RequiredArgsConstructor; @@ -17,6 +18,7 @@ public class TaskFileJobDispatchCoordinator { private final TaskFileJobPublisher taskFileJobPublisher; private final TaskFileJobLocalDispatcher taskFileJobLocalDispatcher; + private final InstanceMetadata instanceMetadata; @Qualifier("taskFileJobDispatchExecutor") private final TaskExecutor taskFileJobDispatchExecutor; @@ -43,6 +45,28 @@ public class TaskFileJobDispatchCoordinator { } private void dispatchAfterCommit(TaskFileJobDispatchEvent event, TaskFileJobMessage message) { + // owner 非空 ⇔ owner-scoped job(只有 SIMILAR_ASIN/APPEARANCE_PATENT/PUBLISH/ + // SHOP_DATA_CRAWL 会写 "task:{id}:owner:{instance}" 形式的 scopeKey)。 + // 这类 job 必须在归属实例上执行,走共享 consumer group 的 MQ 会被轮询投递到 + // 另一台,那边只能跳过,白白多绕一圈。 + String owner = event.owner() == null ? "" : event.owner().trim(); + if (!owner.isBlank()) { + if (owner.equals(currentInstanceId())) { + boolean dispatched = taskFileJobLocalDispatcher.dispatch( + event.jobId(), event.taskId(), event.moduleType(), true); + if (dispatched) { + return; + } + log.warn("[task-file-job] owner-scoped local dispatch declined, fallback to mq jobId={} taskId={} moduleType={} owner={}", + event.jobId(), event.taskId(), event.moduleType(), owner); + } else { + // 归属实例每 15 秒扫一次自己 owner 的 PENDING/FAILED job,这里不发 MQ、 + // 不本地跑,避免消息落到非归属实例后被丢弃。 + log.info("[task-file-job] skip dispatch because job belongs to another instance jobId={} taskId={} moduleType={} owner={} current={}", + event.jobId(), event.taskId(), event.moduleType(), owner, currentInstanceId()); + return; + } + } boolean published = taskFileJobPublisher.publish(message); if (published) { return; @@ -56,4 +80,9 @@ public class TaskFileJobDispatchCoordinator { event.jobId(), event.taskId(), event.moduleType()); } } + + private String currentInstanceId() { + String instanceId = instanceMetadata == null ? null : instanceMetadata.getInstanceId(); + return instanceId == null || instanceId.isBlank() ? "unknown-instance" : instanceId; + } } diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobLocalDispatcher.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobLocalDispatcher.java index 4d4ee937..66d88ad5 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobLocalDispatcher.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobLocalDispatcher.java @@ -42,7 +42,19 @@ public class TaskFileJobLocalDispatcher { return dispatch(jobId, taskId, moduleType, false); } + /** + * 调用方已把 job claim 成 RUNNING 时使用:复用同样的背压/去重/异步执行, + * 但落到 worker.processClaimed,避免二次 claimRunning 失败导致 job 被静默丢弃。 + */ + public boolean dispatchClaimed(Long jobId, Long taskId, String moduleType) { + return submit(jobId, taskId, moduleType, false, true); + } + public boolean dispatch(Long jobId, Long taskId, String moduleType, boolean force) { + return submit(jobId, taskId, moduleType, force, false); + } + + private boolean submit(Long jobId, Long taskId, String moduleType, boolean force, boolean alreadyClaimed) { if (jobId == null || jobId <= 0) { return false; } @@ -63,7 +75,7 @@ public class TaskFileJobLocalDispatcher { return true; } try { - taskFileJobDispatchExecutor.execute(() -> processLocally(jobId, taskId, moduleType)); + taskFileJobDispatchExecutor.execute(() -> processLocally(jobId, taskId, moduleType, alreadyClaimed)); return true; } catch (RuntimeException ex) { inflightJobIds.remove(jobId); @@ -71,7 +83,7 @@ public class TaskFileJobLocalDispatcher { jobId, taskId, moduleType, ex.getMessage(), ex); if (force) { try { - processLocally(jobId, taskId, moduleType); + processLocally(jobId, taskId, moduleType, alreadyClaimed); return true; } finally { inflightJobIds.remove(jobId); @@ -81,7 +93,7 @@ public class TaskFileJobLocalDispatcher { } } - private void processLocally(Long jobId, Long taskId, String moduleType) { + private void processLocally(Long jobId, Long taskId, String moduleType, boolean alreadyClaimed) { try { TaskFileJobEntity job = taskFileJobMapper.selectById(jobId); if (job == null) { @@ -95,7 +107,11 @@ public class TaskFileJobLocalDispatcher { jobId, taskId, moduleType); return; } - taskResultFileJobWorker.process(job); + if (alreadyClaimed) { + taskResultFileJobWorker.processClaimed(job); + } else { + taskResultFileJobWorker.process(job); + } } catch (Exception ex) { log.warn("[task-file-job] local dispatch failed jobId={} taskId={} moduleType={} msg={}", jobId, taskId, moduleType, ex.getMessage(), ex); diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobService.java index 4a401d36..490e0bce 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskFileJobService.java @@ -57,7 +57,14 @@ public class TaskFileJobService { if (existing == null) { return null; } - if ("SUCCESS".equals(existing.getStatus()) || "RUNNING".equals(existing.getStatus()) || "PENDING".equals(existing.getStatus())) { + if ("SUCCESS".equals(existing.getStatus()) || "RUNNING".equals(existing.getStatus())) { + return existing; + } + if ("PENDING".equals(existing.getStatus())) { + // PENDING 僵尸兜底:job 入队时若 MQ 发送失败且本地兜底未生效(或消息被 + // 非 owner 实例静默吞掉),job 会永远停留在 PENDING、无人处理。 + // 补发一次 dispatch 事件,让调度链重试(幂等:consumer 侧 claim 有状态翻转保护)。 + publishDispatchEvent(existing); return existing; } // 终态失败:retryCount 已达上限的 FAILED job 不再重置为 PENDING, @@ -237,6 +244,13 @@ public class TaskFileJobService { .lt(TaskFileJobEntity::getUpdatedAt, threshold) .or() .isNull(TaskFileJobEntity::getUpdatedAt))) + .or(stuckPending -> stuckPending + .eq(TaskFileJobEntity::getStatus, "PENDING") + .ne(TaskFileJobEntity::getRetryCount, MAX_RETRY_COUNT) + .and(age -> age + .lt(TaskFileJobEntity::getUpdatedAt, threshold) + .or() + .isNull(TaskFileJobEntity::getUpdatedAt))) .or(exhausted -> exhausted .eq(TaskFileJobEntity::getStatus, "FAILED") .ge(TaskFileJobEntity::getRetryCount, MAX_RETRY_COUNT) @@ -251,6 +265,24 @@ public class TaskFileJobService { exhaustedJobs.add(job); continue; } + if ("PENDING".equals(job.getStatus())) { + // PENDING 僵尸兜底:补发 dispatch 事件(worker claim 幂等), + // 并累加 retryCount 防无限重发;注意不翻转状态, + // 避免与真实 PENDING 竞争 worker 的 claim 条件更新冲突。 + int retryCount = job.getRetryCount() == null ? 0 : job.getRetryCount(); + int nextRetryCount = Math.min(MAX_RETRY_COUNT - 1, retryCount + 1); + int updated = taskFileJobMapper.update(null, new LambdaUpdateWrapper() + .eq(TaskFileJobEntity::getId, job.getId()) + .eq(TaskFileJobEntity::getStatus, "PENDING") + .eq(job.getUpdatedAt() != null, TaskFileJobEntity::getUpdatedAt, job.getUpdatedAt()) + .isNull(job.getUpdatedAt() == null, TaskFileJobEntity::getUpdatedAt) + .set(TaskFileJobEntity::getRetryCount, nextRetryCount) + .set(TaskFileJobEntity::getUpdatedAt, LocalDateTime.now())); + if (updated > 0) { + publishDispatchEvent(taskFileJobMapper.selectById(job.getId())); + } + continue; + } int retryCount = job.getRetryCount() == null ? 0 : job.getRetryCount(); int nextRetryCount = Math.min(MAX_RETRY_COUNT, retryCount + 1); LocalDateTime now = LocalDateTime.now(); @@ -642,6 +674,7 @@ public class TaskFileJobService { entity.getResultId(), entity.getScopeKey(), entity.getJobType(), + entity.getOwner(), LocalDateTime.now() )); } diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorker.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorker.java index 011ca333..7fef078d 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorker.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorker.java @@ -1,6 +1,7 @@ package com.nanri.aiimage.modules.task.service; import com.nanri.aiimage.common.exception.BusinessException; +import com.nanri.aiimage.common.exception.TaskOwnerMismatchException; import com.nanri.aiimage.modules.appearancepatent.service.AppearancePatentTaskService; import com.nanri.aiimage.config.InstanceMetadata; import com.nanri.aiimage.modules.brand.service.BrandTaskService; @@ -76,13 +77,16 @@ public class TaskResultFileJobWorker { } log.info("[task-file-job] scheduled worker claimed jobs count={}", jobs.size()); for (TaskFileJobEntity job : jobs) { - boolean dispatched = taskFileJobLocalDispatcher.dispatch( + // claimRunnableJobsForOwner 已把 job 原子翻转成 RUNNING,这里必须走"已 claim" + // 通道:若再走 process() 的 claimRunning,条件 UPDATE 会因状态已是 RUNNING 而 + // 失败,job 被静默丢弃并挂到 stuck 扫描才回收,定时兜底等于失效。 + boolean dispatched = taskFileJobLocalDispatcher.dispatchClaimed( job.getId(), job.getTaskId(), job.getModuleType() ); if (!dispatched) { - process(job); + processClaimed(job); } } } @@ -104,18 +108,39 @@ public class TaskResultFileJobWorker { } public void process(TaskFileJobEntity job) { - if (isOwnerScopedJob(job) && !isOwnedByCurrentInstance(job)) { - log.debug("[task-file-job] skip owner-scoped job because owner is another instance jobId={} taskId={} moduleType={} owner={} current={}", - job.getId(), job.getTaskId(), job.getModuleType(), ownerFromScopeKey(job.getScopeKey()), currentInstanceId()); + if (job == null || job.getId() == null) { return; } - if (job == null || job.getId() == null) { + if (isOwnerScopedJob(job) && !isOwnedByCurrentInstance(job)) { + // 非归属实例一律不碰 owner-scoped job:模块层 ensureTaskOwnedByCurrentInstance + // 会直接抛 TaskOwnerMismatchException,抢跑只会白白消耗 job 的重试次数, + // 重试耗尽后把归属机上跑得好好的任务判成失败。 + // job 不会因此卡死:归属实例的 runPendingJobs 每 15 秒扫一次自己 owner 的 + // PENDING/FAILED job,enqueueAssembleResult / stuck 扫描也会补发 dispatch。 + log.info("[task-file-job] skip owner-scoped job because owner is another instance jobId={} taskId={} moduleType={} status={} owner={} current={}", + job.getId(), job.getTaskId(), job.getModuleType(), job.getStatus(), + ownerFromScopeKey(job.getScopeKey()), currentInstanceId()); return; } TaskFileJobEntity claim = taskFileJobService.claimRunning(job.getId()); if (claim == null) { return; } + offloadOrRun(job, claim); + } + + /** + * 供"调用方已持有 RUNNING claim"的路径使用(定时 worker 的 claimRunnableJobsForOwner): + * 跳过 claimRunning,直接按 handler 能力异步/内联执行。 + */ + public void processClaimed(TaskFileJobEntity claim) { + if (claim == null || claim.getId() == null) { + return; + } + offloadOrRun(claim, claim); + } + + private void offloadOrRun(TaskFileJobEntity job, TaskFileJobEntity claim) { ResultFileJobHandler handler = handlerRegistry.asMap().get(job.getModuleType()); if (handler != null && handler.supportsAsyncOffload()) { try { @@ -206,6 +231,12 @@ public class TaskResultFileJobWorker { if (finalizeWithdraw) { withdrawTaskService.tryFinalizeTask(job.getTaskId(), false); } + } catch (TaskOwnerMismatchException ex) { + // 归属校验失败不是"生成失败":不能计入 retryCount,否则任务会在归属机 + // 一切正常的情况下被本机重试耗尽判死。退回 PENDING 交还归属实例。 + log.warn("[task-file-job] process rejected because task belongs to another instance jobId={} taskId={} moduleType={} owner={} current={}", + job.getId(), job.getTaskId(), job.getModuleType(), ex.getOwnerInstanceId(), ex.getCurrentInstanceId()); + taskFileJobService.deferRunning(job.getId(), ex.getMessage()); } catch (Exception ex) { TaskFileJobEntity latest = taskFileJobService.findById(job.getId()); if (latest != null && "SUCCESS".equals(latest.getStatus())) { diff --git a/backend-java/src/main/resources/db/V103__admin_menu_two_level_groups.sql b/backend-java/src/main/resources/db/V103__admin_menu_two_level_groups.sql new file mode 100644 index 00000000..1a589e6f --- /dev/null +++ b/backend-java/src/main/resources/db/V103__admin_menu_two_level_groups.sql @@ -0,0 +1,63 @@ +-- V103: 后台管理菜单改为两级结构(一级分组 + 子页面) +-- 新增 4 个后台一级分组(不映射真实页面,route_path 为分组虚拟值), +-- 并把这 16 个后台页面菜单挂到对应分组下。 +-- 分组只用于权限树层级与「上级菜单」候选;column_key / route_path 校验逻辑不变。 +-- 对应 Java 端清单见 PermissionMenuSchemaInitializer.DEFAULT_ADMIN_GROUPS / DEFAULT_ADMIN_MENUS。 + +-- ---------- 1. 新增一级分组 ---------- + +INSERT INTO `columns` (`name`, `column_key`, `menu_type`, `route_path`, `parent_id`, `sort_order`) +SELECT '账号与权限', 'admin_group_account', 'admin', 'group-account', NULL, 5 +WHERE NOT EXISTS (SELECT 1 FROM `columns` WHERE `column_key` = 'admin_group_account'); + +INSERT INTO `columns` (`name`, `column_key`, `menu_type`, `route_path`, `parent_id`, `sort_order`) +SELECT '数据管理', 'admin_group_data', 'admin', 'group-data', NULL, 15 +WHERE NOT EXISTS (SELECT 1 FROM `columns` WHERE `column_key` = 'admin_group_data'); + +INSERT INTO `columns` (`name`, `column_key`, `menu_type`, `route_path`, `parent_id`, `sort_order`) +SELECT '店铺管理', 'admin_group_shop', 'admin', 'group-shop', NULL, 45 +WHERE NOT EXISTS (SELECT 1 FROM `columns` WHERE `column_key` = 'admin_group_shop'); + +INSERT INTO `columns` (`name`, `column_key`, `menu_type`, `route_path`, `parent_id`, `sort_order`) +SELECT '记录与版本', 'admin_group_record', 'admin', 'group-record', NULL, 55 +WHERE NOT EXISTS (SELECT 1 FROM `columns` WHERE `column_key` = 'admin_group_record'); + +-- ---------- 2. 页面菜单挂到分组下,并同步显示名 ---------- +-- 只更新当前没有父级的页面(保留管理员手工调整过的层级), +-- COALESCE 保证已挂接的叶子不重复指向。 + +UPDATE `columns` child +JOIN `columns` parent ON parent.`column_key` = 'admin_group_account' +SET child.`parent_id` = COALESCE(child.`parent_id`, parent.`id`) +WHERE child.`column_key` IN ('admin_users', 'admin_columns', 'admin_group_manage'); + +UPDATE `columns` child +JOIN `columns` parent ON parent.`column_key` = 'admin_group_data' +SET child.`parent_id` = COALESCE(child.`parent_id`, parent.`id`) +WHERE child.`column_key` IN ('admin_dedupe_total_data', 'admin_invalid_asin_data', + 'admin_query_asin', 'admin_product_categories'); + +UPDATE `columns` child +JOIN `columns` parent ON parent.`column_key` = 'admin_group_shop' +SET child.`parent_id` = COALESCE(child.`parent_id`, parent.`id`) +WHERE child.`column_key` IN ('admin_shop_keys', 'admin_shop_manage', 'admin_skip_price_asin', + 'admin_shop_data_crawl_tasks'); + +UPDATE `columns` child +JOIN `columns` parent ON parent.`column_key` = 'admin_group_record' +SET child.`parent_id` = COALESCE(child.`parent_id`, parent.`id`) +WHERE child.`column_key` IN ('admin_history', 'admin_version', 'digital_human_version', + 'admin_image_video_tasks'); + +-- ---------- 3. 显示名对齐前端界面 --- +UPDATE `columns` SET `name` = '品牌数据库' +WHERE `menu_type` = 'admin' AND `column_key` = 'admin_invalid_asin_data'; + +UPDATE `columns` SET `name` = '最低价ASIN设置' +WHERE `menu_type` = 'admin' AND `column_key` = 'admin_skip_price_asin'; + +UPDATE `columns` SET `name` = '视频任务记录' +WHERE `menu_type` = 'admin' AND `column_key` = 'admin_image_video_tasks'; + +UPDATE `columns` SET `name` = '店铺数据记录' +WHERE `menu_type` = 'admin' AND `column_key` = 'admin_shop_data_crawl_tasks'; diff --git a/backend-java/src/test/java/com/nanri/aiimage/config/ExplainIndexAuditDocTest.java b/backend-java/src/test/java/com/nanri/aiimage/config/ExplainIndexAuditDocTest.java index df54fbe1..1b50a366 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/config/ExplainIndexAuditDocTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/config/ExplainIndexAuditDocTest.java @@ -106,8 +106,17 @@ class ExplainIndexAuditDocTest { assertTrue(audit.contains("不做 DDL") || audit.contains("本任务无 DDL") || audit.contains("不新增索引") || audit.contains("只审计"), "审计文档明确本任务不做 DDL"); - assertFalse(audit.contains("ALTER TABLE"), - "审计文档不包含 ALTER TABLE(DDL 属于后续任务 116)"); + // 任务 116 落地 V101 索引后,审计文档会把 ALTER TABLE 作为候选索引的 + // 回滚说明引用。审计任务自身依然不执行 DDL——判定口径改为「ALTER TABLE + // 只能出现在候选/回滚等说明性上下文」,而不是文档里不得出现该字符串。 + for (String line : audit.split("\r?\n")) { + if (!line.contains("ALTER TABLE")) { + continue; + } + assertTrue(line.contains("回滚") || line.contains("候选") || line.contains("建议") + || line.contains("风险") || line.trim().startsWith(">"), + "审计文档中的 ALTER TABLE 必须是候选/回滚说明,实际行:" + line); + } assertTrue(spec.contains("只读") || spec.contains("巡检"), "12 spec 巡检报表模式与审计文档对应"); } diff --git a/backend-java/src/test/java/com/nanri/aiimage/metrics/ExternalCallMetricsRecorderTest.java b/backend-java/src/test/java/com/nanri/aiimage/metrics/ExternalCallMetricsRecorderTest.java index 54dee77f..a04f1775 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/metrics/ExternalCallMetricsRecorderTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/metrics/ExternalCallMetricsRecorderTest.java @@ -27,6 +27,7 @@ import java.util.List; import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -268,6 +269,7 @@ class ExternalCallMetricsRecorderTest { } } finally { pool.shutdown(); + assertTrue(pool.awaitTermination(60, TimeUnit.SECONDS), "并发任务在超时前全部结束"); } awaitMetric("aiimage.external-call.total", "client", "llm", "result", "success"); @@ -275,6 +277,12 @@ class ExternalCallMetricsRecorderTest { Thread.sleep(10); } assertEquals(20, llmSubmitCount.get(), "并发 20 请求全部真实发出"); + // 指标在响应处理完成之后才落,submit 计数到 20 不代表指标已收敛: + // 直接断言会读到中间值(实测 18.0)。等指标计数自身收敛再断言。 + for (int i = 0; i < 2000 && counterCount("aiimage.external-call.total", + "client", "llm", "result", "success") < 20.0; i++) { + Thread.sleep(10); + } assertEquals(20.0, counterCount("aiimage.external-call.total", "client", "llm", "result", "success"), "20 次成功全部记录,无重复"); } diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/collectdata/service/CollectDataExcelCleanupTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/collectdata/service/CollectDataExcelCleanupTest.java index 5b19ddcc..d8dd28a6 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/collectdata/service/CollectDataExcelCleanupTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/collectdata/service/CollectDataExcelCleanupTest.java @@ -7,6 +7,7 @@ import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataResultRowVo; import org.apache.poi.ss.usermodel.Sheet; import org.apache.poi.ss.usermodel.Workbook; import org.apache.poi.ss.usermodel.WorkbookFactory; +import org.apache.poi.util.DefaultTempFileCreationStrategy; import org.apache.poi.util.TempFile; import org.apache.poi.util.TempFileCreationStrategy; import org.junit.jupiter.api.AfterEach; @@ -39,20 +40,36 @@ class CollectDataExcelCleanupTest { private final List trackedTempFiles = new ArrayList<>(); + private Path sxssfTempDir; + @AfterEach + void restoreDefaultTempFileStrategy() { + // 还原 POI 的进程级全局策略:trackSxssfTempFiles 改写的是单例, + // 不还原会泄漏到同一 JVM 内后续所有用到 SXSSF 的测试。 + // (原来这个 @AfterEach 挂在 private 的 trackSxssfTempFiles 上, + // 既没有还原策略,方法本身还是被各用例手工调用的 setup 辅助。) + TempFile.setTempFileCreationStrategy(new DefaultTempFileCreationStrategy()); + } private void trackSxssfTempFiles() { + // 把 SXSSF 滚动窗口临时文件重定向到本用例独占目录,使残留探测只覆盖 + // 本用例自己产生的文件。 + try { + sxssfTempDir = Files.createDirectory(tempDir.resolve("sxssf-probe")); + } catch (IOException ex) { + throw new IllegalStateException(ex); + } TempFile.setTempFileCreationStrategy(new TempFileCreationStrategy() { @Override public File createTempFile(String prefix, String suffix) throws IOException { - File file = File.createTempFile(prefix, suffix); + File file = Files.createTempFile(sxssfTempDir, prefix, suffix).toFile(); trackedTempFiles.add(file); return file; } @Override public File createTempDirectory(String prefix) throws IOException { - File dir = Files.createTempDirectory(prefix).toFile(); + File dir = Files.createTempDirectory(sxssfTempDir, prefix).toFile(); trackedTempFiles.add(dir); return dir; } @@ -62,10 +79,13 @@ class CollectDataExcelCleanupTest { private void assertSxssfTempFilesCleaned() { assertThat(trackedTempFiles).as("SXSSF 滚动窗口临时文件均已删除") .allSatisfy(f -> assertThat(f).doesNotExist()); - // POI 5.2.5 无法读取当前策略,改用系统临时目录中 SXSSF 固定前缀的残留探测。 - try (Stream paths = Files.list(Path.of(System.getProperty("java.io.tmpdir")))) { + // 只探测本用例独占目录。原实现扫描系统共享临时目录并要求 poi-sxssf* + // 前缀零命中,会把其他进程、其他用例和历史运行遗留的文件算成本用例 + // 的泄漏——本机共享临时目录里就躺着十几个历史 poi-sxssf-sheet*.xml, + // 导致实现完全正确时这 7 个用例照样红。 + try (Stream paths = Files.list(sxssfTempDir)) { assertThat(paths.filter(p -> p.getFileName().toString().startsWith("poi-sxssf"))) - .as("系统临时目录无 SXSSF 滚动窗口残留") + .as("本用例 SXSSF 临时目录无滚动窗口残留") .isEmpty(); } catch (IOException ex) { throw new IllegalStateException(ex); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlCleanupTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlCleanupTest.java index 3d7180bb..da1d4c0b 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlCleanupTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/shopdatacrawl/service/ShopDataCrawlCleanupTest.java @@ -149,6 +149,7 @@ class ShopDataCrawlCleanupTest { private final AtomicLong memberIdSeq = new AtomicLong(1000); private final AtomicLong fileSeq = new AtomicLong(5000); private Long lastJobTaskId; + private Set preExistingWorkRoots = Set.of(); @BeforeEach void configureStorage() { @@ -408,6 +409,12 @@ class ShopDataCrawlCleanupTest { .deleteTaskItems(anyLong(), eq(MODULE_TYPE)); lenient().doNothing().when(taskProgressSnapshotService) .delete(anyLong(), eq(MODULE_TYPE)); + // 共享临时目录里的既有 daily-* 目录不属于本用例,先快照、断言时做差集 + try { + preExistingWorkRoots = Set.copyOf(scanWorkRoots()); + } catch (Exception ex) { + throw new IllegalStateException(ex); + } } // ---- 1. 删除路径 ---- @@ -1261,6 +1268,16 @@ class ShopDataCrawlCleanupTest { } private List workRoots() throws Exception { + // 只统计本用例新产生的工作目录。shop-data-crawl-result 落在机器共享的 + // 系统临时目录下,其他用例、历史运行、本机跑过的应用都会往里留 daily-* + // 目录(本机实测已堆了 200+ 个),原实现把这些全算作本用例的泄漏, + // 于是实现完全正确时这些用例照样红。改为与用例开始时的快照做差集。 + List current = new ArrayList<>(scanWorkRoots()); + current.removeAll(preExistingWorkRoots); + return current; + } + + private List scanWorkRoots() throws Exception { Path root = Path.of(System.getProperty("java.io.tmpdir"), "shop-data-crawl-result"); if (!Files.exists(root)) { return List.of(); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceGroupRefTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceGroupRefTest.java index 9b9e6fc7..22778ed2 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceGroupRefTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceGroupRefTest.java @@ -336,6 +336,56 @@ class SimilarAsinTaskServiceGroupRefTest { assertEquals(80, coverage); } + @Test + void test_task_006_group_payload_embeds_items_for_python() throws Exception { + // Python 旧链路(normalize_groups 直接返回 groups 并按 gp["items"] 取行): + // 写入存储的 payload 中每组必须内嵌完整行,startIndex/endIndex 保留供 Java 内部展开。 + java.util.concurrent.atomic.AtomicReference captured = new java.util.concurrent.atomic.AtomicReference<>(); + when(transientPayloadStorageService.storeParsedPayloadFast( + eq(SimilarAsinTaskService.MODULE_TYPE), any(), anyString(), anyString(), eq(false))) + .thenAnswer(invocation -> { + captured.set(invocation.getArgument(3)); + return "rustfs:task-parsed/similar-asin/30000/payload.json"; + }); + File workbook = buildWorkbook(12); + SimilarAsinParseVo vo = parse(workbook, "uploads/20260829/gr-python.xlsx"); + assertEquals(12, vo.getAcceptedRows()); + SimilarAsinParsedPayloadDto stored = objectMapper.readValue(captured.get(), SimilarAsinParsedPayloadDto.class); + assertNotNull(stored.getGroups()); + assertEquals(12, stored.getGroups().size()); + // 每组内嵌 items 与引用区间宽度一致,且行内容对应 items[startIndex] + for (SimilarAsinParsedGroupVo group : stored.getGroups()) { + assertNotNull(group.getItems()); + assertEquals(group.getEndIndex() - group.getStartIndex(), group.getItems().size(), + "Python 兼容:组内嵌 items 必须与引用区间宽度一致"); + int start = group.getStartIndex(); + assertEquals(stored.getItems().get(start).getRowToken(), group.getItems().get(0).getRowToken(), + "组 items 首行必须是 items[startIndex]"); + } + // Java 内部展开路径不受内嵌 items 影响:仍按引用切片 + List expanded = SimilarAsinTaskService.expandGroupRefs(stored); + assertEquals(12, expanded.size(), "Java 内部展开仍按引用,不受内嵌 items 影响"); + } + + @Test + void test_task_006_group_hydrate_embeds_items_for_legacy_payload() throws Exception { + // 存量任务 payload 只有 startIndex/endIndex 引用(无内嵌 items): + // parsedPayload 读取时 hydrate 必须回填组内嵌 items,Python 才能取到行。 + String legacy = "{\"items\":[" + + "{\"rowToken\":\"t1\",\"asin\":\"B0GRP00001\",\"sourceFileKey\":\"f.xlsx\",\"country\":\"英国\"}," + + "{\"rowToken\":\"t2\",\"asin\":\"B0GRP00002\",\"sourceFileKey\":\"f.xlsx\",\"country\":\"英国\"}" + + "],\"groups\":[{\"groupKey\":\"g1\",\"startIndex\":0,\"endIndex\":2,\"itemCount\":2}]}"; + SimilarAsinParsedPayloadDto payload = objectMapper.readValue(legacy, SimilarAsinParsedPayloadDto.class); + SimilarAsinParsedPayloadDto hydrated = SimilarAsinTaskService.hydrateForPythonCompat(payload); + assertEquals(1, hydrated.getGroups().size()); + assertNotNull(hydrated.getGroups().get(0).getItems()); + assertEquals(2, hydrated.getGroups().get(0).getItems().size(), "hydrate 后组内嵌 items 必须为 2 行"); + assertEquals("B0GRP00001", hydrated.getGroups().get(0).getItems().get(0).getAsin(), + "hydrate 后组内嵌 items 首行必须对应 items[0]"); + assertEquals("B0GRP00002", hydrated.getGroups().get(0).getItems().get(1).getAsin(), + "hydrate 后组内嵌 items 第二行必须对应 items[1]"); + } + // ---- helpers ---- private static String groupRefJson(int groupCount, int[][] ranges) { diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskFileJobDispatchCoordinatorOwnerRoutingTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskFileJobDispatchCoordinatorOwnerRoutingTest.java new file mode 100644 index 00000000..03755e60 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskFileJobDispatchCoordinatorOwnerRoutingTest.java @@ -0,0 +1,85 @@ +package com.nanri.aiimage.modules.task.service; + +import com.nanri.aiimage.config.InstanceMetadata; +import com.nanri.aiimage.modules.task.model.dto.TaskFileJobDispatchEvent; +import com.nanri.aiimage.modules.task.model.dto.TaskFileJobMessage; +import org.junit.jupiter.api.Test; +import org.springframework.core.task.TaskExecutor; + +import java.time.LocalDateTime; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +/** + * owner-scoped job 的派发路由(2026-09-01 事故):两台实例共用一个 RocketMQ consumer group, + * 消息被轮询投递到非归属实例后只能丢弃。归属实例改为直接本地派发,非归属实例不再发消息。 + */ +class TaskFileJobDispatchCoordinatorOwnerRoutingTest { + + private final TaskFileJobPublisher publisher = mock(TaskFileJobPublisher.class); + private final TaskFileJobLocalDispatcher localDispatcher = mock(TaskFileJobLocalDispatcher.class); + private final InstanceMetadata instanceMetadata = mock(InstanceMetadata.class); + private final TaskExecutor inlineExecutor = Runnable::run; + + private TaskFileJobDispatchCoordinator buildCoordinator() { + return new TaskFileJobDispatchCoordinator(publisher, localDispatcher, instanceMetadata, inlineExecutor); + } + + private static TaskFileJobDispatchEvent event(String owner) { + return new TaskFileJobDispatchEvent(18524L, 26402L, "SIMILAR_ASIN", 29405L, + owner == null ? "task:26402" : "task:26402:owner:" + owner, + TaskFileJobService.JOB_TYPE_ASSEMBLE_RESULT, owner, LocalDateTime.now()); + } + + @Test + void ownerJobDispatchedLocallyWithoutMq() { + when(instanceMetadata.getInstanceId()).thenReturn("server-121"); + when(localDispatcher.dispatch(18524L, 26402L, "SIMILAR_ASIN", true)).thenReturn(true); + + buildCoordinator().onTaskFileJobDispatch(event("server-121")); + + verify(localDispatcher).dispatch(18524L, 26402L, "SIMILAR_ASIN", true); + verifyNoInteractions(publisher); + } + + @Test + void ownerJobFallsBackToMqWhenLocalDispatchDeclined() { + when(instanceMetadata.getInstanceId()).thenReturn("server-121"); + when(localDispatcher.dispatch(18524L, 26402L, "SIMILAR_ASIN", true)).thenReturn(false, true); + when(publisher.publish(any(TaskFileJobMessage.class))).thenReturn(false); + + buildCoordinator().onTaskFileJobDispatch(event("server-121")); + + verify(publisher).publish(any(TaskFileJobMessage.class)); + } + + @Test + void otherInstanceJobIsNeitherPublishedNorRun() { + when(instanceMetadata.getInstanceId()).thenReturn("server-110"); + + buildCoordinator().onTaskFileJobDispatch(event("server-121")); + + verifyNoInteractions(publisher); + verify(localDispatcher, never()).dispatch(anyLong(), anyLong(), any(), anyBoolean()); + verify(localDispatcher, never()).dispatch(anyLong(), anyLong(), any()); + } + + @Test + void jobWithoutOwnerStillGoesThroughMq() { + when(instanceMetadata.getInstanceId()).thenReturn("server-110"); + when(publisher.publish(any(TaskFileJobMessage.class))).thenReturn(true); + + buildCoordinator().onTaskFileJobDispatch(event(null)); + + verify(publisher).publish(any(TaskFileJobMessage.class)); + verify(localDispatcher, never()).dispatch(anyLong(), anyLong(), eq("SIMILAR_ASIN"), anyBoolean()); + } +} diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorkerOwnerRoutingTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorkerOwnerRoutingTest.java new file mode 100644 index 00000000..1df61898 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorkerOwnerRoutingTest.java @@ -0,0 +1,163 @@ +package com.nanri.aiimage.modules.task.service; + +import com.nanri.aiimage.common.exception.TaskOwnerMismatchException; +import com.nanri.aiimage.config.InstanceMetadata; +import com.nanri.aiimage.modules.appearancepatent.service.AppearancePatentTaskService; +import com.nanri.aiimage.modules.brand.service.BrandTaskService; +import com.nanri.aiimage.modules.collectdata.service.CollectDataService; +import com.nanri.aiimage.modules.deletebrand.service.DeleteBrandRunService; +import com.nanri.aiimage.modules.patroldelete.service.PatrolDeleteTaskService; +import com.nanri.aiimage.modules.pricetrack.service.PriceTrackTaskService; +import com.nanri.aiimage.modules.productrisk.service.ProductRiskTaskService; +import com.nanri.aiimage.modules.publish.service.PublishTaskService; +import com.nanri.aiimage.modules.queryasin.service.QueryAsinTaskService; +import com.nanri.aiimage.modules.shopdatacrawl.service.ShopDataCrawlTaskService; +import com.nanri.aiimage.modules.shopmatch.service.ShopMatchTaskService; +import com.nanri.aiimage.modules.similarasin.service.SimilarAsinTaskService; +import com.nanri.aiimage.modules.task.mapper.FileResultMapper; +import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity; +import com.nanri.aiimage.modules.withdraw.service.WithdrawTaskService; +import org.junit.jupiter.api.Test; + +import java.lang.reflect.Field; +import java.time.LocalDateTime; +import java.util.List; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * 双实例归属路由回归测试(2026-09-01 货源查询任务被非归属实例误杀): + * 1) 定时 worker 已 claim 的 job 不能再走 claimRunning,否则条件 UPDATE 失败被静默丢弃; + * 2) 处理中撞上归属校验时只退回 PENDING,不得计入 retryCount(markFailed)。 + */ +class TaskResultFileJobWorkerOwnerRoutingTest { + + private final ShopMatchTaskService shopMatch = mock(ShopMatchTaskService.class); + private final PriceTrackTaskService priceTrack = mock(PriceTrackTaskService.class); + private final ProductRiskTaskService productRisk = mock(ProductRiskTaskService.class); + private final PublishTaskService publish = mock(PublishTaskService.class); + private final QueryAsinTaskService queryAsin = mock(QueryAsinTaskService.class); + private final ShopDataCrawlTaskService shopDataCrawl = mock(ShopDataCrawlTaskService.class); + private final WithdrawTaskService withdraw = mock(WithdrawTaskService.class); + private final PatrolDeleteTaskService patrolDelete = mock(PatrolDeleteTaskService.class); + private final AppearancePatentTaskService appearance = mock(AppearancePatentTaskService.class); + private final SimilarAsinTaskService similar = mock(SimilarAsinTaskService.class); + private final DeleteBrandRunService deleteBrand = mock(DeleteBrandRunService.class); + private final BrandTaskService brand = mock(BrandTaskService.class); + private final CollectDataService collectData = mock(CollectDataService.class); + private final TaskResultPayloadService payload = mock(TaskResultPayloadService.class); + private final TaskFileJobService taskFileJobService = mock(TaskFileJobService.class); + private final TaskDistributedLockService taskDistributedLockService = mock(TaskDistributedLockService.class); + private final TaskDistributedLockService.LockHandle lock = mock(TaskDistributedLockService.LockHandle.class); + private final TaskFileJobLocalDispatcher localDispatcher = mock(TaskFileJobLocalDispatcher.class); + private final InstanceMetadata instanceMetadata = mock(InstanceMetadata.class); + + private TaskResultFileJobWorker buildWorker() throws Exception { + List handlers = List.of( + new ShopMatchResultFileJobHandler(shopMatch, payload), + new PriceTrackResultFileJobHandler(priceTrack, payload), + new ProductRiskResultFileJobHandler(productRisk, payload), + new PublishResultFileJobHandler(publish), + new QueryAsinResultFileJobHandler(queryAsin, payload), + new ShopDataCrawlResultFileJobHandler(shopDataCrawl, payload), + new WithdrawResultFileJobHandler(withdraw, payload), + new PatrolDeleteResultFileJobHandler(patrolDelete, payload), + new AppearancePatentResultFileJobHandler(appearance), + new SimilarAsinResultFileJobHandler(similar), + new DeleteBrandResultFileJobHandler(deleteBrand), + new BrandResultFileJobHandler(brand, payload), + new CollectDataResultFileJobHandler(collectData)); + TaskResultFileJobWorker worker = new TaskResultFileJobWorker( + taskFileJobService, + taskDistributedLockService, + mock(FileResultMapper.class), + localDispatcher, + instanceMetadata, + withdraw, brand, + new ResultFileJobHandlerRegistry(handlers)); + set(worker, "taskQueueExecutor", (org.springframework.core.task.TaskExecutor) Runnable::run); + set(worker, "localWorkerEnabled", true); + set(worker, "batchSize", 20); + return worker; + } + + private static void set(TaskResultFileJobWorker worker, String fieldName, Object value) throws Exception { + Field field = TaskResultFileJobWorker.class.getDeclaredField(fieldName); + field.setAccessible(true); + field.set(worker, value); + } + + private static TaskFileJobEntity claimedJob(String moduleType, long jobId, long taskId, String scopeKey) { + TaskFileJobEntity entity = new TaskFileJobEntity(); + entity.setId(jobId); + entity.setTaskId(taskId); + entity.setResultId(taskId + 1000L); + entity.setModuleType(moduleType); + entity.setScopeKey(scopeKey); + entity.setStatus("RUNNING"); + entity.setUpdatedAt(LocalDateTime.now()); + return entity; + } + + @Test + void scheduledWorkerProcessesAlreadyClaimedJobWithoutReclaiming() throws Exception { + TaskResultFileJobWorker worker = buildWorker(); + when(instanceMetadata.getInstanceId()).thenReturn("instance-a"); + TaskFileJobEntity claim = claimedJob(PublishTaskService.MODULE_TYPE, 18524L, 26402L, "task:26402:owner:instance-a"); + when(taskFileJobService.claimRunnableJobsForOwner(20, "instance-a")).thenReturn(List.of(claim)); + when(taskFileJobService.activateRunningClaim(claim)).thenReturn(true); + when(localDispatcher.dispatchClaimed(claim.getId(), claim.getTaskId(), claim.getModuleType())).thenReturn(false); + when(taskDistributedLockService.acquire(PublishTaskService.MODULE_TYPE, claim.getTaskId(), + TaskDistributedLockService.DEFAULT_WAIT_MILLIS)).thenReturn(lock); + + worker.runPendingJobs(); + + verify(taskFileJobService, never()).claimRunning(anyLong()); + verify(publish).processResultFileJob(claim); + verify(taskFileJobService).markSuccess(eq(claim), any()); + } + + @Test + void scheduledWorkerPrefersClaimedDispatch() throws Exception { + TaskResultFileJobWorker worker = buildWorker(); + when(instanceMetadata.getInstanceId()).thenReturn("instance-a"); + TaskFileJobEntity claim = claimedJob(PublishTaskService.MODULE_TYPE, 18525L, 26403L, "task:26403:owner:instance-a"); + when(taskFileJobService.claimRunnableJobsForOwner(20, "instance-a")).thenReturn(List.of(claim)); + when(localDispatcher.dispatchClaimed(claim.getId(), claim.getTaskId(), claim.getModuleType())).thenReturn(true); + + worker.runPendingJobs(); + + verify(localDispatcher).dispatchClaimed(claim.getId(), claim.getTaskId(), claim.getModuleType()); + verify(localDispatcher, never()).dispatch(anyLong(), anyLong(), any()); + verify(taskFileJobService, never()).claimRunning(anyLong()); + verify(publish, never()).processResultFileJob(any()); + } + + @Test + void ownerMismatchDuringProcessingDefersInsteadOfFailing() throws Exception { + TaskResultFileJobWorker worker = buildWorker(); + when(instanceMetadata.getInstanceId()).thenReturn("instance-a"); + // scopeKey 归属本机(通过 process 的前置校验),但模块层读到的 task owner 已经变成他机 + TaskFileJobEntity job = claimedJob(PublishTaskService.MODULE_TYPE, 18526L, 26404L, "task:26404:owner:instance-a"); + when(taskFileJobService.claimRunning(job.getId())).thenReturn(job); + when(taskFileJobService.activateRunningClaim(job)).thenReturn(true); + when(taskDistributedLockService.acquire(PublishTaskService.MODULE_TYPE, job.getTaskId(), + TaskDistributedLockService.DEFAULT_WAIT_MILLIS)).thenReturn(lock); + doThrow(new TaskOwnerMismatchException(job.getTaskId(), "assemble result file", "instance-b", "instance-a")) + .when(publish).processResultFileJob(job); + + worker.process(job); + + verify(taskFileJobService).deferRunning(eq(job.getId()), any()); + verify(taskFileJobService, never()).markFailed(any(), any()); + verify(taskFileJobService, never()).markFailedPermanent(any(), any()); + verify(taskFileJobService, never()).isRetryExhausted(anyLong()); + } +} diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorkerOwnerScopedTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorkerOwnerScopedTest.java index 68ebfb33..11933b96 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorkerOwnerScopedTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskResultFileJobWorkerOwnerScopedTest.java @@ -128,6 +128,20 @@ class TaskResultFileJobWorkerOwnerScopedTest { TaskResultFileJobWorker worker = buildWorker(); when(instanceMetadata.getInstanceId()).thenReturn("instance-a"); TaskFileJobEntity job = job("PUBLISH", 2L, 12L, "task:12:owner:instance-b"); + job.setStatus("RUNNING"); + worker.process(job); + verifyNoInteractions(taskFileJobService, taskDistributedLockService, publish); + } + + @Test + void ownerScopedOtherOwnerPendingSkip() { + // PENDING + owner 不匹配也必须跳过:模块层 ensureTaskOwnedByCurrentInstance 会抛 + // TaskOwnerMismatchException,本机抢跑只会耗尽 job 重试次数并误杀归属机上的任务。 + // job 由归属实例的 runPendingJobs(每 15 秒扫自己 owner 的 PENDING/FAILED)接手。 + TaskResultFileJobWorker worker = buildWorker(); + when(instanceMetadata.getInstanceId()).thenReturn("instance-a"); + TaskFileJobEntity job = job("PUBLISH", 2L, 12L, "task:12:owner:instance-b"); + job.setStatus("PENDING"); worker.process(job); verifyNoInteractions(taskFileJobService, taskDistributedLockService, publish); }