fix(personal): isolate personal objects and dedupe under lock

This commit is contained in:
2026-07-12 03:14:57 +08:00
parent eb1d91cede
commit 10c04eaf67
4 changed files with 363 additions and 300 deletions
@@ -1,5 +1,6 @@
package org.dromara.aihr.personal.service;
import com.baomidou.mybatisplus.core.toolkit.IdWorker;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
@@ -8,12 +9,13 @@ 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.core.OssClient;
import org.dromara.common.oss.entity.UploadResult;
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.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;
@@ -21,16 +23,10 @@ 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;
@@ -38,12 +34,15 @@ import java.util.Locale;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.function.LongSupplier;
import java.util.regex.Pattern;
@Slf4j
@Service
public class PersonalIngestionService {
private static final String ITEM_NOT_FOUND = "PERSONAL_ITEM_NOT_FOUND";
private static final Pattern SAFE_TENANT = Pattern.compile("[A-Za-z0-9_-]{1,20}");
private static final Set<String> SUPPORTED_FILE_SUFFIXES = Set.of(
"txt", "md", "markdown", "pdf", "doc", "docx", "xls", "xlsx", "ppt", "pptx"
);
@@ -53,55 +52,54 @@ public class PersonalIngestionService {
private final PersonalKnowledgeProperties properties;
private final ISysOssService ossService;
private final ObjectMapper objectMapper;
private final PersonalObjectStore objectStore;
private final LongSupplier itemIdSupplier;
@Autowired
public PersonalIngestionService(JdbcTemplate jdbcTemplate, PersonalSpaceService spaceService,
PersonalKnowledgeProperties properties, ISysOssService ossService,
ObjectMapper objectMapper) {
this(jdbcTemplate, spaceService, properties, ossService, objectMapper,
new DefaultPersonalObjectStore(jdbcTemplate), IdWorker::getId);
}
private PersonalIngestionService(JdbcTemplate jdbcTemplate, PersonalSpaceService spaceService,
PersonalKnowledgeProperties properties, ISysOssService ossService,
ObjectMapper objectMapper, PersonalObjectStore objectStore,
LongSupplier itemIdSupplier) {
this.jdbcTemplate = jdbcTemplate;
this.spaceService = spaceService;
this.properties = properties;
this.ossService = ossService;
this.objectMapper = objectMapper;
this.objectStore = objectStore;
this.itemIdSupplier = itemIdSupplier;
}
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);
}
@Transactional
public ItemCreatedResponse createText(PersonalOwner owner, TextItemRequest request) {
validateOwner(owner);
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);
return create(owner, "TEXT", cleanTitle(request.title(), "文字资料"), "txt", "text/plain", bytes,
request.capturedAt(), request.tags());
}
@Transactional
public ItemCreatedResponse createFile(PersonalOwner owner, MultipartFile file, String title,
LocalDateTime capturedAt) {
validateOwner(owner);
validateFile(file);
byte[] bytes;
try {
@@ -109,23 +107,14 @@ public class PersonalIngestionService {
} 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);
return create(owner, "FILE", cleanTitle(title, originalName), suffix(originalName),
cleanMime(file.getContentType()), bytes, capturedAt, List.of());
}
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,
@@ -138,37 +127,50 @@ public class PersonalIngestionService {
}
}
private ItemCreatedResponse persist(PersonalOwner owner, long spaceId, String sourceType, String title,
String mimeType, long size, String hash, LocalDateTime capturedAt,
List<String> tags, SysOssVo uploaded) {
private ItemCreatedResponse create(PersonalOwner owner, String sourceType, String title, String suffix,
String mimeType, byte[] bytes, LocalDateTime capturedAt, List<String> 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 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);
}
private ItemCreatedResponse persist(PersonalOwner owner, long spaceId, long itemId, String sourceType,
String title, String mimeType, long size, String hash,
LocalDateTime capturedAt, List<String> tags, String objectKey,
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) {
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("""
@@ -179,24 +181,24 @@ public class PersonalIngestionService {
if (counterUpdated != 1) {
throw new ServiceException("PERSONAL_SPACE_NOT_AVAILABLE");
}
return new ItemCreatedResponse(key.longValue(), "QUEUED", null);
return new ItemCreatedResponse(itemId, "QUEUED", null);
} catch (RuntimeException ex) {
if (!deferredCleanup) {
cleanupOss(uploaded.getOssId());
cleanupOss(uploaded);
}
throw ex;
}
}
private ItemCreatedResponse duplicate(PersonalOwner owner, String hash) {
private ItemCreatedResponse duplicate(PersonalOwner owner, long spaceId, 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 = ?
where tenant_id = ? and owner_user_id = ? and space_id = ? and content_hash = ?
and status <> 'DELETED'
order by id desc
limit 1
""", owner.tenantId(), owner.userId(), hash);
""", owner.tenantId(), owner.userId(), spaceId, hash);
if (rows.isEmpty()) {
return null;
}
@@ -210,8 +212,7 @@ public class PersonalIngestionService {
throw new ServiceException("PERSONAL_FILE_EMPTY");
}
validateSize(file.getSize());
String suffix = suffix(file.getOriginalFilename());
if (!SUPPORTED_FILE_SUFFIXES.contains(suffix)) {
if (!SUPPORTED_FILE_SUFFIXES.contains(suffix(file.getOriginalFilename()))) {
throw new ServiceException("PERSONAL_FILE_UNSUPPORTED");
}
}
@@ -236,32 +237,31 @@ public class PersonalIngestionService {
@Override
public void afterCompletion(int status) {
if (status != STATUS_COMMITTED) {
cleanupRolledBackOss(uploaded);
cleanupOss(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());
}
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());
}
cleanupOss(uploaded.getOssId());
}
private void cleanupOss(long ossId) {
private String personalOssExt(long itemId) {
try {
ossService.deleteWithValidByIds(List.of(ossId), false);
} catch (RuntimeException cleanupError) {
log.warn("Unable to clean personal OSS object id={}", ossId);
return objectMapper.writeValueAsString(Map.of("source", "personal", "itemId", itemId));
} catch (JsonProcessingException ex) {
throw new ServiceException("PERSONAL_OSS_BIND_FAILED");
}
}
@@ -280,6 +280,32 @@ public class PersonalIngestionService {
}
}
private static String objectKey(PersonalOwner owner, long itemId, String suffix) {
validateOwner(owner);
long safeItemId = 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;
}
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");
}
}
private static long positiveId(long id) {
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));
@@ -313,11 +339,6 @@ public class PersonalIngestionService {
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";
@@ -327,55 +348,62 @@ public class PersonalIngestionService {
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;
public interface PersonalObjectStore {
SysOssVo upload(PersonalOwner owner, long itemId, String objectKey, String suffix, String mimeType,
byte[] bytes);
private SafeMultipartFile(String fileName, String contentType, byte[] bytes) {
this.fileName = fileName;
this.contentType = contentType;
this.bytes = bytes.clone();
void deletePhysical(SysOssVo uploaded);
}
private static final class DefaultPersonalObjectStore implements PersonalObjectStore {
private final JdbcTemplate jdbcTemplate;
private DefaultPersonalObjectStore(JdbcTemplate jdbcTemplate) {
this.jdbcTemplate = jdbcTemplate;
}
@Override
public String getName() {
return "file";
public SysOssVo upload(PersonalOwner owner, long itemId, String objectKey, String suffix, String mimeType,
byte[] bytes) {
OssClient storage = OssFactory.instance();
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 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);
public void deletePhysical(SysOssVo uploaded) {
if (uploaded.getService() != null && uploaded.getFileName() != null) {
OssFactory.instance(uploaded.getService()).delete(uploaded.getFileName());
}
}
}
}
@@ -90,7 +90,8 @@ public class PersonalIngestionWorker {
}
try {
StoredObject stored = objectReader.read(item.ossId(), maxInputBytes);
StoredObject stored = objectReader.read(
item.ossId(), ownerObjectPrefix(item), item.ownerUserId(), maxInputBytes);
ParsedDocument document = parser.parse(stored.fileName(), item.mimeType(), stored.bytes());
List<String> chunks = document.chunks(CHUNK_SIZE, CHUNK_OVERLAP);
if (chunks.isEmpty()) {
@@ -156,10 +157,12 @@ public class PersonalIngestionWorker {
}
private static StoredObjectReader defaultReader(ISysOssService ossService) {
return (ossId, maxBytes) -> {
return (ossId, expectedPrefix, ownerUserId, maxBytes) -> {
SysOssVo object = ossService.getById(ossId);
if (object == null || object.getFileName() == null || object.getFileName().isBlank()
|| object.getService() == null || object.getService().isBlank()) {
|| object.getService() == null || object.getService().isBlank()
|| object.getCreateBy() == null || object.getCreateBy() != ownerUserId
|| !object.getFileName().startsWith(expectedPrefix)) {
throw new IOException("personal source object is unavailable");
}
OssClient storage = OssFactory.instance(object.getService());
@@ -246,9 +249,13 @@ public class PersonalIngestionWorker {
return safe.isBlank() ? "personal-object" : safe;
}
private static String ownerObjectPrefix(Item item) {
return "personal/" + item.tenantId() + "/" + item.ownerUserId() + "/" + item.id() + "/";
}
@FunctionalInterface
public interface StoredObjectReader {
StoredObject read(long ossId, long maxBytes) throws Exception;
StoredObject read(long ossId, String expectedPrefix, long ownerUserId, long maxBytes) throws Exception;
}
public record StoredObject(String fileName, byte[] bytes) {
@@ -1,9 +1,11 @@
package org.dromara.aihr.personal;
import com.fasterxml.jackson.databind.JsonNode;
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.PersonalIngestionService.PersonalObjectStore;
import org.dromara.aihr.personal.service.PersonalSpaceService;
import org.dromara.aihr.personal.support.PersonalKnowledgeProperties;
import org.dromara.aihr.personal.support.PersonalOwner;
@@ -16,10 +18,7 @@ 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;
@@ -32,17 +31,20 @@ import java.time.LocalDateTime;
import java.util.List;
import java.util.Map;
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.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.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;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
@@ -54,224 +56,240 @@ class PersonalIngestionServiceTest {
@Test
void failedItemRetriesThroughOwnerScopedQueuedState() {
JdbcTemplate jdbc = mock(JdbcTemplate.class);
when(jdbc.update(contains("status = 'QUEUED'"), eq("000000"), eq(101L), eq(9L)))
Fixture fixture = fixture(100L);
when(fixture.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);
fixture.service.retry(OWNER, 9L);
verify(jdbc).update(contains("status = 'QUEUED'"), eq("000000"), eq(101L), eq(9L));
verify(fixture.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));
Fixture fixture = fixture(100L);
ServiceException error = assertThrows(ServiceException.class, () -> service.retry(OWNER, 9L));
ServiceException error = assertThrows(ServiceException.class, () -> fixture.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(
void invalidFilesAreRejectedBeforeObjectStorage() {
Fixture fixture = fixture(100L);
MockMultipartFile large = 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(
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(oss);
verifyNoInteractions(jdbc);
verifyNoInteractions(fixture.store, fixture.oss, fixture.jdbc, fixture.spaces);
}
@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);
void textCreationLocksBeforeDedupeAndUsesIsolatedObjectKeyAndExtBinding() throws Exception {
Fixture fixture = fixture(100L);
stubSuccessfulCreate(fixture, OWNER, 7L, 18L, 81L);
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);
ItemCreatedResponse response = fixture.service.createText(OWNER,
new TextItemRequest("周报", "保洁巡检记录", LocalDateTime.of(2026, 7, 12, 9, 0), List.of("保洁")));
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));
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<String> key = ArgumentCaptor.forClass(String.class);
verify(fixture.store).upload(eq(OWNER), eq(100L), key.capture(), eq("txt"), eq("text/plain"),
any(byte[].class));
assertTrue(key.getValue().matches("personal/000000/101/100/[0-9a-f]{32}\\.txt"));
ArgumentCaptor<String> 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));
}
@Test
void duplicateForSameOwnerReturnsExistingItemWithoutReservationOrUpload() {
JdbcTemplate jdbc = mock(JdbcTemplate.class);
PersonalSpaceService spaces = mock(PersonalSpaceService.class);
ISysOssService oss = mock(ISysOssService.class);
void sameOwnerDuplicateIsCheckedUnderLockAndDoesNotUploadOrIncrementCounters() {
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(Map.of("id", 77L, "status", "READY")));
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);
ItemCreatedResponse response = fixture.service.createFile(OWNER, file, null, null);
assertEquals(77L, response.itemId());
assertEquals(77L, response.duplicateOf());
assertEquals("READY", response.status());
verifyNoInteractions(spaces, oss);
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 +"),
any(), any(), any(), any());
}
@Test
void failedPersistenceRegistersBestEffortOssCleanup() {
JdbcTemplate jdbc = mock(JdbcTemplate.class);
PersonalSpaceService spaces = mock(PersonalSpaceService.class);
ISysOssService oss = mock(ISysOssService.class);
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.<Long>getArgument(1),
invocation.getArgument(2)));
stubPersistence(fixture.jdbc);
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)))
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()))
.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 service = transactionalProxy(service(jdbc, spaces, oss), transactions);
PersonalIngestionService proxy = transactionalProxy(fixture.service, transactions);
MockMultipartFile file = new MockMultipartFile(
"file", "notes.txt", "text/plain", "same".getBytes(StandardCharsets.UTF_8));
assertThrows(IllegalStateException.class, () -> service.createFile(OWNER, file, null, null));
assertThrows(IllegalStateException.class, () -> proxy.createFile(OWNER, file, null, null));
verify(oss).deleteWithValidByIds(eq(List.of(82L)), eq(false));
assertEquals(1, transactions.rollbacks);
verify(fixture.store).deletePhysical(uploaded);
verify(fixture.oss).deleteWithValidByIds(List.of(81L), false);
}
@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));
void reservationAndCountersShareOneOuterTransaction() {
Fixture fixture = fixture(100L);
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 -> {
when(fixture.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)))
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 service = transactionalProxy(service(jdbc, spaces, oss), transactions);
PersonalIngestionService proxy = transactionalProxy(fixture.service, transactions);
MockMultipartFile file = new MockMultipartFile(
"file", "notes.txt", "text/plain", "same".getBytes(StandardCharsets.UTF_8));
service.createFile(OWNER, file, null, null);
proxy.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());
@Test
void unsafeTenantIsRejectedBeforeStorage() {
Fixture fixture = fixture(100L);
PersonalOwner unsafe = new PersonalOwner("../000000", 101L, 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);
}
private static PersonalKnowledgeProperties properties() {
return new PersonalKnowledgeProperties();
private static Fixture fixture(long itemId) {
return fixture(() -> itemId);
}
private static SysOssVo oss(long id) {
private static Fixture fixture(java.util.function.LongSupplier itemIds) {
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);
}
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 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 SysOssVo oss(long id, String fileName) {
SysOssVo result = new SysOssVo();
result.setOssId(id);
result.setFileName("personal/random.txt");
result.setOriginalName("random.txt");
result.setFileName(fileName);
result.setOriginalName(fileName.substring(fileName.lastIndexOf('/') + 1));
result.setService("minio");
result.setUrl("https://private.invalid/random.txt");
result.setUrl("https://private.invalid/" + fileName);
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);
@@ -284,6 +302,10 @@ class PersonalIngestionServiceTest {
return (PersonalIngestionService) factory.getProxy();
}
private record Fixture(JdbcTemplate jdbc, PersonalSpaceService spaces, ISysOssService oss,
PersonalObjectStore store, PersonalIngestionService service) {
}
private static final class TestTransactionManager extends AbstractPlatformTransactionManager {
private int commits;
private int rollbacks;
@@ -44,8 +44,12 @@ class PersonalIngestionWorkerTest {
.thenReturn(1);
PersonalIngestionWorker worker = PersonalIngestionWorker.forTest(
jdbc, mock(ISysOssService.class), parser, transactions,
(ossId, maxBytes) -> new PersonalIngestionWorker.StoredObject(
"notes.txt", "same".getBytes(StandardCharsets.UTF_8)));
(ossId, prefix, ownerUserId, maxBytes) -> {
assertTrue(prefix.equals("personal/000000/101/9/"));
assertTrue(ownerUserId == 101L);
return new PersonalIngestionWorker.StoredObject(
"notes.txt", "same".getBytes(StandardCharsets.UTF_8));
});
assertTrue(worker.processNext());
@@ -61,7 +65,8 @@ class PersonalIngestionWorkerTest {
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}));
(ossId, prefix, ownerUserId, maxBytes) ->
new PersonalIngestionWorker.StoredObject("notes.txt", new byte[]{1}));
assertFalse(worker.processNext());
@@ -80,7 +85,8 @@ class PersonalIngestionWorkerTest {
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}));
(ossId, prefix, ownerUserId, maxBytes) ->
new PersonalIngestionWorker.StoredObject("notes.txt", new byte[]{1}));
assertTrue(worker.processNext());