diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/config/PersonalSchedulingConfig.java b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/config/PersonalSchedulingConfig.java index 9909f832..644a0c5e 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/config/PersonalSchedulingConfig.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/config/PersonalSchedulingConfig.java @@ -1,9 +1,22 @@ package org.dromara.aihr.personal.config; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Bean; import org.springframework.scheduling.annotation.EnableScheduling; +import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; @Configuration(proxyBeanMethods = false) @EnableScheduling public class PersonalSchedulingConfig { + + @Bean(name = "personalTaskScheduler") + public ThreadPoolTaskScheduler personalTaskScheduler() { + ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler(); + scheduler.setPoolSize(2); + scheduler.setThreadNamePrefix("personal-ingestion-"); + scheduler.setRemoveOnCancelPolicy(true); + scheduler.setWaitForTasksToCompleteOnShutdown(true); + scheduler.setAwaitTerminationSeconds(30); + return scheduler; + } } diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionService.java b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionService.java index 9342c8ff..5dbe0e37 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionService.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionService.java @@ -249,7 +249,8 @@ public class PersonalIngestionService { return activated; } - @Scheduled(fixedDelayString = "${aihr.personal.upload-cleanup-delay-ms:60000}") + @Scheduled(fixedDelayString = "${aihr.personal.upload-cleanup-delay-ms:60000}", + scheduler = "personalTaskScheduler") public void recoverStaleUploadIntents() { LocalDateTime now = LocalDateTime.now(); LocalDateTime pendingCutoff = now.minusMinutes(properties.getUploadCleanupAgeMinutes()); @@ -260,11 +261,18 @@ public class PersonalIngestionService { for (Map row : staleUploadRows("CLEANING", cleaningCutoff)) { finalizeStaleCleanup(intent(row), cleaningCutoff); } - LocalDateTime tombstoneCutoff = now.minusMinutes(properties.getUploadTombstoneRetentionMinutes()); - for (Map row : tombstoneRows()) { + LocalDateTime deleteIntervalCutoff = now.minusMinutes(properties.getTombstoneDeleteIntervalMinutes()); + long retentionCutoffEpoch = java.time.Instant.now() + .minusSeconds(properties.getUploadTombstoneRetentionMinutes() * 60L).getEpochSecond(); + for (Map row : tombstoneRows(deleteIntervalCutoff)) { UploadIntent intent = intent(row); - if (deleteKnownObject(intent) && uploadUpdatedAt(row).isBefore(tombstoneCutoff)) { + if (!deleteKnownObject(intent)) { + continue; + } + if (tombstonedAt(row) <= retentionCutoffEpoch) { phaseTransaction.executeWithoutResult(status -> deleteTombstoneMetadata(intent)); + } else { + phaseTransaction.executeWithoutResult(status -> touchTombstone(intent)); } } } @@ -282,15 +290,16 @@ public class PersonalIngestionService { and json_unquote(json_extract(o.ext1, '$.uploadState')) = ? and json_unquote(json_extract(o.ext1, '$.uploadToken')) is not null order by o.update_time - limit 20 - """, cutoff, state); + limit ? + """, cutoff, state, properties.getCleanupBatchSize()); } - private List> tombstoneRows() { + private List> tombstoneRows(LocalDateTime cutoff) { return jdbcTemplate.queryForList(""" select i.tenant_id, i.owner_user_id, i.space_id, i.id item_id, i.oss_id, i.size_bytes, i.mime_type, o.file_name, o.service, json_unquote(json_extract(o.ext1, '$.uploadToken')) upload_token, + json_unquote(json_extract(o.ext1, '$.tombstonedAt')) tombstoned_at, o.update_time upload_updated_at from sys_oss o join aihr_personal_item i on i.oss_id = o.oss_id @@ -298,8 +307,10 @@ public class PersonalIngestionService { where i.status = 'DELETED' and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'TOMBSTONE' and json_unquote(json_extract(o.ext1, '$.uploadToken')) is not null + and o.update_time < ? order by o.update_time - """); + limit ? + """, cutoff, properties.getCleanupBatchSize()); } private void beginCleanup(UploadIntent intent, LocalDateTime cutoff) { @@ -414,7 +425,8 @@ public class PersonalIngestionService { intent.sizeBytes()); int tombstoned = jdbcTemplate.update(""" update sys_oss - set ext1 = json_set(ext1, '$.uploadState', 'TOMBSTONE'), url = '', update_time = now() + set ext1 = json_set(ext1, '$.uploadState', 'TOMBSTONE', + '$.tombstonedAt', unix_timestamp(now())), url = '', update_time = now() where tenant_id = ? and oss_id = ? and create_by = ? and file_name = ? and json_unquote(json_extract(ext1, '$.itemId')) = ? and json_unquote(json_extract(ext1, '$.uploadToken')) = ? @@ -439,6 +451,19 @@ public class PersonalIngestionService { String.valueOf(intent.itemId()), intent.uploadToken(), intent.serviceKey()); } + private void touchTombstone(UploadIntent intent) { + jdbcTemplate.update(""" + update sys_oss + set update_time = now() + where tenant_id = ? and oss_id = ? and create_by = ? and file_name = ? + and json_unquote(json_extract(ext1, '$.itemId')) = ? + and json_unquote(json_extract(ext1, '$.uploadToken')) = ? + and json_unquote(json_extract(ext1, '$.uploadState')) = 'TOMBSTONE' + and service = ? + """, intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(), + String.valueOf(intent.itemId()), intent.uploadToken(), intent.serviceKey()); + } + private ItemCreatedResponse duplicate(PersonalOwner owner, long spaceId, String hash) { List> rows = jdbcTemplate.queryForList(""" select i.id, i.status @@ -491,11 +516,14 @@ public class PersonalIngestionService { return number.longValue(); } - private static LocalDateTime uploadUpdatedAt(Map row) { - Object value = row.get("upload_updated_at"); - if (value instanceof LocalDateTime dateTime) return dateTime; - if (value instanceof java.sql.Timestamp timestamp) return timestamp.toLocalDateTime(); - throw new ServiceException("PERSONAL_UPLOAD_CLEANUP_FAILED"); + private static long tombstonedAt(Map row) { + Object value = row.get("tombstoned_at"); + if (value instanceof Number number) return number.longValue(); + try { + return Long.parseLong(String.valueOf(value)); + } catch (NumberFormatException ex) { + throw new ServiceException("PERSONAL_UPLOAD_CLEANUP_FAILED"); + } } private void validateFile(MultipartFile file) { @@ -525,7 +553,11 @@ public class PersonalIngestionService { if (properties.getUploadCleanupAgeMinutes() < MIN_UPLOAD_CLEANUP_AGE_MINUTES || properties.getCleanupFinalizeGraceMinutes() <= 0 || properties.getUploadTombstoneRetentionMinutes() < 60 - || properties.getUploadTombstoneRetentionMinutes() <= cleanupWindow) { + || properties.getUploadTombstoneRetentionMinutes() <= cleanupWindow + || properties.getCleanupBatchSize() <= 0 + || properties.getTombstoneDeleteIntervalMinutes() <= 0 + || properties.getTombstoneDeleteIntervalMinutes() + >= properties.getUploadTombstoneRetentionMinutes()) { throw new IllegalArgumentException("invalid personal upload recovery windows"); } } diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionWorker.java b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionWorker.java index 17e64bf0..8461a3c9 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionWorker.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionWorker.java @@ -79,7 +79,8 @@ public class PersonalIngestionWorker { return defaultReader(ossService, clientProvider); } - @Scheduled(fixedDelayString = "${aihr.personal.ingestion-delay-ms:2000}") + @Scheduled(fixedDelayString = "${aihr.personal.ingestion-delay-ms:2000}", + scheduler = "personalTaskScheduler") public void poll() { recoverStaleParsing(); processNext(); diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/support/PersonalKnowledgeProperties.java b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/support/PersonalKnowledgeProperties.java index f5dd4fd5..db085111 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/support/PersonalKnowledgeProperties.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/support/PersonalKnowledgeProperties.java @@ -24,4 +24,6 @@ public class PersonalKnowledgeProperties { private int uploadCleanupAgeMinutes = 15; private int cleanupFinalizeGraceMinutes = 15; private int uploadTombstoneRetentionMinutes = 1440; + private int cleanupBatchSize = 20; + private int tombstoneDeleteIntervalMinutes = 10; } diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionServiceTest.java b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionServiceTest.java index e4f5740a..194705e3 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionServiceTest.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionServiceTest.java @@ -209,11 +209,14 @@ class PersonalIngestionServiceTest { void cleaningFinalizeCompensatesOnceAndRetainsTombstoneIntent() { Fixture fixture = fixture(); Map stale = staleIntent(); - when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("PENDING"))) + when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("PENDING"), + eq(20))) .thenReturn(List.of(stale), List.of()); - when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("CLEANING"))) + when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("CLEANING"), + eq(20))) .thenReturn(List.of(), List.of(stale)); - when(fixture.jdbc.queryForList(contains("uploadState')) = 'TOMBSTONE'"))).thenReturn(List.of()); + when(fixture.jdbc.queryForList(contains("uploadState')) = 'TOMBSTONE'"), any(LocalDateTime.class), + eq(20))).thenReturn(List.of()); when(fixture.jdbc.update(contains("json_set"), eq("CLEANING"), eq("000000"), eq(101L), eq(101L), eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("PENDING"), any(LocalDateTime.class))) .thenReturn(1); @@ -238,14 +241,20 @@ class PersonalIngestionServiceTest { @Test void tombstoneIsDeletedOnEveryScanAndMetadataRemovedOnlyAfterRetention() { Fixture fixture = fixture(); - Map fresh = staleIntent(LocalDateTime.now()); - Map expired = staleIntent(LocalDateTime.now().minusDays(2)); - when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("PENDING"))) + Map fresh = staleIntent(LocalDateTime.now(), java.time.Instant.now().getEpochSecond()); + Map expired = staleIntent( + LocalDateTime.now(), java.time.Instant.now().minusSeconds(172800).getEpochSecond()); + when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("PENDING"), + eq(20))) .thenReturn(List.of()); - when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("CLEANING"))) + when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("CLEANING"), + eq(20))) .thenReturn(List.of()); - when(fixture.jdbc.queryForList(contains("uploadState')) = 'TOMBSTONE'"))) + when(fixture.jdbc.queryForList(contains("uploadState')) = 'TOMBSTONE'"), any(LocalDateTime.class), + eq(20))) .thenReturn(List.of(fresh), List.of(fresh), List.of(expired)); + when(fixture.jdbc.update(contains("set update_time = now()"), eq("000000"), eq(101L), eq(101L), + eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("personal-private"))).thenReturn(1); when(fixture.jdbc.update(contains("delete from sys_oss"), eq("000000"), eq(101L), eq(101L), eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("personal-private"))).thenReturn(1); @@ -254,10 +263,17 @@ class PersonalIngestionServiceTest { fixture.service.recoverStaleUploadIntents(); verify(fixture.store, times(3)).deletePhysical("personal-private", "personal/key.txt"); + verify(fixture.jdbc, times(2)).update(contains("set update_time = now()"), eq("000000"), eq(101L), + eq(101L), eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("personal-private")); verify(fixture.jdbc, times(1)).update(contains("delete from sys_oss"), eq("000000"), eq(101L), eq(101L), eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("personal-private")); verify(fixture.jdbc, never()).update(contains("used_bytes = used_bytes -"), any(), any(), any(), any(), any()); + ArgumentCaptor tombstoneSql = ArgumentCaptor.forClass(String.class); + verify(fixture.jdbc, times(3)).queryForList(tombstoneSql.capture(), any(LocalDateTime.class), eq(20)); + assertTrue(tombstoneSql.getAllValues().stream().allMatch(sql -> sql.contains("o.update_time < ?"))); + assertTrue(tombstoneSql.getAllValues().stream().allMatch(sql -> sql.contains("order by o.update_time"))); + assertTrue(tombstoneSql.getAllValues().stream().allMatch(sql -> sql.contains("limit ?"))); } @Test @@ -415,6 +431,14 @@ class PersonalIngestionServiceTest { properties.setCleanupFinalizeGraceMinutes(1); properties.setUploadTombstoneRetentionMinutes(59); assertThrows(IllegalArgumentException.class, () -> fixture(properties)); + properties.setUploadTombstoneRetentionMinutes(1440); + properties.setCleanupBatchSize(0); + assertThrows(IllegalArgumentException.class, () -> fixture(properties)); + properties.setCleanupBatchSize(20); + properties.setTombstoneDeleteIntervalMinutes(0); + assertThrows(IllegalArgumentException.class, () -> fixture(properties)); + properties.setTombstoneDeleteIntervalMinutes(1440); + assertThrows(IllegalArgumentException.class, () -> fixture(properties)); } private static Fixture fixture() { @@ -458,16 +482,16 @@ class PersonalIngestionServiceTest { } private static Map staleIntent() { - return staleIntent(LocalDateTime.now()); + return staleIntent(LocalDateTime.now(), java.time.Instant.now().getEpochSecond()); } - private static Map staleIntent(LocalDateTime updatedAt) { + private static Map staleIntent(LocalDateTime updatedAt, long tombstonedAt) { return Map.ofEntries( Map.entry("tenant_id", "000000"), Map.entry("owner_user_id", 101L), Map.entry("space_id", 7L), Map.entry("item_id", 100L), Map.entry("oss_id", 101L), Map.entry("size_bytes", 4L), Map.entry("mime_type", "text/plain"), Map.entry("file_name", "personal/key.txt"), Map.entry("service", "personal-private"), Map.entry("upload_token", "upload-token"), - Map.entry("upload_updated_at", updatedAt) + Map.entry("upload_updated_at", updatedAt), Map.entry("tombstoned_at", String.valueOf(tombstonedAt)) ); } diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalSpaceServiceTest.java b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalSpaceServiceTest.java index 65f5fc7c..88138cc6 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalSpaceServiceTest.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalSpaceServiceTest.java @@ -1,6 +1,8 @@ package org.dromara.aihr.personal; import org.dromara.aihr.personal.config.PersonalSchedulingConfig; +import org.dromara.aihr.personal.service.PersonalIngestionService; +import org.dromara.aihr.personal.service.PersonalIngestionWorker; import org.dromara.aihr.personal.service.PersonalSpaceService; import org.dromara.aihr.personal.support.PersonalKnowledgeProperties; import org.dromara.aihr.personal.support.PersonalOwner; @@ -18,6 +20,9 @@ import org.springframework.transaction.support.AbstractPlatformTransactionManage import org.springframework.transaction.support.DefaultTransactionStatus; import org.springframework.transaction.support.TransactionSynchronizationManager; import org.springframework.scheduling.annotation.EnableScheduling; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; +import org.springframework.test.util.ReflectionTestUtils; import org.springframework.transaction.support.TransactionTemplate; import org.mockito.InOrder; @@ -57,11 +62,25 @@ class PersonalSpaceServiceTest { assertEquals(15, properties.getUploadCleanupAgeMinutes()); assertEquals(15, properties.getCleanupFinalizeGraceMinutes()); assertEquals(1440, properties.getUploadTombstoneRetentionMinutes()); + assertEquals(20, properties.getCleanupBatchSize()); + assertEquals(10, properties.getTombstoneDeleteIntervalMinutes()); } @Test - void personalSchedulingIsExplicitlyEnabled() { + void personalSchedulingUsesBoundedDedicatedScheduler() throws Exception { assertTrue(PersonalSchedulingConfig.class.isAnnotationPresent(EnableScheduling.class)); + ThreadPoolTaskScheduler scheduler = new PersonalSchedulingConfig().personalTaskScheduler(); + assertEquals(2, scheduler.getPoolSize()); + assertEquals("personal-ingestion-", scheduler.getThreadNamePrefix()); + assertTrue(scheduler.isRemoveOnCancelPolicy()); + assertEquals(true, ReflectionTestUtils.getField(scheduler, "waitForTasksToCompleteOnShutdown")); + assertEquals(30000L, ReflectionTestUtils.getField(scheduler, "awaitTerminationMillis")); + + Scheduled poll = PersonalIngestionWorker.class.getMethod("poll").getAnnotation(Scheduled.class); + Scheduled cleanup = PersonalIngestionService.class.getMethod("recoverStaleUploadIntents") + .getAnnotation(Scheduled.class); + assertEquals("personalTaskScheduler", poll.scheduler()); + assertEquals("personalTaskScheduler", cleanup.scheduler()); } @Test