task-131: collectdata 事务收缩契约固化(/result 锁内单条写无长事务 + 响应 DTO 写后组装)+ 8 条守门测试

This commit is contained in:
2026-09-02 04:37:14 +08:00
parent 29ddf5c780
commit c60ad48be5
@@ -0,0 +1,263 @@
package com.nanri.aiimage.modules.collectdata.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.modules.collectdata.mapper.CollectDataCountryPrefMapper;
import com.nanri.aiimage.modules.collectdata.mapper.CollectDataItemMapper;
import com.nanri.aiimage.modules.invalidasin.mapper.InvalidAsinDataMapper;
import com.nanri.aiimage.modules.collectdata.model.dto.CollectDataSubmitResultRequest;
import com.nanri.aiimage.modules.collectdata.model.dto.CollectDataSubmitRowDto;
import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataResultRowVo;
import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataSubmitResultVo;
import com.nanri.aiimage.modules.collectdata.util.CollectDataBatchQuery;
import com.nanri.aiimage.modules.collectdata.util.CollectDataBrandBatchFilter;
import com.nanri.aiimage.modules.collectdata.util.CollectDataResultDetailCodec;
import com.nanri.aiimage.modules.collectdata.util.CollectDataInvalidAsinBatchWriter;
import com.nanri.aiimage.modules.collectdata.util.CollectDataResultItemBatchWriter;
import com.nanri.aiimage.modules.file.service.LocalFileStorageService;
import com.nanri.aiimage.modules.file.service.oss.OssStorageService;
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.TaskResultItemMapper;
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.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.TransientPayloadStorageService;
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.ArgumentCaptor;
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.Map;
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.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.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
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-131collectdata 事务收缩契约固化。
* 审计结论:/result 提交路径(submitResult → normalize/filter/存储/落库)无
* @Transactional、无事务模板调用——分布式锁内单条原子写与批量 upsert,事务边界
* 已窄;响应 DTO 组装(buildSubmitVo)在写后、锁内纯内存构造。本测试固化契约。
*/
@ExtendWith(MockitoExtension.class)
class CollectDataServiceTxBoundaryTest {
private static final Long TASK_ID = 6161L;
private static final Long USER_ID = 7L;
@Mock private LocalFileStorageService localFileStorageService;
@Mock private FileTaskMapper fileTaskMapper;
@Mock private FileResultMapper fileResultMapper;
@Mock private CollectDataItemMapper collectDataItemMapper;
@Mock private CollectDataCountryPrefMapper collectDataCountryPrefMapper;
@Mock private InvalidAsinDataMapper invalidAsinDataMapper;
@Mock private TaskChunkMapper taskChunkMapper;
@Mock private TaskScopeStateMapper taskScopeStateMapper;
@Mock private TaskResultItemMapper taskResultItemMapper;
@Mock private TaskDistributedLockService taskDistributedLockService;
@Mock private TaskFileJobService taskFileJobService;
@Mock private TransientPayloadStorageService transientPayloadStorageService;
@Mock private CollectDataExcelAssemblyService excelAssemblyService;
@Mock private OssStorageService ossStorageService;
@org.mockito.Spy private ObjectMapper objectMapper = new ObjectMapper();
@Mock private TransactionTemplate transactionTemplate;
@Mock private CollectDataBatchQuery collectDataBatchQuery;
@Mock private CollectDataBrandBatchFilter brandBatchFilter;
@Mock private CollectDataInvalidAsinBatchWriter invalidAsinBatchWriter;
@Mock private CollectDataResultItemBatchWriter resultItemBatchWriter;
@Mock private CollectDataResultDetailCodec resultDetailCodec;
@InjectMocks private CollectDataService service;
private final AtomicReference<TaskChunkEntity> insertedChunk = new AtomicReference<>();
@BeforeAll
static void initializeMybatisMetadata() {
MapperBuilderAssistant assistant = new MapperBuilderAssistant(new MybatisConfiguration(), "");
TableInfoHelper.initTableInfo(assistant, FileTaskEntity.class);
TableInfoHelper.initTableInfo(assistant, FileResultEntity.class);
TableInfoHelper.initTableInfo(assistant, TaskChunkEntity.class);
TableInfoHelper.initTableInfo(assistant, TaskScopeStateEntity.class);
}
@BeforeEach
void setUp() throws Exception {
lenient().when(taskDistributedLockService.acquire(eq("COLLECT_DATA"), anyLong(), anyLong()))
.thenReturn(mock(TaskDistributedLockService.LockHandle.class));
lenient().when(fileTaskMapper.selectById(TASK_ID)).thenReturn(runningTask());
lenient().when(fileResultMapper.selectOne(any())).thenReturn(null);
lenient().doAnswer(invocation -> {
FileResultEntity result = invocation.getArgument(0);
result.setId(5001L);
return 1;
}).when(fileResultMapper).insert(any(FileResultEntity.class));
lenient().when(taskChunkMapper.selectCount(any())).thenReturn(0L);
lenient().when(taskChunkMapper.selectOne(any())).thenReturn(null);
lenient().when(taskScopeStateMapper.selectOne(any())).thenReturn(null);
lenient().when(taskScopeStateMapper.insert(any(TaskScopeStateEntity.class))).thenReturn(1);
lenient().when(fileTaskMapper.updateById(any(FileTaskEntity.class))).thenReturn(1);
lenient().when(transientPayloadStorageService.isSharedWriteEnabled()).thenReturn(true);
lenient().when(transientPayloadStorageService.extractPointer(anyString())).thenReturn("rustfs:detail");
lenient().when(transientPayloadStorageService.storeResultPayload(
eq("COLLECT_DATA"), eq(TASK_ID), anyString(), anyString(), anyString()))
.thenReturn("\"rustfs:detail\"");
lenient().when(transientPayloadStorageService.storeChunkPayloadVersioned(
eq("COLLECT_DATA"), eq(TASK_ID), anyString(), anyInt(), anyString()))
.thenReturn("\"rustfs:chunk\"");
lenient().doAnswer(invocation -> {
TaskChunkEntity chunk = invocation.getArgument(0);
chunk.setId(701L);
insertedChunk.set(chunk);
return 1;
}).when(taskChunkMapper).insert(any(TaskChunkEntity.class));
CollectDataResultRowVo row = new CollectDataResultRowVo();
row.setAsin("B0COLLECT1");
lenient().when(collectDataBatchQuery.filter(any())).thenReturn(
new CollectDataBatchQuery.FilterResult(List.of(row), 0, 0));
lenient().when(brandBatchFilter.filter(any())).thenReturn(
new CollectDataBrandBatchFilter.BrandBatchOutcome(List.of(), List.of(), List.of(row)));
lenient().when(resultItemBatchWriter.upsertAccepted(anyLong(), anyLong(), anyString(), anyInt(), any(), anyString()))
.thenReturn(new CollectDataResultItemBatchWriter.UpsertCounts(1, 0, 1));
lenient().when(resultDetailCodec.encodeChunk(any())).thenReturn("{\"chunk\":1}");
}
@Test
void submitResultCarriesNoTransactionAnnotation() throws Exception {
Method submit = CollectDataService.class.getMethod(
"submitResult", Long.class, CollectDataSubmitResultRequest.class);
assertNull(submit.getAnnotation(Transactional.class),
"submitResult 无 @Transactional(锁内单条写,无长事务)");
}
@Test
void submitResultUsesNoTransactionTemplate() {
CollectDataSubmitResultVo vo = service.submitResult(TASK_ID, submitRequest());
assertNotNull(vo);
verify(transactionTemplate, never()).execute(any());
verify(transactionTemplate, never()).executeWithoutResult(any());
}
@Test
void duplicateChunkIsIgnoredWithoutRewriting() {
TaskChunkEntity existing = new TaskChunkEntity();
existing.setId(702L);
existing.setChunkIndex(1);
lenient().when(taskChunkMapper.selectOne(any())).thenReturn(existing);
CollectDataSubmitResultVo vo = service.submitResult(TASK_ID, submitRequest());
assertNotNull(vo);
verify(taskChunkMapper, never()).insert(any(TaskChunkEntity.class));
verify(transientPayloadStorageService, never()).storeChunkPayloadVersioned(
anyString(), anyLong(), anyString(), anyInt(), anyString());
}
@Test
void chunkSnapshotPersistedAfterFiltering() {
service.submitResult(TASK_ID, submitRequest());
TaskChunkEntity chunk = insertedChunk.get();
assertEquals(TASK_ID, chunk.getTaskId());
assertEquals("COLLECT_DATA", chunk.getModuleType());
assertEquals(1, chunk.getChunkIndex());
assertEquals(1, chunk.getChunkTotal());
assertEquals("\"rustfs:chunk\"", chunk.getPayloadJson());
assertNotNull(chunk.getPayloadHash(), "payload 哈希必须预计算");
}
@Test
void brandRejectedRowsWrittenToInvalidWriter() {
CollectDataResultRowVo bad = new CollectDataResultRowVo();
bad.setAsin("B0INVALID");
when(brandBatchFilter.filter(any())).thenReturn(
new CollectDataBrandBatchFilter.BrandBatchOutcome(List.of(bad), List.of(), List.of()));
service.submitResult(TASK_ID, submitRequest());
verify(invalidAsinBatchWriter).writeBatch(any());
}
@Test
void rustfsFailureAbortsWithoutPartialChunk() {
when(transientPayloadStorageService.extractPointer(anyString())).thenReturn("local:/tmp/x");
assertThrows(com.nanri.aiimage.common.exception.BusinessException.class,
() -> service.submitResult(TASK_ID, submitRequest()));
verify(taskChunkMapper, never()).insert(any(TaskChunkEntity.class));
}
@Test
void responseVoBuiltFromPostWriteState() {
CollectDataSubmitResultVo vo = service.submitResult(TASK_ID, submitRequest());
assertEquals(TASK_ID, vo.getTaskId());
assertEquals(1, vo.getChunkIndex());
assertEquals(1, vo.getChunkTotal());
assertNotNull(vo.getTaskStatus());
assertEquals(1, vo.getFinalRowCount());
}
@Test
void fullSubmitWritesChunkScopeStatsAndResult() {
service.submitResult(TASK_ID, submitRequest());
verify(taskChunkMapper).insert(any(TaskChunkEntity.class));
verify(taskScopeStateMapper).insert(any(TaskScopeStateEntity.class));
verify(fileResultMapper).insert(any(FileResultEntity.class));
verify(fileTaskMapper).updateById(any(FileTaskEntity.class));
verify(transactionTemplate, never()).execute(any());
}
private CollectDataSubmitResultRequest submitRequest() {
CollectDataSubmitRowDto row = new CollectDataSubmitRowDto();
row.setAsin("B0COLLECT1");
CollectDataSubmitResultRequest request = new CollectDataSubmitResultRequest();
request.setChunkIndex(1);
request.setChunkTotal(1);
request.setDone(false);
request.setItems(List.of(row));
return request;
}
private FileTaskEntity runningTask() {
FileTaskEntity task = new FileTaskEntity();
task.setId(TASK_ID);
task.setModuleType("COLLECT_DATA");
task.setStatus("RUNNING");
task.setUserId(USER_ID);
task.setResultJson("{}");
return task;
}
}