fix(personal): fence ingestion workers and cleanup

This commit is contained in:
2026-07-12 03:55:15 +08:00
parent 4a35696d60
commit f4847f9ef7
6 changed files with 378 additions and 119 deletions
@@ -44,6 +44,7 @@ import java.util.regex.Pattern;
public class PersonalIngestionService { public class PersonalIngestionService {
private static final String ITEM_NOT_FOUND = "PERSONAL_ITEM_NOT_FOUND"; private static final String ITEM_NOT_FOUND = "PERSONAL_ITEM_NOT_FOUND";
private static final int MIN_UPLOAD_CLEANUP_AGE_MINUTES = 5;
private static final Pattern SAFE_TENANT = Pattern.compile("[A-Za-z0-9_-]{1,20}"); private static final Pattern SAFE_TENANT = Pattern.compile("[A-Za-z0-9_-]{1,20}");
private static final Set<String> SUPPORTED_FILE_SUFFIXES = Set.of( private static final Set<String> SUPPORTED_FILE_SUFFIXES = Set.of(
"txt", "md", "markdown", "pdf", "doc", "docx", "xls", "xlsx", "ppt", "pptx" "txt", "md", "markdown", "pdf", "doc", "docx", "xls", "xlsx", "ppt", "pptx"
@@ -77,6 +78,7 @@ public class PersonalIngestionService {
this.objectStore = objectStore; this.objectStore = objectStore;
this.idSupplier = idSupplier; this.idSupplier = idSupplier;
this.phaseTransaction = phaseTransaction; this.phaseTransaction = phaseTransaction;
validateRecoveryWindows(properties);
} }
public static PersonalIngestionService forTest(JdbcTemplate jdbcTemplate, PersonalSpaceService spaceService, public static PersonalIngestionService forTest(JdbcTemplate jdbcTemplate, PersonalSpaceService spaceService,
@@ -144,8 +146,9 @@ public class PersonalIngestionService {
long itemId = positiveId(idSupplier.getAsLong()); long itemId = positiveId(idSupplier.getAsLong());
long ossId = positiveId(idSupplier.getAsLong()); long ossId = positiveId(idSupplier.getAsLong());
String objectKey = objectKey(owner, itemId, suffix); String objectKey = objectKey(owner, itemId, suffix);
String uploadToken = UUID.randomUUID().toString();
UploadIntent draft = new UploadIntent(owner, 0L, itemId, ossId, objectKey, suffix, mimeType, UploadIntent draft = new UploadIntent(owner, 0L, itemId, ossId, objectKey, suffix, mimeType,
bytes.length, serviceKey); bytes.length, serviceKey, uploadToken);
String hash = sha256(bytes); String hash = sha256(bytes);
PhaseOne phaseOne = phaseTransaction.execute(status -> phaseOne( PhaseOne phaseOne = phaseTransaction.execute(status -> phaseOne(
@@ -158,13 +161,18 @@ public class PersonalIngestionService {
} }
UploadIntent intent = phaseOne.intent(); UploadIntent intent = phaseOne.intent();
String url;
try { try {
String url = objectStore.uploadPhysical(intent.serviceKey(), intent.objectKey(), mimeType, bytes); url = objectStore.uploadPhysical(intent.serviceKey(), intent.objectKey(), mimeType, bytes);
UploadIntent finalIntent = intent; } catch (RuntimeException ex) {
phaseTransaction.executeWithoutResult(status -> activate(finalIntent, url)); beginCleanup(intent, null);
throw ex;
}
try {
activate(intent, url);
return new ItemCreatedResponse(intent.itemId(), "QUEUED", null); return new ItemCreatedResponse(intent.itemId(), "QUEUED", null);
} catch (RuntimeException ex) { } catch (RuntimeException ex) {
cleanupIntent(intent, false, null); reconcileActivationFailure(intent);
throw ex; throw ex;
} }
} }
@@ -184,7 +192,7 @@ public class PersonalIngestionService {
create_time, create_by, update_time, update_by, service) create_time, create_by, update_time, update_by, service)
values (?, ?, ?, ?, ?, '', ?, now(), ?, now(), ?, ?) values (?, ?, ?, ?, ?, '', ?, now(), ?, now(), ?, ?)
""", intent.ossId(), intent.owner().tenantId(), intent.objectKey(), safeName, """, intent.ossId(), intent.owner().tenantId(), intent.objectKey(), safeName,
"." + intent.suffix(), uploadExt(intent.itemId(), "PENDING"), intent.owner().userId(), "." + intent.suffix(), uploadExt(intent, "PENDING"), intent.owner().userId(),
intent.owner().userId(), intent.serviceKey()); intent.owner().userId(), intent.serviceKey());
int itemInserted = jdbcTemplate.update(""" int itemInserted = jdbcTemplate.update("""
insert into aihr_personal_item insert into aihr_personal_item
@@ -209,6 +217,23 @@ public class PersonalIngestionService {
if (url == null || url.isBlank()) { if (url == null || url.isBlank()) {
throw new ServiceException("PERSONAL_OSS_UPLOAD_FAILED"); throw new ServiceException("PERSONAL_OSS_UPLOAD_FAILED");
} }
for (int attempt = 0; attempt < 2; attempt++) {
Integer activated = phaseTransaction.execute(status -> activateOnce(intent, url));
if (activated != null && activated == 1) {
return;
}
UploadState current = uploadState(intent);
if (current == UploadState.READY) {
return;
}
if (current != UploadState.PENDING) {
throw new ServiceException("PERSONAL_OSS_ACTIVATION_FAILED");
}
}
throw new ServiceException("PERSONAL_OSS_ACTIVATION_FAILED");
}
private int activateOnce(UploadIntent intent, String url) {
int activated = jdbcTemplate.update(""" int activated = jdbcTemplate.update("""
update sys_oss o 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 join aihr_personal_item i on i.oss_id = o.oss_id and binary i.tenant_id = binary o.tenant_id
@@ -216,43 +241,63 @@ public class PersonalIngestionService {
set o.url = ?, o.ext1 = ?, o.update_time = now(), o.update_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 = ? where o.tenant_id = ? and o.oss_id = ? and o.create_by = ? and o.file_name = ?
and i.id = ? and i.status = 'QUEUED' and i.id = ? and i.status = 'QUEUED'
and json_unquote(json_extract(o.ext1, '$.uploadToken')) = ?
and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'PENDING' and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'PENDING'
""", url, uploadExt(intent.itemId(), "READY"), intent.owner().userId(), """, url, uploadExt(intent, "READY"), intent.owner().userId(),
intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(), intent.itemId()); intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(), intent.itemId(),
if (activated != 1) { intent.uploadToken());
throw new ServiceException("PERSONAL_OSS_ACTIVATION_FAILED"); return activated;
}
} }
@Scheduled(fixedDelayString = "${aihr.personal.upload-cleanup-delay-ms:60000}") @Scheduled(fixedDelayString = "${aihr.personal.upload-cleanup-delay-ms:60000}")
public void recoverStaleUploadIntents() { public void recoverStaleUploadIntents() {
LocalDateTime cutoff = LocalDateTime.now().minusMinutes(properties.getUploadCleanupAgeMinutes()); LocalDateTime now = LocalDateTime.now();
List<Map<String, Object>> rows = jdbcTemplate.queryForList(""" LocalDateTime pendingCutoff = now.minusMinutes(properties.getUploadCleanupAgeMinutes());
for (Map<String, Object> row : staleUploadRows("PENDING", pendingCutoff)) {
beginCleanup(intent(row), pendingCutoff);
}
LocalDateTime cleaningCutoff = now.minusMinutes(properties.getCleanupFinalizeGraceMinutes());
for (Map<String, Object> row : staleUploadRows("CLEANING", cleaningCutoff)) {
finalizeStaleCleanup(intent(row), cleaningCutoff);
}
}
private List<Map<String, Object>> staleUploadRows(String state, LocalDateTime cutoff) {
return jdbcTemplate.queryForList("""
select i.tenant_id, i.owner_user_id, i.space_id, i.id item_id, i.oss_id, 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 i.size_bytes, i.mime_type, o.file_name, o.service,
json_unquote(json_extract(o.ext1, '$.uploadToken')) upload_token
from aihr_personal_item i 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 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 and o.create_by = i.owner_user_id
where i.status = 'QUEUED' and o.update_time < ? where i.status = 'QUEUED' and o.update_time < ?
and json_unquote(json_extract(o.ext1, '$.uploadState')) in ('PENDING', 'CLEANING') and json_unquote(json_extract(o.ext1, '$.uploadState')) = ?
and json_unquote(json_extract(o.ext1, '$.uploadToken')) is not null
order by o.update_time order by o.update_time
limit 20 limit 20
""", cutoff); """, cutoff, state);
for (Map<String, Object> row : rows) { }
UploadIntent intent = intent(row);
cleanupIntent(intent, true, cutoff); private void beginCleanup(UploadIntent intent, LocalDateTime cutoff) {
Integer claimed = phaseTransaction.execute(status -> claimState(
intent, UploadState.PENDING, UploadState.CLEANING, cutoff));
if (claimed != null && claimed == 1) {
deleteKnownObject(intent);
return;
}
UploadState current = uploadState(intent);
if (current == UploadState.CLEANING || current == UploadState.MISSING) {
deleteKnownObject(intent);
} }
} }
private void cleanupIntent(UploadIntent intent, boolean includeCleaning, LocalDateTime cutoff) { private void finalizeStaleCleanup(UploadIntent intent, LocalDateTime cutoff) {
Integer claimed = phaseTransaction.execute(status -> claimCleanup(intent, includeCleaning, cutoff)); Integer claimed = phaseTransaction.execute(status -> claimState(
intent, UploadState.CLEANING, UploadState.CLEANING, cutoff));
if (claimed == null || claimed != 1) { if (claimed == null || claimed != 1) {
return; return;
} }
try { if (!deleteKnownObject(intent)) {
objectStore.deletePhysical(intent.serviceKey(), intent.objectKey());
} catch (RuntimeException ex) {
log.warn("Personal upload-intent physical cleanup failed itemId={}", intent.itemId());
return; return;
} }
try { try {
@@ -262,22 +307,63 @@ public class PersonalIngestionService {
} }
} }
private int claimCleanup(UploadIntent intent, boolean includeCleaning, LocalDateTime cutoff) { private int claimState(UploadIntent intent, UploadState expected, UploadState target, LocalDateTime cutoff) {
String states = includeCleaning ? "('PENDING', 'CLEANING')" : "('PENDING')";
String cutoffClause = cutoff == null ? "" : " and update_time < ?"; String cutoffClause = cutoff == null ? "" : " and update_time < ?";
String sql = """ String sql = """
update sys_oss update sys_oss
set ext1 = json_set(ext1, '$.uploadState', 'CLEANING'), update_time = now() set ext1 = json_set(ext1, '$.uploadState', ?), update_time = now()
where tenant_id = ? and oss_id = ? and create_by = ? and file_name = ? 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, '$.itemId')) = ?
and json_unquote(json_extract(ext1, '$.uploadState')) in %s%s and json_unquote(json_extract(ext1, '$.uploadToken')) = ?
""".formatted(states, cutoffClause); and json_unquote(json_extract(ext1, '$.uploadState')) = ?%s
""".formatted(cutoffClause);
if (cutoff == null) { if (cutoff == null) {
return jdbcTemplate.update(sql, intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), return jdbcTemplate.update(sql, target.name(), intent.owner().tenantId(), intent.ossId(),
intent.objectKey(), String.valueOf(intent.itemId())); intent.owner().userId(), intent.objectKey(), String.valueOf(intent.itemId()), intent.uploadToken(),
expected.name());
}
return jdbcTemplate.update(sql, target.name(), intent.owner().tenantId(), intent.ossId(),
intent.owner().userId(), intent.objectKey(), String.valueOf(intent.itemId()), intent.uploadToken(),
expected.name(), cutoff);
}
private void reconcileActivationFailure(UploadIntent intent) {
UploadState current = uploadState(intent);
if (current == UploadState.PENDING) {
beginCleanup(intent, null);
} else if (current == UploadState.CLEANING || current == UploadState.MISSING) {
deleteKnownObject(intent);
}
}
private UploadState uploadState(UploadIntent intent) {
List<Map<String, Object>> rows = jdbcTemplate.queryForList("""
select json_unquote(json_extract(ext1, '$.uploadState')) upload_state
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, '$.uploadToken')) = ?
limit 1
""", intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(),
String.valueOf(intent.itemId()), intent.uploadToken());
if (rows.isEmpty()) {
return UploadState.MISSING;
}
try {
return UploadState.valueOf(String.valueOf(rows.get(0).get("upload_state")));
} catch (IllegalArgumentException ex) {
throw new ServiceException("PERSONAL_UPLOAD_STATE_INVALID");
}
}
private boolean deleteKnownObject(UploadIntent intent) {
try {
objectStore.deletePhysical(intent.serviceKey(), intent.objectKey());
return true;
} catch (RuntimeException ex) {
log.warn("Personal upload-intent physical cleanup failed itemId={}", intent.itemId());
return false;
} }
return jdbcTemplate.update(sql, intent.owner().tenantId(), intent.ossId(), intent.owner().userId(),
intent.objectKey(), String.valueOf(intent.itemId()), cutoff);
} }
private void finalizeCleanup(UploadIntent intent) { private void finalizeCleanup(UploadIntent intent) {
@@ -304,9 +390,10 @@ public class PersonalIngestionService {
delete from sys_oss delete from sys_oss
where tenant_id = ? and oss_id = ? and create_by = ? and file_name = ? 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, '$.itemId')) = ?
and json_unquote(json_extract(ext1, '$.uploadToken')) = ?
and json_unquote(json_extract(ext1, '$.uploadState')) = 'CLEANING' and json_unquote(json_extract(ext1, '$.uploadState')) = 'CLEANING'
""", intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(), """, intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(),
String.valueOf(intent.itemId())); String.valueOf(intent.itemId()), intent.uploadToken());
if (itemDeleted != 1 || counterUpdated != 1 || ossDeleted != 1) { if (itemDeleted != 1 || counterUpdated != 1 || ossDeleted != 1) {
throw new ServiceException("PERSONAL_UPLOAD_CLEANUP_FAILED"); throw new ServiceException("PERSONAL_UPLOAD_CLEANUP_FAILED");
} }
@@ -323,10 +410,11 @@ public class PersonalIngestionService {
return new ItemCreatedResponse(id, String.valueOf(rows.get(0).get("status")), id); return new ItemCreatedResponse(id, String.valueOf(rows.get(0).get("status")), id);
} }
private String uploadExt(long itemId, String state) { private String uploadExt(UploadIntent intent, String state) {
try { try {
return objectMapper.writeValueAsString(Map.of( return objectMapper.writeValueAsString(Map.of(
"source", "personal", "itemId", itemId, "uploadState", state)); "source", "personal", "itemId", intent.itemId(), "uploadState", state,
"uploadToken", intent.uploadToken()));
} catch (JsonProcessingException ex) { } catch (JsonProcessingException ex) {
throw new ServiceException("PERSONAL_OSS_BIND_FAILED"); throw new ServiceException("PERSONAL_OSS_BIND_FAILED");
} }
@@ -346,7 +434,8 @@ public class PersonalIngestionService {
PersonalOwner owner = new PersonalOwner(String.valueOf(row.get("tenant_id")), number(row, "owner_user_id"), null); 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"), 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("file_name")), suffix(String.valueOf(row.get("file_name"))),
String.valueOf(row.get("mime_type")), number(row, "size_bytes"), String.valueOf(row.get("service"))); String.valueOf(row.get("mime_type")), number(row, "size_bytes"), String.valueOf(row.get("service")),
String.valueOf(row.get("upload_token")));
} }
private static long number(Map<String, Object> row, String key) { private static long number(Map<String, Object> row, String key) {
@@ -376,6 +465,13 @@ public class PersonalIngestionService {
return template; return template;
} }
private static void validateRecoveryWindows(PersonalKnowledgeProperties properties) {
if (properties.getUploadCleanupAgeMinutes() < MIN_UPLOAD_CLEANUP_AGE_MINUTES
|| properties.getCleanupFinalizeGraceMinutes() <= 0) {
throw new IllegalArgumentException("invalid personal upload recovery windows");
}
}
private static String objectKey(PersonalOwner owner, long itemId, String suffix) { private static String objectKey(PersonalOwner owner, long itemId, String suffix) {
validateOwner(owner); positiveId(itemId); validateOwner(owner); positiveId(itemId);
String safeSuffix = suffix == null ? "" : suffix.toLowerCase(Locale.ROOT); String safeSuffix = suffix == null ? "" : suffix.toLowerCase(Locale.ROOT);
@@ -447,7 +543,9 @@ public class PersonalIngestionService {
return result.getUrl(); return result.getUrl();
} }
@Override public void deletePhysical(String serviceKey, String objectKey) { @Override public void deletePhysical(String serviceKey, String objectKey) {
OssClient storage = clients.get(serviceKey); requirePrivate(storage); storage.delete(objectKey); OssClient storage = clients.get(serviceKey);
if (storage == null) throw new ServiceException("PERSONAL_OSS_UNAVAILABLE");
storage.delete(objectKey);
} }
} }
@@ -459,10 +557,14 @@ public class PersonalIngestionService {
} }
private record PhaseOne(ItemCreatedResponse duplicate, UploadIntent intent) {} private record PhaseOne(ItemCreatedResponse duplicate, UploadIntent intent) {}
private enum UploadState { PENDING, READY, CLEANING, MISSING }
private record UploadIntent(PersonalOwner owner, long spaceId, long itemId, long ossId, String objectKey, private record UploadIntent(PersonalOwner owner, long spaceId, long itemId, long ossId, String objectKey,
String suffix, String mimeType, long sizeBytes, String serviceKey) { String suffix, String mimeType, long sizeBytes, String serviceKey,
String uploadToken) {
private UploadIntent withSpaceId(long value) { private UploadIntent withSpaceId(long value) {
return new UploadIntent(owner, value, itemId, ossId, objectKey, suffix, mimeType, sizeBytes, serviceKey); return new UploadIntent(owner, value, itemId, ossId, objectKey, suffix, mimeType, sizeBytes, serviceKey,
uploadToken);
} }
} }
} }
@@ -87,29 +87,44 @@ public class PersonalIngestionWorker {
public void recoverStaleParsing() { public void recoverStaleParsing() {
java.time.LocalDateTime cutoff = java.time.LocalDateTime.now().minusMinutes(parsingLeaseMinutes); java.time.LocalDateTime cutoff = java.time.LocalDateTime.now().minusMinutes(parsingLeaseMinutes);
jdbcTemplate.update(""" List<Map<String, Object>> stale = jdbcTemplate.queryForList("""
update aihr_personal_item i select i.id, i.tenant_id, i.owner_user_id, i.attempt_count
join sys_oss o on o.oss_id = i.oss_id and binary o.tenant_id = binary i.tenant_id 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 and o.create_by = i.owner_user_id
set i.status = 'FAILED', i.error_code = 'PERSONAL_PARSE_RETRY_EXHAUSTED', where i.status = 'PARSING' and i.update_time < ?
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' and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'READY'
""", maxParseAttempts, cutoff); order by i.update_time
jdbcTemplate.update(""" limit 100
update aihr_personal_item i """, cutoff);
join sys_oss o on o.oss_id = i.oss_id and binary o.tenant_id = binary i.tenant_id for (Map<String, Object> row : stale) {
and o.create_by = i.owner_user_id long id = number(row, "id");
set i.status = 'QUEUED', i.error_code = null, i.error_message = null, i.update_time = now() String tenantId = String.valueOf(row.get("tenant_id"));
where i.status = 'PARSING' and i.attempt_count < ? and i.update_time < ? long ownerUserId = number(row, "owner_user_id");
and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'READY' int attemptVersion = Math.toIntExact(number(row, "attempt_count"));
""", maxParseAttempts, cutoff); if (attemptVersion >= maxParseAttempts) {
jdbcTemplate.update("""
update aihr_personal_item
set status = 'FAILED', error_code = 'PERSONAL_PARSE_RETRY_EXHAUSTED',
error_message = '资料处理重试次数已用尽', update_time = now()
where tenant_id = ? and owner_user_id = ? and id = ?
and status = 'PARSING' and attempt_count = ? and update_time < ?
""", tenantId, ownerUserId, id, attemptVersion, cutoff);
} else {
jdbcTemplate.update("""
update aihr_personal_item
set status = 'QUEUED', error_code = null, error_message = null, update_time = now()
where tenant_id = ? and owner_user_id = ? and id = ?
and status = 'PARSING' and attempt_count = ? and update_time < ?
""", tenantId, ownerUserId, id, attemptVersion, cutoff);
}
}
} }
public boolean processNext() { public boolean processNext() {
List<Map<String, Object>> queued = jdbcTemplate.queryForList(""" List<Map<String, Object>> queued = jdbcTemplate.queryForList("""
select i.id, i.tenant_id, i.space_id, i.owner_user_id, i.source_type, i.title, select i.id, i.tenant_id, i.space_id, i.owner_user_id, i.source_type, i.title,
i.oss_id, i.mime_type, i.tags_json, i.captured_at i.oss_id, i.mime_type, i.tags_json, i.captured_at, i.attempt_count
from aihr_personal_item i 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 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 and o.create_by = i.owner_user_id
@@ -127,15 +142,16 @@ public class PersonalIngestionWorker {
set status = 'PARSING', attempt_count = attempt_count + 1, set status = 'PARSING', attempt_count = attempt_count + 1,
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 attempt_count = ?
and exists (select 1 from sys_oss o where o.oss_id = aihr_personal_item.oss_id 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 binary o.tenant_id = binary aihr_personal_item.tenant_id
and o.create_by = aihr_personal_item.owner_user_id and o.create_by = aihr_personal_item.owner_user_id
and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'READY') and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'READY')
""", item.tenantId(), item.ownerUserId(), item.id()); """, item.tenantId(), item.ownerUserId(), item.id(), item.attemptCount());
if (claimed != 1) { if (claimed != 1) {
return false; return false;
} }
int attemptVersion = Math.addExact(item.attemptCount(), 1);
try { try {
StoredObject stored = objectReader.read( StoredObject stored = objectReader.read(
@@ -147,7 +163,7 @@ public class PersonalIngestionWorker {
KnowledgeDocumentParser.Failure.EMPTY, "document contains no text"); KnowledgeDocumentParser.Failure.EMPTY, "document contains no text");
} }
transactionTemplate.execute(status -> { transactionTemplate.execute(status -> {
persistSuccess(item, document, chunks); persistSuccess(item, attemptVersion, document, chunks);
return null; return null;
}); });
} catch (Exception ex) { } catch (Exception ex) {
@@ -156,15 +172,26 @@ public class PersonalIngestionWorker {
update aihr_personal_item update aihr_personal_item
set status = 'FAILED', error_code = ?, error_message = ?, update_time = now() set status = 'FAILED', error_code = ?, error_message = ?, update_time = now()
where tenant_id = ? and owner_user_id = ? and id = ? where tenant_id = ? and owner_user_id = ? and id = ?
and status = 'PARSING' and status = 'PARSING' and attempt_count = ?
""", failure.code(), failure.message(), item.tenantId(), item.ownerUserId(), item.id()); """, failure.code(), failure.message(), item.tenantId(), item.ownerUserId(), item.id(),
attemptVersion);
log.warn("Personal ingestion failed itemId={} ownerUserId={} code={}", log.warn("Personal ingestion failed itemId={} ownerUserId={} code={}",
item.id(), item.ownerUserId(), failure.code()); item.id(), item.ownerUserId(), failure.code());
} }
return true; return true;
} }
private void persistSuccess(Item item, ParsedDocument document, List<String> chunks) { private void persistSuccess(Item item, int attemptVersion, ParsedDocument document, List<String> chunks) {
Map<String, Object> locked = jdbcTemplate.queryForMap("""
select status, attempt_count
from aihr_personal_item
where tenant_id = ? and owner_user_id = ? and id = ?
for update
""", item.tenantId(), item.ownerUserId(), item.id());
if (!"PARSING".equals(String.valueOf(locked.get("status")))
|| number(locked, "attempt_count") != attemptVersion) {
throw new IllegalStateException("personal item attempt changed while parsing");
}
jdbcTemplate.update(""" jdbcTemplate.update("""
delete from aihr_personal_fragment delete from aihr_personal_fragment
where tenant_id = ? and owner_user_id = ? and item_id = ? where tenant_id = ? and owner_user_id = ? and item_id = ?
@@ -197,8 +224,9 @@ public class PersonalIngestionWorker {
set status = 'READY', parsed_at = now(), summary = ?, tags_json = ?, set status = 'READY', parsed_at = now(), summary = ?, tags_json = ?,
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 = 'PARSING' and status = 'PARSING' and attempt_count = ?
""", summary(document.text()), item.tagsJson(), item.tenantId(), item.ownerUserId(), item.id()); """, summary(document.text()), item.tagsJson(), item.tenantId(), item.ownerUserId(), item.id(),
attemptVersion);
if (updated != 1) { if (updated != 1) {
throw new IllegalStateException("personal item state changed while parsing"); throw new IllegalStateException("personal item state changed while parsing");
} }
@@ -239,7 +267,8 @@ public class PersonalIngestionWorker {
String.valueOf(row.get("title")), String.valueOf(row.get("title")),
number(row, "oss_id"), number(row, "oss_id"),
String.valueOf(row.get("mime_type")), String.valueOf(row.get("mime_type")),
tags == null ? "[]" : String.valueOf(tags) tags == null ? "[]" : String.valueOf(tags),
Math.toIntExact(number(row, "attempt_count"))
); );
} }
@@ -330,7 +359,7 @@ public class PersonalIngestionWorker {
} }
private record Item(long id, String tenantId, long spaceId, long ownerUserId, String sourceType, private record Item(long id, String tenantId, long spaceId, long ownerUserId, String sourceType,
String title, long ossId, String mimeType, String tagsJson) { String title, long ossId, String mimeType, String tagsJson, int attemptCount) {
} }
private record Failure(String code, String message) { private record Failure(String code, String message) {
@@ -22,4 +22,5 @@ public class PersonalKnowledgeProperties {
private int parsingLeaseMinutes = 15; private int parsingLeaseMinutes = 15;
private int maxParseAttempts = 3; private int maxParseAttempts = 3;
private int uploadCleanupAgeMinutes = 15; private int uploadCleanupAgeMinutes = 15;
private int cleanupFinalizeGraceMinutes = 15;
} }
@@ -66,7 +66,7 @@ class PersonalIngestionServiceTest {
return "https://private.invalid/" + invocation.<String>getArgument(1); return "https://private.invalid/" + invocation.<String>getArgument(1);
}); });
when(fixture.jdbc.update(contains("set o.url ="), anyString(), anyString(), eq(101L), eq("000000"), when(fixture.jdbc.update(contains("set o.url ="), anyString(), anyString(), eq(101L), eq("000000"),
eq(101L), eq(101L), anyString(), eq(100L))).thenReturn(1); eq(101L), eq(101L), anyString(), eq(100L), anyString())).thenReturn(1);
ItemCreatedResponse response = fixture.service.createText(OWNER, ItemCreatedResponse response = fixture.service.createText(OWNER,
new TextItemRequest("周报", "保洁巡检记录", null, List.of("保洁"))); new TextItemRequest("周报", "保洁巡检记录", null, List.of("保洁")));
@@ -81,11 +81,13 @@ class PersonalIngestionServiceTest {
eq(key.getValue()), originalName.capture(), eq(".txt"), pendingExt.capture(), eq(101L), eq(101L), eq(key.getValue()), originalName.capture(), eq(".txt"), pendingExt.capture(), eq(101L), eq(101L),
eq("personal-private")); eq("personal-private"));
assertTrue(originalName.getValue().matches("[0-9a-f]{32}\\.txt")); assertTrue(originalName.getValue().matches("[0-9a-f]{32}\\.txt"));
assertUploadState(pendingExt.getValue(), "PENDING"); JsonNode pending = uploadExt(pendingExt.getValue(), "PENDING");
String uploadToken = pending.path("uploadToken").asText();
assertTrue(uploadToken.matches("[0-9a-f-]{36}"));
ArgumentCaptor<String> readyExt = ArgumentCaptor.forClass(String.class); ArgumentCaptor<String> readyExt = ArgumentCaptor.forClass(String.class);
verify(fixture.jdbc).update(contains("set o.url ="), anyString(), readyExt.capture(), eq(101L), verify(fixture.jdbc).update(contains("set o.url ="), anyString(), readyExt.capture(), eq(101L),
eq("000000"), eq(101L), eq(101L), eq(key.getValue()), eq(100L)); eq("000000"), eq(101L), eq(101L), eq(key.getValue()), eq(100L), eq(uploadToken));
assertUploadState(readyExt.getValue(), "READY"); assertEquals(uploadToken, uploadExt(readyExt.getValue(), "READY").path("uploadToken").asText());
assertEquals(2, fixture.transactions.commits); assertEquals(2, fixture.transactions.commits);
assertEquals(0, fixture.transactions.rollbacks); assertEquals(0, fixture.transactions.rollbacks);
} }
@@ -97,7 +99,7 @@ class PersonalIngestionServiceTest {
when(fixture.store.uploadPhysical(eq("personal-private"), anyString(), eq("text/plain"), any(byte[].class))) when(fixture.store.uploadPhysical(eq("personal-private"), anyString(), eq("text/plain"), any(byte[].class)))
.thenReturn("https://private.invalid/object"); .thenReturn("https://private.invalid/object");
when(fixture.jdbc.update(contains("set o.url ="), anyString(), anyString(), eq(101L), eq("000000"), when(fixture.jdbc.update(contains("set o.url ="), anyString(), anyString(), eq(101L), eq("000000"),
eq(101L), eq(101L), anyString(), eq(100L))).thenReturn(1); eq(101L), eq(101L), anyString(), eq(100L), anyString())).thenReturn(1);
MockMultipartFile file = new MockMultipartFile("file", "13800138000-secret.txt", "text/plain", MockMultipartFile file = new MockMultipartFile("file", "13800138000-secret.txt", "text/plain",
"secret".getBytes(StandardCharsets.UTF_8)); "secret".getBytes(StandardCharsets.UTF_8));
@@ -132,16 +134,14 @@ class PersonalIngestionServiceTest {
} }
@Test @Test
void uploadFailureClaimsCleaningDeletesPhysicalAndReversesCountersOnce() { void uploadFailureClaimsCleaningAndDeletesPhysicalButRetainsDurableIntent() {
Fixture fixture = fixture(); Fixture fixture = fixture();
stubPhaseOne(fixture, 4L); stubPhaseOne(fixture, 4L);
AtomicBoolean deleteInTransaction = new AtomicBoolean(true); AtomicBoolean deleteInTransaction = new AtomicBoolean(true);
when(fixture.store.uploadPhysical(eq("personal-private"), anyString(), eq("text/plain"), any(byte[].class))) when(fixture.store.uploadPhysical(eq("personal-private"), anyString(), eq("text/plain"), any(byte[].class)))
.thenThrow(new ServiceException("PERSONAL_OSS_UPLOAD_FAILED")); .thenThrow(new ServiceException("PERSONAL_OSS_UPLOAD_FAILED"));
when(fixture.jdbc.update(contains("json_set"), eq("000000"), eq(101L), eq(101L), anyString(), eq("100"))) when(fixture.jdbc.update(contains("json_set"), eq("CLEANING"), eq("000000"), eq(101L), eq(101L),
.thenReturn(1); anyString(), eq("100"), anyString(), eq("PENDING"))).thenReturn(1);
when(fixture.spaces.lockForUpdate(any(PersonalOwner.class))).thenReturn(7L);
stubFinalizeCleanup(fixture);
org.mockito.Mockito.doAnswer(invocation -> { org.mockito.Mockito.doAnswer(invocation -> {
deleteInTransaction.set(TransactionSynchronizationManager.isActualTransactionActive()); deleteInTransaction.set(TransactionSynchronizationManager.isActualTransactionActive());
return null; return null;
@@ -152,13 +152,12 @@ class PersonalIngestionServiceTest {
assertEquals("PERSONAL_OSS_UPLOAD_FAILED", error.getMessage()); assertEquals("PERSONAL_OSS_UPLOAD_FAILED", error.getMessage());
assertFalse(deleteInTransaction.get()); assertFalse(deleteInTransaction.get());
verify(fixture.jdbc).update(contains("status = 'DELETED'"), eq("000000"), eq(101L), eq(7L), eq(100L), verify(fixture.jdbc, never()).update(contains("status = 'DELETED'"), any(), any(), any(), any(), any());
eq(101L)); verify(fixture.jdbc, never()).update(contains("used_bytes = used_bytes -"), any(), any(), any(), any(),
verify(fixture.jdbc).update(contains("used_bytes = used_bytes -"), eq(4L), eq("000000"), eq(101L), any());
eq(7L), eq(4L)); verify(fixture.jdbc, never()).update(contains("delete from sys_oss"), any(), any(), any(), any(), any(),
verify(fixture.jdbc).update(contains("delete from sys_oss"), eq("000000"), eq(101L), eq(101L), any());
anyString(), eq("100")); assertEquals(2, fixture.transactions.commits);
assertEquals(3, fixture.transactions.commits);
assertEquals(0, fixture.transactions.rollbacks); assertEquals(0, fixture.transactions.rollbacks);
} }
@@ -166,17 +165,24 @@ class PersonalIngestionServiceTest {
void staleCleanupIsIdempotentAndDoesNotDecrementCountersTwice() { void staleCleanupIsIdempotentAndDoesNotDecrementCountersTwice() {
Fixture fixture = fixture(); Fixture fixture = fixture();
Map<String, Object> stale = staleIntent(); Map<String, Object> stale = staleIntent();
when(fixture.jdbc.queryForList(contains("uploadState')) in ('PENDING', 'CLEANING')"), when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("PENDING")))
any(LocalDateTime.class))).thenReturn(List.of(stale)); .thenReturn(List.of(stale), List.of(), List.of());
when(fixture.jdbc.update(contains("json_set"), eq("000000"), eq(101L), eq(101L), eq("personal/key.txt"), when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("CLEANING")))
eq("100"), any(LocalDateTime.class))).thenReturn(1, 0); .thenReturn(List.of(), List.of(stale), List.of(stale));
when(fixture.jdbc.update(contains("json_set"), eq("CLEANING"), eq("000000"), eq(101L), eq(101L),
eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("PENDING"), any(LocalDateTime.class)))
.thenReturn(1);
when(fixture.jdbc.update(contains("json_set"), eq("CLEANING"), eq("000000"), eq(101L), eq(101L),
eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("CLEANING"), any(LocalDateTime.class)))
.thenReturn(1, 0);
when(fixture.spaces.lockForUpdate(any(PersonalOwner.class))).thenReturn(7L); when(fixture.spaces.lockForUpdate(any(PersonalOwner.class))).thenReturn(7L);
stubFinalizeCleanup(fixture); stubFinalizeCleanup(fixture);
fixture.service.recoverStaleUploadIntents();
fixture.service.recoverStaleUploadIntents(); fixture.service.recoverStaleUploadIntents();
fixture.service.recoverStaleUploadIntents(); fixture.service.recoverStaleUploadIntents();
verify(fixture.store, times(1)).deletePhysical("personal-private", "personal/key.txt"); verify(fixture.store, times(2)).deletePhysical("personal-private", "personal/key.txt");
verify(fixture.jdbc, times(1)).update(contains("used_bytes = used_bytes -"), eq(4L), eq("000000"), verify(fixture.jdbc, times(1)).update(contains("used_bytes = used_bytes -"), eq(4L), eq("000000"),
eq(101L), eq(7L), eq(4L)); eq(101L), eq(7L), eq(4L));
} }
@@ -223,7 +229,7 @@ class PersonalIngestionServiceTest {
} }
@Test @Test
void publicStorageIsRejectedAndPrivateStorageUploadsByPhysicalKey() { void uploadRequiresPrivatePolicyButCleanupSurvivesPolicyDrift() {
PersonalKnowledgeProperties properties = new PersonalKnowledgeProperties(); PersonalKnowledgeProperties properties = new PersonalKnowledgeProperties();
PersonalIngestionService.OssClientProvider clients = mock(PersonalIngestionService.OssClientProvider.class); PersonalIngestionService.OssClientProvider clients = mock(PersonalIngestionService.OssClientProvider.class);
OssClient publicClient = mock(OssClient.class); OssClient publicClient = mock(OssClient.class);
@@ -232,8 +238,8 @@ class PersonalIngestionServiceTest {
PersonalObjectStore publicStore = PersonalIngestionService.objectStoreForTest(properties, clients); PersonalObjectStore publicStore = PersonalIngestionService.objectStoreForTest(properties, clients);
assertEquals("PERSONAL_OSS_NOT_PRIVATE", assertEquals("PERSONAL_OSS_NOT_PRIVATE",
assertThrows(ServiceException.class, publicStore::requirePrivateService).getMessage()); assertThrows(ServiceException.class, publicStore::requirePrivateService).getMessage());
assertEquals("PERSONAL_OSS_NOT_PRIVATE", assertThrows(ServiceException.class, publicStore.deletePhysical("", "personal/key.txt");
() -> publicStore.deletePhysical("", "personal/key.txt")).getMessage()); verify(publicClient).delete("personal/key.txt");
verify(publicClient, never()).upload(any(java.io.InputStream.class), anyString(), anyLong(), anyString()); verify(publicClient, never()).upload(any(java.io.InputStream.class), anyString(), anyLong(), anyString());
properties.setOssConfigKey(" personal-private "); properties.setOssConfigKey(" personal-private ");
@@ -251,7 +257,74 @@ class PersonalIngestionServiceTest {
verify(privateClient).upload(any(java.io.InputStream.class), eq("personal/key.txt"), eq(1L), eq("text/plain")); verify(privateClient).upload(any(java.io.InputStream.class), eq("personal/key.txt"), eq(1L), eq("text/plain"));
} }
@Test
void alreadyReadyActivationIsIdempotentAndNeverDeletesConfirmedObject() {
Fixture fixture = fixture();
stubPhaseOne(fixture, 4L);
when(fixture.store.uploadPhysical(eq("personal-private"), anyString(), eq("text/plain"), any(byte[].class)))
.thenReturn("https://private.invalid/object");
when(fixture.jdbc.queryForList(contains("select json_unquote"), eq("000000"), eq(101L), eq(101L),
anyString(), eq("100"), anyString())).thenReturn(List.of(Map.of("upload_state", "READY")));
ItemCreatedResponse response = fixture.service.createFile(OWNER,
new MockMultipartFile("file", "notes.txt", "text/plain", new byte[]{1, 2, 3, 4}), null, null);
assertEquals(100L, response.itemId());
verify(fixture.store, never()).deletePhysical(anyString(), anyString());
}
@Test
void pendingActivationRetriesWithSameTokenBeforeCleanup() {
Fixture fixture = fixture();
stubPhaseOne(fixture, 4L);
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), anyString())).thenReturn(0, 1);
when(fixture.jdbc.queryForList(contains("select json_unquote"), eq("000000"), eq(101L), eq(101L),
anyString(), eq("100"), anyString())).thenReturn(List.of(Map.of("upload_state", "PENDING")));
ItemCreatedResponse response = fixture.service.createFile(OWNER,
new MockMultipartFile("file", "notes.txt", "text/plain", new byte[]{1, 2, 3, 4}), null, null);
assertEquals(100L, response.itemId());
verify(fixture.jdbc, times(2)).update(contains("set o.url ="), anyString(), anyString(), eq(101L),
eq("000000"), eq(101L), eq(101L), anyString(), eq(100L), anyString());
verify(fixture.store, never()).deletePhysical(anyString(), anyString());
}
@Test
void cleaningActivationFailureDeletesKnownObjectWithoutReclaimingIntent() {
Fixture fixture = fixture();
stubPhaseOne(fixture, 4L);
when(fixture.store.uploadPhysical(eq("personal-private"), anyString(), eq("text/plain"), any(byte[].class)))
.thenReturn("https://private.invalid/object");
when(fixture.jdbc.queryForList(contains("select json_unquote"), eq("000000"), eq(101L), eq(101L),
anyString(), eq("100"), anyString())).thenReturn(List.of(Map.of("upload_state", "CLEANING")));
assertThrows(ServiceException.class, () -> fixture.service.createFile(OWNER,
new MockMultipartFile("file", "notes.txt", "text/plain", new byte[]{1, 2, 3, 4}), null, null));
verify(fixture.store).deletePhysical(eq("personal-private"), anyString());
verify(fixture.jdbc, never()).update(contains("json_set"), any(), any(), any(), any(), any(), any(), any(),
any());
}
@Test
void unsafeCleanupWindowsAreRejected() {
PersonalKnowledgeProperties properties = new PersonalKnowledgeProperties();
properties.setUploadCleanupAgeMinutes(4);
assertThrows(IllegalArgumentException.class, () -> fixture(properties));
properties.setUploadCleanupAgeMinutes(15);
properties.setCleanupFinalizeGraceMinutes(0);
assertThrows(IllegalArgumentException.class, () -> fixture(properties));
}
private static Fixture fixture() { private static Fixture fixture() {
return fixture(new PersonalKnowledgeProperties());
}
private static Fixture fixture(PersonalKnowledgeProperties properties) {
JdbcTemplate jdbc = mock(JdbcTemplate.class); JdbcTemplate jdbc = mock(JdbcTemplate.class);
PersonalSpaceService spaces = mock(PersonalSpaceService.class); PersonalSpaceService spaces = mock(PersonalSpaceService.class);
PersonalObjectStore store = mock(PersonalObjectStore.class); PersonalObjectStore store = mock(PersonalObjectStore.class);
@@ -260,7 +333,7 @@ class PersonalIngestionServiceTest {
TransactionTemplate template = new TransactionTemplate(transactions); TransactionTemplate template = new TransactionTemplate(transactions);
template.setPropagationBehavior(TransactionDefinition.PROPAGATION_REQUIRES_NEW); template.setPropagationBehavior(TransactionDefinition.PROPAGATION_REQUIRES_NEW);
PersonalIngestionService service = PersonalIngestionService.forTest(jdbc, spaces, PersonalIngestionService service = PersonalIngestionService.forTest(jdbc, spaces,
new PersonalKnowledgeProperties(), new ObjectMapper(), store, properties, new ObjectMapper(), store,
new AtomicLong(100L)::getAndIncrement, template); new AtomicLong(100L)::getAndIncrement, template);
return new Fixture(jdbc, spaces, store, service, transactions); return new Fixture(jdbc, spaces, store, service, transactions);
} }
@@ -284,22 +357,24 @@ class PersonalIngestionServiceTest {
when(fixture.jdbc.update(contains("used_bytes = used_bytes -"), eq(4L), eq("000000"), eq(101L), when(fixture.jdbc.update(contains("used_bytes = used_bytes -"), eq(4L), eq("000000"), eq(101L),
eq(7L), eq(4L))).thenReturn(1); eq(7L), eq(4L))).thenReturn(1);
when(fixture.jdbc.update(contains("delete from sys_oss"), eq("000000"), eq(101L), eq(101L), when(fixture.jdbc.update(contains("delete from sys_oss"), eq("000000"), eq(101L), eq(101L),
anyString(), eq("100"))).thenReturn(1); anyString(), eq("100"), eq("upload-token"))).thenReturn(1);
} }
private static Map<String, Object> staleIntent() { private static Map<String, Object> staleIntent() {
return Map.of( return Map.of(
"tenant_id", "000000", "owner_user_id", 101L, "space_id", 7L, "item_id", 100L, "tenant_id", "000000", "owner_user_id", 101L, "space_id", 7L, "item_id", 100L,
"oss_id", 101L, "size_bytes", 4L, "mime_type", "text/plain", "oss_id", 101L, "size_bytes", 4L, "mime_type", "text/plain",
"file_name", "personal/key.txt", "service", "personal-private" "file_name", "personal/key.txt", "service", "personal-private",
"upload_token", "upload-token"
); );
} }
private static void assertUploadState(String ext, String state) throws Exception { private static JsonNode uploadExt(String ext, String state) throws Exception {
JsonNode json = new ObjectMapper().readTree(ext); JsonNode json = new ObjectMapper().readTree(ext);
assertEquals("personal", json.path("source").asText()); assertEquals("personal", json.path("source").asText());
assertEquals(100L, json.path("itemId").asLong()); assertEquals(100L, json.path("itemId").asLong());
assertEquals(state, json.path("uploadState").asText()); assertEquals(state, json.path("uploadState").asText());
return json;
} }
private record Fixture(JdbcTemplate jdbc, PersonalSpaceService spaces, PersonalObjectStore store, private record Fixture(JdbcTemplate jdbc, PersonalSpaceService spaces, PersonalObjectStore store,
@@ -25,6 +25,7 @@ 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;
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.mock; import static org.mockito.Mockito.mock;
@@ -54,6 +55,12 @@ class PersonalIngestionWorkerTest {
@Test @Test
void staleParsingUsesLeaseAndExhaustionThreshold() { void staleParsingUsesLeaseAndExhaustionThreshold() {
JdbcTemplate jdbc = mock(JdbcTemplate.class); JdbcTemplate jdbc = mock(JdbcTemplate.class);
Map<String, Object> exhausted = Map.of(
"id", 9L, "tenant_id", "000000", "owner_user_id", 101L, "attempt_count", 3);
Map<String, Object> retryable = Map.of(
"id", 10L, "tenant_id", "000000", "owner_user_id", 101L, "attempt_count", 2);
when(jdbc.queryForList(contains("i.status = 'PARSING'"), any(LocalDateTime.class)))
.thenReturn(List.of(exhausted, retryable));
PersonalIngestionWorker worker = PersonalIngestionWorker.forTest( PersonalIngestionWorker worker = PersonalIngestionWorker.forTest(
jdbc, mock(ISysOssService.class), mock(KnowledgeDocumentParser.class), immediateTransactions(), jdbc, mock(ISysOssService.class), mock(KnowledgeDocumentParser.class), immediateTransactions(),
(ossId, prefix, ownerUserId, maxBytes) -> (ossId, prefix, ownerUserId, maxBytes) ->
@@ -61,8 +68,10 @@ class PersonalIngestionWorkerTest {
worker.recoverStaleParsing(); worker.recoverStaleParsing();
verify(jdbc).update(contains("PERSONAL_PARSE_RETRY_EXHAUSTED"), eq(3), any(LocalDateTime.class)); verify(jdbc).update(contains("PERSONAL_PARSE_RETRY_EXHAUSTED"), eq("000000"), eq(101L), eq(9L),
verify(jdbc).update(contains("set i.status = 'QUEUED'"), eq(3), any(LocalDateTime.class)); eq(3), any(LocalDateTime.class));
verify(jdbc).update(contains("set status = 'QUEUED'"), eq("000000"), eq(101L), eq(10L), eq(2),
any(LocalDateTime.class));
} }
@Test @Test
@@ -72,10 +81,12 @@ class PersonalIngestionWorkerTest {
TransactionTemplate transactions = immediateTransactions(); TransactionTemplate transactions = immediateTransactions();
Map<String, Object> item = item(); Map<String, Object> item = item();
when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item)); when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item));
when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L))).thenReturn(1); when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L), eq(0))).thenReturn(1);
when(jdbc.queryForMap(contains("for update"), eq("000000"), eq(101L), eq(9L)))
.thenReturn(Map.of("status", "PARSING", "attempt_count", 1));
when(parser.parse(eq("notes.txt"), eq("text/plain"), any(byte[].class))) when(parser.parse(eq("notes.txt"), eq("text/plain"), any(byte[].class)))
.thenReturn(new ParsedDocument("一二三四五六", "text/plain", Map.of())); .thenReturn(new ParsedDocument("一二三四五六", "text/plain", Map.of()));
when(jdbc.update(contains("status = 'READY'"), any(), eq("[]"), eq("000000"), eq(101L), eq(9L))) when(jdbc.update(contains("status = 'READY'"), any(), eq("[]"), eq("000000"), eq(101L), eq(9L), eq(1)))
.thenReturn(1); .thenReturn(1);
PersonalIngestionWorker worker = PersonalIngestionWorker.forTest( PersonalIngestionWorker worker = PersonalIngestionWorker.forTest(
jdbc, mock(ISysOssService.class), parser, transactions, jdbc, mock(ISysOssService.class), parser, transactions,
@@ -88,10 +99,11 @@ class PersonalIngestionWorkerTest {
assertTrue(worker.processNext()); assertTrue(worker.processNext());
verify(jdbc).update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L)); verify(jdbc).update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L), eq(0));
verify(jdbc).update(contains("delete from aihr_personal_fragment"), eq("000000"), eq(101L), eq(9L)); verify(jdbc).update(contains("delete from aihr_personal_fragment"), eq("000000"), eq(101L), eq(9L));
verify(jdbc).batchUpdate(contains("insert into aihr_personal_fragment"), any(BatchPreparedStatementSetter.class)); verify(jdbc).batchUpdate(contains("insert into aihr_personal_fragment"), any(BatchPreparedStatementSetter.class));
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),
eq(1));
} }
@Test @Test
@@ -99,11 +111,13 @@ class PersonalIngestionWorkerTest {
JdbcTemplate jdbc = mock(JdbcTemplate.class); JdbcTemplate jdbc = mock(JdbcTemplate.class);
KnowledgeDocumentParser parser = mock(KnowledgeDocumentParser.class); KnowledgeDocumentParser parser = mock(KnowledgeDocumentParser.class);
when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item())); when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item()));
when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L))).thenReturn(1); when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L), eq(0))).thenReturn(1);
when(jdbc.queryForMap(contains("for update"), eq("000000"), eq(101L), eq(9L)))
.thenReturn(Map.of("status", "PARSING", "attempt_count", 1));
String content = "字".repeat(900); String content = "字".repeat(900);
when(parser.parse(any(), any(), any(byte[].class))) when(parser.parse(any(), any(), any(byte[].class)))
.thenReturn(new ParsedDocument(content, "text/plain", Map.of())); .thenReturn(new ParsedDocument(content, "text/plain", Map.of()));
when(jdbc.update(contains("status = 'READY'"), any(), eq("[]"), eq("000000"), eq(101L), eq(9L))) when(jdbc.update(contains("status = 'READY'"), any(), eq("[]"), eq("000000"), eq(101L), eq(9L), eq(1)))
.thenReturn(1); .thenReturn(1);
PersonalIngestionWorker worker = PersonalIngestionWorker.forTest( PersonalIngestionWorker worker = PersonalIngestionWorker.forTest(
jdbc, mock(ISysOssService.class), parser, immediateTransactions(), jdbc, mock(ISysOssService.class), parser, immediateTransactions(),
@@ -133,12 +147,54 @@ class PersonalIngestionWorkerTest {
verify(jdbc, never()).update(contains("status = 'READY'"), any(), any(), any(), any(), any()); verify(jdbc, never()).update(contains("status = 'READY'"), any(), any(), any(), any(), any());
} }
@Test
void expiredWorkerSuccessCannotOverwriteNewAttemptOrFragments() 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), eq(0))).thenReturn(1);
when(parser.parse(any(), any(), any(byte[].class)))
.thenReturn(new ParsedDocument("A worker parsed this", "text/plain", Map.of()));
when(jdbc.queryForMap(contains("for update"), eq("000000"), eq(101L), eq(9L)))
.thenReturn(Map.of("status", "PARSING", "attempt_count", 2));
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());
verify(jdbc, never()).update(contains("delete from aihr_personal_fragment"), any(), any(), any());
verify(jdbc, never()).batchUpdate(any(String.class), any(BatchPreparedStatementSetter.class));
verify(jdbc).update(contains("status = 'FAILED'"), eq("PERSONAL_PARSE_FAILED"), anyString(),
eq("000000"), eq(101L), eq(9L), eq(1));
}
@Test
void expiredWorkerFailureCannotOverwriteNewAttempt() 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), eq(0))).thenReturn(1);
when(parser.parse(any(), any(), any(byte[].class))).thenThrow(new IllegalStateException("late failure"));
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());
verify(jdbc).update(contains("status = 'FAILED'"), eq("PERSONAL_PARSE_FAILED"), anyString(),
eq("000000"), eq(101L), eq(9L), eq(1));
verify(jdbc, never()).update(contains("delete from aihr_personal_fragment"), any(), any(), any());
}
@Test @Test
void workerPersistsOnlyStablePublicFailure() throws Exception { void workerPersistsOnlyStablePublicFailure() throws Exception {
JdbcTemplate jdbc = mock(JdbcTemplate.class); JdbcTemplate jdbc = mock(JdbcTemplate.class);
KnowledgeDocumentParser parser = mock(KnowledgeDocumentParser.class); KnowledgeDocumentParser parser = mock(KnowledgeDocumentParser.class);
when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item())); when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item()));
when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L))).thenReturn(1); when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L), eq(0))).thenReturn(1);
when(parser.parse(any(), any(), any(byte[].class))) when(parser.parse(any(), any(), any(byte[].class)))
.thenThrow(new KnowledgeDocumentParser.ParseException( .thenThrow(new KnowledgeDocumentParser.ParseException(
KnowledgeDocumentParser.Failure.INVALID, "secret parser detail")); KnowledgeDocumentParser.Failure.INVALID, "secret parser detail"));
@@ -150,7 +206,7 @@ class PersonalIngestionWorkerTest {
assertTrue(worker.processNext()); assertTrue(worker.processNext());
verify(jdbc).update(contains("status = 'FAILED'"), eq("PERSONAL_PARSE_INVALID"), verify(jdbc).update(contains("status = 'FAILED'"), eq("PERSONAL_PARSE_INVALID"),
eq("资料解析失败,请检查文件后重试"), eq("000000"), eq(101L), eq(9L)); eq("资料解析失败,请检查文件后重试"), eq("000000"), eq(101L), eq(9L), eq(1));
} }
@Test @Test
@@ -169,7 +225,7 @@ class PersonalIngestionWorkerTest {
mock(PersonalIngestionWorker.OssClientProvider.class); mock(PersonalIngestionWorker.OssClientProvider.class);
when(clients.get("public-client")).thenReturn(publicClient); when(clients.get("public-client")).thenReturn(publicClient);
when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item())); when(jdbc.queryForList(contains("status = 'QUEUED'"))).thenReturn(List.of(item()));
when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L))).thenReturn(1); when(jdbc.update(contains("status = 'PARSING'"), eq("000000"), eq(101L), eq(9L), eq(0))).thenReturn(1);
PersonalIngestionWorker worker = PersonalIngestionWorker.forTest( PersonalIngestionWorker worker = PersonalIngestionWorker.forTest(
jdbc, ossService, mock(KnowledgeDocumentParser.class), immediateTransactions(), jdbc, ossService, mock(KnowledgeDocumentParser.class), immediateTransactions(),
PersonalIngestionWorker.objectReaderForTest(ossService, clients)); PersonalIngestionWorker.objectReaderForTest(ossService, clients));
@@ -177,22 +233,17 @@ class PersonalIngestionWorkerTest {
assertTrue(worker.processNext()); assertTrue(worker.processNext());
verify(jdbc).update(contains("status = 'FAILED'"), eq("PERSONAL_OSS_NOT_PRIVATE"), verify(jdbc).update(contains("status = 'FAILED'"), eq("PERSONAL_OSS_NOT_PRIVATE"),
eq("个人资料存储策略不可用"), eq("000000"), eq(101L), eq(9L)); eq("个人资料存储策略不可用"), eq("000000"), eq(101L), eq(9L), eq(1));
verify(publicClient, never()).getObjectContent(any(String.class)); verify(publicClient, never()).getObjectContent(any(String.class));
} }
private static Map<String, Object> item() { private static Map<String, Object> item() {
return Map.of( return Map.ofEntries(
"id", 9L, Map.entry("id", 9L), Map.entry("tenant_id", "000000"), Map.entry("space_id", 7L),
"tenant_id", "000000", Map.entry("owner_user_id", 101L), Map.entry("source_type", "TEXT"), Map.entry("title", "周报"),
"space_id", 7L, Map.entry("oss_id", 81L), Map.entry("mime_type", "text/plain"), Map.entry("tags_json", "[]"),
"owner_user_id", 101L, Map.entry("captured_at", Timestamp.valueOf(LocalDateTime.of(2026, 7, 12, 9, 0))),
"source_type", "TEXT", Map.entry("attempt_count", 0)
"title", "周报",
"oss_id", 81L,
"mime_type", "text/plain",
"tags_json", "[]",
"captured_at", Timestamp.valueOf(LocalDateTime.of(2026, 7, 12, 9, 0))
); );
} }
@@ -55,6 +55,7 @@ class PersonalSpaceServiceTest {
assertEquals(15, properties.getParsingLeaseMinutes()); assertEquals(15, properties.getParsingLeaseMinutes());
assertEquals(3, properties.getMaxParseAttempts()); assertEquals(3, properties.getMaxParseAttempts());
assertEquals(15, properties.getUploadCleanupAgeMinutes()); assertEquals(15, properties.getUploadCleanupAgeMinutes());
assertEquals(15, properties.getCleanupFinalizeGraceMinutes());
} }
@Test @Test