From 035ec9ac7970133ded99ab8498bc3bfb4b90bdd9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=BB=84=E8=87=AA=E8=BE=BE?= <980324341@qq.com> Date: Wed, 2 Sep 2026 07:19:56 +0800 Subject: [PATCH] =?UTF-8?q?task-161:=20=E5=B7=A1=E6=A3=80=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E5=BC=80=E5=85=B3=EF=BC=88aiimage.inspection.enabled?= =?UTF-8?q?=20=E9=BB=98=E8=AE=A4=20disabled=E3=80=81cron=20=E8=B0=83?= =?UTF-8?q?=E5=BA=A6=E3=80=81=E5=88=86=E5=B8=83=E5=BC=8F=E9=94=81=E5=8F=8C?= =?UTF-8?q?=E5=AE=9E=E4=BE=8B=E5=8E=BB=E9=87=8D=E3=80=81=E5=8D=95=E5=B7=A1?= =?UTF-8?q?=E6=A3=80=E5=BC=82=E5=B8=B8=E9=9A=94=E7=A6=BB=E3=80=81limit=20?= =?UTF-8?q?=E9=80=8F=E4=BC=A0=EF=BC=89+=208=20=E6=9D=A1=E6=B5=8B=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../aiimage/config/InspectionProperties.java | 24 +++ .../task/service/InspectionScheduler.java | 81 ++++++++ .../task/service/InspectionSchedulerTest.java | 188 ++++++++++++++++++ 3 files changed, 293 insertions(+) create mode 100644 backend-java/src/main/java/com/nanri/aiimage/config/InspectionProperties.java create mode 100644 backend-java/src/main/java/com/nanri/aiimage/modules/task/service/InspectionScheduler.java create mode 100644 backend-java/src/test/java/com/nanri/aiimage/modules/task/service/InspectionSchedulerTest.java diff --git a/backend-java/src/main/java/com/nanri/aiimage/config/InspectionProperties.java b/backend-java/src/main/java/com/nanri/aiimage/config/InspectionProperties.java new file mode 100644 index 00000000..f46d5764 --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/config/InspectionProperties.java @@ -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; +} diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/InspectionScheduler.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/InspectionScheduler.java new file mode 100644 index 00000000..1de5a58a --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/InspectionScheduler.java @@ -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); + } + } +} diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/InspectionSchedulerTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/InspectionSchedulerTest.java new file mode 100644 index 00000000..f970dbc0 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/InspectionSchedulerTest.java @@ -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()); + } +}