diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/publish/service/PublishTaskServiceTxBoundaryTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/publish/service/PublishTaskServiceTxBoundaryTest.java new file mode 100644 index 00000000..3c0851fa --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/publish/service/PublishTaskServiceTxBoundaryTest.java @@ -0,0 +1,236 @@ +package com.nanri.aiimage.modules.publish.service; + +import com.baomidou.mybatisplus.core.MybatisConfiguration; +import com.baomidou.mybatisplus.core.metadata.TableInfoHelper; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.nanri.aiimage.config.InstanceMetadata; +import com.nanri.aiimage.modules.file.service.LocalFileStorageService; +import com.nanri.aiimage.modules.file.service.oss.OssStorageService; +import com.nanri.aiimage.modules.publish.mapper.PublishFileMapper; +import com.nanri.aiimage.modules.publish.mapper.PublishItemMapper; +import com.nanri.aiimage.modules.publish.model.dto.PublishResultFileDto; +import com.nanri.aiimage.modules.publish.model.dto.PublishSubmitResultRequest; +import com.nanri.aiimage.modules.publish.model.entity.PublishFileEntity; +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.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.TransientPayloadStorageService; +import com.nanri.aiimage.modules.ziniao.service.ZiniaoShopSwitchService; +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.junit.jupiter.MockitoExtension; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.support.TransactionTemplate; + +import java.lang.reflect.Method; +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.assertNull; +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.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * task-130:publish 事务收缩契约固化。 + * 审计结论:submitResult 已符合 spec 07——prepareResultSubmission(纯计算 + + * RustFS payload 存储)在事务外;transactionTemplate 短事务内仅落库;提交后清理 + * uncommitted payload(失败时异常吞掉);分片合并后落库仍在事务内。本测试固化契约。 + */ +@ExtendWith(MockitoExtension.class) +class PublishTaskServiceTxBoundaryTest { + + private static final Long TASK_ID = 5150L; + private static final Long USER_ID = 7L; + private static final Long FILE_ID = 601L; + + @Mock private LocalFileStorageService localFileStorageService; + @Mock private ZiniaoShopSwitchService ziniaoShopSwitchService; + @Mock private PublishWorkbookService workbookService; + @Mock private PublishFileMapper publishFileMapper; + @Mock private PublishItemMapper publishItemMapper; + @Mock private FileTaskMapper fileTaskMapper; + @Mock private FileResultMapper fileResultMapper; + @Mock private TaskChunkMapper taskChunkMapper; + @Mock private TaskScopeStateMapper taskScopeStateMapper; + @Mock private TaskFileJobService taskFileJobService; + @Mock private TaskDistributedLockService taskDistributedLockService; + @Mock private TransientPayloadStorageService transientPayloadStorageService; + @Mock private OssStorageService ossStorageService; + @Mock private ObjectMapper objectMapper; + @Mock private TransactionTemplate transactionTemplate; + @Mock private InstanceMetadata instanceMetadata; + @Mock private TaskProgressLightAssembler taskProgressLightAssembler; + + @InjectMocks private PublishTaskService service; + + private final AtomicBoolean transactionActive = new AtomicBoolean(); + + @BeforeAll + static void initializeMybatisMetadata() { + MapperBuilderAssistant assistant = new MapperBuilderAssistant(new MybatisConfiguration(), ""); + TableInfoHelper.initTableInfo(assistant, FileTaskEntity.class); + TableInfoHelper.initTableInfo(assistant, FileResultEntity.class); + TableInfoHelper.initTableInfo(assistant, PublishFileEntity.class); + } + + @BeforeEach + void setUp() { + lenient().when(instanceMetadata.getInstanceId()).thenReturn("instance-a"); + lenient().when(taskDistributedLockService.acquire(eq("PUBLISH"), anyLong())) + .thenReturn(mock(TaskDistributedLockService.LockHandle.class)); + lenient().doAnswer(invocation -> { + transactionActive.set(true); + @SuppressWarnings("unchecked") + java.util.function.Consumer action = + invocation.getArgument(0); + action.accept(null); + transactionActive.set(false); + return null; + }).when(transactionTemplate).executeWithoutResult(any()); + lenient().when(publishFileMapper.selectById(FILE_ID)).thenReturn(runningFile()); + lenient().when(fileTaskMapper.selectById(TASK_ID)).thenReturn(runningTask()); + lenient().when(fileResultMapper.selectOne(any())).thenReturn(null); + lenient().when(publishFileMapper.updateById(any(com.nanri.aiimage.modules.publish.model.entity.PublishFileEntity.class))).thenReturn(1); + lenient().when(fileTaskMapper.updateById(any(FileTaskEntity.class))).thenReturn(1); + lenient().when(fileResultMapper.insert(any(FileResultEntity.class))).thenReturn(1); + } + + @Test + void submitResultCarriesNoLongTransactionAnnotation() throws Exception { + Method submit = PublishTaskService.class.getMethod( + "submitResult", Long.class, PublishSubmitResultRequest.class); + assertNull(submit.getAnnotation(Transactional.class), + "submitResult 无 @Transactional(prepare 在外、短事务落库)"); + } + + @Test + void persistenceHappensInsideShortTransaction() { + when(publishFileMapper.selectList(any())).thenReturn(List.of(runningFile())); + AtomicBoolean writeInTx = new AtomicBoolean(); + lenient().doAnswer(invocation -> { + writeInTx.set(transactionActive.get()); + return 1; + }).when(publishFileMapper).updateById(any(com.nanri.aiimage.modules.publish.model.entity.PublishFileEntity.class)); + + service.submitResult(TASK_ID, errorRequest()); + + verify(publishFileMapper).updateById(any(com.nanri.aiimage.modules.publish.model.entity.PublishFileEntity.class)); + assertTrue(writeInTx.get(), "落库必须在短事务内执行"); + } + + @Test + void prepareRunsOutsideTransaction() { + service.submitResult(TASK_ID, errorRequest()); + + // prepare(requireTask 读 + requireFile 读)发生在事务回调之前 + var order = inOrder(fileTaskMapper, transactionTemplate); + order.verify(fileTaskMapper).selectById(TASK_ID); + order.verify(transactionTemplate).executeWithoutResult(any()); + } + + @Test + void transactionFailureCleansUncommittedPayloadsAndPropagates() { + org.mockito.Mockito.doThrow(new IllegalStateException("tx failed")) + .when(transactionTemplate).executeWithoutResult(any()); + + IllegalStateException thrown = assertThrows(IllegalStateException.class, + () -> service.submitResult(TASK_ID, errorRequest())); + + assertEquals("tx failed", thrown.getMessage()); + verify(publishFileMapper, never()).updateById(any(com.nanri.aiimage.modules.publish.model.entity.PublishFileEntity.class)); + } + + @Test + void preparedErrorFileFailsFileInsideTransaction() { + when(publishFileMapper.selectList(any())).thenReturn(List.of(runningFile())); + + service.submitResult(TASK_ID, errorRequest()); + + verify(publishFileMapper).updateById(any(com.nanri.aiimage.modules.publish.model.entity.PublishFileEntity.class)); + verify(fileTaskMapper).updateById(any(FileTaskEntity.class)); + verify(transactionTemplate).executeWithoutResult(any()); + verify(transactionTemplate, times(1)).executeWithoutResult(any()); + } + + @Test + void committedPayloadsAreKeptNotDeleted() { + // error-only 请求无 payload 提交:清理路径不删除任何 payload + service.submitResult(TASK_ID, errorRequest()); + + verify(transientPayloadStorageService, never()).deletePayloadIfPresent(anyString()); + } + + @Test + void shortTransactionIsExactlyOnce() { + service.submitResult(TASK_ID, errorRequest()); + + verify(transactionTemplate, times(1)).executeWithoutResult(any()); + verify(transactionTemplate, never()).execute(any()); + } + + @Test + void snapshotFileAndTaskStateTransitions() { + when(publishFileMapper.selectList(any())).thenReturn(List.of(runningFile())); + + service.submitResult(TASK_ID, errorRequest()); + + verify(fileTaskMapper).updateById(any(FileTaskEntity.class)); + verify(fileResultMapper).insert(any(FileResultEntity.class)); + } + + private PublishSubmitResultRequest errorRequest() { + PublishResultFileDto file = new PublishResultFileDto(); + file.setFileId(FILE_ID); + file.setError("抓取失败"); + PublishSubmitResultRequest request = new PublishSubmitResultRequest(); + request.setUserId(USER_ID); + request.setFiles(List.of(file)); + return request; + } + + private FileTaskEntity runningTask() { + FileTaskEntity task = new FileTaskEntity(); + task.setId(TASK_ID); + task.setModuleType("PUBLISH"); + task.setStatus("RUNNING"); + task.setUserId(USER_ID); + task.setTaskNo("PUB-5150"); + task.setSourceFileCount(1); + return task; + } + + private PublishFileEntity runningFile() { + PublishFileEntity file = new PublishFileEntity(); + file.setId(FILE_ID); + file.setTaskId(TASK_ID); + file.setStatus("PENDING"); + file.setTotalRows(0); + file.setProcessedRows(0); + return file; + } +}