From ddefcbed56d4241078158198447e3ce7a3ddd7ac Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=BB=84=E8=87=AA=E8=BE=BE?= <980324341@qq.com> Date: Thu, 17 Sep 2026 11:04:38 +0800 Subject: [PATCH] =?UTF-8?q?feat(=E4=BB=BB=E5=8A=A1=E6=81=A2=E5=A4=8D):=20?= =?UTF-8?q?=E5=A4=96=E8=A7=82=E4=B8=93=E5=88=A9=E6=8E=A5=E4=B8=8A=E3=80=8C?= =?UTF-8?q?=E8=A1=A5=E4=BC=A0=E6=81=A2=E5=A4=8D=E3=80=8D=EF=BC=9B=E5=90=88?= =?UTF-8?q?=E5=B9=B6=E5=86=B2=E7=AA=81=E4=BF=9D=E7=95=99=E5=85=9C=E5=BA=95?= =?UTF-8?q?=E5=AF=B9=E8=B1=A1=E5=B9=B6=E5=8F=AF=E8=AF=8A=E6=96=AD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 28459 的两个遗留项: 1) 外观专利缺「补传恢复」入口——分片缺失致组装 job 重试耗尽后,客户端补传缺口 也无法自动重跑,只能人工重置 job。TaskFileJobService 早有 resetTerminalFailedForRecovery,但只被删除品牌模块接了。现按同一口径在分片 提交成功后检查并恢复;best-effort,恢复失败不影响补传本身。 2) 合并 CAS 冲突不可诊断且会删掉唯一的兜底对象——原实现每次冲突都删掉刚写入的 版本化对象,重试耗尽即抛异常、行仍指向旧指针;一旦旧对象也不在,该分片永久 读不到(28459 的 chunk-462/473/484 正是这个形态)。现在冲突时读回行上当前 哈希并写进异常与日志;终局失败保留最后一个对象,作为读路径「同槽位兄弟对象」 兜底的恢复源。外观专利与相似ASIN 同一口径。 测试:新增 17 个用例(补传恢复 8 / 外观专利合并冲突 5 / 相似ASIN 合并冲突 4), RED 均已确认;全量 mvn test 3099 个 0 失败。 --- .../service/AppearancePatentTaskService.java | 78 +++++- .../support/SimilarAsinPipelineSupport.java | 31 ++- ...ppearancePatentChunkMergeConflictTest.java | 218 ++++++++++++++++ ...rancePatentTerminalFailedRecoveryTest.java | 243 ++++++++++++++++++ ...larAsinTaskServiceChunkMergeLimitTest.java | 97 +++++++ 5 files changed, 659 insertions(+), 8 deletions(-) create mode 100644 backend-java/src/test/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentChunkMergeConflictTest.java create mode 100644 backend-java/src/test/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTerminalFailedRecoveryTest.java 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 a0e492e0..267615db 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 @@ -424,6 +424,53 @@ public class AppearancePatentTaskService { public void submitResult(Long taskId, AppearancePatentSubmitResultRequest request) { submitResultLocked(taskId, request); + // 分片补传后自动恢复终态失败的组装 job —— 走与删除品牌同一套口径, + // 此前外观专利没有接该入口,分片补齐后只能人工重置 job(线上任务 28459 即如此)。 + maybeRecoverTerminalFailedAssemble(taskId); + } + + /** + * 补传恢复:此前因分片缺失导致组装 job 重试耗尽(终态失败),客户端补传缺口后 + * 把失败的组装 job 重置为 PENDING 重新派发({@code resetTerminalFailedForRecovery} 自行补发 dispatch 事件)。 + * + *
best-effort:恢复失败不得影响补传本身——分片已经落库,恢复只是让后续组装继续推进。
+ * 常态(无终态失败 job)下只查两次即返回,不触发分片扫描。
+ */
+ private void maybeRecoverTerminalFailedAssemble(Long taskId) {
+ if (taskId == null || taskId <= 0) {
+ return;
+ }
+ try {
+ FileResultEntity result = findResultRecord(taskId);
+ if (result == null) {
+ return;
+ }
+ if (taskFileJobService.hasSuccessfulAssembleJob(taskId, MODULE_TYPE, result.getId())) {
+ return;
+ }
+ if (!taskFileJobService.isTerminalFailedAssembleJob(taskId, MODULE_TYPE, result.getId())) {
+ return;
+ }
+ if (!isResultSubmissionComplete(taskId)) {
+ log.info("[appearance-patent] 组装 job 终态失败但分片仍未补传完整,暂不恢复 taskId={} resultId={}",
+ taskId, result.getId());
+ return;
+ }
+ log.info("[appearance-patent] 分片已补传完整,恢复终态失败的组装 job taskId={} resultId={}",
+ taskId, result.getId());
+ taskFileJobService.resetTerminalFailedForRecovery(taskId, MODULE_TYPE, result.getId());
+ } catch (Exception ex) {
+ log.warn("[appearance-patent] 补传恢复检查失败(不影响本次补传)taskId={} err={}", taskId, ex.getMessage(), ex);
+ }
+ }
+
+ /** 按 task 取结果行(不创建);不存在返回 null。 */
+ private FileResultEntity findResultRecord(Long taskId) {
+ List 原实现在每次 CAS 冲突后都删掉刚写入的版本化对象,重试耗尽即抛异常、行仍指向旧指针——
+ * 一旦旧对象也不在,该分片就永久读不到(28459 的 chunk-462/473/484 正是这个形态)。
+ * 现在:冲突时读回行上的当前哈希以便定位;**终局失败保留最后一个兜底对象**,
+ * 让读路径的「同槽位兄弟对象」兜底仍有数据可取。
+ */
+@ExtendWith(MockitoExtension.class)
+class AppearancePatentChunkMergeConflictTest {
+
+ private static final Long TASK_ID = 28459L;
+ private static final String SCOPE_HASH = "2248d39710545b47b7c7035fc60c4e33924a16174abf65dd1c5a230918e6f61e";
+ private static final int CHUNK_INDEX = 462;
+
+ @Mock private LocalFileStorageService localFileStorageService;
+ @Mock private OssStorageService ossStorageService;
+ @Mock private StorageProperties storageProperties;
+ @Mock private FileTaskMapper fileTaskMapper;
+ @Mock private FileResultMapper fileResultMapper;
+ @Mock private TaskScopeStateMapper taskScopeStateMapper;
+ @Mock private TaskChunkMapper taskChunkMapper;
+ @Spy private ObjectMapper objectMapper = new ObjectMapper();
+ @Mock private AppearancePatentLlmClient llmClient;
+ @Mock private AppearancePatentTaskCacheService taskCacheService;
+ @Mock private AppearancePatentProperties properties;
+ @Mock private TaskFileJobService taskFileJobService;
+ @Mock private TaskProgressSnapshotService taskProgressSnapshotService;
+ @Mock private TransientPayloadStorageService transientPayloadStorageService;
+ @Mock private PlatformTransactionManager transactionManager;
+ @Mock private DistributedJobLockService distributedJobLockService;
+ @Mock private TaskDistributedLockService taskDistributedLockService;
+ @Mock private InstanceMetadata instanceMetadata;
+ @Mock private TaskProgressLightAssembler taskProgressLightAssembler;
+
+ @InjectMocks private AppearancePatentTaskService service;
+
+ @BeforeAll
+ static void initializeMybatisMetadata() {
+ MapperBuilderAssistant assistant = new MapperBuilderAssistant(new MybatisConfiguration(), "");
+ TableInfoHelper.initTableInfo(assistant, FileTaskEntity.class);
+ TableInfoHelper.initTableInfo(assistant, FileResultEntity.class);
+ TableInfoHelper.initTableInfo(assistant, TaskChunkEntity.class);
+ TableInfoHelper.initTableInfo(assistant, TaskScopeStateEntity.class);
+ }
+
+ @BeforeEach
+ void setUp() {
+ lenient().when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
+ // 共享写开启才会走版本化对象存储;否则落到本地兜底路径直接报「RustFS 未配置」
+ lenient().when(transientPayloadStorageService.isSharedWriteEnabled()).thenReturn(true);
+ }
+
+ /** 首次冲突后重试成功:中间那次写的对象要删(会被下次重写),且最终行被改到新对象。 */
+ @Test
+ void retryAfterConflictDeletesSupersededObject() {
+ stubChunk("hash-old", "ptr-chunk-462.json");
+ AtomicInteger stores = new AtomicInteger();
+ stubVersionedStore(stores);
+ when(taskChunkMapper.update(any(), any())).thenReturn(0, 1);
+ when(taskChunkMapper.selectById(anyLong())).thenReturn(chunk("hash-other", "ptr-chunk-462.json"));
+
+ invokeMerge();
+
+ verify(transientPayloadStorageService, times(2)).storeChunkPayloadVersioned(anyString(), any(), anyString(), any(), anyString());
+ // 第 1 次冲突写的对象被删;第 2 次成功,走的是「替换旧对象」而不是删新对象
+ verify(transientPayloadStorageService).deletePayloadIfPresent("sibling-1");
+ verify(transientPayloadStorageService, never()).deletePayloadIfPresent("sibling-2");
+ verify(transientPayloadStorageService).deleteReplacedPayloadIfNeeded(eq("ptr-chunk-462.json"), eq("sibling-2"));
+ }
+
+ /** 冲突重试耗尽:**保留**最后一次写入的兜底对象(本次要修的形态),并抛异常带出两个哈希。 */
+ @Test
+ void exhaustedConflictKeepsLastStoredObjectAsFallback() {
+ stubChunk("hash-old", "ptr-chunk-462.json");
+ AtomicInteger stores = new AtomicInteger();
+ stubVersionedStore(stores);
+ when(taskChunkMapper.update(any(), any())).thenReturn(0);
+ when(taskChunkMapper.selectById(anyLong())).thenReturn(chunk("hash-other", "ptr-chunk-462.json"));
+
+ IllegalStateException ex = assertThrows(IllegalStateException.class, this::invokeMerge);
+
+ verify(transientPayloadStorageService, times(3)).storeChunkPayloadVersioned(anyString(), any(), anyString(), any(), anyString());
+ // 前两次冲突的对象照旧删除;第三次(终局)的对象必须保留
+ verify(transientPayloadStorageService).deletePayloadIfPresent("sibling-1");
+ verify(transientPayloadStorageService).deletePayloadIfPresent("sibling-2");
+ verify(transientPayloadStorageService, never()).deletePayloadIfPresent("sibling-3");
+ assertTrue(ex.getMessage().contains("hash-old"), "异常需带出期望哈希,实际: " + ex.getMessage());
+ assertTrue(ex.getMessage().contains("hash-other"), "异常需带出当前哈希,实际: " + ex.getMessage());
+ }
+
+ /** 冲突时读回行上的当前哈希,供定位(此前只有一句 update conflict,线上无法定位)。 */
+ @Test
+ void conflictReadsBackCurrentHashForDiagnostics() {
+ stubChunk("hash-old", "ptr-chunk-462.json");
+ stubVersionedStore(new AtomicInteger());
+ when(taskChunkMapper.update(any(), any())).thenReturn(0);
+ when(taskChunkMapper.selectById(7L)).thenReturn(chunk("hash-changed-by-other-writer", "ptr-chunk-462.json"));
+
+ assertThrows(IllegalStateException.class, this::invokeMerge);
+
+ verify(taskChunkMapper, times(3)).selectById(7L);
+ }
+
+ /** 读回行失败不得掩盖原始冲突:异常信息里给出可读标记。 */
+ @Test
+ void readBackFailureDoesNotMaskConflict() {
+ stubChunk("hash-old", "ptr-chunk-462.json");
+ stubVersionedStore(new AtomicInteger());
+ when(taskChunkMapper.update(any(), any())).thenReturn(0);
+ when(taskChunkMapper.selectById(anyLong())).thenThrow(new IllegalStateException("db down"));
+
+ IllegalStateException ex = assertThrows(IllegalStateException.class, this::invokeMerge);
+
+ assertTrue(ex.getMessage().contains("读取失败"), "实际: " + ex.getMessage());
+ }
+
+ // ===== 辅助 =====
+
+ private void invokeMerge() {
+ AppearancePatentResultRowDto row = new AppearancePatentResultRowDto();
+ row.setRowToken("r1");
+ row.setAsin("B0A0000001");
+ ReflectionTestUtils.invokeMethod(service, "mergeChunkPayload", TASK_ID, SCOPE_HASH, CHUNK_INDEX, List.of(row));
+ }
+
+ private void stubChunk(String payloadHash, String payloadPointer) {
+ when(taskChunkMapper.selectOne(any())).thenReturn(chunk(payloadHash, payloadPointer));
+ lenient().when(transientPayloadStorageService.resolvePayload(eq(payloadPointer), anyString()))
+ .thenReturn("[{\"rowToken\":\"r1\",\"asin\":\"B0A0000001\"}]");
+ }
+
+ private void stubVersionedStore(AtomicInteger stores) {
+ lenient().when(transientPayloadStorageService.storeChunkPayloadVersioned(
+ anyString(), anyLong(), anyString(), any(), anyString()))
+ .thenAnswer(inv -> "sibling-" + stores.incrementAndGet());
+ }
+
+ private TaskChunkEntity chunk(String payloadHash, String payloadPointer) {
+ TaskChunkEntity chunk = new TaskChunkEntity();
+ chunk.setId(7L);
+ chunk.setTaskId(TASK_ID);
+ chunk.setModuleType(AppearancePatentTaskService.MODULE_TYPE);
+ chunk.setScopeHash(SCOPE_HASH);
+ chunk.setChunkIndex(CHUNK_INDEX);
+ chunk.setPayloadJson(payloadPointer);
+ chunk.setPayloadHash(payloadHash);
+ return chunk;
+ }
+
+ /** 断言辅助:确认存的对象数(避免误用未使用的 import)。 */
+ @Test
+ void storeIsCalledOnceWhenUpdateSucceeds() {
+ stubChunk("hash-old", "ptr-chunk-462.json");
+ AtomicInteger stores = new AtomicInteger();
+ stubVersionedStore(stores);
+ when(taskChunkMapper.update(any(), any())).thenReturn(1);
+
+ invokeMerge();
+
+ assertEquals(1, stores.get());
+ verify(taskChunkMapper, never()).selectById(anyLong());
+ }
+}
diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTerminalFailedRecoveryTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTerminalFailedRecoveryTest.java
new file mode 100644
index 00000000..690b27bc
--- /dev/null
+++ b/backend-java/src/test/java/com/nanri/aiimage/modules/appearancepatent/service/AppearancePatentTerminalFailedRecoveryTest.java
@@ -0,0 +1,243 @@
+package com.nanri.aiimage.modules.appearancepatent.service;
+
+import com.baomidou.mybatisplus.core.MybatisConfiguration;
+import com.baomidou.mybatisplus.core.metadata.TableInfoHelper;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.nanri.aiimage.common.service.DistributedJobLockService;
+import com.nanri.aiimage.config.AppearancePatentProperties;
+import com.nanri.aiimage.config.InstanceMetadata;
+import com.nanri.aiimage.config.StorageProperties;
+import com.nanri.aiimage.modules.appearancepatent.client.AppearancePatentLlmClient;
+import com.nanri.aiimage.modules.appearancepatent.model.dto.AppearancePatentSubmitResultRequest;
+import com.nanri.aiimage.modules.file.service.LocalFileStorageService;
+import com.nanri.aiimage.modules.file.service.oss.OssStorageService;
+import com.nanri.aiimage.modules.task.mapper.FileResultMapper;
+import com.nanri.aiimage.modules.task.mapper.FileTaskMapper;
+import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper;
+import com.nanri.aiimage.modules.task.mapper.TaskScopeStateMapper;
+import com.nanri.aiimage.modules.task.model.entity.FileResultEntity;
+import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
+import com.nanri.aiimage.modules.task.model.entity.TaskChunkEntity;
+import com.nanri.aiimage.modules.task.model.entity.TaskScopeStateEntity;
+import com.nanri.aiimage.modules.task.service.TaskDistributedLockService;
+import com.nanri.aiimage.modules.task.service.TaskFileJobService;
+import com.nanri.aiimage.modules.task.service.TaskProgressLightAssembler;
+import com.nanri.aiimage.modules.task.service.TaskProgressSnapshotService;
+import com.nanri.aiimage.modules.task.service.TransientPayloadStorageService;
+import org.apache.ibatis.builder.MapperBuilderAssistant;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+import org.mockito.Spy;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.test.util.ReflectionTestUtils;
+import org.springframework.transaction.PlatformTransactionManager;
+import org.springframework.transaction.TransactionDefinition;
+import org.springframework.transaction.TransactionStatus;
+
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.lenient;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * 外观专利「补传恢复」入口(2026-09-17 线上任务 28459 的遗留项)。
+ *
+ * 分片缺失导致组装 job 重试耗尽后,客户端补传缺口应能自动把该 job 重置重跑。
+ * 该能力在 {@code TaskFileJobService} 里早就有了,但只有删除品牌模块接了入口,
+ * 外观专利没有——28459 补齐分片后仍需人工重置 job 才能出结果。
+ */
+@ExtendWith(MockitoExtension.class)
+class AppearancePatentTerminalFailedRecoveryTest {
+
+ private static final Long TASK_ID = 28459L;
+ private static final Long RESULT_ID = 31490L;
+ private static final String MODULE_TYPE = "APPEARANCE_PATENT";
+
+ @Mock private LocalFileStorageService localFileStorageService;
+ @Mock private OssStorageService ossStorageService;
+ @Mock private StorageProperties storageProperties;
+ @Mock private FileTaskMapper fileTaskMapper;
+ @Mock private FileResultMapper fileResultMapper;
+ @Mock private TaskScopeStateMapper taskScopeStateMapper;
+ @Mock private TaskChunkMapper taskChunkMapper;
+ @Spy private ObjectMapper objectMapper = new ObjectMapper();
+ @Mock private AppearancePatentLlmClient llmClient;
+ @Mock private AppearancePatentTaskCacheService taskCacheService;
+ @Mock private AppearancePatentProperties properties;
+ @Mock private TaskFileJobService taskFileJobService;
+ @Mock private TaskProgressSnapshotService taskProgressSnapshotService;
+ @Mock private TransientPayloadStorageService transientPayloadStorageService;
+ @Mock private PlatformTransactionManager transactionManager;
+ @Mock private DistributedJobLockService distributedJobLockService;
+ @Mock private TaskDistributedLockService taskDistributedLockService;
+ @Mock private InstanceMetadata instanceMetadata;
+ @Mock private TaskProgressLightAssembler taskProgressLightAssembler;
+ @Mock private TransactionStatus transactionStatus;
+
+ @InjectMocks private AppearancePatentTaskService service;
+
+ @BeforeAll
+ static void initializeMybatisMetadata() {
+ MapperBuilderAssistant assistant = new MapperBuilderAssistant(new MybatisConfiguration(), "");
+ TableInfoHelper.initTableInfo(assistant, FileTaskEntity.class);
+ TableInfoHelper.initTableInfo(assistant, FileResultEntity.class);
+ TableInfoHelper.initTableInfo(assistant, TaskChunkEntity.class);
+ TableInfoHelper.initTableInfo(assistant, TaskScopeStateEntity.class);
+ }
+
+ @BeforeEach
+ void setUp() {
+ lenient().when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
+ lenient().when(transactionManager.getTransaction(any(TransactionDefinition.class))).thenReturn(transactionStatus);
+ lenient().doAnswer(inv -> null).when(transactionManager).commit(transactionStatus);
+ lenient().doAnswer(inv -> null).when(transactionManager).rollback(transactionStatus);
+ lenient().when(taskChunkMapper.selectList(any())).thenReturn(List.of());
+ lenient().when(taskChunkMapper.selectCount(any())).thenReturn(0L);
+ lenient().when(taskChunkMapper.selectOne(any())).thenReturn(null);
+ lenient().when(taskScopeStateMapper.selectOne(any())).thenReturn(null);
+ lenient().when(fileTaskMapper.updateById(any(FileTaskEntity.class))).thenReturn(1);
+ lenient().when(transientPayloadStorageService.isSharedWriteEnabled()).thenReturn(true);
+ lenient().when(transientPayloadStorageService.storeChunkPayload(
+ eq(MODULE_TYPE), eq(TASK_ID), anyString(), any(), anyString()))
+ .thenReturn("\"rustfs:task-chunk/appearance_patent/28459/hash/chunk-11.json\"");
+ }
+
+ // ===== 直接覆盖恢复判定 =====
+
+ /** 已有成功的组装 job → 不恢复。 */
+ @Test
+ void successfulAssembleJobSkipsRecovery() {
+ stubResultRow();
+ when(taskFileJobService.hasSuccessfulAssembleJob(TASK_ID, MODULE_TYPE, RESULT_ID)).thenReturn(true);
+
+ invokeRecovery();
+
+ verify(taskFileJobService, never()).resetTerminalFailedForRecovery(anyLong(), anyString(), anyLong());
+ }
+
+ /** 组装 job 不是「重试耗尽的终态失败」→ 不恢复。 */
+ @Test
+ void nonTerminalFailedJobSkipsRecovery() {
+ stubResultRow();
+ when(taskFileJobService.hasSuccessfulAssembleJob(TASK_ID, MODULE_TYPE, RESULT_ID)).thenReturn(false);
+ when(taskFileJobService.isTerminalFailedAssembleJob(TASK_ID, MODULE_TYPE, RESULT_ID)).thenReturn(false);
+
+ invokeRecovery();
+
+ verify(taskFileJobService, never()).resetTerminalFailedForRecovery(anyLong(), anyString(), anyLong());
+ }
+
+ /** 终态失败 + 分片已补传完整 → 重置 job 重新派发(本次要修的场景)。 */
+ @Test
+ void terminalFailedJobWithCompleteChunksIsReset() {
+ stubResultRow();
+ when(taskFileJobService.hasSuccessfulAssembleJob(TASK_ID, MODULE_TYPE, RESULT_ID)).thenReturn(false);
+ when(taskFileJobService.isTerminalFailedAssembleJob(TASK_ID, MODULE_TYPE, RESULT_ID)).thenReturn(true);
+ when(taskScopeStateMapper.selectCount(any())).thenReturn(1L);
+
+ invokeRecovery();
+
+ verify(taskFileJobService).resetTerminalFailedForRecovery(TASK_ID, MODULE_TYPE, RESULT_ID);
+ }
+
+ /** 终态失败但分片尚未补齐 → 不恢复(否则又会读到缺失分片再失败一次)。 */
+ @Test
+ void terminalFailedJobWithIncompleteChunksIsNotReset() {
+ stubResultRow();
+ when(taskFileJobService.hasSuccessfulAssembleJob(TASK_ID, MODULE_TYPE, RESULT_ID)).thenReturn(false);
+ when(taskFileJobService.isTerminalFailedAssembleJob(TASK_ID, MODULE_TYPE, RESULT_ID)).thenReturn(true);
+ when(taskScopeStateMapper.selectCount(any())).thenReturn(0L);
+
+ invokeRecovery();
+
+ verify(taskFileJobService, never()).resetTerminalFailedForRecovery(anyLong(), anyString(), anyLong());
+ }
+
+ /** 还没有结果行 → 不恢复。 */
+ @Test
+ void missingResultRowSkipsRecovery() {
+ when(fileResultMapper.selectList(any())).thenReturn(List.of());
+
+ invokeRecovery();
+
+ verify(taskFileJobService, never()).isTerminalFailedAssembleJob(anyLong(), anyString(), anyLong());
+ }
+
+ /** taskId 非法 → 直接返回,不查库。 */
+ @Test
+ void invalidTaskIdSkipsLookup() {
+ ReflectionTestUtils.invokeMethod(service, "maybeRecoverTerminalFailedAssemble", 0L);
+ ReflectionTestUtils.invokeMethod(service, "maybeRecoverTerminalFailedAssemble", (Object) null);
+
+ verify(fileResultMapper, never()).selectList(any());
+ }
+
+ /** 恢复检查自身抛异常 → 吞掉,不影响补传结果(best-effort)。 */
+ @Test
+ void recoveryFailureDoesNotPropagate() {
+ when(fileResultMapper.selectList(any())).thenThrow(new IllegalStateException("db down"));
+
+ assertDoesNotThrow(this::invokeRecovery);
+ }
+
+ // ===== 走完整提交路径 =====
+
+ /** 分片提交成功后触发恢复检查(接线正确)。 */
+ @Test
+ void submitResultTriggersRecoveryCheck() {
+ stubRunningTask();
+ stubResultRow();
+ when(taskFileJobService.hasSuccessfulAssembleJob(TASK_ID, MODULE_TYPE, RESULT_ID)).thenReturn(false);
+ when(taskFileJobService.isTerminalFailedAssembleJob(TASK_ID, MODULE_TYPE, RESULT_ID)).thenReturn(true);
+ when(taskScopeStateMapper.selectCount(any())).thenReturn(1L);
+
+ service.submitResult(TASK_ID, request());
+
+ verify(taskFileJobService).resetTerminalFailedForRecovery(TASK_ID, MODULE_TYPE, RESULT_ID);
+ }
+
+ // ===== 辅助 =====
+
+ private void invokeRecovery() {
+ ReflectionTestUtils.invokeMethod(service, "maybeRecoverTerminalFailedAssemble", TASK_ID);
+ }
+
+ private void stubResultRow() {
+ FileResultEntity result = new FileResultEntity();
+ result.setId(RESULT_ID);
+ result.setTaskId(TASK_ID);
+ result.setModuleType(MODULE_TYPE);
+ lenient().when(fileResultMapper.selectList(any())).thenReturn(List.of(result));
+ }
+
+ private void stubRunningTask() {
+ FileTaskEntity task = new FileTaskEntity();
+ task.setId(TASK_ID);
+ task.setModuleType(MODULE_TYPE);
+ task.setStatus("RUNNING");
+ task.setUserId(1121L);
+ task.setResultJson("{\"parsedPayloadRef\":\"rustfs:task-parsed/x.json\",\"ownerInstanceId\":\"instance-a\"}");
+ when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task);
+ }
+
+ private AppearancePatentSubmitResultRequest request() {
+ AppearancePatentSubmitResultRequest request = new AppearancePatentSubmitResultRequest();
+ request.setSubmissionId("appearance-patent-" + TASK_ID);
+ request.setChunkIndex(11);
+ request.setChunkTotal(500);
+ request.setDone(false);
+ return request;
+ }
+}
diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceChunkMergeLimitTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceChunkMergeLimitTest.java
index d70c3854..4b28ce67 100644
--- a/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceChunkMergeLimitTest.java
+++ b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceChunkMergeLimitTest.java
@@ -33,6 +33,7 @@ import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
@@ -42,10 +43,12 @@ import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.lenient;
+import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -397,4 +400,98 @@ class SimilarAsinTaskServiceChunkMergeLimitTest {
verify(taskChunkMapper, times(1)).update(any(), any());
verify(taskScopeStateMapper, times(0)).insert(any(TaskScopeStateEntity.class));
}
+
+ // ===== CAS 冲突处理(2026-09-17 线上任务 28459 遗留项,与外观专利同一口径) =====
+
+ /** 冲突重试耗尽:**保留**最后一次写入的对象作读兜底,异常带出期望/当前哈希。 */
+ @Test
+ void casConflictExhaustedKeepsLastStoredObjectAsFallback() throws Exception {
+ TaskChunkEntity chunk = chunk(7L, "hashA", 1, "ptr:chunk-A");
+ chunk.setPayloadHash("hashA");
+ stubConflictMerge(chunk, 0);
+
+ FileTaskEntity task = new FileTaskEntity();
+ task.setId(9004L);
+ Throwable ex = invokeMergeExpectingFailure(task, List.of(row("r1", "B0A0000001", "标题1")));
+
+ assertTrue(ex instanceof IllegalStateException, "实际: " + ex);
+ verify(transientPayloadStorageService).deletePayloadIfPresent("sibling-1");
+ verify(transientPayloadStorageService).deletePayloadIfPresent("sibling-2");
+ verify(transientPayloadStorageService, never()).deletePayloadIfPresent("sibling-3");
+ assertTrue(ex.getMessage().contains("hashA"), "需带出期望哈希,实际: " + ex.getMessage());
+ assertTrue(ex.getMessage().contains("hashB"), "需带出当前哈希,实际: " + ex.getMessage());
+ }
+
+ /** 冲突后重试成功:被顶替的中间对象删除,最终对象写入行(不删)。 */
+ @Test
+ void casConflictThenSuccessDeletesSupersededObjects() throws Exception {
+ TaskChunkEntity chunk = chunk(7L, "hashA", 1, "ptr:chunk-A");
+ chunk.setPayloadHash("hashA");
+ stubConflictMerge(chunk, 0, 0, 1);
+
+ FileTaskEntity task = new FileTaskEntity();
+ task.setId(9004L);
+ invokeMerge(service, task, "hashA", 1, List.of(row("r1", "B0A0000001", "标题1")));
+
+ verify(transientPayloadStorageService).deletePayloadIfPresent("sibling-1");
+ verify(transientPayloadStorageService).deletePayloadIfPresent("sibling-2");
+ verify(transientPayloadStorageService, never()).deletePayloadIfPresent("sibling-3");
+ verify(transientPayloadStorageService).deleteReplacedPayloadIfNeeded(eq("ptr:chunk-A"), eq("sibling-3"));
+ }
+
+ /** 冲突时读回行上的当前哈希供定位(此前只有一句 conflict,线上无法定位)。 */
+ @Test
+ void casConflictReadsBackCurrentHashForDiagnostics() throws Exception {
+ TaskChunkEntity chunk = chunk(7L, "hashA", 1, "ptr:chunk-A");
+ chunk.setPayloadHash("hashA");
+ stubConflictMerge(chunk, 0);
+
+ FileTaskEntity task = new FileTaskEntity();
+ task.setId(9004L);
+ invokeMergeExpectingFailure(task, List.of(row("r1", "B0A0000001", "标题1")));
+
+ verify(taskChunkMapper, times(3)).selectById(7L);
+ }
+
+ /** 读回行失败不得掩盖原始冲突:异常信息里给出可读标记。 */
+ @Test
+ void casConflictReadBackFailureDoesNotMaskConflict() throws Exception {
+ TaskChunkEntity chunk = chunk(7L, "hashA", 1, "ptr:chunk-A");
+ chunk.setPayloadHash("hashA");
+ stubConflictMerge(chunk, 0);
+ when(taskChunkMapper.selectById(anyLong())).thenThrow(new IllegalStateException("db down"));
+
+ FileTaskEntity task = new FileTaskEntity();
+ task.setId(9004L);
+ Throwable ex = invokeMergeExpectingFailure(task, List.of(row("r1", "B0A0000001", "标题1")));
+
+ assertTrue(ex.getMessage().contains("读取失败"), "实际: " + ex.getMessage());
+ }
+
+ /** 反射调用会包一层 InvocationTargetException,取根因以便断言业务异常。 */
+ private Throwable invokeMergeExpectingFailure(FileTaskEntity task, List