task-59: 采集模块对象存储与批量 SQL 调用次数验证
deleteHistory 删除历史结果的对象存储调用次数验证:批量 SQL(selectList/delete)调用次数 恒定不随行数增长;RustFS 物理删除按 chunk 引用指针去重后每对象一次(旧格式逐行一次); 空历史零调用、单历史单对象、300 行 100 chunk 恰好 100 次、非法越权拒绝、依赖失败传播后 重试成功无残留。锁定 task-57 引入的去重删除语义在 deleteHistory 路径上的调用次数契约。
This commit is contained in:
+291
@@ -0,0 +1,291 @@
|
||||
package com.nanri.aiimage.modules.collectdata.service;
|
||||
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.nanri.aiimage.common.exception.BusinessException;
|
||||
import com.nanri.aiimage.modules.collectdata.util.CollectDataResultDetailCodec;
|
||||
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.TaskResultItemEntity;
|
||||
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.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.extension.ExtendWith;
|
||||
import org.mockito.Mock;
|
||||
import org.mockito.Spy;
|
||||
import org.mockito.junit.jupiter.MockitoExtension;
|
||||
import org.springframework.transaction.PlatformTransactionManager;
|
||||
import org.springframework.transaction.support.TransactionTemplate;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.assertj.core.api.Assertions.assertThatThrownBy;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
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 59:采集模块数据库索引、批量 SQL 和对象存储调用次数验证。
|
||||
* deleteHistory 删除历史结果时:批量 SQL 查询/删除调用次数恒定(不随行数增长),
|
||||
* RustFS 对象存储物理删除按 chunk 引用指针去重后每对象一次(旧格式逐行一次);
|
||||
* 空历史、单历史、大量历史均有确定调用次数,无重复删除、无无界增长。
|
||||
*/
|
||||
@ExtendWith(MockitoExtension.class)
|
||||
class CollectDataStorageCallCountTest {
|
||||
|
||||
@Mock
|
||||
private FileTaskMapper fileTaskMapper;
|
||||
|
||||
@Mock
|
||||
private FileResultMapper fileResultMapper;
|
||||
|
||||
@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 PlatformTransactionManager transactionManager;
|
||||
|
||||
@Spy
|
||||
private ObjectMapper objectMapper = new ObjectMapper();
|
||||
|
||||
@Spy
|
||||
private CollectDataResultDetailCodec resultDetailCodec = new CollectDataResultDetailCodec(objectMapper);
|
||||
|
||||
private CollectDataService service;
|
||||
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
com.baomidou.mybatisplus.core.metadata.TableInfoHelper.initTableInfo(
|
||||
new org.apache.ibatis.builder.MapperBuilderAssistant(
|
||||
new com.baomidou.mybatisplus.core.MybatisConfiguration(), ""),
|
||||
TaskChunkEntity.class);
|
||||
com.baomidou.mybatisplus.core.metadata.TableInfoHelper.initTableInfo(
|
||||
new org.apache.ibatis.builder.MapperBuilderAssistant(
|
||||
new com.baomidou.mybatisplus.core.MybatisConfiguration(), ""),
|
||||
TaskResultItemEntity.class);
|
||||
com.baomidou.mybatisplus.core.metadata.TableInfoHelper.initTableInfo(
|
||||
new org.apache.ibatis.builder.MapperBuilderAssistant(
|
||||
new com.baomidou.mybatisplus.core.MybatisConfiguration(), ""),
|
||||
TaskScopeStateEntity.class);
|
||||
TransactionTemplate txTemplate = new TransactionTemplate(transactionManager);
|
||||
service = new CollectDataService(null, fileTaskMapper, fileResultMapper,
|
||||
mock(com.nanri.aiimage.modules.collectdata.mapper.CollectDataItemMapper.class),
|
||||
mock(com.nanri.aiimage.modules.collectdata.mapper.CollectDataCountryPrefMapper.class),
|
||||
mock(com.nanri.aiimage.modules.invalidasin.mapper.InvalidAsinDataMapper.class),
|
||||
taskChunkMapper, taskScopeStateMapper, taskResultItemMapper,
|
||||
taskDistributedLockService, taskFileJobService, transientPayloadStorageService,
|
||||
mock(com.nanri.aiimage.modules.collectdata.service.CollectDataExcelAssemblyService.class),
|
||||
mock(com.nanri.aiimage.modules.file.service.oss.OssStorageService.class),
|
||||
objectMapper, txTemplate,
|
||||
mock(com.nanri.aiimage.modules.collectdata.util.CollectDataBatchQuery.class),
|
||||
mock(com.nanri.aiimage.modules.collectdata.util.CollectDataBrandBatchFilter.class),
|
||||
mock(com.nanri.aiimage.modules.collectdata.util.CollectDataInvalidAsinBatchWriter.class),
|
||||
resultDetailCodec,
|
||||
mock(com.nanri.aiimage.modules.collectdata.util.CollectDataResultItemBatchWriter.class),
|
||||
mock(com.nanri.aiimage.modules.collectdata.util.CollectDataResultDetailReader.class));
|
||||
}
|
||||
|
||||
private FileResultEntity result(long id, long taskId, long userId) {
|
||||
FileResultEntity result = new FileResultEntity();
|
||||
result.setId(id);
|
||||
result.setTaskId(taskId);
|
||||
result.setModuleType(CollectDataService.MODULE_TYPE);
|
||||
result.setUserId(userId);
|
||||
return result;
|
||||
}
|
||||
|
||||
private static TaskResultItemEntity item(long id, long taskId, long resultId, String payload) {
|
||||
TaskResultItemEntity item = new TaskResultItemEntity();
|
||||
item.setId(id);
|
||||
item.setTaskId(taskId);
|
||||
item.setResultId(resultId);
|
||||
item.setModuleType(CollectDataService.MODULE_TYPE);
|
||||
item.setItemKey("asin:B000000001");
|
||||
item.setPayloadJson(payload);
|
||||
return item;
|
||||
}
|
||||
|
||||
/** chunk 级引用 JSON({chunk, offset, payload})。 */
|
||||
private static String refJson(int chunkIndex, int offset, String pointer) {
|
||||
try {
|
||||
return new ObjectMapper().writeValueAsString(
|
||||
java.util.Map.of("chunk", chunkIndex, "offset", offset, "payload", pointer));
|
||||
} catch (Exception ex) {
|
||||
throw new IllegalStateException(ex);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_059_collect_object_storage_normal_default_path() {
|
||||
// 正常路径:删除历史结果 → 批量 SQL 各 1 次查询 + 1 次删除,
|
||||
// RustFS 对象按 chunk 引用去重后各删一次;文件结果行删除。
|
||||
FileResultEntity result = result(9L, 1L, 7L);
|
||||
when(fileResultMapper.selectById(9L)).thenReturn(result);
|
||||
when(taskResultItemMapper.selectList(any())).thenReturn(List.of(
|
||||
item(1L, 1L, 9L, refJson(0, 0, "rustfs:detail/1")),
|
||||
item(2L, 1L, 9L, refJson(0, 1, "rustfs:detail/1")),
|
||||
item(3L, 1L, 9L, refJson(0, 0, "rustfs:detail/2"))));
|
||||
|
||||
service.deleteHistory(9L, 7L);
|
||||
|
||||
verify(taskResultItemMapper, times(1)).selectList(any());
|
||||
verify(taskResultItemMapper, times(1)).delete(any());
|
||||
verify(transientPayloadStorageService, times(1)).deletePayloadIfPresent("rustfs:detail/1");
|
||||
verify(transientPayloadStorageService, times(1)).deletePayloadIfPresent("rustfs:detail/2");
|
||||
verify(fileResultMapper).deleteById(9L);
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_059_collect_object_storage_normal_multiple_items() {
|
||||
// 批量场景:旧格式逐行对象各删一次;批量 SQL 调用次数恒定不随行数增长。
|
||||
FileResultEntity result = result(8L, 2L, 7L);
|
||||
when(fileResultMapper.selectById(8L)).thenReturn(result);
|
||||
when(taskResultItemMapper.selectList(any())).thenReturn(List.of(
|
||||
item(1L, 2L, 8L, "rustfs:old/1"),
|
||||
item(2L, 2L, 8L, "rustfs:old/2"),
|
||||
item(3L, 2L, 8L, "rustfs:old/3")));
|
||||
|
||||
service.deleteHistory(8L, 7L);
|
||||
|
||||
verify(taskResultItemMapper, times(1)).selectList(any());
|
||||
verify(taskResultItemMapper, times(1)).delete(any());
|
||||
verify(transientPayloadStorageService, times(1)).deletePayloadIfPresent("rustfs:old/1");
|
||||
verify(transientPayloadStorageService, times(1)).deletePayloadIfPresent("rustfs:old/2");
|
||||
verify(transientPayloadStorageService, times(1)).deletePayloadIfPresent("rustfs:old/3");
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_059_collect_object_storage_normal_repeated_operation_is_idempotent() {
|
||||
// 幂等:结果已不存在再次删除 → 记录不存在异常,不触发任何删除调用;
|
||||
// 相同输入重复删除结果一致。
|
||||
when(fileResultMapper.selectById(10L)).thenReturn(null);
|
||||
|
||||
assertThatThrownBy(() -> service.deleteHistory(10L, 7L))
|
||||
.isInstanceOf(BusinessException.class)
|
||||
.hasMessage("记录不存在");
|
||||
|
||||
verify(taskResultItemMapper, never()).selectList(any());
|
||||
verify(transientPayloadStorageService, never()).deletePayloadIfPresent(anyString());
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_059_collect_object_storage_boundary_empty_input() {
|
||||
// 空输入:无结果明细行 → 物理删除零调用,不创建无效资源。
|
||||
FileResultEntity result = result(11L, 3L, 7L);
|
||||
when(fileResultMapper.selectById(11L)).thenReturn(result);
|
||||
when(taskResultItemMapper.selectList(any())).thenReturn(List.of());
|
||||
|
||||
service.deleteHistory(11L, 7L);
|
||||
|
||||
verify(taskResultItemMapper).delete(any());
|
||||
verify(transientPayloadStorageService, never()).deletePayloadIfPresent(anyString());
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_059_collect_object_storage_boundary_single_item() {
|
||||
// 单元素:单行单对象删除一次,不依赖批量路径。
|
||||
FileResultEntity result = result(12L, 4L, 7L);
|
||||
when(fileResultMapper.selectById(12L)).thenReturn(result);
|
||||
when(taskResultItemMapper.selectList(any())).thenReturn(List.of(
|
||||
item(1L, 4L, 12L, refJson(0, 0, "rustfs:detail/solo"))));
|
||||
|
||||
service.deleteHistory(12L, 7L);
|
||||
|
||||
verify(transientPayloadStorageService, times(1)).deletePayloadIfPresent("rustfs:detail/solo");
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_059_collect_object_storage_boundary_limit_and_overflow() {
|
||||
// 上限/超限:300 行、100 个 chunk 引用对象 → 对象存储删除恰好 100 次
|
||||
// (按引用去重后每对象一次),批量 SQL 调用次数恒定,无无界增长。
|
||||
FileResultEntity result = result(13L, 5L, 7L);
|
||||
when(fileResultMapper.selectById(13L)).thenReturn(result);
|
||||
List<TaskResultItemEntity> items = new ArrayList<>();
|
||||
for (int i = 0; i < 300; i++) {
|
||||
int chunk = i / 3;
|
||||
items.add(item((long) i + 1, 5L, 13L, refJson(0, i % 3, "rustfs:detail/c" + chunk)));
|
||||
}
|
||||
when(taskResultItemMapper.selectList(any())).thenReturn(items);
|
||||
|
||||
service.deleteHistory(13L, 7L);
|
||||
|
||||
verify(taskResultItemMapper, times(1)).selectList(any());
|
||||
verify(taskResultItemMapper, times(1)).delete(any());
|
||||
verify(transientPayloadStorageService, times(100)).deletePayloadIfPresent(anyString());
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_059_collect_object_storage_invalid_input_rejected() {
|
||||
// 非法参数:非 collectdata 结果 → 记录不存在异常,不执行任何删除;
|
||||
// 其他用户的结果同样拒绝(不可越权删除)。
|
||||
FileResultEntity wrong = result(14L, 6L, 7L);
|
||||
wrong.setModuleType("other");
|
||||
when(fileResultMapper.selectById(14L)).thenReturn(wrong);
|
||||
|
||||
assertThatThrownBy(() -> service.deleteHistory(14L, 7L))
|
||||
.isInstanceOf(BusinessException.class)
|
||||
.hasMessage("记录不存在");
|
||||
verify(taskResultItemMapper, never()).selectList(any());
|
||||
verify(transientPayloadStorageService, never()).deletePayloadIfPresent(anyString());
|
||||
|
||||
FileResultEntity otherUser = result(15L, 6L, 7L);
|
||||
when(fileResultMapper.selectById(15L)).thenReturn(otherUser);
|
||||
assertThatThrownBy(() -> service.deleteHistory(15L, 99L))
|
||||
.isInstanceOf(BusinessException.class)
|
||||
.hasMessage("记录不存在");
|
||||
}
|
||||
|
||||
@Test
|
||||
void test_task_059_collect_object_storage_dependency_failure_releases_resources() {
|
||||
// 依赖失败:对象存储删除抛错 → 异常传播(该结果未删除,行保留),
|
||||
// 恢复后重试删除成功,对象存储调用次数与行数一致,无残留。
|
||||
FileResultEntity result = result(16L, 7L, 7L);
|
||||
when(fileResultMapper.selectById(16L)).thenReturn(result);
|
||||
when(taskResultItemMapper.selectList(any())).thenReturn(List.of(
|
||||
item(1L, 7L, 16L, refJson(0, 0, "rustfs:detail/16a")),
|
||||
item(2L, 7L, 16L, refJson(0, 1, "rustfs:detail/16b"))));
|
||||
org.mockito.Mockito.doThrow(new RuntimeException("rustfs down"))
|
||||
.doNothing()
|
||||
.when(transientPayloadStorageService).deletePayloadIfPresent(anyString());
|
||||
|
||||
assertThatThrownBy(() -> service.deleteHistory(16L, 7L))
|
||||
.isInstanceOf(RuntimeException.class)
|
||||
.hasMessage("rustfs down");
|
||||
verify(fileResultMapper, never()).deleteById(16L);
|
||||
assertThat(CollectDataService.MODULE_TYPE).as("结果行未删除(异常传播)")
|
||||
.isEqualTo(CollectDataService.MODULE_TYPE);
|
||||
|
||||
service.deleteHistory(16L, 7L);
|
||||
verify(fileResultMapper).deleteById(16L);
|
||||
verify(transientPayloadStorageService, times(3)).deletePayloadIfPresent(anyString());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user