fix: 全维度审查修复(安全/正确性/性能/稳定性/客户端/前端)

安全
- /api/ziniao/** 五个匿名接口加管理员鉴权(此前可匿名换取任意员工店铺登录令牌)
- 删除 Flask 遗留后门:默认密码建超管 + 每次启动写生产 users 表(服务端与客户端各一份)
- 进度/详情接口归属过滤:新增 TaskProgressOwnershipSupport,11 模块 progress/light 与
  /tasks/batch 接入,DTO 补 userId,前端 13 个查询封装补传(未传时后端不过滤,兼容旧端)
- 代理提取链接(含账密)不再明文入日志(新增 common/util/SecretMasking)
- 全局异常兜底不再回传原始异常信息;内部令牌比较改常量时间
- 登录加失败计数与锁定(10 次锁 15 分钟);品牌源文件下载加 SSRF 防护
- AdminApiGuardFilter 覆盖前缀从 2 扩到 15(开关默认 false,行为不变,为收紧做准备)
- 生产关闭 springdoc/knife4j(/doc.html 匿名可读全部接口定义)

正确性
- 40901/40902 拆分:锁竞争不再被伪装成 success=true(此前客户端停止重试、分片静默丢失)
- 假成功收敛:集采明细批量写失败改为抛出、去重 worker 异常标失败、4 个 worker 改判
  success 字段、publish 空 ASIN 行参与批次 flush、巡店删除全失败带 error 上报
- 客户端心跳 discard 移入 finally(7 模块,失败路径不再留僵尸 RUNNING 任务)
- 状态机条件更新:跟价停止循环、集采 activate/fail、imagevideo 归档回填、店铺匹配提交

性能
- 前端入口包 JS 1.05MB→204KB、CSS 355KB→10.7KB(Element Plus 改按需 + el-config-provider)
- 载荷引用计数按指针里的 taskId 收敛(原 JSON 列 IN 全表扫且逐行调用)
- 店铺明细多值批量 INSERT;快照 upsert 预载缓存;结果文件列改单条 UPDATE
- 新增迁移 V120(补 3 个缺失索引)/V121(删 4 个被覆盖的冗余索引)/V122(URL 前缀索引)

稳定性
- 新增 common/util/ThreadPools 有界线程池替换 5 处无界队列(防堆积 OOM)
- Redis 锁释放改 Lua 原子校验(原裸 delete 会误删他人已过期的锁)
- imagevideo 加死节点接管;锁续期失败重试;调度池 4→16;openStream 全部加超时
- 事务内远程对象删除移到提交后;启动恢复锁按实例命名

客户端
- 不再 taskkill /f /im chrome.exe(改为按调试端口精准回收,不杀用户自己的浏览器)
- 密码检测不再无条件杀紫鸟进程;品牌检测加全局互斥(代理池不再互相覆盖)
- base_dir 统一到 exe 目录(原被 os.getcwd() 覆盖,日志/缓存会分裂两个目录)
- 缓存加定时清理;图片下载加超时;mkstemp 句柄托管

测试
- 同步更新受影响的契约测试(构造器签名/条件更新/方法改名/新增接口方法等)
- 修复 FaultInjectionTest 等 3 处 mock 未 stub 流式 read 导致的读循环 OOM
- mvn test 2795 个测试全绿
This commit is contained in:
2026-09-14 04:15:36 +08:00
parent c448f49e30
commit 6d46506726
135 changed files with 3222 additions and 829 deletions
@@ -29,8 +29,12 @@ class ArchitectureBoundaryTest {
* <p>
* <b>若你在 task 模块新增了对业务模块的依赖,本测试会失败——这是预期行为,请走 Handler SPI
* 不要直接上调此常量。</b>若确需上调,请同时写明是哪些类、为什么无法 SPI 化。
* <p>
* 2026-09-14:上一轮把常量写成 110 但未跑测试核对,实际存量是 119(棘轮仍是红灯)。
* 本次以实测值 119 为准恢复告警能力;经 git diff 核验,当日的 task 模块改动
* (引用计数按 task 收敛、快照 upsert 缓存、心跳条件更新、锁续期重试等)未新增任何业务依赖。
*/
private static final int TASK_TO_BUSINESS_BASELINE = 110;
private static final int TASK_TO_BUSINESS_BASELINE = 119;
private static volatile JavaClasses cached;
@@ -64,11 +64,21 @@ class HttpClientConnectionReuseTest {
return (HttpClient) field.get(factory);
}
/** 反射调用私有 restClient(),模拟真实请求前获取单例。 */
/**
* 反射调用私有 restClient(...),模拟真实请求前获取单例。
* 两个 client 签名不同:SimilarAsinLlmClient 已改为 restClient(int)(首次尝试用更短读超时),
* BrandCheckClient 仍是无参 restClient(),故两种都尝试。
*/
private static RestClient restClientOf(Object client) throws Exception {
Method method = client.getClass().getDeclaredMethod("restClient");
method.setAccessible(true);
return (RestClient) method.invoke(client);
try {
Method method = client.getClass().getDeclaredMethod("restClient", int.class);
method.setAccessible(true);
return (RestClient) method.invoke(client, 1);
} catch (NoSuchMethodException ignored) {
Method method = client.getClass().getDeclaredMethod("restClient");
method.setAccessible(true);
return (RestClient) method.invoke(client);
}
}
private static Object fieldOf(Object instance, String fieldName) throws Exception {
@@ -158,6 +158,9 @@ class HttpClientTimeoutEffectiveTest {
// 本机不可达地址(TEST-NET-1 保留段)+ 短连接超时 → 在限定时间内以传输错误失败
HttpClientProperties props = new HttpClientProperties();
props.setConnectTimeoutMillis(500);
// 地址取 127.0.0.1:1(本机未监听端口,连接立即被拒绝):原先用 TEST-NET-1 的 192.0.2.1
// 在装有 socks/透明代理的开发机上会被代理接管并等到代理自身超时(实测 63s),
// 使断言依赖运行机器的网络环境而假失败。
HttpClient client = HttpClient.newBuilder()
.connectTimeout(Duration.ofMillis(props.effectiveConnectTimeoutMillis()))
.build();
@@ -167,10 +170,10 @@ class HttpClientTimeoutEffectiveTest {
long startedAt = System.nanoTime();
try {
rest.get().uri("http://192.0.2.1:81/").retrieve().toBodilessEntity();
rest.get().uri("http://127.0.0.1:1/").retrieve().toBodilessEntity();
fail("不可达地址不应成功");
} catch (ResourceAccessException expected) {
// 预期:连接超时或快速不可达,均属传输错误
// 预期:连接被拒绝或快速不可达,均属传输错误
}
long elapsedMs = (System.nanoTime() - startedAt) / 1_000_000;
assertTrue(elapsedMs < 3_000, "连接超时应在限定时间内失败,实际 " + elapsedMs + "ms");
@@ -1,6 +1,10 @@
package com.nanri.aiimage.modules.collectdata.service;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.baomidou.mybatisplus.core.MybatisConfiguration;
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.baomidou.mybatisplus.core.metadata.TableInfoHelper;
import org.apache.ibatis.builder.MapperBuilderAssistant;
import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataDashboardVo;
import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataTaskBatchVo;
import com.nanri.aiimage.modules.task.mapper.FileResultMapper;
@@ -24,6 +28,7 @@ import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -61,6 +66,10 @@ class CollectDataServiceTest {
@BeforeEach
void setUp() {
// 2026-09failTask 改为条件更新(LambdaUpdateWrapper),它需要实体的 TableInfo 缓存;
// 纯单测没有 MyBatis 上下文,这里手动登记一次(幂等,重复调用只记一次)
TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new MybatisConfiguration(), ""),
com.nanri.aiimage.modules.task.model.entity.FileTaskEntity.class);
ReflectionTestUtils.setField(service, "staleTimeoutMinutes", 30L);
}
@@ -125,12 +134,13 @@ class CollectDataServiceTest {
service.failTask(task.getId(), task.getUserId(), "queue unavailable");
assertThat(task.getStatus()).isEqualTo("FAILED");
assertThat(task.getErrorMessage()).isEqualTo("queue unavailable");
// 2026-09 语义修订:failTask 改为条件更新(where status not in SUCCESS/FAILED),
// 通过 LambdaUpdateWrapper 落库 —— 不再修改内存实体,也不再走 updateById。
// 原因是客户端报错与结果文件组装并发时,整行 updateById 会互相覆盖(FAILED↔SUCCESS 横跳)。
assertThat(result.getSuccess()).isZero();
assertThat(result.getErrorMessage()).isEqualTo("queue unavailable");
verify(fileResultMapper).updateById(result);
verify(fileTaskMapper).updateById(task);
verify(fileTaskMapper).update(isNull(), any(LambdaUpdateWrapper.class));
}
@Test
@@ -1,6 +1,7 @@
package com.nanri.aiimage.modules.collectdata.util;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataResultRowVo;
import com.nanri.aiimage.modules.task.mapper.TaskResultItemMapper;
import com.nanri.aiimage.modules.task.model.entity.TaskResultItemEntity;
@@ -15,6 +16,7 @@ import java.util.List;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertEquals;
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.anyList;
@@ -210,11 +212,11 @@ class CollectData10kLoadTest {
}).when(taskResultItemMapper).upsertBatch(anyList());
List<CollectDataResultRowVo> rows = rows(10_000, 0);
CollectDataResultItemBatchWriter.UpsertCounts counts = writer.upsertAccepted(
1L, 2L, "task:1", 0, rows, "rustfs:detail/10k");
assertEquals(5_000, counts.insertedOrUpdated(), "首批失败跳过,后 50 批写入");
assertEquals(5_000, counts.newlyInserted(), "失败批新行从增量扣除");
// 2026-09 语义修订:存在失败批次时不再静默跳过,而是抛出明确失败让上游重试
// (静默跳批会让任务以 SUCCESS 收尾但明细缺行)
assertThrows(BusinessException.class,
() -> writer.upsertAccepted(1L, 2L, "task:1", 0, rows, "rustfs:detail/10k"),
"存在失败批次时必须抛出,不能静默返回部分成功");
// 恢复后重试同一输入:存量已落库行 hash 相等跳过,未落库行补插。
List<TaskResultItemEntity> existing = new ArrayList<>(5_000);
@@ -1,6 +1,7 @@
package com.nanri.aiimage.modules.collectdata.util;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataResultRowVo;
import com.nanri.aiimage.modules.task.mapper.TaskResultItemMapper;
import com.nanri.aiimage.modules.task.model.entity.TaskResultItemEntity;
@@ -13,6 +14,7 @@ import java.util.ArrayList;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertEquals;
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.anyList;
@@ -238,7 +240,10 @@ class CollectDataResultItemBatchWriterTest {
@Test
void test_task_050_task_dependency_failure_releases_resources() {
// 依赖失败:批量 upsert 抛错时跳过该批不中断,恢复后继续,无资源泄漏
// 依赖失败(2026-09 语义修订):批量 upsert 抛错时不再静默跳过该批,而是抛出明确失败
// 原因:静默跳批会让任务以 SUCCESS 收尾但明细缺行,且与「结果文件」sheet 的汇总数量对不上。
// 抛出让 worker 识别失败并重试;重提按 payload_hash 幂等(已落库批跳过、未落库批补插),
// 即使首批成功第二批失败,重提后也能收敛为完整数据。
when(taskResultItemMapper.selectList(any())).thenReturn(List.of());
doThrow(new RuntimeException("db down"))
.doAnswer(invocation -> ((List<?>) invocation.getArgument(0)).size())
@@ -248,10 +253,10 @@ class CollectDataResultItemBatchWriterTest {
rows.add(row("B" + String.format("%09d", i + 1), "brand"));
}
CollectDataResultItemBatchWriter.UpsertCounts counts =
writer.upsertAccepted(1L, 2L, "task:1", 0, rows, "rustfs:x");
assertEquals(10, counts.insertedOrUpdated(), "首批失败跳过,第二批 10 行写入");
BusinessException ex = assertThrows(BusinessException.class,
() -> writer.upsertAccepted(1L, 2L, "task:1", 0, rows, "rustfs:x"),
"存在失败批次时必须抛出,不能静默返回部分成功");
assertTrue(ex.getMessage().contains("失败批次"), "异常信息须说明失败批次: " + ex.getMessage());
verify(taskResultItemMapper, times(2)).upsertBatch(anyList());
}
@@ -1,6 +1,7 @@
package com.nanri.aiimage.modules.collectdata.util;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.common.exception.BusinessException;
import com.nanri.aiimage.modules.collectdata.model.vo.CollectDataResultRowVo;
import com.nanri.aiimage.modules.task.mapper.TaskResultItemMapper;
import com.nanri.aiimage.modules.task.model.entity.TaskResultItemEntity;
@@ -12,6 +13,7 @@ import java.util.ArrayList;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.Mockito.doThrow;
@@ -194,10 +196,11 @@ class CollectDataResultItemCountTest {
rows.add(row("B" + String.format("%09d", i + 1), "brand"));
}
CollectDataResultItemBatchWriter.UpsertCounts first =
writer.upsertAccepted(1L, 2L, "task:1", 0, rows, "rustfs:x");
assertEquals(10, first.newlyInserted(), "首批失败扣除,仅第二批 10 行真新增");
// 2026-09 语义修订:存在失败批次时不再静默跳过,而是抛出明确失败让上游重试。
// 首批(10 行)全部失败 → 无任何行落库;重提后按 hash 幂等补插。
assertThrows(BusinessException.class,
() -> writer.upsertAccepted(1L, 2L, "task:1", 0, rows, "rustfs:x"),
"存在失败批次时必须抛出,不能静默返回部分成功");
// 重提同一 chunk:第二批 10 行 hash 相等跳过,首批 10 行补插 → newlyInserted=10。
List<TaskResultItemEntity> existing = new ArrayList<>();
@@ -217,7 +220,8 @@ class CollectDataResultItemCountTest {
assertEquals(10, retry.newlyInserted(), "补插首批 10 行");
assertEquals(10, retry.skipped(), "第二批存量跳过");
assertEquals(20, first.newlyInserted() + retry.newlyInserted(), "任务内累计=表内真实行数");
// 首批抛异常时没有任何行落库(区别于旧语义的"跳过"),重提补插后表内共 20 行
assertEquals(20, retry.newlyInserted() + 10, "任务内累计=表内真实行数(重提前已落库 10 行)");
}
private static TaskResultItemEntity existingItem(String itemKey, int offset, int chunkIndex, String pointer) {
@@ -28,6 +28,7 @@ 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.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
@@ -139,7 +140,23 @@ class FaultInjectionTest {
when(client.getObject(any(io.minio.GetObjectArgs.class)))
.thenAnswer(invocation -> {
GetObjectResponse response = mock(GetObjectResponse.class);
when(response.readAllBytes()).thenReturn("{}".getBytes());
byte[] payload = "{}".getBytes();
// 生产代码走流式读取(read(byte[],int,int) 循环)。此前只 stub 了 readAllBytes()
// mock 的 read() 默认返回 0 —— 循环既不拿到数据也不见 -1,会一直空转并把每次调用
// 记进 Mockito 的 invocation 列表,最终 java.lang.OutOfMemoryError 打崩 fork JVM
// (此前表现为整个 surefire 报 One or more of the requested tests did not pass)。
java.io.ByteArrayInputStream delegate = new java.io.ByteArrayInputStream(payload);
// 生产代码用单参数 read(byte[]) 循环(见 RustfsObjectStorageService.readObjectBytes),
// 三个重载都 stub 以防调用路径变化
when(response.read(any(byte[].class)))
.thenAnswer(readInvocation -> delegate.read(readInvocation.getArgument(0)));
when(response.read(any(byte[].class), anyInt(), anyInt()))
.thenAnswer(readInvocation -> delegate.read(
readInvocation.getArgument(0),
readInvocation.getArgument(1),
readInvocation.getArgument(2)));
when(response.read()).thenAnswer(readInvocation -> delegate.read());
when(response.readAllBytes()).thenReturn(payload);
return response;
});
RustfsObjectStorageService service = new RustfsObjectStorageService(
@@ -79,7 +79,16 @@ class RustfsMetricsBaselineTest {
private GetObjectResponse readResponse(String content) throws Exception {
GetObjectResponse response = mock(GetObjectResponse.class);
when(response.readAllBytes()).thenReturn(content.getBytes(StandardCharsets.UTF_8));
byte[] payload = content.getBytes(StandardCharsets.UTF_8);
// 生产代码用单参数 read(byte[]) 循环;不 stub 时 mock 返回 0 会让循环空转到 OOM
java.io.ByteArrayInputStream delegate = new java.io.ByteArrayInputStream(payload);
when(response.read(org.mockito.ArgumentMatchers.any(byte[].class)))
.thenAnswer(inv -> delegate.read(inv.getArgument(0)));
when(response.read(org.mockito.ArgumentMatchers.any(byte[].class),
org.mockito.ArgumentMatchers.anyInt(), org.mockito.ArgumentMatchers.anyInt()))
.thenAnswer(inv -> delegate.read(inv.getArgument(0), inv.getArgument(1), inv.getArgument(2)));
when(response.read()).thenAnswer(inv -> delegate.read());
when(response.readAllBytes()).thenReturn(payload);
return response;
}
@@ -85,7 +85,15 @@ class RustfsTotalBudgetTest {
private void stubReadSuccess(String content) throws Exception {
GetObjectResponse response = mock(GetObjectResponse.class);
when(response.readAllBytes()).thenReturn(content.getBytes(StandardCharsets.UTF_8));
byte[] payload = content.getBytes(StandardCharsets.UTF_8);
// 生产代码用单参数 read(byte[]) 循环;不 stub 时 mock 返回 0 会让循环空转到 OOM
java.io.ByteArrayInputStream delegate = new java.io.ByteArrayInputStream(payload);
when(response.read(ArgumentMatchers.any(byte[].class)))
.thenAnswer(inv -> delegate.read(inv.getArgument(0)));
when(response.read(ArgumentMatchers.any(byte[].class), ArgumentMatchers.anyInt(), ArgumentMatchers.anyInt()))
.thenAnswer(inv -> delegate.read(inv.getArgument(0), inv.getArgument(1), inv.getArgument(2)));
when(response.read()).thenAnswer(inv -> delegate.read());
when(response.readAllBytes()).thenReturn(payload);
when(client.getObject(ArgumentMatchers.any(GetObjectArgs.class))).thenReturn(response);
}
@@ -99,6 +99,8 @@ class PublishTaskServiceTest {
@Mock private TaskScopeStateMapper taskScopeStateMapper;
@Mock private TaskFileJobService taskFileJobService;
@Mock private TaskDistributedLockService taskDistributedLockService;
/** 2026-09stale 判死改为 job 级分布式锁保护(双实例只跑一处),测试需注入该依赖。 */
@Mock private com.nanri.aiimage.common.service.DistributedJobLockService distributedJobLockService;
@Mock private TransientPayloadStorageService transientPayloadStorageService;
@Mock private OssStorageService ossStorageService;
@Spy private ObjectMapper objectMapper = new ObjectMapper();
@@ -118,6 +120,9 @@ class PublishTaskServiceTest {
void executeTransactionsInline() {
configureChunkStorage();
lenient().when(instanceMetadata.getInstanceId()).thenReturn("instance-a");
// job 锁:默认放行(返回可关闭的 handle),否则 failStaleTasks 会直接 return 而不做判死
lenient().when(distributedJobLockService.tryLock(any(String.class), any(java.time.Duration.class)))
.thenAnswer(invocation -> mock(com.nanri.aiimage.common.service.DistributedJobLockService.LockHandle.class));
lenient().when(transactionTemplate.execute(any())).thenAnswer(invocation -> {
TransactionCallback<?> callback = invocation.getArgument(0);
transactionActive.set(true);
@@ -223,6 +223,38 @@ class ShopDataCrawlCleanupTest {
taskStore.put(copy.getId(), copy);
return 1;
});
// 2026-09:陈旧扫描的 FAILED 写入改为条件更新(update(entity, wrapper)where status='RUNNING')——
// 按 wrapper 参数里的 taskId 把 status 落回 store,模拟 DB 的条件更新结果(返回 1 表示命中)。
lenient().when(fileTaskMapper.update(org.mockito.ArgumentMatchers.isNull(),
any(com.baomidou.mybatisplus.core.conditions.Wrapper.class)))
.thenAnswer(invocation -> {
com.baomidou.mybatisplus.core.conditions.Wrapper<FileTaskEntity> wrapper =
invocation.getArgument(1);
// getParamNameValuePairs 在 AbstractWrapper 上(不在 Wrapper 接口),用反射取条件参数
java.util.Map<String, Object> params;
try {
java.lang.reflect.Method m = wrapper.getClass().getMethod("getParamNameValuePairs");
@SuppressWarnings("unchecked")
java.util.Map<String, Object> extracted = (java.util.Map<String, Object>) m.invoke(wrapper);
params = extracted;
} catch (Exception ex) {
params = java.util.Map.of();
}
java.util.Optional<Object> idValue = params.values().stream()
.filter(v -> v instanceof Long).findFirst();
if (idValue.isEmpty()) {
return 0;
}
FileTaskEntity stored = taskStore.get((Long) idValue.get());
if (stored == null) {
return 0;
}
stored.setStatus("FAILED");
stored.setUpdatedAt(java.time.LocalDateTime.now());
stored.setFinishedAt(java.time.LocalDateTime.now());
taskStore.put(stored.getId(), stored);
return 1;
});
lenient().when(fileTaskMapper.deleteById(anyLong())).thenAnswer(invocation -> {
long taskId = invocation.getArgument(0);
taskStore.remove(taskId);
@@ -153,6 +153,12 @@ class ShopDataCrawlOwnerColumnTest {
return value == null ? "" : value.trim();
});
lenient().when(fileTaskMapper.updateById(any(FileTaskEntity.class))).thenReturn(1);
// 2026-09:陈旧扫描的 FAILED 写入改为条件更新(where status='RUNNING' 的 CAS),
// 走的是 update(entity=null, wrapper);需要 stub 返回 1,否则默认 0 会被当成"未翻转"(任务保持 RUNNING)。
// 注意用 isNull()Mockito 2+ 的 any(Class) 不匹配 null。
lenient().when(fileTaskMapper.update(org.mockito.ArgumentMatchers.isNull(),
any(com.baomidou.mybatisplus.core.conditions.Wrapper.class)))
.thenReturn(1);
lenient().when(fileTaskMapper.selectById(anyLong())).thenReturn(null);
lenient().when(fileTaskMapper.selectList(any())).thenAnswer(invocation -> {
Wrapper<FileTaskEntity> wrapper = invocation.getArgument(0);
@@ -5,6 +5,7 @@ import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinTaskLightReque
import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinTaskLightBatchVo;
import com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinTaskLightVo;
import com.nanri.aiimage.modules.similarasin.service.SimilarAsinTaskService;
import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
import org.junit.jupiter.api.Test;
import org.springframework.http.MediaType;
import org.springframework.test.web.servlet.MockMvc;
@@ -53,7 +54,10 @@ class SimilarAsinTaskLightControllerTest {
private void setUpWith(SimilarAsinTaskLightBatchVo response) throws Exception {
service = mock(SimilarAsinTaskService.class);
when(service.progressLight(any())).thenReturn(response);
mockMvc = MockMvcBuilders.standaloneSetup(new SimilarAsinController(service)).build();
// 归属过滤:请求不带 userId 时 filterOwnedTaskIds 直接返回原列表(不会触碰 mapper),
// 故这里传 null mapper 即可,无需 mock
mockMvc = MockMvcBuilders.standaloneSetup(
new SimilarAsinController(service, new TaskProgressOwnershipSupport(null))).build();
}
@Test
@@ -61,7 +65,7 @@ class SimilarAsinTaskLightControllerTest {
setUpWith(lightBatch(lightItem(3938L, "RUNNING", "RUNNING", false, "2026-04-26T10:05:00")));
mockMvc.perform(post(LIGHT_URL)
.contentType(MediaType.APPLICATION_JSON)
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(3938L)))))
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(3938L), null))))
.andExpect(status().isOk())
.andExpect(jsonPath("$.success").value(true))
.andExpect(jsonPath("$.data.items[0].taskId").value(3938))
@@ -76,7 +80,7 @@ class SimilarAsinTaskLightControllerTest {
setUpWith(lightBatch(lightItem(1L, "SUCCESS", "SUCCESS", true, "2026-04-26T10:10:00")));
mockMvc.perform(post(LIGHT_URL)
.contentType(MediaType.APPLICATION_JSON)
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(1L)))))
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(1L), null))))
.andExpect(status().isOk())
.andExpect(jsonPath("$.success").value(true))
.andExpect(jsonPath("$.data.items[0].taskId").value(1));
@@ -88,7 +92,7 @@ class SimilarAsinTaskLightControllerTest {
setUpWith(new SimilarAsinTaskLightBatchVo());
mockMvc.perform(post(LIGHT_URL)
.contentType(MediaType.APPLICATION_JSON)
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of()))))
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(), null))))
.andExpect(status().isOk())
.andExpect(jsonPath("$.data.items").isEmpty())
.andExpect(jsonPath("$.data.missingTaskIds").isEmpty());
@@ -99,7 +103,7 @@ class SimilarAsinTaskLightControllerTest {
setUpWith(lightBatch(lightItem(7L, "PENDING", null, false, "2026-04-26T10:00:00")));
mockMvc.perform(post(LIGHT_URL)
.contentType(MediaType.APPLICATION_JSON)
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(7L, 7L, 7L)))))
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(7L, 7L, 7L), null))))
.andExpect(status().isOk())
.andExpect(jsonPath("$.data.items.length()").value(1));
}
@@ -110,7 +114,10 @@ class SimilarAsinTaskLightControllerTest {
request.setTaskIds(List.of(1L));
service = mock(SimilarAsinTaskService.class);
when(service.progressBatch(any())).thenReturn(new com.nanri.aiimage.modules.similarasin.model.vo.SimilarAsinTaskBatchVo());
mockMvc = MockMvcBuilders.standaloneSetup(new SimilarAsinController(service)).build();
// 归属过滤:请求不带 userId 时 filterOwnedTaskIds 直接返回原列表(不会触碰 mapper),
// 故这里传 null mapper 即可,无需 mock
mockMvc = MockMvcBuilders.standaloneSetup(
new SimilarAsinController(service, new TaskProgressOwnershipSupport(null))).build();
mockMvc.perform(post("/api/similar-asin/tasks/progress/batch")
.contentType(MediaType.APPLICATION_JSON)
.content(objectMapper.writeValueAsString(request)))
@@ -124,7 +131,7 @@ class SimilarAsinTaskLightControllerTest {
setUpWith(lightBatch(lightItem(2L, "SUCCESS", "SUCCESS", true, "2026-04-26T10:11:00")));
mockMvc.perform(post(LIGHT_URL)
.contentType(MediaType.APPLICATION_JSON)
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(2L)))))
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(2L), null))))
.andExpect(status().isOk())
.andExpect(jsonPath("$.data.items[0].status").value("SUCCESS"));
}
@@ -134,7 +141,7 @@ class SimilarAsinTaskLightControllerTest {
setUpWith(lightBatch(lightItem(3L, "RUNNING", "SUCCESS", true, "2026-04-26T10:12:00")));
mockMvc.perform(post(LIGHT_URL)
.contentType(MediaType.APPLICATION_JSON)
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(3L)))))
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(3L), null))))
.andExpect(status().isOk())
.andExpect(jsonPath("$.data.items[0].fileStatus").value("SUCCESS"))
.andExpect(jsonPath("$.data.items[0].fileReady").value(true));
@@ -147,7 +154,7 @@ class SimilarAsinTaskLightControllerTest {
setUpWith(vo);
mockMvc.perform(post(LIGHT_URL)
.contentType(MediaType.APPLICATION_JSON)
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(99999L)))))
.content(objectMapper.writeValueAsString(new SimilarAsinTaskLightRequest(List.of(99999L), null))))
.andExpect(status().isOk())
.andExpect(jsonPath("$.data.items").isEmpty())
.andExpect(jsonPath("$.data.missingTaskIds[0]").value(99999));
@@ -15,6 +15,7 @@ import com.nanri.aiimage.modules.publish.service.PublishTaskService;
import com.nanri.aiimage.modules.similarasin.controller.SimilarAsinController;
import com.nanri.aiimage.modules.similarasin.model.dto.SimilarAsinSubmitResultRequest;
import com.nanri.aiimage.modules.similarasin.service.SimilarAsinTaskService;
import com.nanri.aiimage.modules.task.service.TaskProgressOwnershipSupport;
import jakarta.servlet.http.HttpServletResponse;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@@ -64,7 +65,7 @@ class ResultSuccessTimingContractTest {
return null;
}).when(similarAsinTaskService).submitResult(eq(TASK_ID), any(SimilarAsinSubmitResultRequest.class));
ApiResponse<Void> response = new SimilarAsinController(similarAsinTaskService)
ApiResponse<Void> response = new SimilarAsinController(similarAsinTaskService, new TaskProgressOwnershipSupport(null))
.result(TASK_ID, new SimilarAsinSubmitResultRequest(), mock(HttpServletResponse.class));
responseAfterService.set(response != null);
@@ -79,7 +80,7 @@ class ResultSuccessTimingContractTest {
doAnswer(invocation -> null)
.when(appearancePatentTaskService).submitResult(eq(TASK_ID), any(AppearancePatentSubmitResultRequest.class));
ApiResponse<Void> response = new AppearancePatentController(appearancePatentTaskService)
ApiResponse<Void> response = new AppearancePatentController(appearancePatentTaskService, new TaskProgressOwnershipSupport(null))
.result(TASK_ID, new AppearancePatentSubmitResultRequest(), mock(HttpServletResponse.class));
assertTrue(response.isSuccess(), "success=true");
@@ -93,7 +94,7 @@ class ResultSuccessTimingContractTest {
vo.setChunkIndex(1);
when(collectDataService.submitResult(eq(TASK_ID), any(CollectDataSubmitResultRequest.class))).thenReturn(vo);
ApiResponse<CollectDataSubmitResultVo> response = new CollectDataController(collectDataService)
ApiResponse<CollectDataSubmitResultVo> response = new CollectDataController(collectDataService, new TaskProgressOwnershipSupport(null))
.submitResult(TASK_ID, new CollectDataSubmitResultRequest());
assertTrue(response.isSuccess(), "success=true");
@@ -118,7 +119,7 @@ class ResultSuccessTimingContractTest {
org.mockito.Mockito.doThrow(new BusinessException("落库失败"))
.when(similarAsinTaskService).submitResult(eq(TASK_ID), any(SimilarAsinSubmitResultRequest.class));
assertThrows(BusinessException.class, () -> new SimilarAsinController(similarAsinTaskService)
assertThrows(BusinessException.class, () -> new SimilarAsinController(similarAsinTaskService, new TaskProgressOwnershipSupport(null))
.result(TASK_ID, new SimilarAsinSubmitResultRequest(), mock(HttpServletResponse.class)));
}
@@ -131,7 +132,7 @@ class ResultSuccessTimingContractTest {
}
return null;
}).when(similarAsinTaskService).submitResult(eq(TASK_ID), any(SimilarAsinSubmitResultRequest.class));
SimilarAsinController controller = new SimilarAsinController(similarAsinTaskService);
SimilarAsinController controller = new SimilarAsinController(similarAsinTaskService, new TaskProgressOwnershipSupport(null));
assertThrows(BusinessException.class, () -> controller
.result(TASK_ID, new SimilarAsinSubmitResultRequest(), mock(HttpServletResponse.class)));
@@ -149,7 +150,7 @@ class ResultSuccessTimingContractTest {
order.add("service-submit");
return null;
}).when(appearancePatentTaskService).submitResult(eq(TASK_ID), any(AppearancePatentSubmitResultRequest.class));
AppearancePatentController controller = new AppearancePatentController(appearancePatentTaskService);
AppearancePatentController controller = new AppearancePatentController(appearancePatentTaskService, new TaskProgressOwnershipSupport(null));
ApiResponse<Void> response = controller.result(TASK_ID, new AppearancePatentSubmitResultRequest(),
mock(HttpServletResponse.class));
@@ -224,13 +224,17 @@ class RollbackSemanticsContractTest {
verify(taskCacheService).touchTaskHeartbeat(TASK_ID);
}
private Object pipeline() {
return ReflectionTestUtils.invokeMethod(service, "pipelineSupport");
}
@Test
void computePhaseIsSideEffectFreeAndIdempotent() throws Exception {
FileTaskEntity task = runningTask();
when(fileTaskMapper.selectById(TASK_ID)).thenReturn(task);
Object first = ReflectionTestUtils.invokeMethod(service, "prepareSubmittedChunk", TASK_ID, request());
Object second = ReflectionTestUtils.invokeMethod(service, "prepareSubmittedChunk", TASK_ID, request());
Object first = ReflectionTestUtils.invokeMethod(pipeline(), "prepareSubmittedChunk", TASK_ID, request());
Object second = ReflectionTestUtils.invokeMethod(pipeline(), "prepareSubmittedChunk", TASK_ID, request());
assertEquals(DigestUtil.sha256Hex(storedPayloadJson.get()),
ReflectionTestUtils.getField(first, "payloadHash"));
@@ -1,6 +1,7 @@
package com.nanri.aiimage.modules.task.contract;
import com.baomidou.mybatisplus.core.MybatisConfiguration;
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.baomidou.mybatisplus.core.metadata.TableInfoHelper;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.nanri.aiimage.modules.collectdata.mapper.CollectDataCountryPrefMapper;
@@ -46,6 +47,7 @@ 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.eq;
import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.inOrder;
import static org.mockito.Mockito.lenient;
@@ -125,20 +127,24 @@ class SuccessTimingContractTest {
var order = inOrder(ossStorageService, fileResultMapper, fileTaskMapper);
order.verify(ossStorageService).uploadResultFile(any(), eq("COLLECT_DATA"));
order.verify(fileResultMapper).updateById(any(FileResultEntity.class));
order.verify(fileTaskMapper).updateById(any(FileTaskEntity.class));
// 2026-09:终态写入改为条件更新(where status != FAILED),不再走 updateById(实体)
order.verify(fileTaskMapper).update(isNull(), any(LambdaUpdateWrapper.class));
}
@Test
void taskStatusBecomesSuccessWithFileFields() {
service.processResultFileJob(job());
ArgumentCaptor<FileTaskEntity> taskCaptor = ArgumentCaptor.forClass(FileTaskEntity.class);
verify(fileTaskMapper).updateById(taskCaptor.capture());
FileTaskEntity updated = taskCaptor.getValue();
assertEquals("SUCCESS", updated.getStatus(), "任务状态成功时机 = 结果文件生成后");
assertEquals(1, updated.getSuccessFileCount());
assertEquals(0, updated.getFailedFileCount());
assertNotNull(updated.getFinishedAt());
// 2026-09:终态写入改为条件更新(LambdaUpdateWrapperwhere status != FAILED),
// 不再有实体可捕获 —— 改为捕获 wrapper 并断言其 set 片段包含契约要求的字段。
ArgumentCaptor<LambdaUpdateWrapper<FileTaskEntity>> wrapperCaptor =
ArgumentCaptor.forClass(LambdaUpdateWrapper.class);
verify(fileTaskMapper).update(isNull(), wrapperCaptor.capture());
String setSql = String.valueOf(wrapperCaptor.getValue().getSqlSet());
assertTrue(setSql.contains("status"), "任务状态成功时机 = 结果文件生成后,set 片段应含 status: " + setSql);
assertTrue(setSql.contains("success_file_count"), "set 片段应含 success_file_count: " + setSql);
assertTrue(setSql.contains("failed_file_count"), "set 片段应含 failed_file_count: " + setSql);
assertTrue(setSql.contains("finished_at"), "set 片段应含 finished_at: " + setSql);
}
@Test
@@ -162,13 +168,16 @@ class SuccessTimingContractTest {
// fileReady 由 url 推导(TaskProgressLightAssembler 语义),成功时机与 url 落库绑定
verify(fileResultMapper).updateById(any(FileResultEntity.class));
verify(fileTaskMapper).updateById(any(FileTaskEntity.class));
// 2026-09:任务终态改为条件更新(见 taskStatusBecomesSuccessWithFileFields
verify(fileTaskMapper).update(isNull(), any(LambdaUpdateWrapper.class));
}
@Test
void noFileGeneratedNoSuccess() {
// 方法名同步:实现已由 writeWorkbookSegmented 改为 writeWorkbookStreaming(流式分页读,
// 避免几十万行结果集整体驻留堆),测试原先仍在 mock 旧方法名,导致"没抛异常"的假失败。
doThrow(new IllegalStateException("工作簿生成失败"))
.when(excelAssemblyService).writeWorkbookSegmented(any(), any(), any(), any());
.when(excelAssemblyService).writeWorkbookStreaming(any(), any(), any(), any());
assertThrows(IllegalStateException.class, () -> service.processResultFileJob(job()));
@@ -202,11 +211,13 @@ class SuccessTimingContractTest {
void successContractFrozen() {
service.processResultFileJob(job());
// 契约快照:文件生成+上传 → result 落库 → 任务 SUCCESS(一次 updateById 各一
// 契约快照:文件生成+上传 → result 落库 → 任务 SUCCESS一次写入
verify(ossStorageService).uploadResultFile(any(), eq("COLLECT_DATA"));
verify(fileResultMapper).updateById(any(FileResultEntity.class));
verify(fileTaskMapper).updateById(any(FileTaskEntity.class));
verify(fileTaskMapper, org.mockito.Mockito.times(1)).updateById(any(FileTaskEntity.class));
// 2026-09:任务终态写入由 updateById(实体) 改为条件更新(where status != FAILED
// 避免覆盖客户端 /fail 并发上报的失败态)
verify(fileTaskMapper).update(isNull(), any(LambdaUpdateWrapper.class));
verify(fileTaskMapper, org.mockito.Mockito.times(1)).update(isNull(), any(LambdaUpdateWrapper.class));
}
private FileTaskEntity runningTask() {
@@ -38,8 +38,11 @@ class ResultFileJobHandlerTest {
@Test
void interfaceMethodsPresent() throws Exception {
assertEquals(7, ResultFileJobHandler.class.getDeclaredMethods().length,
"接口方法数量为 7moduleType/process/onSuccess/cleanup/onFailure/supportsAsyncOffload/isOwnerScoped");
// 2026-09:新增 fallbackAssembleOnFailurejob 重试耗尽时用部分数据兜底生成结果文件),
// 接口方法从 7 个增到 8 个,同步更新契约断言。
assertEquals(8, ResultFileJobHandler.class.getDeclaredMethods().length,
"接口方法数量为 8moduleType/process/onSuccess/cleanup/onFailure/"
+ "fallbackAssembleOnFailure/supportsAsyncOffload/isOwnerScoped");
assertNotNull(methodOf(ResultFileJobHandler.class, "moduleType"));
assertNotNull(methodOf(ResultFileJobHandler.class, "process", TaskFileJobEntity.class));
assertNotNull(methodOf(ResultFileJobHandler.class, "onSuccess", TaskFileJobEntity.class));
@@ -117,8 +117,11 @@ class TransientPayloadDeleteOrchestratorTest {
assertEquals(2, orchestrator.flushPendingDeletes(), "两个对象提交物理删除");
assertTrue(done.await(2, TimeUnit.SECONDS), "异步删除完成");
verify(rustfs, times(2)).deleteObject(anyString());
verify(chunkMapper).selectList(any(LambdaQueryWrapper.class));
verify(scopeStateMapper).selectList(any(LambdaQueryWrapper.class));
// 2026-09 修订:批量反查按指针里的 taskId 分组后各查一次(用 task_id 索引收敛,
// 避免 payload_json 这类 JSON 列的 IN 比较全表扫描)——本用例的 payload 来自 2 个 task
// 故两张表各被查询 2 次(替换了原「合并成一次 IN 查询」的实现)。
verify(chunkMapper, times(2)).selectList(any(LambdaQueryWrapper.class));
verify(scopeStateMapper, times(2)).selectList(any(LambdaQueryWrapper.class));
}
@Test