fix(personal): make ingestion crash recoverable

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