From 28f61eb70d230ffb9cd1668eab02521b7a4398c1 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:44:38 +0800 Subject: [PATCH] =?UTF-8?q?task-127:=20similarasin=20=E4=B8=B4=E6=97=B6?= =?UTF-8?q?=E6=96=87=E4=BB=B6=E6=B8=85=E7=90=86=E6=97=B6=E6=9C=BA=E5=A5=91?= =?UTF-8?q?=E7=BA=A6=E5=9B=BA=E5=8C=96=EF=BC=88=E6=8F=90=E4=BA=A4=E5=90=8E?= =?UTF-8?q?=E6=B8=85=E7=90=86/=E5=9B=9E=E6=BB=9A=E4=B8=8D=E6=B8=85?= =?UTF-8?q?=E7=90=86/=E5=BC=82=E5=B8=B8=E5=90=9E=E6=8E=89/=E5=8F=AF?= =?UTF-8?q?=E9=87=8D=E8=AF=95/=E5=B9=82=E7=AD=89=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 --- ...AsinTaskServiceCleanupAfterCommitTest.java | 295 ++++++++++++++++++ 1 file changed, 295 insertions(+) create mode 100644 backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceCleanupAfterCommitTest.java diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceCleanupAfterCommitTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceCleanupAfterCommitTest.java new file mode 100644 index 00000000..1a7f9514 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/similarasin/service/SimilarAsinTaskServiceCleanupAfterCommitTest.java @@ -0,0 +1,295 @@ +package com.nanri.aiimage.modules.similarasin.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.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.dao.DuplicateKeyException; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.TransactionDefinition; +import org.springframework.transaction.TransactionStatus; + +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.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.doThrow; +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.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * task-127:similarasin 临时文件清理时机契约固化。 + * 审计结论:提交后清理(cleanupPreparedSubmittedChunkIfUnreferenced 在事务 + * execute 返回后执行)已符合"挂 afterCommit"语义——回滚不清理、清理异常吞掉、 + * 失败可重试;本测试固化该契约。 + */ +@ExtendWith(MockitoExtension.class) +class SimilarAsinTaskServiceCleanupAfterCommitTest { + + 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 AtomicBoolean cleanupOutsideTx = new AtomicBoolean(true); + + @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(0L); + lenient().when(taskScopeStateMapper.selectCount(any())).thenReturn(0L); + lenient().when(taskChunkMapper.selectOne(any())).thenReturn(null); + lenient().when(transientPayloadStorageService.storeChunkPayloadVersioned( + eq(SimilarAsinTaskService.MODULE_TYPE), eq(TASK_ID), anyString(), any(), anyString())) + .thenReturn(STORED_CHUNK_POINTER); + lenient().when(transientPayloadStorageService.wasLastStoreLocalFallback()).thenReturn(false); + lenient().when(transientPayloadStorageService.extractPointer(STORED_CHUNK_POINTER)) + .thenReturn(CHUNK_POINTER); + lenient().doAnswer(invocation -> { + TaskScopeStateEntity scope = invocation.getArgument(0); + scope.setId(401L); + return 1; + }).when(taskScopeStateMapper).insert(any(TaskScopeStateEntity.class)); + lenient().doAnswer(invocation -> { + cleanupOutsideTx.set(!transactionActive.get()); + return null; + }).when(transientPayloadStorageService).deletePayloadIfPresent(anyString()); + } + + @AfterEach + void shutdownExecutors() { + service.shutdownAssembleExecutor(); + } + + /** 并发重复 chunk(insert 撞唯一键)→ payload 未落库 → 提交后清理。 */ + private void configureDuplicateChunkFlow() { + doThrow(new DuplicateKeyException("duplicate chunk")) + .when(taskChunkMapper).insert(any(TaskChunkEntity.class)); + } + + @Test + void cleanupRunsAfterCommitForUnpersistedPayload() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + configureDuplicateChunkFlow(); + + service.submitResult(TASK_ID, request(false)); + + verify(transientPayloadStorageService).deletePayloadIfPresent(STORED_CHUNK_POINTER); + var order = inOrder(transactionManager, transientPayloadStorageService); + order.verify(transactionManager).commit(transactionStatus); + order.verify(transientPayloadStorageService).deletePayloadIfPresent(STORED_CHUNK_POINTER); + } + + @Test + void rollbackSkipsCleanup() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + configureDuplicateChunkFlow(); + doThrow(new IllegalStateException("scope upsert failed")) + .when(taskScopeStateMapper).insert(any(TaskScopeStateEntity.class)); + + org.junit.jupiter.api.Assertions.assertThrows(IllegalStateException.class, + () -> service.submitResult(TASK_ID, request(false))); + + verify(transactionManager).rollback(transactionStatus); + verify(transientPayloadStorageService, never()).deletePayloadIfPresent(anyString()); + } + + @Test + void cleanupErrorIsSwallowedAndResponseSucceeds() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + configureDuplicateChunkFlow(); + doThrow(new IllegalStateException("cleanup failed")) + .when(transientPayloadStorageService).deletePayloadIfPresent(anyString()); + + service.submitResult(TASK_ID, request(false)); + + verify(taskCacheService).touchTaskHeartbeat(TASK_ID); + verify(transientPayloadStorageService).deletePayloadIfPresent(STORED_CHUNK_POINTER); + } + + @Test + void cleanupObservesTransactionAlreadyFinished() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + configureDuplicateChunkFlow(); + + service.submitResult(TASK_ID, request(false)); + + verify(transientPayloadStorageService).deletePayloadIfPresent(STORED_CHUNK_POINTER); + org.junit.jupiter.api.Assertions.assertTrue(cleanupOutsideTx.get(), + "清理必须在事务结束后执行(等价 afterCommit)"); + } + + @Test + void cleanupFailureKeepsPayloadRetryableOnNextSubmission() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + configureDuplicateChunkFlow(); + doThrow(new IllegalStateException("cleanup failed")) + .when(transientPayloadStorageService).deletePayloadIfPresent(anyString()); + + service.submitResult(TASK_ID, request(false)); + // 恢复后再次提交,同一 payload 可再次进入清理 + service.submitResult(TASK_ID, request(false)); + + verify(transientPayloadStorageService, times(2)).deletePayloadIfPresent(STORED_CHUNK_POINTER); + } + + @Test + void cleanupIsIdempotentAcrossRepeatedSubmissions() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + configureDuplicateChunkFlow(); + + service.submitResult(TASK_ID, request(false)); + service.submitResult(TASK_ID, request(false)); + + verify(transientPayloadStorageService, times(2)).deletePayloadIfPresent(STORED_CHUNK_POINTER); + } + + @Test + void cleanupChecksReferencesBeforeDelete() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + configureDuplicateChunkFlow(); + + service.submitResult(TASK_ID, request(false)); + + var order = inOrder(transientPayloadStorageService, taskChunkMapper, taskScopeStateMapper, + transientPayloadStorageService); + order.verify(transientPayloadStorageService).extractPointer(STORED_CHUNK_POINTER); + order.verify(taskChunkMapper).selectCount(any()); + order.verify(taskScopeStateMapper).selectCount(any()); + order.verify(transientPayloadStorageService).deletePayloadIfPresent(STORED_CHUNK_POINTER); + } + + @Test + void referencedPayloadIsKeptNotDeleted() { + FileTaskEntity task = runningTask("instance-a"); + when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task); + configureDuplicateChunkFlow(); + when(taskChunkMapper.selectCount(any())).thenReturn(1L); + + service.submitResult(TASK_ID, request(false)); + + verify(transientPayloadStorageService, never()).deletePayloadIfPresent(STORED_CHUNK_POINTER); + } + + 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; + } +}