task-161: 巡检任务开关(aiimage.inspection.enabled 默认 disabled、cron 调度、分布式锁双实例去重、单巡检异常隔离、limit 透传)+ 8 条测试
This commit is contained in:
@@ -0,0 +1,24 @@
|
|||||||
|
package com.nanri.aiimage.config;
|
||||||
|
|
||||||
|
import lombok.Data;
|
||||||
|
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||||
|
import org.springframework.stereotype.Component;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 巡检任务配置(task-161)。
|
||||||
|
* 全部巡检默认 disabled;启用后按调度间隔执行;limit 为单次巡检行数上限。
|
||||||
|
*/
|
||||||
|
@Data
|
||||||
|
@Component
|
||||||
|
@ConfigurationProperties(prefix = "aiimage.inspection")
|
||||||
|
public class InspectionProperties {
|
||||||
|
|
||||||
|
/** 巡检总开关(默认关闭)。 */
|
||||||
|
private boolean enabled = false;
|
||||||
|
|
||||||
|
/** 调度 cron(默认凌晨 3 点)。 */
|
||||||
|
private String cron = "0 0 3 * * *";
|
||||||
|
|
||||||
|
/** 单次巡检行数上限。 */
|
||||||
|
private int limit = 200;
|
||||||
|
}
|
||||||
+81
@@ -0,0 +1,81 @@
|
|||||||
|
package com.nanri.aiimage.modules.task.service;
|
||||||
|
|
||||||
|
import com.nanri.aiimage.common.service.DistributedJobLockService;
|
||||||
|
import com.nanri.aiimage.config.InspectionProperties;
|
||||||
|
import com.nanri.aiimage.config.StorageProperties;
|
||||||
|
import com.nanri.aiimage.modules.file.service.TempOrphanInspector;
|
||||||
|
import lombok.RequiredArgsConstructor;
|
||||||
|
import lombok.extern.slf4j.Slf4j;
|
||||||
|
import org.springframework.scheduling.annotation.Scheduled;
|
||||||
|
import org.springframework.stereotype.Service;
|
||||||
|
|
||||||
|
import java.io.File;
|
||||||
|
import java.time.Duration;
|
||||||
|
import java.time.Instant;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 巡检任务调度(task-161)。
|
||||||
|
*
|
||||||
|
* 全部巡检由 aiimage.inspection.enabled 开关控制(默认 disabled);
|
||||||
|
* 启用后按 cron 调度执行;分布式锁防双实例重复输出;单个巡检异常不阻断其余。
|
||||||
|
*/
|
||||||
|
@Slf4j
|
||||||
|
@Service
|
||||||
|
@RequiredArgsConstructor
|
||||||
|
public class InspectionScheduler {
|
||||||
|
|
||||||
|
private static final Duration INSPECTION_LOCK_TTL = Duration.ofMinutes(10);
|
||||||
|
|
||||||
|
private final InspectionProperties inspectionProperties;
|
||||||
|
private final DistributedJobLockService distributedJobLockService;
|
||||||
|
private final StorageProperties storageProperties;
|
||||||
|
private final TempOrphanInspector tempOrphanInspector;
|
||||||
|
private final OrphanJobInspector orphanJobInspector;
|
||||||
|
private final TaskResultMissingInspector taskResultMissingInspector;
|
||||||
|
private final ResultFileMissingInspector resultFileMissingInspector;
|
||||||
|
private final CompletedTaskActiveJobInspector completedTaskActiveJobInspector;
|
||||||
|
|
||||||
|
@Scheduled(cron = "${aiimage.inspection.cron:0 0 3 * * *}")
|
||||||
|
public void runInspections() {
|
||||||
|
if (!inspectionProperties.isEnabled()) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
DistributedJobLockService.LockHandle lockHandle =
|
||||||
|
distributedJobLockService.tryLock("temp-file-inspection", INSPECTION_LOCK_TTL);
|
||||||
|
if (lockHandle == null) {
|
||||||
|
log.info("[inspection] skip because another instance holds the inspection lock");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
try (lockHandle) {
|
||||||
|
runAllSafely();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private void runAllSafely() {
|
||||||
|
int limit = inspectionProperties.getLimit();
|
||||||
|
runSafely("orphan-file", () -> {
|
||||||
|
File tempDir = cn.hutool.core.io.FileUtil.file(storageProperties.getLocalTempDir());
|
||||||
|
if (tempDir.isDirectory()) {
|
||||||
|
var report = tempOrphanInspector.inspectOrphanFiles(
|
||||||
|
tempDir, Instant.now().minus(Duration.ofHours(24)), name -> false);
|
||||||
|
log.info("[inspection] orphan-file report count={}", report.entries().size());
|
||||||
|
}
|
||||||
|
});
|
||||||
|
runSafely("orphan-job", () -> log.info("[inspection] orphan-job report count={}",
|
||||||
|
orphanJobInspector.inspectOrphanJobs(limit).entries().size()));
|
||||||
|
runSafely("task-missing-result", () -> log.info("[inspection] task-missing-result report count={}",
|
||||||
|
taskResultMissingInspector.inspectTasksMissingResult(limit).entries().size()));
|
||||||
|
runSafely("result-missing-file", () -> log.info("[inspection] result-missing-file report count={}",
|
||||||
|
resultFileMissingInspector.inspectResultsMissingFile(limit).entries().size()));
|
||||||
|
runSafely("terminal-active-job", () -> log.info("[inspection] terminal-active-job report count={}",
|
||||||
|
completedTaskActiveJobInspector.inspectTerminalTasksWithActiveJobs(limit).entries().size()));
|
||||||
|
}
|
||||||
|
|
||||||
|
private void runSafely(String name, Runnable action) {
|
||||||
|
try {
|
||||||
|
action.run();
|
||||||
|
} catch (Exception ex) {
|
||||||
|
log.warn("[inspection] {} failed msg={}", name, ex.getMessage(), ex);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
+188
@@ -0,0 +1,188 @@
|
|||||||
|
package com.nanri.aiimage.modules.task.service;
|
||||||
|
|
||||||
|
import com.nanri.aiimage.common.service.DistributedJobLockService;
|
||||||
|
import com.nanri.aiimage.config.InspectionProperties;
|
||||||
|
import com.nanri.aiimage.config.StorageProperties;
|
||||||
|
import com.nanri.aiimage.modules.file.service.TempOrphanInspector;
|
||||||
|
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.time.Duration;
|
||||||
|
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||||
|
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;
|
||||||
|
import static org.mockito.Mockito.verify;
|
||||||
|
import static org.mockito.Mockito.when;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* task-161:巡检任务开关契约(plan 09)。
|
||||||
|
* 默认 disabled 不执行;启用后执行并持分布式锁(双实例不重复);单个巡检异常
|
||||||
|
* 不阻断其余;limit 传入各巡检。
|
||||||
|
*/
|
||||||
|
@ExtendWith(MockitoExtension.class)
|
||||||
|
class InspectionSchedulerTest {
|
||||||
|
|
||||||
|
@Mock private DistributedJobLockService distributedJobLockService;
|
||||||
|
@Mock private StorageProperties storageProperties;
|
||||||
|
@Mock private TempOrphanInspector tempOrphanInspector;
|
||||||
|
@Mock private OrphanJobInspector orphanJobInspector;
|
||||||
|
@Mock private TaskResultMissingInspector taskResultMissingInspector;
|
||||||
|
@Mock private ResultFileMissingInspector resultFileMissingInspector;
|
||||||
|
@Mock private CompletedTaskActiveJobInspector completedTaskActiveJobInspector;
|
||||||
|
|
||||||
|
private InspectionProperties properties;
|
||||||
|
private InspectionScheduler scheduler;
|
||||||
|
|
||||||
|
@BeforeEach
|
||||||
|
void setUp() {
|
||||||
|
properties = new InspectionProperties();
|
||||||
|
scheduler = new InspectionScheduler(properties, distributedJobLockService, storageProperties,
|
||||||
|
tempOrphanInspector, orphanJobInspector, taskResultMissingInspector,
|
||||||
|
resultFileMissingInspector, completedTaskActiveJobInspector);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void disabledByDefaultDoesNotRun() {
|
||||||
|
assertFalse(properties.isEnabled(), "默认 disabled");
|
||||||
|
|
||||||
|
scheduler.runInspections();
|
||||||
|
|
||||||
|
verify(orphanJobInspector, never()).inspectOrphanJobs(anyInt());
|
||||||
|
verify(distributedJobLockService, never()).tryLock(any(), any());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void enabledRunsInspections() {
|
||||||
|
properties.setEnabled(true);
|
||||||
|
when(distributedJobLockService.tryLock(any(), any()))
|
||||||
|
.thenReturn(mock(DistributedJobLockService.LockHandle.class));
|
||||||
|
when(orphanJobInspector.inspectOrphanJobs(anyInt())).thenReturn(
|
||||||
|
new OrphanJobInspector.OrphanJobReport(java.util.List.of()));
|
||||||
|
when(taskResultMissingInspector.inspectTasksMissingResult(anyInt())).thenReturn(
|
||||||
|
new TaskResultMissingInspector.MissingResultReport(java.util.List.of()));
|
||||||
|
when(resultFileMissingInspector.inspectResultsMissingFile(anyInt())).thenReturn(
|
||||||
|
new ResultFileMissingInspector.MissingFileReport(java.util.List.of()));
|
||||||
|
when(completedTaskActiveJobInspector.inspectTerminalTasksWithActiveJobs(anyInt())).thenReturn(
|
||||||
|
new CompletedTaskActiveJobInspector.ActiveJobReport(java.util.List.of()));
|
||||||
|
when(storageProperties.getLocalTempDir()).thenReturn("target/nonexistent-inspection-dir");
|
||||||
|
|
||||||
|
scheduler.runInspections();
|
||||||
|
|
||||||
|
verify(orphanJobInspector).inspectOrphanJobs(200);
|
||||||
|
verify(taskResultMissingInspector).inspectTasksMissingResult(200);
|
||||||
|
verify(resultFileMissingInspector).inspectResultsMissingFile(200);
|
||||||
|
verify(completedTaskActiveJobInspector).inspectTerminalTasksWithActiveJobs(200);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void scheduleIntervalConfigured() {
|
||||||
|
assertEquals("0 0 3 * * *", properties.getCron(), "默认调度 cron");
|
||||||
|
properties.setCron("0 */30 * * * *");
|
||||||
|
assertEquals("0 */30 * * * *", properties.getCron(), "cron 可配置");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void dualInstanceSkipsWhenLockBusy() {
|
||||||
|
properties.setEnabled(true);
|
||||||
|
when(distributedJobLockService.tryLock(any(), any())).thenReturn(null);
|
||||||
|
|
||||||
|
scheduler.runInspections();
|
||||||
|
|
||||||
|
verify(orphanJobInspector, never()).inspectOrphanJobs(anyInt());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void oneInspectionFailureDoesNotBlockOthers() {
|
||||||
|
properties.setEnabled(true);
|
||||||
|
when(distributedJobLockService.tryLock(any(), any()))
|
||||||
|
.thenReturn(mock(DistributedJobLockService.LockHandle.class));
|
||||||
|
doThrow(new RuntimeException("orphan job inspector down"))
|
||||||
|
.when(orphanJobInspector).inspectOrphanJobs(anyInt());
|
||||||
|
when(taskResultMissingInspector.inspectTasksMissingResult(anyInt())).thenReturn(
|
||||||
|
new TaskResultMissingInspector.MissingResultReport(java.util.List.of()));
|
||||||
|
when(resultFileMissingInspector.inspectResultsMissingFile(anyInt())).thenReturn(
|
||||||
|
new ResultFileMissingInspector.MissingFileReport(java.util.List.of()));
|
||||||
|
when(completedTaskActiveJobInspector.inspectTerminalTasksWithActiveJobs(anyInt())).thenReturn(
|
||||||
|
new CompletedTaskActiveJobInspector.ActiveJobReport(java.util.List.of()));
|
||||||
|
when(storageProperties.getLocalTempDir()).thenReturn("target/nonexistent-inspection-dir");
|
||||||
|
|
||||||
|
scheduler.runInspections();
|
||||||
|
|
||||||
|
verify(taskResultMissingInspector).inspectTasksMissingResult(200);
|
||||||
|
verify(resultFileMissingInspector).inspectResultsMissingFile(200);
|
||||||
|
verify(completedTaskActiveJobInspector).inspectTerminalTasksWithActiveJobs(200);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void limitIsApplied() {
|
||||||
|
properties.setEnabled(true);
|
||||||
|
properties.setLimit(50);
|
||||||
|
when(distributedJobLockService.tryLock(any(), any()))
|
||||||
|
.thenReturn(mock(DistributedJobLockService.LockHandle.class));
|
||||||
|
when(orphanJobInspector.inspectOrphanJobs(anyInt())).thenReturn(
|
||||||
|
new OrphanJobInspector.OrphanJobReport(java.util.List.of()));
|
||||||
|
when(taskResultMissingInspector.inspectTasksMissingResult(anyInt())).thenReturn(
|
||||||
|
new TaskResultMissingInspector.MissingResultReport(java.util.List.of()));
|
||||||
|
when(resultFileMissingInspector.inspectResultsMissingFile(anyInt())).thenReturn(
|
||||||
|
new ResultFileMissingInspector.MissingFileReport(java.util.List.of()));
|
||||||
|
when(completedTaskActiveJobInspector.inspectTerminalTasksWithActiveJobs(anyInt())).thenReturn(
|
||||||
|
new CompletedTaskActiveJobInspector.ActiveJobReport(java.util.List.of()));
|
||||||
|
when(storageProperties.getLocalTempDir()).thenReturn("target/nonexistent-inspection-dir");
|
||||||
|
|
||||||
|
scheduler.runInspections();
|
||||||
|
|
||||||
|
verify(orphanJobInspector).inspectOrphanJobs(50);
|
||||||
|
verify(taskResultMissingInspector).inspectTasksMissingResult(50);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void lockAcquiredWithTtl() {
|
||||||
|
properties.setEnabled(true);
|
||||||
|
when(distributedJobLockService.tryLock("temp-file-inspection", Duration.ofMinutes(10)))
|
||||||
|
.thenReturn(mock(DistributedJobLockService.LockHandle.class));
|
||||||
|
when(orphanJobInspector.inspectOrphanJobs(anyInt())).thenReturn(
|
||||||
|
new OrphanJobInspector.OrphanJobReport(java.util.List.of()));
|
||||||
|
when(taskResultMissingInspector.inspectTasksMissingResult(anyInt())).thenReturn(
|
||||||
|
new TaskResultMissingInspector.MissingResultReport(java.util.List.of()));
|
||||||
|
when(resultFileMissingInspector.inspectResultsMissingFile(anyInt())).thenReturn(
|
||||||
|
new ResultFileMissingInspector.MissingFileReport(java.util.List.of()));
|
||||||
|
when(completedTaskActiveJobInspector.inspectTerminalTasksWithActiveJobs(anyInt())).thenReturn(
|
||||||
|
new CompletedTaskActiveJobInspector.ActiveJobReport(java.util.List.of()));
|
||||||
|
when(storageProperties.getLocalTempDir()).thenReturn("target/nonexistent-inspection-dir");
|
||||||
|
|
||||||
|
scheduler.runInspections();
|
||||||
|
|
||||||
|
verify(distributedJobLockService).tryLock("temp-file-inspection", Duration.ofMinutes(10));
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void rerunIsSafe() {
|
||||||
|
properties.setEnabled(true);
|
||||||
|
when(distributedJobLockService.tryLock(any(), any()))
|
||||||
|
.thenReturn(mock(DistributedJobLockService.LockHandle.class));
|
||||||
|
when(orphanJobInspector.inspectOrphanJobs(anyInt())).thenReturn(
|
||||||
|
new OrphanJobInspector.OrphanJobReport(java.util.List.of()));
|
||||||
|
when(taskResultMissingInspector.inspectTasksMissingResult(anyInt())).thenReturn(
|
||||||
|
new TaskResultMissingInspector.MissingResultReport(java.util.List.of()));
|
||||||
|
when(resultFileMissingInspector.inspectResultsMissingFile(anyInt())).thenReturn(
|
||||||
|
new ResultFileMissingInspector.MissingFileReport(java.util.List.of()));
|
||||||
|
when(completedTaskActiveJobInspector.inspectTerminalTasksWithActiveJobs(anyInt())).thenReturn(
|
||||||
|
new CompletedTaskActiveJobInspector.ActiveJobReport(java.util.List.of()));
|
||||||
|
when(storageProperties.getLocalTempDir()).thenReturn("target/nonexistent-inspection-dir");
|
||||||
|
|
||||||
|
scheduler.runInspections();
|
||||||
|
scheduler.runInspections();
|
||||||
|
|
||||||
|
verify(orphanJobInspector, org.mockito.Mockito.times(2)).inspectOrphanJobs(anyInt());
|
||||||
|
assertTrue(properties.isEnabled());
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user