task-127: similarasin 临时文件清理时机契约固化(提交后清理/回滚不清理/异常吞掉/可重试/幂等)+ 8 条守门测试

This commit is contained in:
2026-09-02 03:44:38 +08:00
parent ef4ea8d24b
commit 28f61eb70d
@@ -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-127similarasin 临时文件清理时机契约固化。
* 审计结论:提交后清理(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();
}
/** 并发重复 chunkinsert 撞唯一键)→ 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;
}
}