From ef4ea8d24b441776b2d31b42caf65cc0b251bae3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=BB=84=E8=87=AA=E8=BE=BE?= <980324341@qq.com> Date: Wed, 2 Sep 2026 03:39:54 +0800 Subject: [PATCH] =?UTF-8?q?task-126:=20similarasin=20DTO=20=E7=BB=84?= =?UTF-8?q?=E8=A3=85=E4=BD=8D=E7=BD=AE=E5=A5=91=E7=BA=A6=E5=9B=BA=E5=8C=96?= =?UTF-8?q?=EF=BC=88=E5=AE=A1=E8=AE=A1=E7=BB=93=E8=AE=BA=EF=BC=9A/result?= =?UTF-8?q?=20=E6=97=A0=E5=93=8D=E5=BA=94=20DTO=E3=80=81=E5=8F=AF=E7=A7=BB?= =?UTF-8?q?=E5=87=BA=E7=BB=84=E8=A3=85=E5=B7=B2=E5=9C=A8=E4=BA=8B=E5=8A=A1?= =?UTF-8?q?=E5=A4=96=E3=80=81=E8=90=BD=E5=BA=93=E8=A1=8C=E7=BB=84=E8=A3=85?= =?UTF-8?q?=E9=A1=BA=E5=BA=8F=E4=B8=8D=E5=8F=98=EF=BC=89+=208=20=E6=9D=A1?= =?UTF-8?q?=E5=AE=88=E9=97=A8=E6=B5=8B=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...SimilarAsinTaskServiceDtoAssemblyTest.java | 299 ++++++++++++++++++ 1 file changed, 299 insertions(+) create mode 100644 backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceDtoAssemblyTest.java diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceDtoAssemblyTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceDtoAssemblyTest.java new file mode 100644 index 00000000..44d62bd9 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceDtoAssemblyTest.java @@ -0,0 +1,299 @@ +package com.nanri.aiimage.modules.similarasin.service; + +import cn.hutool.crypto.digest.DigestUtil; +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.InstanceMetadata; +import com.nanri.aiimage.config.SimilarAsinProperties; +import com.nanri.aiimage.config.StorageProperties; +import com.nanri.aiimage.modules.file.service.LocalFileStorageService; +import com.nanri.aiimage.modules.file.service.oss.OssStorageService; +import com.nanri.aiimage.modules.similarasin.mapper.SimilarAsinFilterConditionMapper; +import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinSubmitResultRequest; +import com.nanri.aiimage.modules.similarasin.util.SimilarAsinImageEmbedder; +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.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.TaskProgressSnapshotService; +import com.nanri.aiimage.modules.task.service.TransientPayloadStorageService; +import org.apache.ibatis.builder.MapperBuilderAssistant; +import org.junit.jupiter.api.AfterEach; +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.transaction.PlatformTransactionManager; +import org.springframework.transaction.TransactionDefinition; +import org.springframework.transaction.TransactionStatus; + +import java.lang.reflect.Method; +import java.lang.reflect.Modifier; +import java.time.Duration; +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertEquals; +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.inOrder; +import static org.mockito.Mockito.lenient; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * task-126:similarasin /result DTO 组装位置契约固化。 + * 审计结论:/result 响应为 ApiResponse.success(null)(无响应 DTO);可移出的 + * DTO 组装(PreparedSubmittedChunk/SubmitContext 前置数据)已全部在 prepare + * 阶段、事务开始前完成(task-125);写事务内仅剩落库行实体构造(必须与 + * insert/upsert 同序,属 spec §2 禁拆顺序的一部分)。 + * 本测试固化上述契约:组装在事务前/提交后、回滚无事务后组装、落库行快照不变。 + */ +@ExtendWith(MockitoExtension.class) +class SimilarAsinTaskServiceDtoAssemblyTest { + + private static final Long TASK_ID = 21879L; + private static final String PARSED_POINTER = "rustfs:task-parsed/similar-asin/21879/payload.json"; + private static final String CHUNK_POINTER = "rustfs:task-chunk/similar-asin/21879/chunk.json"; + private static final String STORED_CHUNK_POINTER = "\"" + CHUNK_POINTER + "\""; + + @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; + @Mock private SimilarAsinFilterConditionMapper filterConditionMapper; + @Spy private ObjectMapper objectMapper = new ObjectMapper(); + @Mock private SimilarAsinTaskCacheService taskCacheService; + @Mock private SimilarAsinProperties properties; + @Mock private TaskFileJobService taskFileJobService; + @Mock private TaskDistributedLockService taskDistributedLockService; + @Mock private TaskProgressSnapshotService taskProgressSnapshotService; + @Mock private TransientPayloadStorageService transientPayloadStorageService; + @Mock private PlatformTransactionManager transactionManager; + @Mock private DistributedJobLockService distributedJobLockService; + @Mock private InstanceMetadata instanceMetadata; + @Mock private SimilarAsinImageEmbedder imageEmbedder; + @Mock private SimilarAsinImagePrefetchService imagePrefetchService; + @Mock private TransactionStatus transactionStatus; + + @InjectMocks private SimilarAsinTaskService service; + + private final AtomicBoolean transactionActive = new AtomicBoolean(); + private final AtomicReference insertedChunk = new AtomicReference<>(); + private final AtomicReference storedPayloadJson = new AtomicReference<>(); + + @BeforeAll + static void initializeMybatisMetadata() { + MapperBuilderAssistant assistant = new MapperBuilderAssistant(new MybatisConfiguration(), ""); + TableInfoHelper.initTableInfo(assistant, FileTaskEntity.class); + TableInfoHelper.initTableInfo(assistant, TaskChunkEntity.class); + TableInfoHelper.initTableInfo(assistant, TaskScopeStateEntity.class); + } + + @BeforeEach + void setUpTransactionAndLock() { + lenient().when(instanceMetadata.getInstanceId()).thenReturn("instance-a"); + lenient().when(taskDistributedLockService.acquire( + eq(SimilarAsinTaskService.MODULE_TYPE), anyLong(), any(Duration.class), eq(10_000L))) + .thenReturn(mock(TaskDistributedLockService.LockHandle.class)); + lenient().when(transactionManager.getTransaction(any(TransactionDefinition.class))) + .thenAnswer(invocation -> { + transactionActive.set(true); + return transactionStatus; + }); + lenient().doAnswer(invocation -> { + transactionActive.set(false); + return null; + }).when(transactionManager).commit(transactionStatus); + lenient().doAnswer(invocation -> { + transactionActive.set(false); + return null; + }).when(transactionManager).rollback(transactionStatus); + lenient().when(taskChunkMapper.selectList(any())).thenReturn(List.of()); + lenient().when(fileResultMapper.selectList(any())).thenReturn(List.of()); + lenient().when(taskChunkMapper.selectCount(any())).thenReturn(1L); + lenient().when(taskChunkMapper.selectOne(any())).thenReturn(null); + lenient().when(transientPayloadStorageService.storeChunkPayloadVersioned( + eq(SimilarAsinTaskService.MODULE_TYPE), eq(TASK_ID), anyString(), any(), anyString())) + .thenAnswer(invocation -> { + storedPayloadJson.set(invocation.getArgument(4)); + return STORED_CHUNK_POINTER; + }); + lenient().when(transientPayloadStorageService.wasLastStoreLocalFallback()).thenReturn(false); + lenient().doAnswer(invocation -> { + TaskChunkEntity chunk = invocation.getArgument(0); + chunk.setId(301L); + insertedChunk.set(chunk); + return 1; + }).when(taskChunkMapper).insert(any(TaskChunkEntity.class)); + lenient().doAnswer(invocation -> { + TaskScopeStateEntity scope = invocation.getArgument(0); + scope.setId(401L); + return 1; + }).when(taskScopeStateMapper).insert(any(TaskScopeStateEntity.class)); + } + + @AfterEach + void shutdownExecutors() { + service.shutdownAssembleExecutor(); + } + + @Test + void dtoAssemblyHappensBeforeTransactionStarts() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + + service.submitResult(TASK_ID, request(false)); + + // prepare 组装(含 payload 存储)→ 事务开始 → 落库 + var order = inOrder(transientPayloadStorageService, transactionManager, taskChunkMapper); + order.verify(transientPayloadStorageService).storeChunkPayloadVersioned( + eq(SimilarAsinTaskService.MODULE_TYPE), eq(TASK_ID), anyString(), any(), anyString()); + order.verify(transactionManager).getTransaction(any(TransactionDefinition.class)); + order.verify(taskChunkMapper).insert(any(TaskChunkEntity.class)); + } + + @Test + void assembledChunkFieldsMatchSnapshot() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + + service.submitResult(TASK_ID, request(false)); + + TaskChunkEntity chunk = insertedChunk.get(); + assertEquals(TASK_ID, chunk.getTaskId()); + assertEquals(SimilarAsinTaskService.MODULE_TYPE, chunk.getModuleType()); + assertEquals("similar-asin-21879", chunk.getScopeKey()); + assertEquals(0, chunk.getChunkIndex()); + assertEquals(1, chunk.getChunkTotal()); + assertEquals(STORED_CHUNK_POINTER, chunk.getPayloadJson()); + assertEquals(DigestUtil.sha256Hex(storedPayloadJson.get()), chunk.getPayloadHash()); + assertNotNull(chunk.getCreatedAt()); + assertNotNull(chunk.getUpdatedAt()); + } + + @Test + void postTransactionSideEffectsRunAfterCommit() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + + service.submitResult(TASK_ID, request(false)); + + var order = inOrder(transactionManager, taskCacheService); + order.verify(transactionManager).commit(transactionStatus); + order.verify(taskCacheService).touchTaskHeartbeat(TASK_ID); + } + + @Test + void rollbackSkipsPostTransactionAssemblyAndSideEffects() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + org.mockito.Mockito.doThrow(new IllegalStateException("scope upsert failed")) + .when(taskScopeStateMapper).insert(any(TaskScopeStateEntity.class)); + + assertThrows(IllegalStateException.class, () -> service.submitResult(TASK_ID, request(false))); + + verify(transactionManager).rollback(transactionStatus); + verify(transactionManager, never()).commit(transactionStatus); + verify(taskCacheService, never()).touchTaskHeartbeat(TASK_ID); + } + + @Test + void assemblyFailureBeforeTransactionLeavesNoPartialWrites() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + when(transientPayloadStorageService.storeChunkPayloadVersioned( + eq(SimilarAsinTaskService.MODULE_TYPE), eq(TASK_ID), anyString(), any(), anyString())) + .thenThrow(new IllegalStateException("payload store failed")); + + assertThrows(IllegalStateException.class, () -> service.submitResult(TASK_ID, request(false))); + + verify(taskChunkMapper, never()).insert(any(TaskChunkEntity.class)); + verify(taskScopeStateMapper, never()).insert(any(TaskScopeStateEntity.class)); + verify(transactionManager, never()).commit(transactionStatus); + } + + @Test + void resultEndpointReturnsNoResponseDto() throws Exception { + Method submitResult = SimilarAsinTaskService.class.getMethod("submitResult", Long.class, SimilarAsinSubmitResultRequest.class); + assertEquals(Void.TYPE, submitResult.getReturnType(), "/result 服务方法必须返回 void(无响应 DTO 组装)"); + assertTrue(Modifier.isPublic(submitResult.getModifiers())); + } + + @Test + void chunkAssemblyAddsNoExtraDatabaseReads() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + + service.submitResult(TASK_ID, request(false)); + + // chunk 落库行组装不引入额外查询:insert 前最后一次 mapper 交互为 + // completeSubmittedChunk 的任务状态重读(与 finalize 判定同序,属现状)。 + var order = inOrder(taskChunkMapper, fileTaskMapper, taskChunkMapper); + order.verify(taskChunkMapper).selectOne(any()); // prepare 查重 + order.verify(taskChunkMapper).selectOne(any()); // persist 查重 + order.verify(taskChunkMapper).insert(any(TaskChunkEntity.class)); + } + + @Test + void chunkSnapshotIsStableAcrossSubmissions() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + + service.submitResult(TASK_ID, request(false)); + TaskChunkEntity first = insertedChunk.get(); + service.submitResult(TASK_ID, request(false)); + TaskChunkEntity second = insertedChunk.get(); + + assertEquals(first.getScopeKey(), second.getScopeKey()); + assertEquals(first.getScopeHash(), second.getScopeHash()); + assertEquals(first.getPayloadHash(), second.getPayloadHash()); + assertEquals(first.getPayloadJson(), second.getPayloadJson()); + } + + private SimilarAsinSubmitResultRequest request(boolean done) { + SimilarAsinSubmitResultRequest request = new SimilarAsinSubmitResultRequest(); + request.setSubmissionId("similar-asin-21879"); + request.setChunkIndex(0); + request.setChunkTotal(1); + request.setDone(done); + return request; + } + + private FileTaskEntity runningTask(String owner) { + FileTaskEntity task = new FileTaskEntity(); + task.setId(TASK_ID); + task.setModuleType(SimilarAsinTaskService.MODULE_TYPE); + task.setStatus("RUNNING"); + task.setUserId(7L); + String resultJson = "{\"parsedPayloadRef\":\"" + PARSED_POINTER + "\""; + if (owner != null) { + resultJson += ",\"ownerInstanceId\":\"" + owner + "\""; + } + task.setResultJson(resultJson + "}"); + return task; + } +}