feat(personal): ingest personal text and files asynchronously

This commit is contained in:
2026-07-12 03:03:25 +08:00
parent 806f90a5de
commit eb1d91cede
4 changed files with 1076 additions and 0 deletions
@@ -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<String> 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<String> 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<Map<String, Object>> 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<String, Object> 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<String> tags) {
List<String> 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);
}
}
}
@@ -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<Map<String, Object>> 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<String> 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<String> 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<String, Object> 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<String, Object> 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) {
}
}
@@ -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<org.springframework.web.multipart.MultipartFile> 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++;
}
}
}
@@ -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<String, Object> 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<String, Object> 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;
}
}