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 new file mode 100644 index 00000000..9909f832 --- /dev/null +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/config/PersonalSchedulingConfig.java @@ -0,0 +1,9 @@ +package org.dromara.aihr.personal.config; + +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.annotation.EnableScheduling; + +@Configuration(proxyBeanMethods = false) +@EnableScheduling +public class PersonalSchedulingConfig { +} 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 df5b1d0c..b91f1e1e 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 @@ -13,14 +13,15 @@ import org.dromara.common.oss.core.OssClient; import org.dromara.common.oss.entity.UploadResult; import org.dromara.common.oss.enums.AccessPolicyType; import org.dromara.common.oss.factory.OssFactory; -import org.dromara.system.domain.vo.SysOssVo; -import org.dromara.system.service.ISysOssService; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.TransactionDefinition; +import org.springframework.transaction.annotation.Propagation; import org.springframework.transaction.annotation.Transactional; -import org.springframework.transaction.support.TransactionSynchronization; -import org.springframework.transaction.support.TransactionSynchronizationManager; +import org.springframework.transaction.support.TransactionTemplate; import org.springframework.web.multipart.MultipartFile; import java.io.ByteArrayInputStream; @@ -51,48 +52,47 @@ public class PersonalIngestionService { private final JdbcTemplate jdbcTemplate; private final PersonalSpaceService spaceService; private final PersonalKnowledgeProperties properties; - private final ISysOssService ossService; private final ObjectMapper objectMapper; private final PersonalObjectStore objectStore; - private final LongSupplier itemIdSupplier; + private final LongSupplier idSupplier; + private final TransactionTemplate phaseTransaction; @Autowired public PersonalIngestionService(JdbcTemplate jdbcTemplate, PersonalSpaceService spaceService, - PersonalKnowledgeProperties properties, ISysOssService ossService, - ObjectMapper objectMapper) { - this(jdbcTemplate, spaceService, properties, ossService, objectMapper, - new DefaultPersonalObjectStore(jdbcTemplate, properties, PersonalIngestionService::ossClient), - IdWorker::getId); + PersonalKnowledgeProperties properties, ObjectMapper objectMapper, + PlatformTransactionManager transactionManager) { + this(jdbcTemplate, spaceService, properties, objectMapper, + new DefaultPersonalObjectStore(properties, PersonalIngestionService::ossClient), + IdWorker::getId, requiresNew(transactionManager)); } private PersonalIngestionService(JdbcTemplate jdbcTemplate, PersonalSpaceService spaceService, - PersonalKnowledgeProperties properties, ISysOssService ossService, - ObjectMapper objectMapper, PersonalObjectStore objectStore, - LongSupplier itemIdSupplier) { + PersonalKnowledgeProperties properties, ObjectMapper objectMapper, + PersonalObjectStore objectStore, + LongSupplier idSupplier, TransactionTemplate phaseTransaction) { this.jdbcTemplate = jdbcTemplate; this.spaceService = spaceService; this.properties = properties; - this.ossService = ossService; this.objectMapper = objectMapper; this.objectStore = objectStore; - this.itemIdSupplier = itemIdSupplier; + this.idSupplier = idSupplier; + this.phaseTransaction = phaseTransaction; } public static PersonalIngestionService forTest(JdbcTemplate jdbcTemplate, PersonalSpaceService spaceService, - PersonalKnowledgeProperties properties, ISysOssService ossService, - ObjectMapper objectMapper, PersonalObjectStore objectStore, - LongSupplier itemIdSupplier) { - return new PersonalIngestionService(jdbcTemplate, spaceService, properties, ossService, objectMapper, - objectStore, itemIdSupplier); + PersonalKnowledgeProperties properties, ObjectMapper objectMapper, + PersonalObjectStore objectStore, + LongSupplier idSupplier, TransactionTemplate phaseTransaction) { + return new PersonalIngestionService(jdbcTemplate, spaceService, properties, objectMapper, + objectStore, idSupplier, phaseTransaction); } - public static PersonalObjectStore objectStoreForTest(JdbcTemplate jdbcTemplate, - PersonalKnowledgeProperties properties, + public static PersonalObjectStore objectStoreForTest(PersonalKnowledgeProperties properties, OssClientProvider clientProvider) { - return new DefaultPersonalObjectStore(jdbcTemplate, properties, clientProvider); + return new DefaultPersonalObjectStore(properties, clientProvider); } - @Transactional + @Transactional(propagation = Propagation.NOT_SUPPORTED) public ItemCreatedResponse createText(PersonalOwner owner, TextItemRequest request) { validateOwner(owner); if (request == null || request.content() == null || request.content().isBlank()) { @@ -104,7 +104,7 @@ public class PersonalIngestionService { request.capturedAt(), request.tags()); } - @Transactional + @Transactional(propagation = Propagation.NOT_SUPPORTED) public ItemCreatedResponse createFile(PersonalOwner owner, MultipartFile file, String title, LocalDateTime capturedAt) { validateOwner(owner); @@ -124,11 +124,14 @@ public class PersonalIngestionService { public void retry(PersonalOwner owner, long itemId) { validateOwner(owner); int updated = jdbcTemplate.update(""" - update aihr_personal_item - set status = 'QUEUED', error_code = null, error_message = null, - parsed_at = null, update_time = now() - where tenant_id = ? and owner_user_id = ? and id = ? - and status = 'FAILED' + update aihr_personal_item i + join sys_oss o on o.oss_id = i.oss_id and binary o.tenant_id = binary i.tenant_id + and o.create_by = i.owner_user_id + set i.status = 'QUEUED', i.error_code = null, i.error_message = null, + i.parsed_at = null, i.update_time = now() + where i.tenant_id = ? and i.owner_user_id = ? and i.id = ? + and i.status = 'FAILED' + and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'READY' """, owner.tenantId(), owner.userId(), itemId); if (updated == 0) { throw new ServiceException(ITEM_NOT_FOUND); @@ -137,88 +140,223 @@ public class PersonalIngestionService { private ItemCreatedResponse create(PersonalOwner owner, String sourceType, String title, String suffix, String mimeType, byte[] bytes, LocalDateTime capturedAt, List tags) { - String hash = sha256(bytes); - // reserve locks the current owner's space row. Dedupe must happen while that lock is held. - long spaceId = spaceService.reserve(owner, bytes.length); - ItemCreatedResponse duplicate = duplicate(owner, spaceId, hash); - if (duplicate != null) { - return duplicate; - } - - long itemId = positiveId(itemIdSupplier.getAsLong()); + String serviceKey = objectStore.requirePrivateService(); + long itemId = positiveId(idSupplier.getAsLong()); + long ossId = positiveId(idSupplier.getAsLong()); String objectKey = objectKey(owner, itemId, suffix); - SysOssVo uploaded = objectStore.upload(owner, itemId, objectKey, suffix, mimeType, bytes); - return persist(owner, spaceId, itemId, sourceType, title, mimeType, bytes.length, hash, capturedAt, - tags, objectKey, uploaded); + UploadIntent draft = new UploadIntent(owner, 0L, itemId, ossId, objectKey, suffix, mimeType, + bytes.length, serviceKey); + String hash = sha256(bytes); + + PhaseOne phaseOne = phaseTransaction.execute(status -> phaseOne( + draft, sourceType, title, hash, capturedAt, tags)); + if (phaseOne == null) { + throw new ServiceException("PERSONAL_ITEM_CREATE_FAILED"); + } + if (phaseOne.duplicate() != null) { + return phaseOne.duplicate(); + } + UploadIntent intent = phaseOne.intent(); + + try { + String url = objectStore.uploadPhysical(intent.serviceKey(), intent.objectKey(), mimeType, bytes); + UploadIntent finalIntent = intent; + phaseTransaction.executeWithoutResult(status -> activate(finalIntent, url)); + return new ItemCreatedResponse(intent.itemId(), "QUEUED", null); + } catch (RuntimeException ex) { + cleanupIntent(intent, false, null); + throw ex; + } } - private ItemCreatedResponse persist(PersonalOwner owner, long spaceId, long itemId, String sourceType, - String title, String mimeType, long size, String hash, - LocalDateTime capturedAt, List tags, String objectKey, - SysOssVo uploaded) { - if (uploaded == null || uploaded.getOssId() == null) { + private PhaseOne phaseOne(UploadIntent draft, String sourceType, String title, String hash, + LocalDateTime capturedAt, List tags) { + long spaceId = spaceService.reserve(draft.owner(), draft.sizeBytes()); + ItemCreatedResponse duplicate = duplicate(draft.owner(), spaceId, hash); + if (duplicate != null) { + return new PhaseOne(duplicate, null); + } + UploadIntent intent = draft.withSpaceId(spaceId); + String safeName = intent.objectKey().substring(intent.objectKey().lastIndexOf('/') + 1); + int ossInserted = jdbcTemplate.update(""" + insert into sys_oss + (oss_id, tenant_id, file_name, original_name, file_suffix, url, ext1, + create_time, create_by, update_time, update_by, service) + values (?, ?, ?, ?, ?, '', ?, now(), ?, now(), ?, ?) + """, intent.ossId(), intent.owner().tenantId(), intent.objectKey(), safeName, + "." + intent.suffix(), uploadExt(intent.itemId(), "PENDING"), intent.owner().userId(), + intent.owner().userId(), intent.serviceKey()); + int itemInserted = jdbcTemplate.update(""" + insert into aihr_personal_item + (id, tenant_id, space_id, owner_user_id, source_type, title, oss_id, mime_type, + size_bytes, content_hash, status, tags_json, captured_at, create_time, update_time) + values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'QUEUED', ?, ?, now(), now()) + """, intent.itemId(), intent.owner().tenantId(), spaceId, intent.owner().userId(), sourceType, + title, intent.ossId(), intent.mimeType(), intent.sizeBytes(), hash, tagsJson(tags), + capturedAt == null ? LocalDateTime.now() : capturedAt); + int counterUpdated = jdbcTemplate.update(""" + update aihr_personal_space + set used_bytes = used_bytes + ?, item_count = item_count + 1, update_time = now() + where tenant_id = ? and owner_user_id = ? and id = ? and status = 'ACTIVE' + """, intent.sizeBytes(), intent.owner().tenantId(), intent.owner().userId(), spaceId); + if (ossInserted != 1 || itemInserted != 1 || counterUpdated != 1) { + throw new ServiceException("PERSONAL_ITEM_CREATE_FAILED"); + } + return new PhaseOne(null, intent); + } + + private void activate(UploadIntent intent, String url) { + if (url == null || url.isBlank()) { throw new ServiceException("PERSONAL_OSS_UPLOAD_FAILED"); } - boolean deferredCleanup = registerRollbackCleanup(uploaded); + int activated = jdbcTemplate.update(""" + update sys_oss o + join aihr_personal_item i on i.oss_id = o.oss_id and binary i.tenant_id = binary o.tenant_id + and i.owner_user_id = o.create_by + set o.url = ?, o.ext1 = ?, o.update_time = now(), o.update_by = ? + where o.tenant_id = ? and o.oss_id = ? and o.create_by = ? and o.file_name = ? + and i.id = ? and i.status = 'QUEUED' + and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'PENDING' + """, url, uploadExt(intent.itemId(), "READY"), intent.owner().userId(), + intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(), intent.itemId()); + if (activated != 1) { + throw new ServiceException("PERSONAL_OSS_ACTIVATION_FAILED"); + } + } + + @Scheduled(fixedDelayString = "${aihr.personal.upload-cleanup-delay-ms:60000}") + public void recoverStaleUploadIntents() { + LocalDateTime cutoff = LocalDateTime.now().minusMinutes(properties.getUploadCleanupAgeMinutes()); + List> rows = 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 + from aihr_personal_item i + join sys_oss o on o.oss_id = i.oss_id and binary o.tenant_id = binary i.tenant_id + and o.create_by = i.owner_user_id + where i.status = 'QUEUED' and o.update_time < ? + and json_unquote(json_extract(o.ext1, '$.uploadState')) in ('PENDING', 'CLEANING') + order by o.update_time + limit 20 + """, cutoff); + for (Map row : rows) { + UploadIntent intent = intent(row); + cleanupIntent(intent, true, cutoff); + } + } + + private void cleanupIntent(UploadIntent intent, boolean includeCleaning, LocalDateTime cutoff) { + Integer claimed = phaseTransaction.execute(status -> claimCleanup(intent, includeCleaning, cutoff)); + if (claimed == null || claimed != 1) { + return; + } try { - int bound = jdbcTemplate.update(""" - update sys_oss - set ext1 = ?, update_time = now(), update_by = ? - where tenant_id = ? and oss_id = ? and create_by = ? and file_name = ? - """, personalOssExt(itemId), owner.userId(), owner.tenantId(), uploaded.getOssId(), - owner.userId(), objectKey); - if (bound != 1) { - throw new ServiceException("PERSONAL_OSS_BIND_FAILED"); - } - int inserted = jdbcTemplate.update(""" - insert into aihr_personal_item - (id, tenant_id, space_id, owner_user_id, source_type, title, oss_id, mime_type, - size_bytes, content_hash, status, tags_json, captured_at, create_time, update_time) - values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'QUEUED', ?, ?, now(), now()) - """, itemId, owner.tenantId(), spaceId, owner.userId(), sourceType, title, - uploaded.getOssId(), mimeType, size, hash, tagsJson(tags), - capturedAt == null ? LocalDateTime.now() : capturedAt); - if (inserted != 1) { - throw new ServiceException("PERSONAL_ITEM_CREATE_FAILED"); - } - int counterUpdated = jdbcTemplate.update(""" - update aihr_personal_space - set used_bytes = used_bytes + ?, item_count = item_count + 1, update_time = now() - where tenant_id = ? and owner_user_id = ? and id = ? and status = 'ACTIVE' - """, size, owner.tenantId(), owner.userId(), spaceId); - if (counterUpdated != 1) { - throw new ServiceException("PERSONAL_SPACE_NOT_AVAILABLE"); - } - return new ItemCreatedResponse(itemId, "QUEUED", null); + objectStore.deletePhysical(intent.serviceKey(), intent.objectKey()); } catch (RuntimeException ex) { - if (!deferredCleanup) { - cleanupOss(uploaded); - } - throw ex; + log.warn("Personal upload-intent physical cleanup failed itemId={}", intent.itemId()); + return; + } + try { + phaseTransaction.executeWithoutResult(status -> finalizeCleanup(intent)); + } catch (RuntimeException ex) { + log.warn("Personal upload-intent database cleanup failed itemId={}", intent.itemId()); + } + } + + private int claimCleanup(UploadIntent intent, boolean includeCleaning, LocalDateTime cutoff) { + String states = includeCleaning ? "('PENDING', 'CLEANING')" : "('PENDING')"; + String cutoffClause = cutoff == null ? "" : " and update_time < ?"; + String sql = """ + update sys_oss + set ext1 = json_set(ext1, '$.uploadState', 'CLEANING'), 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, '$.uploadState')) in %s%s + """.formatted(states, cutoffClause); + if (cutoff == null) { + return jdbcTemplate.update(sql, intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), + intent.objectKey(), String.valueOf(intent.itemId())); + } + return jdbcTemplate.update(sql, intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), + intent.objectKey(), String.valueOf(intent.itemId()), cutoff); + } + + private void finalizeCleanup(UploadIntent intent) { + long lockedSpace = spaceService.lockForUpdate(intent.owner()); + if (lockedSpace != intent.spaceId()) { + throw new ServiceException("PERSONAL_SPACE_NOT_AVAILABLE"); + } + int itemDeleted = jdbcTemplate.update(""" + update aihr_personal_item + set status = 'DELETED', deleted_at = now(), error_code = null, error_message = null, + update_time = now() + where tenant_id = ? and owner_user_id = ? and space_id = ? and id = ? and oss_id = ? + and status = 'QUEUED' and parsed_at is null + """, intent.owner().tenantId(), intent.owner().userId(), intent.spaceId(), intent.itemId(), + intent.ossId()); + int counterUpdated = jdbcTemplate.update(""" + update aihr_personal_space + set used_bytes = used_bytes - ?, item_count = item_count - 1, update_time = now() + where tenant_id = ? and owner_user_id = ? and id = ? + and used_bytes >= ? and item_count > 0 + """, intent.sizeBytes(), intent.owner().tenantId(), intent.owner().userId(), intent.spaceId(), + intent.sizeBytes()); + int ossDeleted = jdbcTemplate.update(""" + delete from sys_oss + 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, '$.uploadState')) = 'CLEANING' + """, intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(), + String.valueOf(intent.itemId())); + if (itemDeleted != 1 || counterUpdated != 1 || ossDeleted != 1) { + throw new ServiceException("PERSONAL_UPLOAD_CLEANUP_FAILED"); } } private ItemCreatedResponse duplicate(PersonalOwner owner, long spaceId, String hash) { List> rows = jdbcTemplate.queryForList(""" - select id, status - from aihr_personal_item + select id, status from aihr_personal_item where tenant_id = ? and owner_user_id = ? and space_id = ? and content_hash = ? - and status <> 'DELETED' - order by id desc - limit 1 + and status <> 'DELETED' order by id desc limit 1 """, owner.tenantId(), owner.userId(), spaceId, hash); - if (rows.isEmpty()) { - return null; + if (rows.isEmpty()) return null; + long id = ((Number) rows.get(0).get("id")).longValue(); + return new ItemCreatedResponse(id, String.valueOf(rows.get(0).get("status")), id); + } + + private String uploadExt(long itemId, String state) { + try { + return objectMapper.writeValueAsString(Map.of( + "source", "personal", "itemId", itemId, "uploadState", state)); + } catch (JsonProcessingException ex) { + throw new ServiceException("PERSONAL_OSS_BIND_FAILED"); } - Map row = rows.get(0); - long id = ((Number) row.get("id")).longValue(); - return new ItemCreatedResponse(id, String.valueOf(row.get("status")), id); + } + + private String tagsJson(List tags) { + List safe = tags == null ? List.of() : tags.stream().filter(t -> t != null && !t.isBlank()) + .map(String::trim).map(t -> t.length() > 50 ? t.substring(0, 50) : t).distinct().limit(20).toList(); + try { + return objectMapper.writeValueAsString(safe); + } catch (JsonProcessingException ex) { + throw new ServiceException("PERSONAL_TAGS_INVALID"); + } + } + + private static UploadIntent intent(Map row) { + PersonalOwner owner = new PersonalOwner(String.valueOf(row.get("tenant_id")), number(row, "owner_user_id"), null); + return new UploadIntent(owner, number(row, "space_id"), number(row, "item_id"), number(row, "oss_id"), + String.valueOf(row.get("file_name")), suffix(String.valueOf(row.get("file_name"))), + String.valueOf(row.get("mime_type")), number(row, "size_bytes"), String.valueOf(row.get("service"))); + } + + private static long number(Map row, String key) { + Object value = row.get(key); + if (!(value instanceof Number number)) throw new ServiceException("PERSONAL_UPLOAD_CLEANUP_FAILED"); + return number.longValue(); } private void validateFile(MultipartFile file) { - if (file == null || file.isEmpty() || file.getSize() <= 0) { - throw new ServiceException("PERSONAL_FILE_EMPTY"); - } + if (file == null || file.isEmpty() || file.getSize() <= 0) throw new ServiceException("PERSONAL_FILE_EMPTY"); validateSize(file.getSize()); if (!SUPPORTED_FILE_SUFFIXES.contains(suffix(file.getOriginalFilename()))) { throw new ServiceException("PERSONAL_FILE_UNSUPPORTED"); @@ -226,219 +364,105 @@ public class PersonalIngestionService { } private void validateSize(long bytes) { - long maxBytes; - try { - maxBytes = Math.multiplyExact(properties.getMaxFileSizeMb(), 1024L * 1024L); - } catch (ArithmeticException ex) { - throw new ServiceException("PERSONAL_FILE_TOO_LARGE"); - } - if (bytes <= 0 || maxBytes <= 0 || bytes > maxBytes) { - throw new ServiceException("PERSONAL_FILE_TOO_LARGE"); - } + long max; + try { max = Math.multiplyExact(properties.getMaxFileSizeMb(), 1024L * 1024L); } + catch (ArithmeticException ex) { throw new ServiceException("PERSONAL_FILE_TOO_LARGE"); } + if (bytes <= 0 || max <= 0 || bytes > max) throw new ServiceException("PERSONAL_FILE_TOO_LARGE"); } - private boolean registerRollbackCleanup(SysOssVo uploaded) { - if (!TransactionSynchronizationManager.isSynchronizationActive()) { - return false; - } - TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { - @Override - public void afterCompletion(int status) { - if (status != STATUS_COMMITTED) { - cleanupOss(uploaded); - } - } - }); - return true; - } - - private void cleanupOss(SysOssVo uploaded) { - try { - objectStore.deletePhysical(uploaded); - } catch (RuntimeException cleanupError) { - log.warn("Unable to clean personal OSS object id={}", uploaded.getOssId()); - } - try { - ossService.deleteWithValidByIds(List.of(uploaded.getOssId()), false); - } catch (RuntimeException cleanupError) { - log.warn("Unable to clean personal OSS metadata id={}", uploaded.getOssId()); - } - } - - private String personalOssExt(long itemId) { - try { - return objectMapper.writeValueAsString(Map.of("source", "personal", "itemId", itemId)); - } catch (JsonProcessingException ex) { - throw new ServiceException("PERSONAL_OSS_BIND_FAILED"); - } - } - - private String tagsJson(List tags) { - List safeTags = tags == null ? List.of() : tags.stream() - .filter(tag -> tag != null && !tag.isBlank()) - .map(String::trim) - .map(tag -> tag.length() > 50 ? tag.substring(0, 50) : tag) - .distinct() - .limit(20) - .toList(); - try { - return objectMapper.writeValueAsString(safeTags); - } catch (JsonProcessingException ex) { - throw new ServiceException("PERSONAL_TAGS_INVALID"); - } + private static TransactionTemplate requiresNew(PlatformTransactionManager manager) { + TransactionTemplate template = new TransactionTemplate(manager); + template.setPropagationBehavior(TransactionDefinition.PROPAGATION_REQUIRES_NEW); + return template; } private static String objectKey(PersonalOwner owner, long itemId, String suffix) { - validateOwner(owner); - long safeItemId = positiveId(itemId); + validateOwner(owner); positiveId(itemId); String safeSuffix = suffix == null ? "" : suffix.toLowerCase(Locale.ROOT); - if (!SUPPORTED_FILE_SUFFIXES.contains(safeSuffix)) { - throw new ServiceException("PERSONAL_FILE_UNSUPPORTED"); - } - String randomName = UUID.randomUUID().toString().replace("-", ""); - return "personal/" + owner.tenantId() + "/" + owner.userId() + "/" + safeItemId + "/" - + randomName + "." + safeSuffix; + if (!SUPPORTED_FILE_SUFFIXES.contains(safeSuffix)) throw new ServiceException("PERSONAL_FILE_UNSUPPORTED"); + return "personal/" + owner.tenantId() + "/" + owner.userId() + "/" + itemId + "/" + + UUID.randomUUID().toString().replace("-", "") + "." + safeSuffix; } private static void validateOwner(PersonalOwner owner) { if (owner == null || owner.userId() <= 0 || owner.tenantId() == null - || !SAFE_TENANT.matcher(owner.tenantId()).matches()) { - throw new ServiceException("PERSONAL_OWNER_INVALID"); - } + || !SAFE_TENANT.matcher(owner.tenantId()).matches()) throw new ServiceException("PERSONAL_OWNER_INVALID"); } private static long positiveId(long id) { - if (id <= 0) { - throw new ServiceException("PERSONAL_ID_INVALID"); - } + if (id <= 0) throw new ServiceException("PERSONAL_ID_INVALID"); return id; } private static String sha256(byte[] bytes) { - try { - return HexFormat.of().formatHex(MessageDigest.getInstance("SHA-256").digest(bytes)); - } catch (NoSuchAlgorithmException ex) { - throw new IllegalStateException("SHA-256 unavailable", ex); - } + try { return HexFormat.of().formatHex(MessageDigest.getInstance("SHA-256").digest(bytes)); } + catch (NoSuchAlgorithmException ex) { throw new IllegalStateException("SHA-256 unavailable", ex); } } private static String cleanTitle(String value, String fallback) { String title = value == null || value.isBlank() ? fallback : value.trim(); title = title.replace('\r', ' ').replace('\n', ' ').trim(); - if (title.isBlank()) { - title = "个人资料"; - } + if (title.isBlank()) title = "个人资料"; return title.length() > 500 ? title.substring(0, 500) : title; } private static String safeFileName(String value) { String name = value == null ? "personal-file" : value.replace('\\', '/'); - int slash = name.lastIndexOf('/'); - if (slash >= 0) { - name = name.substring(slash + 1); - } + int slash = name.lastIndexOf('/'); if (slash >= 0) name = name.substring(slash + 1); name = name.replace('\r', '_').replace('\n', '_').trim(); return name.isBlank() ? "personal-file" : name; } private static String suffix(String fileName) { - String safe = safeFileName(fileName); - int dot = safe.lastIndexOf('.'); + String safe = safeFileName(fileName); int dot = safe.lastIndexOf('.'); return dot < 0 ? "" : safe.substring(dot + 1).toLowerCase(Locale.ROOT); } private static String cleanMime(String value) { - if (value == null || value.isBlank()) { - return "application/octet-stream"; - } + if (value == null || value.isBlank()) return "application/octet-stream"; String mime = value.replace('\r', ' ').replace('\n', ' ').trim().toLowerCase(Locale.ROOT); - int separator = mime.indexOf(';'); - return separator < 0 ? mime : mime.substring(0, separator).trim(); + int separator = mime.indexOf(';'); return separator < 0 ? mime : mime.substring(0, separator).trim(); } public interface PersonalObjectStore { - SysOssVo upload(PersonalOwner owner, long itemId, String objectKey, String suffix, String mimeType, - byte[] bytes); - - void deletePhysical(SysOssVo uploaded); + String requirePrivateService(); + String uploadPhysical(String serviceKey, String objectKey, String mimeType, byte[] bytes); + void deletePhysical(String serviceKey, String objectKey); } - @FunctionalInterface - public interface OssClientProvider { - OssClient get(String configKey); - } + @FunctionalInterface public interface OssClientProvider { OssClient get(String configKey); } private static final class DefaultPersonalObjectStore implements PersonalObjectStore { - private final JdbcTemplate jdbcTemplate; - private final PersonalKnowledgeProperties properties; - private final OssClientProvider clientProvider; - - private DefaultPersonalObjectStore(JdbcTemplate jdbcTemplate, PersonalKnowledgeProperties properties, - OssClientProvider clientProvider) { - this.jdbcTemplate = jdbcTemplate; - this.properties = properties; - this.clientProvider = clientProvider; + private final PersonalKnowledgeProperties properties; private final OssClientProvider clients; + private DefaultPersonalObjectStore(PersonalKnowledgeProperties properties, OssClientProvider clients) { + this.properties = properties; this.clients = clients; } - - @Override - public SysOssVo upload(PersonalOwner owner, long itemId, String objectKey, String suffix, String mimeType, - byte[] bytes) { - OssClient storage = clientProvider.get(normalizedConfigKey(properties.getOssConfigKey())); - requirePrivate(storage); - UploadResult result = storage.upload( - new ByteArrayInputStream(bytes), objectKey, (long) bytes.length, mimeType); - long ossId = IdWorker.getId(); - String safeName = objectKey.substring(objectKey.lastIndexOf('/') + 1); - try { - int inserted = jdbcTemplate.update(""" - insert into sys_oss - (oss_id, tenant_id, file_name, original_name, file_suffix, url, ext1, - create_time, create_by, update_time, update_by, service) - values (?, ?, ?, ?, ?, ?, null, now(), ?, now(), ?, ?) - """, ossId, owner.tenantId(), result.getFilename(), safeName, "." + suffix, - result.getUrl(), owner.userId(), owner.userId(), storage.getConfigKey()); - if (inserted != 1) { - throw new ServiceException("PERSONAL_OSS_METADATA_FAILED"); - } - } catch (RuntimeException ex) { - try { - storage.delete(result.getFilename()); - } catch (RuntimeException cleanupError) { - log.warn("Unable to clean personal object after metadata failure"); - } - throw ex; - } - SysOssVo uploaded = new SysOssVo(); - uploaded.setOssId(ossId); - uploaded.setFileName(result.getFilename()); - uploaded.setOriginalName(safeName); - uploaded.setFileSuffix("." + suffix); - uploaded.setUrl(result.getUrl()); - uploaded.setService(storage.getConfigKey()); - return uploaded; + @Override public String requirePrivateService() { + OssClient storage = clients.get(normalize(properties.getOssConfigKey())); requirePrivate(storage); + return storage.getConfigKey(); } - - @Override - public void deletePhysical(SysOssVo uploaded) { - if (uploaded.getService() != null && uploaded.getFileName() != null) { - OssFactory.instance(uploaded.getService()).delete(uploaded.getFileName()); - } + @Override public String uploadPhysical(String serviceKey, String objectKey, String mimeType, byte[] bytes) { + OssClient storage = clients.get(serviceKey); requirePrivate(storage); + UploadResult result = storage.upload(new ByteArrayInputStream(bytes), objectKey, (long) bytes.length, mimeType); + return result.getUrl(); + } + @Override public void deletePhysical(String serviceKey, String objectKey) { + OssClient storage = clients.get(serviceKey); requirePrivate(storage); storage.delete(objectKey); } } - private static OssClient ossClient(String configKey) { - return configKey == null || configKey.isBlank() - ? OssFactory.instance() - : OssFactory.instance(configKey); - } - - private static String normalizedConfigKey(String value) { - return value == null ? "" : value.trim(); - } - + private static OssClient ossClient(String key) { return key == null || key.isBlank() ? OssFactory.instance() : OssFactory.instance(key); } + private static String normalize(String value) { return value == null ? "" : value.trim(); } private static void requirePrivate(OssClient storage) { - if (storage == null || storage.getAccessPolicy() != AccessPolicyType.PRIVATE) { + if (storage == null || storage.getAccessPolicy() != AccessPolicyType.PRIVATE) throw new ServiceException("PERSONAL_OSS_NOT_PRIVATE"); + } + + private record PhaseOne(ItemCreatedResponse duplicate, UploadIntent intent) {} + private record UploadIntent(PersonalOwner owner, long spaceId, long itemId, long ossId, String objectKey, + String suffix, String mimeType, long sizeBytes, String serviceKey) { + private UploadIntent withSpaceId(long value) { + return new UploadIntent(owner, value, itemId, ossId, objectKey, suffix, mimeType, sizeBytes, serviceKey); } } } 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 1ed9328e..e90d0c31 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 @@ -28,30 +28,42 @@ import java.util.Map; @Service public class PersonalIngestionWorker { - private static final int CHUNK_SIZE = 800; - private static final int CHUNK_OVERLAP = 100; - private final JdbcTemplate jdbcTemplate; private final KnowledgeDocumentParser parser; private final TransactionTemplate transactionTemplate; private final StoredObjectReader objectReader; private final long maxInputBytes; + private final int chunkSize; + private final int chunkOverlap; + private final int parsingLeaseMinutes; + private final int maxParseAttempts; public PersonalIngestionWorker(JdbcTemplate jdbcTemplate, ISysOssService ossService, KnowledgeDocumentParser parser, PersonalKnowledgeProperties properties, PlatformTransactionManager transactionManager) { this(jdbcTemplate, parser, new TransactionTemplate(transactionManager), - defaultReader(ossService, PersonalIngestionWorker::ossClient), configuredMaxBytes(properties)); + defaultReader(ossService, PersonalIngestionWorker::ossClient), configuredMaxBytes(properties), + properties.getChunkSize(), properties.getChunkOverlap(), properties.getParsingLeaseMinutes(), + properties.getMaxParseAttempts()); } private PersonalIngestionWorker(JdbcTemplate jdbcTemplate, KnowledgeDocumentParser parser, TransactionTemplate transactionTemplate, StoredObjectReader objectReader, - long maxInputBytes) { + long maxInputBytes, int chunkSize, int chunkOverlap, + int parsingLeaseMinutes, int maxParseAttempts) { this.jdbcTemplate = jdbcTemplate; this.parser = parser; this.transactionTemplate = transactionTemplate; this.objectReader = objectReader; this.maxInputBytes = maxInputBytes; + if (chunkSize <= 0 || chunkOverlap < 0 || chunkOverlap >= chunkSize + || parsingLeaseMinutes <= 0 || maxParseAttempts <= 0) { + throw new IllegalArgumentException("invalid personal ingestion worker settings"); + } + this.chunkSize = chunkSize; + this.chunkOverlap = chunkOverlap; + this.parsingLeaseMinutes = parsingLeaseMinutes; + this.maxParseAttempts = maxParseAttempts; } public static PersonalIngestionWorker forTest(JdbcTemplate jdbcTemplate, ISysOssService ossService, @@ -59,7 +71,7 @@ public class PersonalIngestionWorker { TransactionTemplate transactionTemplate, StoredObjectReader objectReader) { return new PersonalIngestionWorker(jdbcTemplate, parser, transactionTemplate, objectReader, - 20L * 1024 * 1024); + 20L * 1024 * 1024, 800, 120, 15, 3); } public static StoredObjectReader objectReaderForTest(ISysOssService ossService, @@ -69,16 +81,41 @@ public class PersonalIngestionWorker { @Scheduled(fixedDelayString = "${aihr.personal.ingestion-delay-ms:2000}") public void poll() { + recoverStaleParsing(); processNext(); } + public void recoverStaleParsing() { + java.time.LocalDateTime cutoff = java.time.LocalDateTime.now().minusMinutes(parsingLeaseMinutes); + jdbcTemplate.update(""" + update aihr_personal_item i + join sys_oss o on o.oss_id = i.oss_id and binary o.tenant_id = binary i.tenant_id + and o.create_by = i.owner_user_id + set i.status = 'FAILED', i.error_code = 'PERSONAL_PARSE_RETRY_EXHAUSTED', + i.error_message = '资料处理重试次数已用尽', i.update_time = now() + where i.status = 'PARSING' and i.attempt_count >= ? and i.update_time < ? + and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'READY' + """, maxParseAttempts, cutoff); + jdbcTemplate.update(""" + update aihr_personal_item i + join sys_oss o on o.oss_id = i.oss_id and binary o.tenant_id = binary i.tenant_id + and o.create_by = i.owner_user_id + set i.status = 'QUEUED', i.error_code = null, i.error_message = null, i.update_time = now() + where i.status = 'PARSING' and i.attempt_count < ? and i.update_time < ? + and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'READY' + """, maxParseAttempts, cutoff); + } + public boolean processNext() { List> queued = jdbcTemplate.queryForList(""" - select id, tenant_id, space_id, owner_user_id, source_type, title, - oss_id, mime_type, tags_json, captured_at - from aihr_personal_item - where status = 'QUEUED' - order by id + select i.id, i.tenant_id, i.space_id, i.owner_user_id, i.source_type, i.title, + i.oss_id, i.mime_type, i.tags_json, i.captured_at + from aihr_personal_item i + join sys_oss o on o.oss_id = i.oss_id and binary o.tenant_id = binary i.tenant_id + and o.create_by = i.owner_user_id + where i.status = 'QUEUED' + and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'READY' + order by i.id limit 1 """); if (queued.isEmpty()) { @@ -91,6 +128,10 @@ public class PersonalIngestionWorker { error_code = null, error_message = null, update_time = now() where tenant_id = ? and owner_user_id = ? and id = ? and status = 'QUEUED' + and exists (select 1 from sys_oss o where o.oss_id = aihr_personal_item.oss_id + and binary o.tenant_id = binary aihr_personal_item.tenant_id + and o.create_by = aihr_personal_item.owner_user_id + and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'READY') """, item.tenantId(), item.ownerUserId(), item.id()); if (claimed != 1) { return false; @@ -100,7 +141,7 @@ public class PersonalIngestionWorker { StoredObject stored = objectReader.read( item.ossId(), ownerObjectPrefix(item), item.ownerUserId(), maxInputBytes); ParsedDocument document = parser.parse(stored.fileName(), item.mimeType(), stored.bytes()); - List chunks = document.chunks(CHUNK_SIZE, CHUNK_OVERLAP); + List chunks = document.chunks(chunkSize, chunkOverlap); if (chunks.isEmpty()) { throw new KnowledgeDocumentParser.ParseException( KnowledgeDocumentParser.Failure.EMPTY, "document contains no text"); diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalSpaceService.java b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalSpaceService.java index ce13846f..f15e214d 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalSpaceService.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalSpaceService.java @@ -66,6 +66,12 @@ public class PersonalSpaceService { return spaceId; } + /** Locks an existing owner space for compensating updates without applying quota admission rules. */ + @Transactional(propagation = Propagation.MANDATORY) + public long lockForUpdate(PersonalOwner owner) { + return ((Number) lockSpace(owner).get("id")).longValue(); + } + private Map ensureAndLockSpace(PersonalOwner owner) { jdbcTemplate.update(""" insert into aihr_personal_space 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 59ba4a7e..81e7720a 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 @@ -17,4 +17,9 @@ public class PersonalKnowledgeProperties { private String qdrantCollection = "aihr_personal_knowledge"; /** Optional sys_oss_config key. Blank selects the system default client. */ private String ossConfigKey = ""; + private int chunkSize = 800; + private int chunkOverlap = 120; + private int parsingLeaseMinutes = 15; + private int maxParseAttempts = 3; + private int uploadCleanupAgeMinutes = 15; } 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 ab8bd1d9..f1c41628 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 @@ -13,21 +13,18 @@ import org.dromara.common.core.exception.ServiceException; import org.dromara.common.oss.core.OssClient; import org.dromara.common.oss.entity.UploadResult; import org.dromara.common.oss.enums.AccessPolicyType; -import org.dromara.system.domain.vo.SysOssVo; -import org.dromara.system.service.ISysOssService; import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; import org.mockito.ArgumentCaptor; -import org.mockito.InOrder; -import org.springframework.aop.framework.ProxyFactory; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.mock.web.MockMultipartFile; import org.springframework.transaction.TransactionDefinition; -import org.springframework.transaction.annotation.AnnotationTransactionAttributeSource; -import org.springframework.transaction.interceptor.TransactionInterceptor; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.support.AbstractPlatformTransactionManager; import org.springframework.transaction.support.DefaultTransactionStatus; import org.springframework.transaction.support.TransactionSynchronizationManager; +import org.springframework.transaction.support.TransactionTemplate; import java.nio.charset.StandardCharsets; import java.time.LocalDateTime; @@ -37,6 +34,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; @@ -44,7 +42,6 @@ import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.contains; import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.Mockito.inOrder; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; @@ -58,308 +55,255 @@ class PersonalIngestionServiceTest { private static final PersonalOwner OWNER = new PersonalOwner("000000", 101L, "ext-101"); @Test - void failedItemRetriesThroughOwnerScopedQueuedState() { - Fixture fixture = fixture(100L); - when(fixture.jdbc.update(contains("status = 'QUEUED'"), eq("000000"), eq(101L), eq(9L))) - .thenReturn(1); - - fixture.service.retry(OWNER, 9L); - - verify(fixture.jdbc).update(contains("status = 'QUEUED'"), eq("000000"), eq(101L), eq(9L)); - } - - @Test - void retryDoesNotDiscloseMissingOrForeignItem() { - Fixture fixture = fixture(100L); - - ServiceException error = assertThrows(ServiceException.class, () -> fixture.service.retry(OWNER, 9L)); - - assertEquals("PERSONAL_ITEM_NOT_FOUND", error.getMessage()); - } - - @Test - void invalidFilesAreRejectedBeforeObjectStorage() { - Fixture fixture = fixture(100L); - MockMultipartFile large = new MockMultipartFile( - "file", "large.pdf", "application/pdf", new byte[21 * 1024 * 1024]); - - assertEquals("PERSONAL_FILE_TOO_LARGE", assertThrows(ServiceException.class, - () -> fixture.service.createFile(OWNER, large, null, null)).getMessage()); - assertThrows(ServiceException.class, () -> fixture.service.createFile( - OWNER, new MockMultipartFile("file", "script.exe", "application/octet-stream", new byte[]{1}), - null, null)); - assertThrows(ServiceException.class, () -> fixture.service.createFile( - OWNER, new MockMultipartFile("file", "empty.txt", "text/plain", new byte[0]), null, null)); - - verifyNoInteractions(fixture.store, fixture.oss, fixture.jdbc, fixture.spaces); - } - - @Test - void textCreationLocksBeforeDedupeAndUsesIsolatedObjectKeyAndExtBinding() throws Exception { - Fixture fixture = fixture(100L); - stubSuccessfulCreate(fixture, OWNER, 7L, 18L, 81L); + void successfulCreatePersistsPendingUploadsOutsideTransactionAndActivatesReady() throws Exception { + Fixture fixture = fixture(); + stubPhaseOne(fixture, 18L); + AtomicBoolean uploadInTransaction = new AtomicBoolean(true); + ArgumentCaptor key = ArgumentCaptor.forClass(String.class); + when(fixture.store.uploadPhysical(eq("personal-private"), key.capture(), eq("text/plain"), any(byte[].class))) + .thenAnswer(invocation -> { + uploadInTransaction.set(TransactionSynchronizationManager.isActualTransactionActive()); + return "https://private.invalid/" + invocation.getArgument(1); + }); + when(fixture.jdbc.update(contains("set o.url ="), anyString(), anyString(), eq(101L), eq("000000"), + eq(101L), eq(101L), anyString(), eq(100L))).thenReturn(1); ItemCreatedResponse response = fixture.service.createText(OWNER, - new TextItemRequest("周报", "保洁巡检记录", LocalDateTime.of(2026, 7, 12, 9, 0), List.of("保洁"))); + new TextItemRequest("周报", "保洁巡检记录", null, List.of("保洁"))); assertEquals(100L, response.itemId()); assertEquals("QUEUED", response.status()); - InOrder lockOrder = inOrder(fixture.spaces, fixture.jdbc); - lockOrder.verify(fixture.spaces).reserve(OWNER, 18L); - lockOrder.verify(fixture.jdbc).queryForList(contains("content_hash"), - eq("000000"), eq(101L), eq(7L), anyString()); - - ArgumentCaptor key = ArgumentCaptor.forClass(String.class); - verify(fixture.store).upload(eq(OWNER), eq(100L), key.capture(), eq("txt"), eq("text/plain"), - any(byte[].class)); + assertFalse(uploadInTransaction.get()); assertTrue(key.getValue().matches("personal/000000/101/100/[0-9a-f]{32}\\.txt")); - - ArgumentCaptor ext = ArgumentCaptor.forClass(String.class); - verify(fixture.jdbc).update(contains("update sys_oss"), ext.capture(), eq(101L), eq("000000"), - eq(81L), eq(101L), eq(key.getValue())); - JsonNode extJson = new ObjectMapper().readTree(ext.getValue()); - assertEquals("personal", extJson.path("source").asText()); - assertEquals(100L, extJson.path("itemId").asLong()); - verify(fixture.jdbc).update(contains("insert into aihr_personal_item"), eq(100L), eq("000000"), - eq(7L), eq(101L), eq("TEXT"), eq("周报"), eq(81L), eq("text/plain"), eq(18L), - anyString(), anyString(), any(LocalDateTime.class)); - verify(fixture.jdbc).update(contains("used_bytes = used_bytes +"), - eq(18L), eq("000000"), eq(101L), eq(7L)); + ArgumentCaptor pendingExt = ArgumentCaptor.forClass(String.class); + ArgumentCaptor originalName = ArgumentCaptor.forClass(String.class); + verify(fixture.jdbc).update(contains("insert into sys_oss"), eq(101L), eq("000000"), + eq(key.getValue()), originalName.capture(), eq(".txt"), pendingExt.capture(), eq(101L), eq(101L), + eq("personal-private")); + assertTrue(originalName.getValue().matches("[0-9a-f]{32}\\.txt")); + assertUploadState(pendingExt.getValue(), "PENDING"); + ArgumentCaptor readyExt = ArgumentCaptor.forClass(String.class); + verify(fixture.jdbc).update(contains("set o.url ="), anyString(), readyExt.capture(), eq(101L), + eq("000000"), eq(101L), eq(101L), eq(key.getValue()), eq(100L)); + assertUploadState(readyExt.getValue(), "READY"); + assertEquals(2, fixture.transactions.commits); + assertEquals(0, fixture.transactions.rollbacks); } @Test - void sameOwnerDuplicateIsCheckedUnderLockAndDoesNotUploadOrIncrementCounters() { - Fixture fixture = fixture(100L); + void originalFileNameIsVisibleOnlyOnPersonalItemNotSystemObjectMetadata() { + Fixture fixture = fixture(); + stubPhaseOne(fixture, 6L); + when(fixture.store.uploadPhysical(eq("personal-private"), anyString(), eq("text/plain"), any(byte[].class))) + .thenReturn("https://private.invalid/object"); + when(fixture.jdbc.update(contains("set o.url ="), anyString(), anyString(), eq(101L), eq("000000"), + eq(101L), eq(101L), anyString(), eq(100L))).thenReturn(1); + MockMultipartFile file = new MockMultipartFile("file", "13800138000-secret.txt", "text/plain", + "secret".getBytes(StandardCharsets.UTF_8)); + + fixture.service.createFile(OWNER, file, null, null); + + ArgumentCaptor safeObjectName = ArgumentCaptor.forClass(String.class); + verify(fixture.jdbc).update(contains("insert into sys_oss"), eq(101L), eq("000000"), anyString(), + safeObjectName.capture(), eq(".txt"), anyString(), eq(101L), eq(101L), eq("personal-private")); + assertFalse(safeObjectName.getValue().contains("13800138000")); + verify(fixture.jdbc).update(contains("insert into aihr_personal_item"), eq(100L), eq("000000"), eq(7L), + eq(101L), eq("FILE"), eq("13800138000-secret.txt"), eq(101L), eq("text/plain"), eq(6L), + anyString(), anyString(), any(LocalDateTime.class)); + } + + @Test + void duplicateIsResolvedUnderOwnerLockWithoutCreatingUploadIntent() { + Fixture fixture = fixture(); when(fixture.spaces.reserve(OWNER, 4L)).thenReturn(7L); - when(fixture.jdbc.queryForList(contains("content_hash"), - eq("000000"), eq(101L), eq(7L), anyString())) + when(fixture.jdbc.queryForList(contains("content_hash"), eq("000000"), eq(101L), eq(7L), anyString())) .thenReturn(List.of(Map.of("id", 77L, "status", "READY"))); - MockMultipartFile file = new MockMultipartFile( - "file", "notes.txt", "text/plain", "same".getBytes(StandardCharsets.UTF_8)); + MockMultipartFile file = new MockMultipartFile("file", "notes.txt", "text/plain", new byte[]{1, 2, 3, 4}); ItemCreatedResponse response = fixture.service.createFile(OWNER, file, null, null); assertEquals(77L, response.itemId()); assertEquals(77L, response.duplicateOf()); - InOrder order = inOrder(fixture.spaces, fixture.jdbc); - order.verify(fixture.spaces).reserve(OWNER, 4L); - order.verify(fixture.jdbc).queryForList(contains("content_hash"), - eq("000000"), eq(101L), eq(7L), anyString()); - verifyNoInteractions(fixture.store, fixture.oss); - verify(fixture.jdbc, never()).update(contains("used_bytes = used_bytes +"), + verify(fixture.store).requirePrivateService(); + verify(fixture.store, never()).uploadPhysical(anyString(), anyString(), anyString(), any(byte[].class)); + verify(fixture.jdbc, never()).update(contains("insert into sys_oss"), any(), any(), any(), any(), any(), any(), any(), any(), any()); + assertEquals(1, fixture.transactions.commits); } @Test - void sameContentAcrossOwnersCreatesSeparateObjects() { - AtomicLong ids = new AtomicLong(100L); - Fixture fixture = fixture(ids::getAndIncrement); - PersonalOwner other = new PersonalOwner("000000", 202L, "ext-202"); - when(fixture.spaces.reserve(OWNER, 4L)).thenReturn(7L); - when(fixture.spaces.reserve(other, 4L)).thenReturn(8L); - when(fixture.jdbc.queryForList(contains("content_hash"), any(), any(), any(), any())) - .thenReturn(List.of()); - when(fixture.store.upload(any(), anyLong(), anyString(), eq("txt"), eq("text/plain"), - any(byte[].class))).thenAnswer(invocation -> oss(80L + invocation.getArgument(1), - invocation.getArgument(2))); - stubPersistence(fixture.jdbc); - MockMultipartFile file = new MockMultipartFile( - "file", "notes.txt", "text/plain", "same".getBytes(StandardCharsets.UTF_8)); - - ItemCreatedResponse first = fixture.service.createFile(OWNER, file, null, null); - ItemCreatedResponse second = fixture.service.createFile(other, file, null, null); - - assertEquals(100L, first.itemId()); - assertEquals(101L, second.itemId()); - verify(fixture.jdbc).queryForList(contains("content_hash"), - eq("000000"), eq(101L), eq(7L), anyString()); - verify(fixture.jdbc).queryForList(contains("content_hash"), - eq("000000"), eq(202L), eq(8L), anyString()); - verify(fixture.store, times(2)).upload(any(), anyLong(), anyString(), eq("txt"), eq("text/plain"), - any(byte[].class)); - } - - @Test - void failedItemPersistenceRollsBackAndCleansPhysicalObject() { - Fixture fixture = fixture(100L); - when(fixture.spaces.reserve(OWNER, 4L)).thenReturn(7L); - when(fixture.jdbc.queryForList(contains("content_hash"), - eq("000000"), eq(101L), eq(7L), anyString())).thenReturn(List.of()); - SysOssVo uploaded = oss(81L, "personal/000000/101/100/a.txt"); - when(fixture.store.upload(any(), anyLong(), anyString(), anyString(), anyString(), any(byte[].class))) - .thenReturn(uploaded); - when(fixture.jdbc.update(contains("update sys_oss"), any(), any(), any(), any(), any(), any())) + void uploadFailureClaimsCleaningDeletesPhysicalAndReversesCountersOnce() { + Fixture fixture = fixture(); + stubPhaseOne(fixture, 4L); + AtomicBoolean deleteInTransaction = new AtomicBoolean(true); + when(fixture.store.uploadPhysical(eq("personal-private"), anyString(), eq("text/plain"), any(byte[].class))) + .thenThrow(new ServiceException("PERSONAL_OSS_UPLOAD_FAILED")); + when(fixture.jdbc.update(contains("json_set"), eq("000000"), eq(101L), eq(101L), anyString(), eq("100"))) .thenReturn(1); - when(fixture.jdbc.update(contains("insert into aihr_personal_item"), - any(), any(), any(), any(), any(), any(), any(), any(), any(), any(), any(), any())) - .thenThrow(new IllegalStateException("db failed")); - TestTransactionManager transactions = new TestTransactionManager(); - PersonalIngestionService proxy = transactionalProxy(fixture.service, transactions); - MockMultipartFile file = new MockMultipartFile( - "file", "notes.txt", "text/plain", "same".getBytes(StandardCharsets.UTF_8)); + when(fixture.spaces.lockForUpdate(any(PersonalOwner.class))).thenReturn(7L); + stubFinalizeCleanup(fixture); + org.mockito.Mockito.doAnswer(invocation -> { + deleteInTransaction.set(TransactionSynchronizationManager.isActualTransactionActive()); + return null; + }).when(fixture.store).deletePhysical(eq("personal-private"), anyString()); - assertThrows(IllegalStateException.class, () -> proxy.createFile(OWNER, file, null, null)); + ServiceException error = assertThrows(ServiceException.class, () -> fixture.service.createFile(OWNER, + new MockMultipartFile("file", "notes.txt", "text/plain", new byte[]{1, 2, 3, 4}), null, null)); - assertEquals(1, transactions.rollbacks); - verify(fixture.store).deletePhysical(uploaded); - verify(fixture.oss).deleteWithValidByIds(List.of(81L), false); + assertEquals("PERSONAL_OSS_UPLOAD_FAILED", error.getMessage()); + assertFalse(deleteInTransaction.get()); + verify(fixture.jdbc).update(contains("status = 'DELETED'"), eq("000000"), eq(101L), eq(7L), eq(100L), + eq(101L)); + verify(fixture.jdbc).update(contains("used_bytes = used_bytes -"), eq(4L), eq("000000"), eq(101L), + eq(7L), eq(4L)); + verify(fixture.jdbc).update(contains("delete from sys_oss"), eq("000000"), eq(101L), eq(101L), + anyString(), eq("100")); + assertEquals(3, fixture.transactions.commits); + assertEquals(0, fixture.transactions.rollbacks); } @Test - void reservationAndCountersShareOneOuterTransaction() { - Fixture fixture = fixture(100L); - AtomicBoolean reserveInTransaction = new AtomicBoolean(); - AtomicBoolean counterInTransaction = new AtomicBoolean(); - when(fixture.spaces.reserve(OWNER, 4L)).thenAnswer(invocation -> { - reserveInTransaction.set(TransactionSynchronizationManager.isActualTransactionActive()); - return 7L; - }); - when(fixture.jdbc.queryForList(contains("content_hash"), any(), any(), any(), any())) - .thenReturn(List.of()); - when(fixture.store.upload(any(), anyLong(), anyString(), anyString(), anyString(), any(byte[].class))) - .thenReturn(oss(81L, "personal/000000/101/100/a.txt")); - when(fixture.jdbc.update(contains("update sys_oss"), any(), any(), any(), any(), any(), any())) - .thenReturn(1); - when(fixture.jdbc.update(contains("insert into aihr_personal_item"), - any(), any(), any(), any(), any(), any(), any(), any(), any(), any(), any(), any())) - .thenReturn(1); - when(fixture.jdbc.update(contains("used_bytes = used_bytes +"), any(), any(), any(), any())) - .thenAnswer(invocation -> { - counterInTransaction.set(TransactionSynchronizationManager.isActualTransactionActive()); - return 1; - }); - TestTransactionManager transactions = new TestTransactionManager(); - PersonalIngestionService proxy = transactionalProxy(fixture.service, transactions); - MockMultipartFile file = new MockMultipartFile( - "file", "notes.txt", "text/plain", "same".getBytes(StandardCharsets.UTF_8)); + void staleCleanupIsIdempotentAndDoesNotDecrementCountersTwice() { + Fixture fixture = fixture(); + Map stale = staleIntent(); + when(fixture.jdbc.queryForList(contains("uploadState')) in ('PENDING', 'CLEANING')"), + any(LocalDateTime.class))).thenReturn(List.of(stale)); + when(fixture.jdbc.update(contains("json_set"), eq("000000"), eq(101L), eq(101L), eq("personal/key.txt"), + eq("100"), any(LocalDateTime.class))).thenReturn(1, 0); + when(fixture.spaces.lockForUpdate(any(PersonalOwner.class))).thenReturn(7L); + stubFinalizeCleanup(fixture); - proxy.createFile(OWNER, file, null, null); + fixture.service.recoverStaleUploadIntents(); + fixture.service.recoverStaleUploadIntents(); - assertTrue(reserveInTransaction.get()); - assertTrue(counterInTransaction.get()); - assertEquals(1, transactions.commits); + verify(fixture.store, times(1)).deletePhysical("personal-private", "personal/key.txt"); + verify(fixture.jdbc, times(1)).update(contains("used_bytes = used_bytes -"), eq(4L), eq("000000"), + eq(101L), eq(7L), eq(4L)); } @Test - void unsafeTenantIsRejectedBeforeStorage() { - Fixture fixture = fixture(100L); - PersonalOwner unsafe = new PersonalOwner("../000000", 101L, null); + void retryRequiresFailedItemWithReadyUpload() { + Fixture fixture = fixture(); + when(fixture.jdbc.update(contains("status = 'QUEUED'"), eq("000000"), eq(101L), eq(9L))).thenReturn(1); + fixture.service.retry(OWNER, 9L); + + ArgumentCaptor sql = ArgumentCaptor.forClass(String.class); + verify(fixture.jdbc).update(sql.capture(), eq("000000"), eq(101L), eq(9L)); + assertTrue(sql.getValue().contains("uploadState')) = 'READY'")); + assertTrue(sql.getValue().contains("binary o.tenant_id = binary i.tenant_id")); + } + + @Test + void invalidFilesAndOwnerAreRejectedBeforeStorage() { + Fixture fixture = fixture(); + MockMultipartFile large = new MockMultipartFile( + "file", "large.pdf", "application/pdf", new byte[21 * 1024 * 1024]); + + assertEquals("PERSONAL_FILE_TOO_LARGE", assertThrows(ServiceException.class, + () -> fixture.service.createFile(OWNER, large, null, null)).getMessage()); + assertThrows(ServiceException.class, () -> fixture.service.createFile(OWNER, + new MockMultipartFile("file", "script.exe", "application/octet-stream", new byte[]{1}), null, null)); assertEquals("PERSONAL_OWNER_INVALID", assertThrows(ServiceException.class, - () -> fixture.service.createText(unsafe, new TextItemRequest("x", "body", null, List.of()))) - .getMessage()); - - verifyNoInteractions(fixture.spaces, fixture.store, fixture.jdbc, fixture.oss); + () -> fixture.service.createText(new PersonalOwner("../bad", 101L, null), + new TextItemRequest("x", "body", null, List.of()))).getMessage()); + verifyNoInteractions(fixture.store, fixture.jdbc, fixture.spaces); } @Test - void defaultPublicOssClientIsRejectedBeforeUploadOrMetadataWrite() { - JdbcTemplate jdbc = mock(JdbcTemplate.class); + void createEndpointsExplicitlySuspendCallerTransactions() throws Exception { + Transactional text = PersonalIngestionService.class + .getMethod("createText", PersonalOwner.class, TextItemRequest.class).getAnnotation(Transactional.class); + Transactional file = PersonalIngestionService.class + .getMethod("createFile", PersonalOwner.class, org.springframework.web.multipart.MultipartFile.class, + String.class, LocalDateTime.class).getAnnotation(Transactional.class); + + assertEquals(Propagation.NOT_SUPPORTED, text.propagation()); + assertEquals(Propagation.NOT_SUPPORTED, file.propagation()); + } + + @Test + void publicStorageIsRejectedAndPrivateStorageUploadsByPhysicalKey() { PersonalKnowledgeProperties properties = new PersonalKnowledgeProperties(); - PersonalIngestionService.OssClientProvider clients = - mock(PersonalIngestionService.OssClientProvider.class); + PersonalIngestionService.OssClientProvider clients = mock(PersonalIngestionService.OssClientProvider.class); OssClient publicClient = mock(OssClient.class); when(clients.get("")).thenReturn(publicClient); when(publicClient.getAccessPolicy()).thenReturn(AccessPolicyType.PUBLIC); - PersonalObjectStore store = PersonalIngestionService.objectStoreForTest(jdbc, properties, clients); - - ServiceException error = assertThrows(ServiceException.class, () -> store.upload( - OWNER, 100L, "personal/000000/101/100/a.txt", "txt", "text/plain", new byte[]{1})); - - assertEquals("PERSONAL_OSS_NOT_PRIVATE", error.getMessage()); - verify(clients).get(""); + PersonalObjectStore publicStore = PersonalIngestionService.objectStoreForTest(properties, clients); + assertEquals("PERSONAL_OSS_NOT_PRIVATE", + assertThrows(ServiceException.class, publicStore::requirePrivateService).getMessage()); + assertEquals("PERSONAL_OSS_NOT_PRIVATE", assertThrows(ServiceException.class, + () -> publicStore.deletePhysical("", "personal/key.txt")).getMessage()); verify(publicClient, never()).upload(any(java.io.InputStream.class), anyString(), anyLong(), anyString()); - verifyNoInteractions(jdbc); - } - @Test - void configuredPrivateOssClientIsSelectedAndWritesSystemMetadata() { - JdbcTemplate jdbc = mock(JdbcTemplate.class); - PersonalKnowledgeProperties properties = new PersonalKnowledgeProperties(); properties.setOssConfigKey(" personal-private "); - PersonalIngestionService.OssClientProvider clients = - mock(PersonalIngestionService.OssClientProvider.class); OssClient privateClient = mock(OssClient.class); when(clients.get("personal-private")).thenReturn(privateClient); when(privateClient.getAccessPolicy()).thenReturn(AccessPolicyType.PRIVATE); when(privateClient.getConfigKey()).thenReturn("personal-private"); - when(privateClient.upload(any(java.io.InputStream.class), anyString(), anyLong(), eq("text/plain"))) - .thenReturn(UploadResult.builder() - .filename("personal/000000/101/100/a.txt") - .url("https://private.invalid/personal/000000/101/100/a.txt") - .build()); - when(jdbc.update(contains("insert into sys_oss"), any(), any(), any(), any(), any(), any(), any(), any(), any())) - .thenReturn(1); - PersonalObjectStore store = PersonalIngestionService.objectStoreForTest(jdbc, properties, clients); + when(privateClient.upload(any(java.io.InputStream.class), anyString(), eq(1L), eq("text/plain"))) + .thenReturn(UploadResult.builder().filename("personal/key.txt").url("https://private/key.txt").build()); + PersonalObjectStore privateStore = PersonalIngestionService.objectStoreForTest(properties, clients); - SysOssVo uploaded = store.upload( - OWNER, 100L, "personal/000000/101/100/a.txt", "txt", "text/plain", new byte[]{1}); - - assertEquals("personal-private", uploaded.getService()); - verify(clients).get("personal-private"); - verify(privateClient).upload(any(java.io.InputStream.class), - eq("personal/000000/101/100/a.txt"), eq(1L), eq("text/plain")); - verify(jdbc).update(contains("insert into sys_oss"), any(), eq("000000"), - eq("personal/000000/101/100/a.txt"), eq("a.txt"), eq(".txt"), - eq("https://private.invalid/personal/000000/101/100/a.txt"), eq(101L), eq(101L), - eq("personal-private")); + assertEquals("personal-private", privateStore.requirePrivateService()); + assertEquals("https://private/key.txt", + privateStore.uploadPhysical("personal-private", "personal/key.txt", "text/plain", new byte[]{1})); + verify(privateClient).upload(any(java.io.InputStream.class), eq("personal/key.txt"), eq(1L), eq("text/plain")); } - private static Fixture fixture(long itemId) { - return fixture(() -> itemId); - } - - private static Fixture fixture(java.util.function.LongSupplier itemIds) { + private static Fixture fixture() { JdbcTemplate jdbc = mock(JdbcTemplate.class); PersonalSpaceService spaces = mock(PersonalSpaceService.class); - ISysOssService ossService = mock(ISysOssService.class); PersonalObjectStore store = mock(PersonalObjectStore.class); - PersonalIngestionService service = PersonalIngestionService.forTest( - jdbc, spaces, new PersonalKnowledgeProperties(), ossService, new ObjectMapper(), store, itemIds); - return new Fixture(jdbc, spaces, ossService, store, service); + when(store.requirePrivateService()).thenReturn("personal-private"); + TestTransactionManager transactions = new TestTransactionManager(); + TransactionTemplate template = new TransactionTemplate(transactions); + template.setPropagationBehavior(TransactionDefinition.PROPAGATION_REQUIRES_NEW); + PersonalIngestionService service = PersonalIngestionService.forTest(jdbc, spaces, + new PersonalKnowledgeProperties(), new ObjectMapper(), store, + new AtomicLong(100L)::getAndIncrement, template); + return new Fixture(jdbc, spaces, store, service, transactions); } - private static void stubSuccessfulCreate(Fixture fixture, PersonalOwner owner, long spaceId, long bytes, - long ossId) { - when(fixture.spaces.reserve(owner, bytes)).thenReturn(spaceId); - when(fixture.jdbc.queryForList(contains("content_hash"), - eq(owner.tenantId()), eq(owner.userId()), eq(spaceId), anyString())).thenReturn(List.of()); - when(fixture.store.upload(eq(owner), anyLong(), anyString(), anyString(), anyString(), any(byte[].class))) - .thenAnswer(invocation -> oss(ossId, invocation.getArgument(2))); - stubPersistence(fixture.jdbc); + private static void stubPhaseOne(Fixture fixture, long bytes) { + when(fixture.spaces.reserve(OWNER, bytes)).thenReturn(7L); + when(fixture.jdbc.queryForList(contains("content_hash"), eq("000000"), eq(101L), eq(7L), anyString())) + .thenReturn(List.of()); + when(fixture.jdbc.update(contains("insert into sys_oss"), eq(101L), eq("000000"), anyString(), + anyString(), anyString(), anyString(), eq(101L), eq(101L), eq("personal-private"))).thenReturn(1); + when(fixture.jdbc.update(contains("insert into aihr_personal_item"), eq(100L), eq("000000"), eq(7L), + eq(101L), anyString(), anyString(), eq(101L), anyString(), eq(bytes), anyString(), anyString(), + any(LocalDateTime.class))).thenReturn(1); + when(fixture.jdbc.update(contains("used_bytes = used_bytes +"), eq(bytes), eq("000000"), eq(101L), + eq(7L))).thenReturn(1); } - private static void stubPersistence(JdbcTemplate jdbc) { - when(jdbc.update(contains("update sys_oss"), any(), any(), any(), any(), any(), any())).thenReturn(1); - when(jdbc.update(contains("insert into aihr_personal_item"), - any(), any(), any(), any(), any(), any(), any(), any(), any(), any(), any(), any())).thenReturn(1); - when(jdbc.update(contains("used_bytes = used_bytes +"), any(), any(), any(), any())).thenReturn(1); + private static void stubFinalizeCleanup(Fixture fixture) { + when(fixture.jdbc.update(contains("status = 'DELETED'"), eq("000000"), eq(101L), eq(7L), eq(100L), + eq(101L))).thenReturn(1); + when(fixture.jdbc.update(contains("used_bytes = used_bytes -"), eq(4L), eq("000000"), eq(101L), + eq(7L), eq(4L))).thenReturn(1); + when(fixture.jdbc.update(contains("delete from sys_oss"), eq("000000"), eq(101L), eq(101L), + anyString(), eq("100"))).thenReturn(1); } - private static SysOssVo oss(long id, String fileName) { - SysOssVo result = new SysOssVo(); - result.setOssId(id); - result.setFileName(fileName); - result.setOriginalName(fileName.substring(fileName.lastIndexOf('/') + 1)); - result.setService("minio"); - result.setUrl("https://private.invalid/" + fileName); - return result; + private static Map staleIntent() { + return Map.of( + "tenant_id", "000000", "owner_user_id", 101L, "space_id", 7L, "item_id", 100L, + "oss_id", 101L, "size_bytes", 4L, "mime_type", "text/plain", + "file_name", "personal/key.txt", "service", "personal-private" + ); } - private static PersonalIngestionService transactionalProxy(PersonalIngestionService target, - TestTransactionManager transactionManager) { - ProxyFactory factory = new ProxyFactory(target); - factory.setProxyTargetClass(true); - TransactionInterceptor interceptor = new TransactionInterceptor(); - interceptor.setTransactionManager(transactionManager); - interceptor.setTransactionAttributeSource(new AnnotationTransactionAttributeSource()); - interceptor.afterPropertiesSet(); - factory.addAdvice(interceptor); - return (PersonalIngestionService) factory.getProxy(); + private static void assertUploadState(String ext, String state) throws Exception { + JsonNode json = new ObjectMapper().readTree(ext); + assertEquals("personal", json.path("source").asText()); + assertEquals(100L, json.path("itemId").asLong()); + assertEquals(state, json.path("uploadState").asText()); } - private record Fixture(JdbcTemplate jdbc, PersonalSpaceService spaces, ISysOssService oss, - PersonalObjectStore store, PersonalIngestionService service) { + private record Fixture(JdbcTemplate jdbc, PersonalSpaceService spaces, PersonalObjectStore store, + PersonalIngestionService service, TestTransactionManager transactions) { } private static final class TestTransactionManager extends AbstractPlatformTransactionManager { diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionWorkerTest.java b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionWorkerTest.java index b1d5df31..4c0b4642 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionWorkerTest.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionWorkerTest.java @@ -9,6 +9,7 @@ import org.dromara.common.oss.core.OssClient; import org.dromara.common.oss.enums.AccessPolicyType; import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; import org.springframework.jdbc.core.BatchPreparedStatementSetter; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.transaction.support.TransactionCallback; @@ -20,6 +21,7 @@ import java.time.LocalDateTime; import java.util.List; import java.util.Map; +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; @@ -33,6 +35,36 @@ import static org.mockito.Mockito.when; @Tag("dev") class PersonalIngestionWorkerTest { + @Test + void workerQueueOnlySelectsReadyUploadIntents() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + PersonalIngestionWorker worker = PersonalIngestionWorker.forTest( + jdbc, mock(ISysOssService.class), mock(KnowledgeDocumentParser.class), immediateTransactions(), + (ossId, prefix, ownerUserId, maxBytes) -> + new PersonalIngestionWorker.StoredObject("notes.txt", new byte[]{1})); + + assertFalse(worker.processNext()); + + ArgumentCaptor sql = ArgumentCaptor.forClass(String.class); + verify(jdbc).queryForList(sql.capture()); + assertTrue(sql.getValue().contains("$.uploadState')) = 'READY'")); + assertTrue(sql.getValue().contains("binary o.tenant_id = binary i.tenant_id")); + } + + @Test + void staleParsingUsesLeaseAndExhaustionThreshold() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + PersonalIngestionWorker worker = PersonalIngestionWorker.forTest( + jdbc, mock(ISysOssService.class), mock(KnowledgeDocumentParser.class), immediateTransactions(), + (ossId, prefix, ownerUserId, maxBytes) -> + new PersonalIngestionWorker.StoredObject("notes.txt", new byte[]{1})); + + worker.recoverStaleParsing(); + + verify(jdbc).update(contains("PERSONAL_PARSE_RETRY_EXHAUSTED"), eq(3), any(LocalDateTime.class)); + verify(jdbc).update(contains("set i.status = 'QUEUED'"), eq(3), any(LocalDateTime.class)); + } + @Test void workerClaimsOwnerScopedItemParsesFragmentsAndMarksReady() throws Exception { JdbcTemplate jdbc = mock(JdbcTemplate.class); @@ -62,6 +94,30 @@ class PersonalIngestionWorkerTest { verify(jdbc).update(contains("status = 'READY'"), any(), eq("[]"), eq("000000"), eq(101L), eq(9L)); } + @Test + void workerChunksWithEightHundredCharactersAndOneHundredTwentyOverlap() throws Exception { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + KnowledgeDocumentParser parser = mock(KnowledgeDocumentParser.class); + when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item())); + when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L))).thenReturn(1); + String content = "字".repeat(900); + when(parser.parse(any(), any(), any(byte[].class))) + .thenReturn(new ParsedDocument(content, "text/plain", Map.of())); + when(jdbc.update(contains("status = 'READY'"), any(), eq("[]"), eq("000000"), eq(101L), eq(9L))) + .thenReturn(1); + PersonalIngestionWorker worker = PersonalIngestionWorker.forTest( + jdbc, mock(ISysOssService.class), parser, immediateTransactions(), + (ossId, prefix, ownerUserId, maxBytes) -> + new PersonalIngestionWorker.StoredObject("notes.txt", new byte[]{1})); + + assertTrue(worker.processNext()); + + ArgumentCaptor setter = + ArgumentCaptor.forClass(BatchPreparedStatementSetter.class); + verify(jdbc).batchUpdate(contains("insert into aihr_personal_fragment"), setter.capture()); + assertEquals(2, setter.getValue().getBatchSize()); + } + @Test void workerDoesNothingWhenClaimLosesRace() { JdbcTemplate jdbc = mock(JdbcTemplate.class); 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 58d309ff..b5b655f6 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,5 +1,6 @@ package org.dromara.aihr.personal; +import org.dromara.aihr.personal.config.PersonalSchedulingConfig; import org.dromara.aihr.personal.service.PersonalSpaceService; import org.dromara.aihr.personal.support.PersonalKnowledgeProperties; import org.dromara.aihr.personal.support.PersonalOwner; @@ -16,6 +17,7 @@ import org.springframework.transaction.interceptor.TransactionInterceptor; import org.springframework.transaction.support.AbstractPlatformTransactionManager; import org.springframework.transaction.support.DefaultTransactionStatus; import org.springframework.transaction.support.TransactionSynchronizationManager; +import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.transaction.support.TransactionTemplate; import org.mockito.InOrder; @@ -23,6 +25,7 @@ import java.util.Map; 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.anyString; import static org.mockito.ArgumentMatchers.contains; import static org.mockito.ArgumentMatchers.eq; @@ -47,6 +50,16 @@ class PersonalSpaceServiceTest { assertEquals(5, properties.getDownloadUrlMinutes()); assertEquals("aihr_personal_knowledge", properties.getQdrantCollection()); assertEquals("", properties.getOssConfigKey()); + assertEquals(800, properties.getChunkSize()); + assertEquals(120, properties.getChunkOverlap()); + assertEquals(15, properties.getParsingLeaseMinutes()); + assertEquals(3, properties.getMaxParseAttempts()); + assertEquals(15, properties.getUploadCleanupAgeMinutes()); + } + + @Test + void personalSchedulingIsExplicitlyEnabled() { + assertTrue(PersonalSchedulingConfig.class.isAnnotationPresent(EnableScheduling.class)); } @Test @@ -116,6 +129,19 @@ class PersonalSpaceServiceTest { verifyNoMoreInteractions(jdbc); } + @Test + void cleanupLockDoesNotUpsertOrApplyAdmissionRules() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + when(jdbc.queryForMap(contains("for update"), eq("000000"), eq(101L))) + .thenReturn(space(7L, 1L, 1L, 1000)); + PersonalSpaceService service = new PersonalSpaceService(jdbc, properties()); + + assertEquals(7L, service.lockForUpdate(new PersonalOwner("000000", 101L, null))); + + verify(jdbc).queryForMap(contains("for update"), eq("000000"), eq(101L)); + verifyNoMoreInteractions(jdbc); + } + @Test void reserveRejectsNegativeBytesBeforeTouchingStorage() { JdbcTemplate jdbc = mock(JdbcTemplate.class);