From eb1d91cede53385926ecdaab34aadb2c24062759 Mon Sep 17 00:00:00 2001 From: let5sne Date: Sun, 12 Jul 2026 03:03:25 +0800 Subject: [PATCH] feat(personal): ingest personal text and files asynchronously --- .../service/PersonalIngestionService.java | 381 ++++++++++++++++++ .../service/PersonalIngestionWorker.java | 271 +++++++++++++ .../PersonalIngestionServiceTest.java | 310 ++++++++++++++ .../personal/PersonalIngestionWorkerTest.java | 114 ++++++ 4 files changed, 1076 insertions(+) create mode 100644 backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionService.java create mode 100644 backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionWorker.java create mode 100644 backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionServiceTest.java create mode 100644 backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionWorkerTest.java 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 new file mode 100644 index 00000000..6268b97a --- /dev/null +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionService.java @@ -0,0 +1,381 @@ +package org.dromara.aihr.personal.service; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.extern.slf4j.Slf4j; +import org.dromara.aihr.personal.domain.PersonalAssistantDto.ItemCreatedResponse; +import org.dromara.aihr.personal.domain.PersonalAssistantDto.TextItemRequest; +import org.dromara.aihr.personal.support.PersonalKnowledgeProperties; +import org.dromara.aihr.personal.support.PersonalOwner; +import org.dromara.common.core.exception.ServiceException; +import org.dromara.common.oss.factory.OssFactory; +import org.dromara.system.domain.vo.SysOssVo; +import org.dromara.system.service.ISysOssService; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.support.GeneratedKeyHolder; +import org.springframework.jdbc.support.KeyHolder; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.support.TransactionSynchronization; +import org.springframework.transaction.support.TransactionSynchronizationManager; +import org.springframework.web.multipart.MultipartFile; + +import java.io.ByteArrayInputStream; +import java.io.File; +import java.io.IOException; +import java.io.InputStream; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.sql.PreparedStatement; +import java.sql.Statement; +import java.time.LocalDateTime; +import java.util.HexFormat; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Set; +import java.util.UUID; + +@Slf4j +@Service +public class PersonalIngestionService { + + private static final String ITEM_NOT_FOUND = "PERSONAL_ITEM_NOT_FOUND"; + private static final Set SUPPORTED_FILE_SUFFIXES = Set.of( + "txt", "md", "markdown", "pdf", "doc", "docx", "xls", "xlsx", "ppt", "pptx" + ); + + private final JdbcTemplate jdbcTemplate; + private final PersonalSpaceService spaceService; + private final PersonalKnowledgeProperties properties; + private final ISysOssService ossService; + private final ObjectMapper objectMapper; + + public PersonalIngestionService(JdbcTemplate jdbcTemplate, PersonalSpaceService spaceService, + PersonalKnowledgeProperties properties, ISysOssService ossService, + ObjectMapper objectMapper) { + this.jdbcTemplate = jdbcTemplate; + this.spaceService = spaceService; + this.properties = properties; + this.ossService = ossService; + this.objectMapper = objectMapper; + } + + @Transactional + public ItemCreatedResponse createText(PersonalOwner owner, TextItemRequest request) { + if (request == null || request.content() == null || request.content().isBlank()) { + throw new ServiceException("PERSONAL_TEXT_EMPTY"); + } + byte[] bytes = request.content().getBytes(StandardCharsets.UTF_8); + validateSize(bytes.length); + String hash = sha256(bytes); + ItemCreatedResponse duplicate = duplicate(owner, hash); + if (duplicate != null) { + return duplicate; + } + + long spaceId = spaceService.reserve(owner, bytes.length); + Path temporary = null; + SysOssVo uploaded; + try { + temporary = Files.createTempFile("personal-" + UUID.randomUUID(), ".txt"); + Files.write(temporary, bytes); + uploaded = ossService.upload(temporary.toFile()); + } catch (IOException ex) { + throw new ServiceException("PERSONAL_TEXT_STORAGE_FAILED"); + } finally { + if (temporary != null) { + try { + Files.deleteIfExists(temporary); + } catch (IOException ex) { + log.warn("Unable to remove temporary personal text object"); + } + } + } + return persist(owner, spaceId, "TEXT", cleanTitle(request.title(), "文字资料"), + "text/plain", bytes.length, hash, request.capturedAt(), request.tags(), uploaded); + } + + @Transactional + public ItemCreatedResponse createFile(PersonalOwner owner, MultipartFile file, String title, + LocalDateTime capturedAt) { + validateFile(file); + byte[] bytes; + try { + bytes = file.getBytes(); + } catch (IOException ex) { + throw new ServiceException("PERSONAL_FILE_READ_FAILED"); + } + String hash = sha256(bytes); + ItemCreatedResponse duplicate = duplicate(owner, hash); + if (duplicate != null) { + return duplicate; + } + + validateSize(bytes.length); + long spaceId = spaceService.reserve(owner, bytes.length); + String safeUploadName = "personal-" + UUID.randomUUID() + suffixWithDot(file.getOriginalFilename()); + SysOssVo uploaded = ossService.upload(new SafeMultipartFile( + safeUploadName, cleanMime(file.getContentType()), bytes)); + String originalName = safeFileName(file.getOriginalFilename()); + return persist(owner, spaceId, "FILE", cleanTitle(title, originalName), + cleanMime(file.getContentType()), bytes.length, hash, capturedAt, List.of(), uploaded); + } + + public void retry(PersonalOwner owner, long itemId) { + 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' + """, owner.tenantId(), owner.userId(), itemId); + if (updated == 0) { + throw new ServiceException(ITEM_NOT_FOUND); + } + } + + private ItemCreatedResponse persist(PersonalOwner owner, long spaceId, String sourceType, String title, + String mimeType, long size, String hash, LocalDateTime capturedAt, + List tags, SysOssVo uploaded) { + if (uploaded == null || uploaded.getOssId() == null) { + throw new ServiceException("PERSONAL_OSS_UPLOAD_FAILED"); + } + boolean deferredCleanup = registerRollbackCleanup(uploaded); + try { + KeyHolder keyHolder = new GeneratedKeyHolder(); + jdbcTemplate.update(connection -> { + PreparedStatement statement = connection.prepareStatement(""" + insert into aihr_personal_item + (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()) + """, Statement.RETURN_GENERATED_KEYS); + statement.setString(1, owner.tenantId()); + statement.setLong(2, spaceId); + statement.setLong(3, owner.userId()); + statement.setString(4, sourceType); + statement.setString(5, title); + statement.setLong(6, uploaded.getOssId()); + statement.setString(7, mimeType); + statement.setLong(8, size); + statement.setString(9, hash); + statement.setString(10, tagsJson(tags)); + statement.setObject(11, capturedAt == null ? LocalDateTime.now() : capturedAt); + return statement; + }, keyHolder); + Number key = keyHolder.getKey(); + if (key == null) { + 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(key.longValue(), "QUEUED", null); + } catch (RuntimeException ex) { + if (!deferredCleanup) { + cleanupOss(uploaded.getOssId()); + } + throw ex; + } + } + + private ItemCreatedResponse duplicate(PersonalOwner owner, String hash) { + List> rows = jdbcTemplate.queryForList(""" + select id, status + from aihr_personal_item + where tenant_id = ? and owner_user_id = ? and content_hash = ? + and status <> 'DELETED' + order by id desc + limit 1 + """, owner.tenantId(), owner.userId(), hash); + if (rows.isEmpty()) { + return null; + } + Map row = rows.get(0); + long id = ((Number) row.get("id")).longValue(); + return new ItemCreatedResponse(id, String.valueOf(row.get("status")), id); + } + + private void validateFile(MultipartFile file) { + if (file == null || file.isEmpty() || file.getSize() <= 0) { + throw new ServiceException("PERSONAL_FILE_EMPTY"); + } + validateSize(file.getSize()); + String suffix = suffix(file.getOriginalFilename()); + if (!SUPPORTED_FILE_SUFFIXES.contains(suffix)) { + throw new ServiceException("PERSONAL_FILE_UNSUPPORTED"); + } + } + + 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"); + } + } + + private boolean registerRollbackCleanup(SysOssVo uploaded) { + if (!TransactionSynchronizationManager.isSynchronizationActive()) { + return false; + } + TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { + @Override + public void afterCompletion(int status) { + if (status != STATUS_COMMITTED) { + cleanupRolledBackOss(uploaded); + } + } + }); + return true; + } + + private void cleanupRolledBackOss(SysOssVo uploaded) { + if (uploaded.getService() != null && !uploaded.getService().isBlank() + && uploaded.getFileName() != null && !uploaded.getFileName().isBlank()) { + try { + // sys_oss participates in the outer transaction and may already be rolled back here, + // so remove the physical object by its private object key before metadata cleanup. + OssFactory.instance(uploaded.getService()).delete(uploaded.getFileName()); + } catch (RuntimeException cleanupError) { + log.warn("Unable to clean rolled-back personal OSS object id={}", uploaded.getOssId()); + } + } + cleanupOss(uploaded.getOssId()); + } + + private void cleanupOss(long ossId) { + try { + ossService.deleteWithValidByIds(List.of(ossId), false); + } catch (RuntimeException cleanupError) { + log.warn("Unable to clean personal OSS object id={}", ossId); + } + } + + 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 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); + } + } + + 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 = "个人资料"; + } + 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); + } + 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('.'); + return dot < 0 ? "" : safe.substring(dot + 1).toLowerCase(Locale.ROOT); + } + + private static String suffixWithDot(String fileName) { + String suffix = suffix(fileName); + return suffix.isBlank() ? "" : "." + suffix; + } + + private static String cleanMime(String value) { + 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(); + } + + private static final class SafeMultipartFile implements MultipartFile { + private final String fileName; + private final String contentType; + private final byte[] bytes; + + private SafeMultipartFile(String fileName, String contentType, byte[] bytes) { + this.fileName = fileName; + this.contentType = contentType; + this.bytes = bytes.clone(); + } + + @Override + public String getName() { + return "file"; + } + + @Override + public String getOriginalFilename() { + return fileName; + } + + @Override + public String getContentType() { + return contentType; + } + + @Override + public boolean isEmpty() { + return bytes.length == 0; + } + + @Override + public long getSize() { + return bytes.length; + } + + @Override + public byte[] getBytes() { + return bytes.clone(); + } + + @Override + public InputStream getInputStream() { + return new ByteArrayInputStream(bytes); + } + + @Override + public void transferTo(File destination) throws IOException { + Files.write(destination.toPath(), bytes); + } + } +} 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 new file mode 100644 index 00000000..0cf73f9f --- /dev/null +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/personal/service/PersonalIngestionWorker.java @@ -0,0 +1,271 @@ +package org.dromara.aihr.personal.service; + +import lombok.extern.slf4j.Slf4j; +import org.dromara.aihr.knowledge.parse.KnowledgeDocumentParser; +import org.dromara.aihr.knowledge.parse.ParsedDocument; +import org.dromara.aihr.personal.support.PersonalKnowledgeProperties; +import org.dromara.common.oss.core.OssClient; +import org.dromara.common.oss.factory.OssFactory; +import org.dromara.system.domain.vo.SysOssVo; +import org.dromara.system.service.ISysOssService; +import org.springframework.jdbc.core.BatchPreparedStatementSetter; +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.support.TransactionTemplate; + +import java.io.IOException; +import java.io.InputStream; +import java.sql.PreparedStatement; +import java.sql.SQLException; +import java.util.List; +import java.util.Map; + +@Slf4j +@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; + + public PersonalIngestionWorker(JdbcTemplate jdbcTemplate, ISysOssService ossService, + KnowledgeDocumentParser parser, PersonalKnowledgeProperties properties, + PlatformTransactionManager transactionManager) { + this(jdbcTemplate, parser, new TransactionTemplate(transactionManager), + defaultReader(ossService), configuredMaxBytes(properties)); + } + + private PersonalIngestionWorker(JdbcTemplate jdbcTemplate, KnowledgeDocumentParser parser, + TransactionTemplate transactionTemplate, StoredObjectReader objectReader, + long maxInputBytes) { + this.jdbcTemplate = jdbcTemplate; + this.parser = parser; + this.transactionTemplate = transactionTemplate; + this.objectReader = objectReader; + this.maxInputBytes = maxInputBytes; + } + + public static PersonalIngestionWorker forTest(JdbcTemplate jdbcTemplate, ISysOssService ossService, + KnowledgeDocumentParser parser, + TransactionTemplate transactionTemplate, + StoredObjectReader objectReader) { + return new PersonalIngestionWorker(jdbcTemplate, parser, transactionTemplate, objectReader, + 20L * 1024 * 1024); + } + + @Scheduled(fixedDelayString = "${aihr.personal.ingestion-delay-ms:2000}") + public void poll() { + processNext(); + } + + 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 + limit 1 + """); + if (queued.isEmpty()) { + return false; + } + Item item = item(queued.get(0)); + int claimed = jdbcTemplate.update(""" + update aihr_personal_item + set status = 'PARSING', attempt_count = attempt_count + 1, + error_code = null, error_message = null, update_time = now() + where tenant_id = ? and owner_user_id = ? and id = ? + and status = 'QUEUED' + """, item.tenantId(), item.ownerUserId(), item.id()); + if (claimed != 1) { + return false; + } + + try { + StoredObject stored = objectReader.read(item.ossId(), maxInputBytes); + ParsedDocument document = parser.parse(stored.fileName(), item.mimeType(), stored.bytes()); + List chunks = document.chunks(CHUNK_SIZE, CHUNK_OVERLAP); + if (chunks.isEmpty()) { + throw new KnowledgeDocumentParser.ParseException( + KnowledgeDocumentParser.Failure.EMPTY, "document contains no text"); + } + transactionTemplate.execute(status -> { + persistSuccess(item, document, chunks); + return null; + }); + } catch (Exception ex) { + Failure failure = publicFailure(ex); + jdbcTemplate.update(""" + update aihr_personal_item + set status = 'FAILED', error_code = ?, error_message = ?, update_time = now() + where tenant_id = ? and owner_user_id = ? and id = ? + and status = 'PARSING' + """, failure.code(), failure.message(), item.tenantId(), item.ownerUserId(), item.id()); + log.warn("Personal ingestion failed itemId={} ownerUserId={} code={}", + item.id(), item.ownerUserId(), failure.code()); + } + return true; + } + + private void persistSuccess(Item item, ParsedDocument document, List chunks) { + jdbcTemplate.update(""" + delete from aihr_personal_fragment + where tenant_id = ? and owner_user_id = ? and item_id = ? + """, item.tenantId(), item.ownerUserId(), item.id()); + jdbcTemplate.batchUpdate(""" + insert into aihr_personal_fragment + (tenant_id, space_id, owner_user_id, item_id, idx, content, token_count, + embedding_json, embedding_model, embedding_time, create_time) + values (?, ?, ?, ?, ?, ?, ?, null, null, null, now()) + """, new BatchPreparedStatementSetter() { + @Override + public void setValues(PreparedStatement statement, int index) throws SQLException { + String chunk = chunks.get(index); + statement.setString(1, item.tenantId()); + statement.setLong(2, item.spaceId()); + statement.setLong(3, item.ownerUserId()); + statement.setLong(4, item.id()); + statement.setInt(5, index); + statement.setString(6, chunk); + statement.setInt(7, estimatedTokens(chunk)); + } + + @Override + public int getBatchSize() { + return chunks.size(); + } + }); + int updated = jdbcTemplate.update(""" + update aihr_personal_item + set status = 'READY', parsed_at = now(), summary = ?, tags_json = ?, + error_code = null, error_message = null, update_time = now() + where tenant_id = ? and owner_user_id = ? and id = ? + and status = 'PARSING' + """, summary(document.text()), item.tagsJson(), item.tenantId(), item.ownerUserId(), item.id()); + if (updated != 1) { + throw new IllegalStateException("personal item state changed while parsing"); + } + } + + private static StoredObjectReader defaultReader(ISysOssService ossService) { + return (ossId, maxBytes) -> { + SysOssVo object = ossService.getById(ossId); + if (object == null || object.getFileName() == null || object.getFileName().isBlank() + || object.getService() == null || object.getService().isBlank()) { + throw new IOException("personal source object is unavailable"); + } + OssClient storage = OssFactory.instance(object.getService()); + try (InputStream input = storage.getObjectContent(object.getFileName())) { + int boundedLimit = (int) Math.min(Integer.MAX_VALUE - 1L, maxBytes); + byte[] bytes = input.readNBytes(boundedLimit + 1); + if (bytes.length > maxBytes) { + throw new IOException("personal source object exceeds limit"); + } + return new StoredObject(safeObjectName(object.getOriginalName()), bytes); + } + }; + } + + private static Item item(Map row) { + Object tags = row.get("tags_json"); + return new Item( + number(row, "id"), + String.valueOf(row.get("tenant_id")), + number(row, "space_id"), + number(row, "owner_user_id"), + String.valueOf(row.get("source_type")), + String.valueOf(row.get("title")), + number(row, "oss_id"), + String.valueOf(row.get("mime_type")), + tags == null ? "[]" : String.valueOf(tags) + ); + } + + private static long number(Map row, String key) { + Object value = row.get(key); + if (!(value instanceof Number number)) { + throw new IllegalStateException("personal item metadata is incomplete"); + } + return number.longValue(); + } + + private static Failure publicFailure(Exception error) { + Throwable candidate = error; + while (candidate != null) { + if (candidate instanceof KnowledgeDocumentParser.ParseException parseError) { + return switch (parseError.failure()) { + case EMPTY -> new Failure("PERSONAL_PARSE_EMPTY", "资料中未识别到可用文字"); + case TOO_LARGE -> new Failure("PERSONAL_PARSE_TOO_LARGE", "资料解析后内容超过限制"); + case INVALID -> new Failure("PERSONAL_PARSE_INVALID", "资料解析失败,请检查文件后重试"); + }; + } + candidate = candidate.getCause(); + } + return new Failure("PERSONAL_PARSE_FAILED", "资料处理失败,请稍后重试"); + } + + private static long configuredMaxBytes(PersonalKnowledgeProperties properties) { + try { + long bytes = Math.multiplyExact(properties.getMaxFileSizeMb(), 1024L * 1024L); + if (bytes <= 0) { + throw new IllegalArgumentException("personal max file size must be positive"); + } + return bytes; + } catch (ArithmeticException ex) { + throw new IllegalArgumentException("personal max file size is invalid", ex); + } + } + + private static String summary(String text) { + int[] codePoints = text.codePoints().limit(300).toArray(); + return new String(codePoints, 0, codePoints.length).trim(); + } + + private static int estimatedTokens(String content) { + return Math.max(1, (content.codePointCount(0, content.length()) + 1) / 2); + } + + private static String safeObjectName(String fileName) { + if (fileName == null || fileName.isBlank()) { + return "personal-object"; + } + String safe = fileName.replace('\\', '/'); + int slash = safe.lastIndexOf('/'); + if (slash >= 0) { + safe = safe.substring(slash + 1); + } + safe = safe.replace('\r', '_').replace('\n', '_').trim(); + return safe.isBlank() ? "personal-object" : safe; + } + + @FunctionalInterface + public interface StoredObjectReader { + StoredObject read(long ossId, long maxBytes) throws Exception; + } + + public record StoredObject(String fileName, byte[] bytes) { + public StoredObject { + bytes = bytes == null ? new byte[0] : bytes.clone(); + } + + @Override + public byte[] bytes() { + return bytes.clone(); + } + } + + private record Item(long id, String tenantId, long spaceId, long ownerUserId, String sourceType, + String title, long ossId, String mimeType, String tagsJson) { + } + + private record Failure(String code, String message) { + } +} 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 new file mode 100644 index 00000000..0a06ec5d --- /dev/null +++ b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionServiceTest.java @@ -0,0 +1,310 @@ +package org.dromara.aihr.personal; + +import com.fasterxml.jackson.databind.ObjectMapper; +import org.dromara.aihr.personal.domain.PersonalAssistantDto.ItemCreatedResponse; +import org.dromara.aihr.personal.domain.PersonalAssistantDto.TextItemRequest; +import org.dromara.aihr.personal.service.PersonalIngestionService; +import org.dromara.aihr.personal.service.PersonalSpaceService; +import org.dromara.aihr.personal.support.PersonalKnowledgeProperties; +import org.dromara.aihr.personal.support.PersonalOwner; +import org.dromara.common.core.exception.ServiceException; +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.jdbc.core.PreparedStatementCreator; +import org.springframework.jdbc.support.KeyHolder; +import org.springframework.mock.web.MockMultipartFile; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.TransactionDefinition; +import org.springframework.transaction.annotation.AnnotationTransactionAttributeSource; +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 java.nio.charset.StandardCharsets; +import java.time.LocalDateTime; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +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.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +@Tag("dev") +class PersonalIngestionServiceTest { + + private static final PersonalOwner OWNER = new PersonalOwner("000000", 101L, "ext-101"); + + @Test + void failedItemRetriesThroughOwnerScopedQueuedState() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + when(jdbc.update(contains("status = 'QUEUED'"), eq("000000"), eq(101L), eq(9L))) + .thenReturn(1); + PersonalIngestionService service = service(jdbc, mock(PersonalSpaceService.class), mock(ISysOssService.class)); + + service.retry(OWNER, 9L); + + verify(jdbc).update(contains("status = 'QUEUED'"), eq("000000"), eq(101L), eq(9L)); + } + + @Test + void retryDoesNotDiscloseMissingOrForeignItem() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + PersonalIngestionService service = service(jdbc, mock(PersonalSpaceService.class), mock(ISysOssService.class)); + + ServiceException error = assertThrows(ServiceException.class, () -> service.retry(OWNER, 9L)); + + assertEquals("PERSONAL_ITEM_NOT_FOUND", error.getMessage()); + } + + @Test + void fileLargerThanConfiguredLimitIsRejectedBeforeOssUpload() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + ISysOssService oss = mock(ISysOssService.class); + MockMultipartFile file = new MockMultipartFile( + "file", "large.pdf", "application/pdf", new byte[21 * 1024 * 1024]); + PersonalIngestionService service = service(jdbc, mock(PersonalSpaceService.class), oss); + + ServiceException error = assertThrows(ServiceException.class, + () -> service.createFile(OWNER, file, null, null)); + + assertEquals("PERSONAL_FILE_TOO_LARGE", error.getMessage()); + verifyNoInteractions(oss); + verifyNoInteractions(jdbc); + } + + @Test + void unsupportedAndEmptyFilesAreRejectedBeforeOssUpload() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + ISysOssService oss = mock(ISysOssService.class); + PersonalIngestionService service = service(jdbc, mock(PersonalSpaceService.class), oss); + + assertThrows(ServiceException.class, () -> service.createFile( + OWNER, new MockMultipartFile("file", "script.exe", "application/octet-stream", new byte[]{1}), null, null)); + assertThrows(ServiceException.class, () -> service.createFile( + OWNER, new MockMultipartFile("file", "empty.txt", "text/plain", new byte[0]), null, null)); + + verifyNoInteractions(oss); + verifyNoInteractions(jdbc); + } + + @Test + void createTextReservesUploadsInsertsAndMutatesCountersOnce() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + PersonalSpaceService spaces = mock(PersonalSpaceService.class); + ISysOssService oss = mock(ISysOssService.class); + when(spaces.reserve(eq(OWNER), eq(18L))).thenReturn(7L); + when(jdbc.queryForList(contains("content_hash"), eq("000000"), eq(101L), any(String.class))) + .thenReturn(List.of()); + when(oss.upload(any(java.io.File.class))).thenReturn(oss(81L)); + generatedId(jdbc, 99L); + when(jdbc.update(contains("used_bytes = used_bytes +"), eq(18L), eq("000000"), eq(101L), eq(7L))) + .thenReturn(1); + PersonalIngestionService service = service(jdbc, spaces, oss); + + ItemCreatedResponse response = service.createText( + OWNER, new TextItemRequest("周报", "保洁巡检记录", LocalDateTime.of(2026, 7, 12, 9, 0), List.of("保洁"))); + + assertEquals(99L, response.itemId()); + assertEquals("QUEUED", response.status()); + assertEquals(null, response.duplicateOf()); + InOrder order = inOrder(spaces, oss, jdbc); + order.verify(spaces).reserve(OWNER, 18L); + order.verify(oss).upload(any(java.io.File.class)); + order.verify(jdbc).update(any(PreparedStatementCreator.class), any(KeyHolder.class)); + order.verify(jdbc).update(contains("used_bytes = used_bytes +"), eq(18L), eq("000000"), eq(101L), eq(7L)); + assertNotNull(PersonalIngestionService.class.getAnnotation(org.springframework.stereotype.Service.class)); + assertNotNull(method("createText", PersonalOwner.class, TextItemRequest.class).getAnnotation(Transactional.class)); + } + + @Test + void createFileDoesNotDeduplicateAcrossOwners() throws Exception { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + PersonalSpaceService spaces = mock(PersonalSpaceService.class); + ISysOssService oss = mock(ISysOssService.class); + MockMultipartFile file = new MockMultipartFile( + "file", "notes.txt", "text/plain", "same".getBytes(StandardCharsets.UTF_8)); + when(jdbc.queryForList(contains("content_hash"), eq("000000"), eq(101L), any(String.class))) + .thenReturn(List.of()); + when(spaces.reserve(OWNER, 4L)).thenReturn(7L); + when(oss.upload(any(org.springframework.web.multipart.MultipartFile.class))).thenReturn(oss(82L)); + generatedId(jdbc, 100L); + when(jdbc.update(contains("used_bytes = used_bytes +"), eq(4L), eq("000000"), eq(101L), eq(7L))) + .thenReturn(1); + PersonalIngestionService service = service(jdbc, spaces, oss); + + ItemCreatedResponse response = service.createFile(OWNER, file, null, null); + + assertEquals(100L, response.itemId()); + ArgumentCaptor uploaded = + ArgumentCaptor.forClass(org.springframework.web.multipart.MultipartFile.class); + verify(oss).upload(uploaded.capture()); + assertTrue(uploaded.getValue().getOriginalFilename().matches("personal-[0-9a-f-]+\\.txt")); + assertEquals("same", new String(uploaded.getValue().getBytes(), StandardCharsets.UTF_8)); + verify(jdbc).queryForList(contains("content_hash"), eq("000000"), eq(101L), any(String.class)); + verify(jdbc, never()).queryForList(contains("content_hash"), eq("000000"), eq(202L), any(String.class)); + } + + @Test + void duplicateForSameOwnerReturnsExistingItemWithoutReservationOrUpload() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + PersonalSpaceService spaces = mock(PersonalSpaceService.class); + ISysOssService oss = mock(ISysOssService.class); + MockMultipartFile file = new MockMultipartFile( + "file", "notes.txt", "text/plain", "same".getBytes(StandardCharsets.UTF_8)); + when(jdbc.queryForList(contains("content_hash"), eq("000000"), eq(101L), any(String.class))) + .thenReturn(List.of(Map.of("id", 77L, "status", "READY"))); + PersonalIngestionService service = service(jdbc, spaces, oss); + + ItemCreatedResponse response = service.createFile(OWNER, file, null, null); + + assertEquals(77L, response.itemId()); + assertEquals(77L, response.duplicateOf()); + assertEquals("READY", response.status()); + verifyNoInteractions(spaces, oss); + } + + @Test + void failedPersistenceRegistersBestEffortOssCleanup() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + PersonalSpaceService spaces = mock(PersonalSpaceService.class); + ISysOssService oss = mock(ISysOssService.class); + MockMultipartFile file = new MockMultipartFile( + "file", "notes.txt", "text/plain", "same".getBytes(StandardCharsets.UTF_8)); + when(jdbc.queryForList(contains("content_hash"), eq("000000"), eq(101L), any(String.class))) + .thenReturn(List.of()); + when(spaces.reserve(OWNER, 4L)).thenReturn(7L); + SysOssVo rolledBack = oss(82L); + rolledBack.setService(null); + when(oss.upload(any(org.springframework.web.multipart.MultipartFile.class))).thenReturn(rolledBack); + when(jdbc.update(any(PreparedStatementCreator.class), any(KeyHolder.class))) + .thenThrow(new IllegalStateException("db failed")); + TestTransactionManager transactions = new TestTransactionManager(); + PersonalIngestionService service = transactionalProxy(service(jdbc, spaces, oss), transactions); + + assertThrows(IllegalStateException.class, () -> service.createFile(OWNER, file, null, null)); + + verify(oss).deleteWithValidByIds(eq(List.of(82L)), eq(false)); + assertEquals(1, transactions.rollbacks); + } + + @Test + void reservationAndCounterMutationRunInsideOneOuterTransaction() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + PersonalSpaceService spaces = mock(PersonalSpaceService.class); + ISysOssService oss = mock(ISysOssService.class); + MockMultipartFile file = new MockMultipartFile( + "file", "notes.txt", "text/plain", "same".getBytes(StandardCharsets.UTF_8)); + AtomicBoolean reserveInTransaction = new AtomicBoolean(); + AtomicBoolean counterInTransaction = new AtomicBoolean(); + when(jdbc.queryForList(contains("content_hash"), eq("000000"), eq(101L), any(String.class))) + .thenReturn(List.of()); + when(spaces.reserve(OWNER, 4L)).thenAnswer(invocation -> { + reserveInTransaction.set(TransactionSynchronizationManager.isActualTransactionActive()); + return 7L; + }); + when(oss.upload(any(org.springframework.web.multipart.MultipartFile.class))).thenReturn(oss(82L)); + generatedId(jdbc, 100L); + when(jdbc.update(contains("used_bytes = used_bytes +"), eq(4L), eq("000000"), eq(101L), eq(7L))) + .thenAnswer(invocation -> { + counterInTransaction.set(TransactionSynchronizationManager.isActualTransactionActive()); + return 1; + }); + TestTransactionManager transactions = new TestTransactionManager(); + PersonalIngestionService service = transactionalProxy(service(jdbc, spaces, oss), transactions); + + service.createFile(OWNER, file, null, null); + + assertTrue(reserveInTransaction.get()); + assertTrue(counterInTransaction.get()); + assertEquals(1, transactions.commits); + assertEquals(0, transactions.rollbacks); + } + + private static PersonalIngestionService service(JdbcTemplate jdbc, PersonalSpaceService spaces, ISysOssService oss) { + return new PersonalIngestionService(jdbc, spaces, properties(), oss, new ObjectMapper()); + } + + private static PersonalKnowledgeProperties properties() { + return new PersonalKnowledgeProperties(); + } + + private static SysOssVo oss(long id) { + SysOssVo result = new SysOssVo(); + result.setOssId(id); + result.setFileName("personal/random.txt"); + result.setOriginalName("random.txt"); + result.setService("minio"); + result.setUrl("https://private.invalid/random.txt"); + return result; + } + + private static void generatedId(JdbcTemplate jdbc, long id) { + when(jdbc.update(any(PreparedStatementCreator.class), any(KeyHolder.class))).thenAnswer(invocation -> { + KeyHolder holder = invocation.getArgument(1); + holder.getKeyList().add(Map.of("GENERATED_KEY", id)); + return 1; + }); + } + + private static java.lang.reflect.Method method(String name, Class... types) { + try { + return PersonalIngestionService.class.getMethod(name, types); + } catch (NoSuchMethodException e) { + throw new AssertionError(e); + } + } + + 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 final class TestTransactionManager extends AbstractPlatformTransactionManager { + private int commits; + private int rollbacks; + + @Override + protected Object doGetTransaction() { + return new Object(); + } + + @Override + protected void doBegin(Object transaction, TransactionDefinition definition) { + } + + @Override + protected void doCommit(DefaultTransactionStatus status) { + commits++; + } + + @Override + protected void doRollback(DefaultTransactionStatus status) { + rollbacks++; + } + } +} 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 new file mode 100644 index 00000000..c844d4a8 --- /dev/null +++ b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/personal/PersonalIngestionWorkerTest.java @@ -0,0 +1,114 @@ +package org.dromara.aihr.personal; + +import org.dromara.aihr.knowledge.parse.KnowledgeDocumentParser; +import org.dromara.aihr.knowledge.parse.ParsedDocument; +import org.dromara.aihr.personal.service.PersonalIngestionWorker; +import org.dromara.system.service.ISysOssService; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.springframework.jdbc.core.BatchPreparedStatementSetter; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.transaction.support.TransactionCallback; +import org.springframework.transaction.support.TransactionTemplate; + +import java.nio.charset.StandardCharsets; +import java.sql.Timestamp; +import java.time.LocalDateTime; +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.contains; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@Tag("dev") +class PersonalIngestionWorkerTest { + + @Test + void workerClaimsOwnerScopedItemParsesFragmentsAndMarksReady() throws Exception { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + KnowledgeDocumentParser parser = mock(KnowledgeDocumentParser.class); + TransactionTemplate transactions = immediateTransactions(); + Map item = item(); + when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item)); + when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L))).thenReturn(1); + when(parser.parse(eq("notes.txt"), eq("text/plain"), any(byte[].class))) + .thenReturn(new ParsedDocument("一二三四五六", "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, transactions, + (ossId, maxBytes) -> new PersonalIngestionWorker.StoredObject( + "notes.txt", "same".getBytes(StandardCharsets.UTF_8))); + + assertTrue(worker.processNext()); + + verify(jdbc).update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L)); + verify(jdbc).update(contains("delete from aihr_personal_fragment"), eq("000000"), eq(101L), eq(9L)); + verify(jdbc).batchUpdate(contains("insert into aihr_personal_fragment"), any(BatchPreparedStatementSetter.class)); + verify(jdbc).update(contains("status = 'READY'"), any(), eq("[]"), eq("000000"), eq(101L), eq(9L)); + } + + @Test + void workerDoesNothingWhenClaimLosesRace() { + JdbcTemplate jdbc = mock(JdbcTemplate.class); + when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item())); + PersonalIngestionWorker worker = PersonalIngestionWorker.forTest( + jdbc, mock(ISysOssService.class), mock(KnowledgeDocumentParser.class), immediateTransactions(), + (ossId, maxBytes) -> new PersonalIngestionWorker.StoredObject("notes.txt", new byte[]{1})); + + assertFalse(worker.processNext()); + + verify(jdbc, never()).batchUpdate(any(String.class), any(BatchPreparedStatementSetter.class)); + verify(jdbc, never()).update(contains("status = 'READY'"), any(), any(), any(), any(), any()); + } + + @Test + void workerPersistsOnlyStablePublicFailure() 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); + when(parser.parse(any(), any(), any(byte[].class))) + .thenThrow(new KnowledgeDocumentParser.ParseException( + KnowledgeDocumentParser.Failure.INVALID, "secret parser detail")); + PersonalIngestionWorker worker = PersonalIngestionWorker.forTest( + jdbc, mock(ISysOssService.class), parser, immediateTransactions(), + (ossId, maxBytes) -> new PersonalIngestionWorker.StoredObject("notes.txt", new byte[]{1})); + + assertTrue(worker.processNext()); + + verify(jdbc).update(contains("status = 'FAILED'"), eq("PERSONAL_PARSE_INVALID"), + eq("资料解析失败,请检查文件后重试"), eq("000000"), eq(101L), eq(9L)); + } + + private static Map item() { + return Map.of( + "id", 9L, + "tenant_id", "000000", + "space_id", 7L, + "owner_user_id", 101L, + "source_type", "TEXT", + "title", "周报", + "oss_id", 81L, + "mime_type", "text/plain", + "tags_json", "[]", + "captured_at", Timestamp.valueOf(LocalDateTime.of(2026, 7, 12, 9, 0)) + ); + } + + private static TransactionTemplate immediateTransactions() { + TransactionTemplate template = mock(TransactionTemplate.class); + when(template.execute(any())).thenAnswer(invocation -> { + TransactionCallback callback = invocation.getArgument(0); + return callback.doInTransaction(null); + }); + return template; + } +}