task-159: 已完成任务活跃 Job 巡检(终态任务 RUNNING/PENDING Job 检出、运行中任务正常、只读、可重复)+ 8 条测试

This commit is contained in:
2026-09-02 07:08:53 +08:00
parent 9998571770
commit 5305d6dbc2
2 changed files with 252 additions and 0 deletions
@@ -0,0 +1,82 @@
package com.nanri.aiimage.modules.task.service;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.nanri.aiimage.modules.task.mapper.FileTaskMapper;
import com.nanri.aiimage.modules.task.mapper.TaskFileJobMapper;
import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.List;
import java.util.Set;
import java.util.stream.Collectors;
/**
* 已完成任务活跃 Job 巡检(task-159)。
*
* 只读巡检:任务已终态(SUCCESS/FAILED)但仍有活跃(RUNNING/PENDINGJob 的
* 清单输出报表;运行中任务的活跃 Job 属正常;绝不修改任何状态。可重复执行。
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class CompletedTaskActiveJobInspector {
private final FileTaskMapper fileTaskMapper;
private final TaskFileJobMapper taskFileJobMapper;
public record ActiveJobEntry(Long taskId, String moduleType, Long jobId, String jobStatus, String jobType) {
}
public record ActiveJobReport(List<ActiveJobEntry> entries) {
public boolean isEmpty() {
return entries == null || entries.isEmpty();
}
}
public ActiveJobReport inspectTerminalTasksWithActiveJobs(int limit) {
List<FileTaskEntity> tasks = fileTaskMapper.selectList(new LambdaQueryWrapper<FileTaskEntity>()
.in(FileTaskEntity::getStatus, List.of("SUCCESS", "FAILED"))
.orderByAsc(FileTaskEntity::getUpdatedAt)
.last("limit " + Math.max(1, Math.min(limit, 500))));
if (tasks == null || tasks.isEmpty()) {
return new ActiveJobReport(List.of());
}
Set<Long> taskIds = tasks.stream()
.map(FileTaskEntity::getId)
.filter(id -> id != null && id > 0)
.collect(Collectors.toSet());
List<TaskFileJobEntity> activeJobs = taskIds.isEmpty() ? List.of()
: taskFileJobMapper.selectList(new LambdaQueryWrapper<TaskFileJobEntity>()
.in(TaskFileJobEntity::getTaskId, taskIds)
.in(TaskFileJobEntity::getStatus, List.of("RUNNING", "PENDING")));
if (activeJobs == null || activeJobs.isEmpty()) {
return new ActiveJobReport(List.of());
}
java.util.Map<Long, FileTaskEntity> taskById = tasks.stream()
.collect(Collectors.toMap(FileTaskEntity::getId, task -> task, (a, b) -> a));
List<ActiveJobEntry> entries = new ArrayList<>();
for (TaskFileJobEntity job : activeJobs) {
FileTaskEntity task = taskById.get(job.getTaskId());
if (task != null) {
entries.add(new ActiveJobEntry(
task.getId(), task.getModuleType(),
job.getId(), job.getStatus(), job.getJobType()));
}
}
ActiveJobReport report = new ActiveJobReport(List.copyOf(entries));
if (!report.isEmpty()) {
log.info("terminal-task-active-job report: count={}", report.entries().size());
for (ActiveJobEntry entry : report.entries()) {
log.info("terminal-task-active-job: taskId={} moduleType={} jobId={} jobStatus={} jobType={}",
entry.taskId(), entry.moduleType(), entry.jobId(), entry.jobStatus(), entry.jobType());
}
}
return report;
}
}
@@ -0,0 +1,170 @@
package com.nanri.aiimage.modules.task.service;
import com.baomidou.mybatisplus.core.MybatisConfiguration;
import com.baomidou.mybatisplus.core.metadata.TableInfoHelper;
import com.nanri.aiimage.modules.task.mapper.FileTaskMapper;
import com.nanri.aiimage.modules.task.mapper.TaskFileJobMapper;
import com.nanri.aiimage.modules.task.model.entity.FileTaskEntity;
import com.nanri.aiimage.modules.task.model.entity.TaskFileJobEntity;
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.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
/**
* task-159:已完成任务活跃 Job 巡检契约(plan 09)。
* 终态任务仍有 RUNNING/PENDING Job → 检出;运行中任务的活跃 Job 正常不报;
* 只读;可重复。
*/
@ExtendWith(MockitoExtension.class)
class CompletedTaskActiveJobInspectorTest {
private static final Long TERMINAL_TASK_ID = 5001L;
private static final Long RUNNING_TASK_ID = 5101L;
@Mock private FileTaskMapper fileTaskMapper;
@Mock private TaskFileJobMapper taskFileJobMapper;
private CompletedTaskActiveJobInspector inspector;
@BeforeAll
static void initializeMybatisMetadata() {
MapperBuilderAssistant assistant = new MapperBuilderAssistant(new MybatisConfiguration(), "");
TableInfoHelper.initTableInfo(assistant, FileTaskEntity.class);
TableInfoHelper.initTableInfo(assistant, TaskFileJobEntity.class);
}
@BeforeEach
void setUp() {
inspector = new CompletedTaskActiveJobInspector(fileTaskMapper, taskFileJobMapper);
lenient().when(taskFileJobMapper.selectList(any())).thenReturn(List.of());
}
private FileTaskEntity task(Long id, String status) {
FileTaskEntity task = new FileTaskEntity();
task.setId(id);
task.setModuleType("SIMILAR_ASIN");
task.setStatus(status);
return task;
}
private TaskFileJobEntity job(Long id, Long taskId, String status) {
TaskFileJobEntity job = new TaskFileJobEntity();
job.setId(id);
job.setTaskId(taskId);
job.setStatus(status);
job.setJobType("ASSEMBLE_RESULT");
return job;
}
@Test
void activeJobAfterTerminalTaskDetected() {
when(fileTaskMapper.selectList(any())).thenReturn(List.of(task(TERMINAL_TASK_ID, "SUCCESS")));
when(taskFileJobMapper.selectList(any()))
.thenReturn(List.of(job(901L, TERMINAL_TASK_ID, "RUNNING")));
var report = inspector.inspectTerminalTasksWithActiveJobs(50);
assertEquals(1, report.entries().size());
assertEquals(TERMINAL_TASK_ID, report.entries().getFirst().taskId());
assertEquals("RUNNING", report.entries().getFirst().jobStatus());
}
@Test
void terminalTaskWithoutJobsIsFine() {
when(fileTaskMapper.selectList(any())).thenReturn(List.of(task(TERMINAL_TASK_ID, "SUCCESS")));
var report = inspector.inspectTerminalTasksWithActiveJobs(50);
assertTrue(report.isEmpty());
}
@Test
void runningTaskWithActiveJobIsNormal() {
// 巡检只查终态任务:RUNNING 任务的活跃 Job 不在报表范围
when(fileTaskMapper.selectList(any())).thenReturn(List.of());
inspector.inspectTerminalTasksWithActiveJobs(50);
var wrapper = org.mockito.ArgumentCaptor.forClass(com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper.class);
verify(fileTaskMapper).selectList(wrapper.capture());
assertTrue(wrapper.getValue().getSqlSegment().contains("IN"), "只查终态任务");
assertFalse(wrapper.getValue().getSqlSegment().contains("RUNNING"));
}
@Test
void stuckJobAfterTerminalReported() {
when(fileTaskMapper.selectList(any())).thenReturn(List.of(task(TERMINAL_TASK_ID, "FAILED")));
when(taskFileJobMapper.selectList(any()))
.thenReturn(List.of(job(902L, TERMINAL_TASK_ID, "PENDING")));
var report = inspector.inspectTerminalTasksWithActiveJobs(50);
assertEquals(1, report.entries().size(), "卡死 PENDING Job 同样检出");
}
@Test
void emptyReportWhenNoTerminalTasks() {
when(fileTaskMapper.selectList(any())).thenReturn(List.of());
var report = inspector.inspectTerminalTasksWithActiveJobs(50);
assertTrue(report.isEmpty());
}
@Test
void readOnlyDoesNotModifyAnything() {
when(fileTaskMapper.selectList(any())).thenReturn(List.of(task(TERMINAL_TASK_ID, "SUCCESS")));
when(taskFileJobMapper.selectList(any()))
.thenReturn(List.of(job(903L, TERMINAL_TASK_ID, "RUNNING")));
inspector.inspectTerminalTasksWithActiveJobs(50);
verify(taskFileJobMapper, never()).update(any(), any());
verify(taskFileJobMapper, never()).updateById(any(TaskFileJobEntity.class));
verify(taskFileJobMapper, never()).delete(any());
}
@Test
void reportIncludesDetails() {
when(fileTaskMapper.selectList(any())).thenReturn(List.of(task(TERMINAL_TASK_ID, "SUCCESS")));
when(taskFileJobMapper.selectList(any()))
.thenReturn(List.of(job(904L, TERMINAL_TASK_ID, "RUNNING")));
var report = inspector.inspectTerminalTasksWithActiveJobs(50);
CompletedTaskActiveJobInspector.ActiveJobEntry entry = report.entries().getFirst();
assertEquals(TERMINAL_TASK_ID, entry.taskId());
assertEquals("SIMILAR_ASIN", entry.moduleType());
assertEquals(904L, entry.jobId());
assertEquals("ASSEMBLE_RESULT", entry.jobType());
assertNotNull(entry.jobStatus());
}
@Test
void rerunIsSafeAndStable() {
when(fileTaskMapper.selectList(any())).thenReturn(List.of(task(TERMINAL_TASK_ID, "SUCCESS")));
when(taskFileJobMapper.selectList(any()))
.thenReturn(List.of(job(905L, TERMINAL_TASK_ID, "PENDING")));
var first = inspector.inspectTerminalTasksWithActiveJobs(50);
var second = inspector.inspectTerminalTasksWithActiveJobs(50);
assertEquals(first.entries().size(), second.entries().size());
assertEquals(first.entries().getFirst().jobId(), second.entries().getFirst().jobId());
}
}