task-130: publish 事务收缩契约固化(prepare 外置/短事务落库/提交后清理已符合 spec 07)+ 8 条守门测试
This commit is contained in:
+236
@@ -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<org.springframework.transaction.TransactionStatus> 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;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user