fix(personal): retain upload cleanup tombstones

This commit is contained in:
2026-07-12 04:03:16 +08:00
parent f4847f9ef7
commit 6b187060de
4 changed files with 143 additions and 24 deletions
@@ -260,13 +260,21 @@ public class PersonalIngestionService {
for (Map<String, Object> row : staleUploadRows("CLEANING", cleaningCutoff)) { for (Map<String, Object> row : staleUploadRows("CLEANING", cleaningCutoff)) {
finalizeStaleCleanup(intent(row), cleaningCutoff); finalizeStaleCleanup(intent(row), cleaningCutoff);
} }
LocalDateTime tombstoneCutoff = now.minusMinutes(properties.getUploadTombstoneRetentionMinutes());
for (Map<String, Object> row : tombstoneRows()) {
UploadIntent intent = intent(row);
if (deleteKnownObject(intent) && uploadUpdatedAt(row).isBefore(tombstoneCutoff)) {
phaseTransaction.executeWithoutResult(status -> deleteTombstoneMetadata(intent));
}
}
} }
private List<Map<String, Object>> staleUploadRows(String state, LocalDateTime cutoff) { private List<Map<String, Object>> staleUploadRows(String state, LocalDateTime cutoff) {
return jdbcTemplate.queryForList(""" 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 json_unquote(json_extract(o.ext1, '$.uploadToken')) upload_token,
o.update_time upload_updated_at
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
@@ -278,6 +286,22 @@ public class PersonalIngestionService {
""", cutoff, state); """, cutoff, state);
} }
private List<Map<String, Object>> tombstoneRows() {
return 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,
json_unquote(json_extract(o.ext1, '$.uploadToken')) upload_token,
o.update_time upload_updated_at
from 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
where i.status = 'DELETED'
and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'TOMBSTONE'
and json_unquote(json_extract(o.ext1, '$.uploadToken')) is not null
order by o.update_time
""");
}
private void beginCleanup(UploadIntent intent, LocalDateTime cutoff) { private void beginCleanup(UploadIntent intent, LocalDateTime cutoff) {
Integer claimed = phaseTransaction.execute(status -> claimState( Integer claimed = phaseTransaction.execute(status -> claimState(
intent, UploadState.PENDING, UploadState.CLEANING, cutoff)); intent, UploadState.PENDING, UploadState.CLEANING, cutoff));
@@ -286,7 +310,8 @@ public class PersonalIngestionService {
return; return;
} }
UploadState current = uploadState(intent); UploadState current = uploadState(intent);
if (current == UploadState.CLEANING || current == UploadState.MISSING) { if (current == UploadState.CLEANING || current == UploadState.TOMBSTONE
|| current == UploadState.MISSING) {
deleteKnownObject(intent); deleteKnownObject(intent);
} }
} }
@@ -331,7 +356,8 @@ public class PersonalIngestionService {
UploadState current = uploadState(intent); UploadState current = uploadState(intent);
if (current == UploadState.PENDING) { if (current == UploadState.PENDING) {
beginCleanup(intent, null); beginCleanup(intent, null);
} else if (current == UploadState.CLEANING || current == UploadState.MISSING) { } else if (current == UploadState.CLEANING || current == UploadState.TOMBSTONE
|| current == UploadState.MISSING) {
deleteKnownObject(intent); deleteKnownObject(intent);
} }
} }
@@ -386,24 +412,43 @@ public class PersonalIngestionService {
and used_bytes >= ? and item_count > 0 and used_bytes >= ? and item_count > 0
""", intent.sizeBytes(), intent.owner().tenantId(), intent.owner().userId(), intent.spaceId(), """, intent.sizeBytes(), intent.owner().tenantId(), intent.owner().userId(), intent.spaceId(),
intent.sizeBytes()); intent.sizeBytes());
int ossDeleted = jdbcTemplate.update(""" int tombstoned = jdbcTemplate.update("""
delete from sys_oss update sys_oss
set ext1 = json_set(ext1, '$.uploadState', 'TOMBSTONE'), url = '', 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, '$.uploadToken')) = ? and json_unquote(json_extract(ext1, '$.uploadToken')) = ?
and json_unquote(json_extract(ext1, '$.uploadState')) = 'CLEANING' and json_unquote(json_extract(ext1, '$.uploadState')) = 'CLEANING'
and service = ?
""", intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(), """, intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(),
String.valueOf(intent.itemId()), intent.uploadToken()); String.valueOf(intent.itemId()), intent.uploadToken(), intent.serviceKey());
if (itemDeleted != 1 || counterUpdated != 1 || ossDeleted != 1) { if (itemDeleted != 1 || counterUpdated != 1 || tombstoned != 1) {
throw new ServiceException("PERSONAL_UPLOAD_CLEANUP_FAILED"); throw new ServiceException("PERSONAL_UPLOAD_CLEANUP_FAILED");
} }
} }
private void deleteTombstoneMetadata(UploadIntent intent) {
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, '$.uploadToken')) = ?
and json_unquote(json_extract(ext1, '$.uploadState')) = 'TOMBSTONE'
and service = ?
""", intent.owner().tenantId(), intent.ossId(), intent.owner().userId(), intent.objectKey(),
String.valueOf(intent.itemId()), intent.uploadToken(), intent.serviceKey());
}
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 from aihr_personal_item select i.id, i.status
where tenant_id = ? and owner_user_id = ? and space_id = ? and content_hash = ? from aihr_personal_item i
and status <> 'DELETED' order by id desc limit 1 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.tenant_id = ? and i.owner_user_id = ? and i.space_id = ? and i.content_hash = ?
and i.status <> 'DELETED'
and json_unquote(json_extract(o.ext1, '$.uploadState')) = 'READY'
order by i.id desc limit 1
""", owner.tenantId(), owner.userId(), spaceId, hash); """, owner.tenantId(), owner.userId(), spaceId, hash);
if (rows.isEmpty()) return null; if (rows.isEmpty()) return null;
long id = ((Number) rows.get(0).get("id")).longValue(); long id = ((Number) rows.get(0).get("id")).longValue();
@@ -444,6 +489,13 @@ public class PersonalIngestionService {
return number.longValue(); return number.longValue();
} }
private static LocalDateTime uploadUpdatedAt(Map<String, Object> row) {
Object value = row.get("upload_updated_at");
if (value instanceof LocalDateTime dateTime) return dateTime;
if (value instanceof java.sql.Timestamp timestamp) return timestamp.toLocalDateTime();
throw new ServiceException("PERSONAL_UPLOAD_CLEANUP_FAILED");
}
private void validateFile(MultipartFile file) { private void validateFile(MultipartFile file) {
if (file == null || file.isEmpty() || file.getSize() <= 0) throw new ServiceException("PERSONAL_FILE_EMPTY"); if (file == null || file.isEmpty() || file.getSize() <= 0) throw new ServiceException("PERSONAL_FILE_EMPTY");
validateSize(file.getSize()); validateSize(file.getSize());
@@ -466,8 +518,12 @@ public class PersonalIngestionService {
} }
private static void validateRecoveryWindows(PersonalKnowledgeProperties properties) { private static void validateRecoveryWindows(PersonalKnowledgeProperties properties) {
long cleanupWindow = (long) properties.getUploadCleanupAgeMinutes()
+ properties.getCleanupFinalizeGraceMinutes();
if (properties.getUploadCleanupAgeMinutes() < MIN_UPLOAD_CLEANUP_AGE_MINUTES if (properties.getUploadCleanupAgeMinutes() < MIN_UPLOAD_CLEANUP_AGE_MINUTES
|| properties.getCleanupFinalizeGraceMinutes() <= 0) { || properties.getCleanupFinalizeGraceMinutes() <= 0
|| properties.getUploadTombstoneRetentionMinutes() < 60
|| properties.getUploadTombstoneRetentionMinutes() <= cleanupWindow) {
throw new IllegalArgumentException("invalid personal upload recovery windows"); throw new IllegalArgumentException("invalid personal upload recovery windows");
} }
} }
@@ -557,7 +613,7 @@ 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 enum UploadState { PENDING, READY, CLEANING, TOMBSTONE, 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,
@@ -23,4 +23,5 @@ public class PersonalKnowledgeProperties {
private int maxParseAttempts = 3; private int maxParseAttempts = 3;
private int uploadCleanupAgeMinutes = 15; private int uploadCleanupAgeMinutes = 15;
private int cleanupFinalizeGraceMinutes = 15; private int cleanupFinalizeGraceMinutes = 15;
private int uploadTombstoneRetentionMinutes = 1440;
} }
@@ -115,7 +115,7 @@ class PersonalIngestionServiceTest {
} }
@Test @Test
void duplicateIsResolvedUnderOwnerLockWithoutCreatingUploadIntent() { void duplicateIsResolvedOnlyFromReadyUploadUnderOwnerLock() {
Fixture fixture = fixture(); 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"), eq("000000"), eq(101L), eq(7L), anyString())) when(fixture.jdbc.queryForList(contains("content_hash"), eq("000000"), eq(101L), eq(7L), anyString()))
@@ -126,6 +126,9 @@ class PersonalIngestionServiceTest {
assertEquals(77L, response.itemId()); assertEquals(77L, response.itemId());
assertEquals(77L, response.duplicateOf()); assertEquals(77L, response.duplicateOf());
ArgumentCaptor<String> dedupeSql = ArgumentCaptor.forClass(String.class);
verify(fixture.jdbc).queryForList(dedupeSql.capture(), eq("000000"), eq(101L), eq(7L), anyString());
assertTrue(dedupeSql.getValue().contains("$.uploadState')) = 'READY'"));
verify(fixture.store).requirePrivateService(); verify(fixture.store).requirePrivateService();
verify(fixture.store, never()).uploadPhysical(anyString(), anyString(), anyString(), any(byte[].class)); verify(fixture.store, never()).uploadPhysical(anyString(), anyString(), anyString(), any(byte[].class));
verify(fixture.jdbc, never()).update(contains("insert into sys_oss"), any(), any(), any(), any(), any(), verify(fixture.jdbc, never()).update(contains("insert into sys_oss"), any(), any(), any(), any(), any(),
@@ -162,13 +165,14 @@ class PersonalIngestionServiceTest {
} }
@Test @Test
void staleCleanupIsIdempotentAndDoesNotDecrementCountersTwice() { void cleaningFinalizeCompensatesOnceAndRetainsTombstoneIntent() {
Fixture fixture = fixture(); Fixture fixture = fixture();
Map<String, Object> stale = staleIntent(); Map<String, Object> stale = staleIntent();
when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("PENDING"))) when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("PENDING")))
.thenReturn(List.of(stale), List.of(), List.of()); .thenReturn(List.of(stale), List.of());
when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("CLEANING"))) when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("CLEANING")))
.thenReturn(List.of(), List.of(stale), List.of(stale)); .thenReturn(List.of(), List.of(stale));
when(fixture.jdbc.queryForList(contains("uploadState')) = 'TOMBSTONE'"))).thenReturn(List.of());
when(fixture.jdbc.update(contains("json_set"), eq("CLEANING"), eq("000000"), eq(101L), eq(101L), 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))) eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("PENDING"), any(LocalDateTime.class)))
.thenReturn(1); .thenReturn(1);
@@ -178,13 +182,41 @@ class PersonalIngestionServiceTest {
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(2)).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));
verify(fixture.jdbc).update(contains("'$.uploadState', 'TOMBSTONE'"), eq("000000"), eq(101L),
eq(101L), eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("personal-private"));
verify(fixture.jdbc, never()).update(contains("delete from sys_oss"), any(), any(), any(), any(), any(),
any(), any());
}
@Test
void tombstoneIsDeletedOnEveryScanAndMetadataRemovedOnlyAfterRetention() {
Fixture fixture = fixture();
Map<String, Object> fresh = staleIntent(LocalDateTime.now());
Map<String, Object> expired = staleIntent(LocalDateTime.now().minusDays(2));
when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("PENDING")))
.thenReturn(List.of());
when(fixture.jdbc.queryForList(contains("uploadState')) = ?"), any(LocalDateTime.class), eq("CLEANING")))
.thenReturn(List.of());
when(fixture.jdbc.queryForList(contains("uploadState')) = 'TOMBSTONE'")))
.thenReturn(List.of(fresh), List.of(fresh), List.of(expired));
when(fixture.jdbc.update(contains("delete from sys_oss"), eq("000000"), eq(101L), eq(101L),
eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("personal-private"))).thenReturn(1);
fixture.service.recoverStaleUploadIntents();
fixture.service.recoverStaleUploadIntents();
fixture.service.recoverStaleUploadIntents();
verify(fixture.store, times(3)).deletePhysical("personal-private", "personal/key.txt");
verify(fixture.jdbc, times(1)).update(contains("delete from sys_oss"), eq("000000"), eq(101L), eq(101L),
eq("personal/key.txt"), eq("100"), eq("upload-token"), eq("personal-private"));
verify(fixture.jdbc, never()).update(contains("used_bytes = used_bytes -"), any(), any(), any(), any(),
any());
} }
@Test @Test
@@ -310,6 +342,23 @@ class PersonalIngestionServiceTest {
any()); any());
} }
@Test
void tombstoneActivationFailureDeletesKnownLateObjectWithoutReclaimingIntent() {
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", "TOMBSTONE")));
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 @Test
void unsafeCleanupWindowsAreRejected() { void unsafeCleanupWindowsAreRejected() {
PersonalKnowledgeProperties properties = new PersonalKnowledgeProperties(); PersonalKnowledgeProperties properties = new PersonalKnowledgeProperties();
@@ -318,6 +367,13 @@ class PersonalIngestionServiceTest {
properties.setUploadCleanupAgeMinutes(15); properties.setUploadCleanupAgeMinutes(15);
properties.setCleanupFinalizeGraceMinutes(0); properties.setCleanupFinalizeGraceMinutes(0);
assertThrows(IllegalArgumentException.class, () -> fixture(properties)); assertThrows(IllegalArgumentException.class, () -> fixture(properties));
properties.setCleanupFinalizeGraceMinutes(15);
properties.setUploadTombstoneRetentionMinutes(30);
assertThrows(IllegalArgumentException.class, () -> fixture(properties));
properties.setUploadCleanupAgeMinutes(5);
properties.setCleanupFinalizeGraceMinutes(1);
properties.setUploadTombstoneRetentionMinutes(59);
assertThrows(IllegalArgumentException.class, () -> fixture(properties));
} }
private static Fixture fixture() { private static Fixture fixture() {
@@ -356,16 +412,21 @@ class PersonalIngestionServiceTest {
eq(101L))).thenReturn(1); eq(101L))).thenReturn(1);
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("'$.uploadState', 'TOMBSTONE'"), eq("000000"), eq(101L), eq(101L),
anyString(), eq("100"), eq("upload-token"))).thenReturn(1); anyString(), eq("100"), eq("upload-token"), eq("personal-private"))).thenReturn(1);
} }
private static Map<String, Object> staleIntent() { private static Map<String, Object> staleIntent() {
return Map.of( return staleIntent(LocalDateTime.now());
"tenant_id", "000000", "owner_user_id", 101L, "space_id", 7L, "item_id", 100L, }
"oss_id", 101L, "size_bytes", 4L, "mime_type", "text/plain",
"file_name", "personal/key.txt", "service", "personal-private", private static Map<String, Object> staleIntent(LocalDateTime updatedAt) {
"upload_token", "upload-token" return Map.ofEntries(
Map.entry("tenant_id", "000000"), Map.entry("owner_user_id", 101L), Map.entry("space_id", 7L),
Map.entry("item_id", 100L), Map.entry("oss_id", 101L), Map.entry("size_bytes", 4L),
Map.entry("mime_type", "text/plain"), Map.entry("file_name", "personal/key.txt"),
Map.entry("service", "personal-private"), Map.entry("upload_token", "upload-token"),
Map.entry("upload_updated_at", updatedAt)
); );
} }
@@ -56,6 +56,7 @@ class PersonalSpaceServiceTest {
assertEquals(3, properties.getMaxParseAttempts()); assertEquals(3, properties.getMaxParseAttempts());
assertEquals(15, properties.getUploadCleanupAgeMinutes()); assertEquals(15, properties.getUploadCleanupAgeMinutes());
assertEquals(15, properties.getCleanupFinalizeGraceMinutes()); assertEquals(15, properties.getCleanupFinalizeGraceMinutes());
assertEquals(1440, properties.getUploadTombstoneRetentionMinutes());
} }
@Test @Test