diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/file/service/object/RustfsObjectStorageService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/file/service/object/RustfsObjectStorageService.java index 3eab9e59..0f95d710 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/file/service/object/RustfsObjectStorageService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/file/service/object/RustfsObjectStorageService.java @@ -5,11 +5,14 @@ import io.micrometer.core.instrument.DistributionSummary; import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.Timer; import io.minio.GetObjectArgs; +import io.minio.ListObjectsArgs; import io.minio.MinioClient; import io.minio.PutObjectArgs; import io.minio.RemoveObjectArgs; +import io.minio.Result; import io.minio.StatObjectArgs; import io.minio.errors.ErrorResponseException; +import io.minio.messages.Item; import lombok.extern.slf4j.Slf4j; import okhttp3.ConnectionPool; import okhttp3.Dispatcher; @@ -21,8 +24,11 @@ import org.springframework.stereotype.Service; import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.List; import java.util.Objects; import java.util.Set; +import java.util.UUID; import java.util.concurrent.Semaphore; import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; @@ -45,6 +51,12 @@ public class RustfsObjectStorageService { "NoSuchKey", "NoSuchBucket", "NoSuchVersion", "AccessDenied", "InvalidAccessKeyId", "SignatureDoesNotMatch", "InvalidBucketName"); + /** 标准 UUID 字符串长度(8-4-4-4-12),用于识别版本化对象 key。 */ + private static final int UUID_STRING_LENGTH = 36; + + /** 兄弟对象兜底列出的上限:只为找回同槽位的版本化对象,不需要列全。 */ + private static final int MAX_SIBLING_LIST_KEYS = 50; + private final TransientStorageProperties properties; private final ObjectProvider meterRegistryProvider; private final ObjectProvider deleteRetryServiceProvider; @@ -106,7 +118,22 @@ public class RustfsObjectStorageService { return uploadBytes(objectKey, bytes, verifyAfterUpload); } + /** + * 三参重载:是否做「上传失败补偿删除」按对象 key 形态自动判定。 + * + *

只有版本化 key(末段以 UUID 结尾)是本次写入独占的;确定性 key 会被重传重写复用, + * 删它就可能删掉别的 DB 行仍在引用的对象(2026-09-17 线上任务 28459 的载荷对象就是这么丢的)。 + */ public String uploadBytes(String objectKey, byte[] content, boolean verifyAfterUpload) { + return uploadBytes(objectKey, content, verifyAfterUpload, isVersionedObjectKey(objectKey)); + } + + /** + * @param compensateDeleteOnFailure put 已完成、但后续可见性校验失败时,是否把该对象排进删除补偿队列。 + * 仅当调用方能确认「该对象不会被其它写入复用时」才可传 true。 + */ + public String uploadBytes(String objectKey, byte[] content, boolean verifyAfterUpload, + boolean compensateDeleteOnFailure) { long deadlineNanos = operationDeadlineNanos(); if (!isConfigured()) { throw new IllegalStateException("transient storage is not configured"); @@ -136,13 +163,45 @@ public class RustfsObjectStorageService { } return uploadedObjectKey; } catch (RuntimeException ex) { - if (putCompleted.get()) { + if (putCompleted.get() && compensateDeleteOnFailure) { enqueueDeleteRetry(objectKey, ex); + } else if (putCompleted.get()) { + // 共享 key 会被重传重写:此处删除可能删掉别的行正在引用的对象,交给保留期清理兜底。 + // 线上任务 28459 的 chunk-462/473/484 就是被这条无条件删除队列删掉的。 + log.warn("[rustfs] 跳过上传失败补偿删除(对象非本次独占,可能被复用)objectKey={} err={}", + objectKey, ex.getMessage()); } throw ex; } } + /** + * 对象 key 是否为「本次写入独占」的版本化 key:末段(去掉 {@code .json} 后缀)以 UUID 结尾。 + * + *

只看末段——UUID 出现在中间段(如 scopeHash)不代表该对象被独占;解析失败一律按共享处理 + * (保守:宁可留下孤儿对象,也不删掉可能仍被引用的对象)。 + */ + static boolean isVersionedObjectKey(String objectKey) { + if (objectKey == null) { + return false; + } + String key = objectKey.trim(); + if (key.endsWith(".json")) { + key = key.substring(0, key.length() - ".json".length()); + } + int slash = key.lastIndexOf('/'); + String lastSegment = slash < 0 ? key : key.substring(slash + 1); + if (lastSegment.length() < UUID_STRING_LENGTH) { + return false; + } + try { + UUID.fromString(lastSegment.substring(lastSegment.length() - UUID_STRING_LENGTH)); + return true; + } catch (IllegalArgumentException ex) { + return false; + } + } + public String readObjectAsString(String objectKey) { byte[] bytes = readObjectBytes(objectKey); return new String(bytes, StandardCharsets.UTF_8); @@ -182,6 +241,43 @@ public class RustfsObjectStorageService { }); } + /** + * 列出前缀下的对象 key,按最后修改时间倒序(最新在前)。 + * + *

只服务于「指针指向的对象已不存在、需要找回同槽位的版本化兄弟对象」这一兜底路径, + * 因此刻意不做重试、不参与失败窗口记账:列出失败直接抛错,由调用方按原错误语义处理。 + */ + public List listObjectKeysNewestFirst(String prefix, int limit) { + if (!isConfigured()) { + throw new IllegalStateException("transient storage is not configured"); + } + int safeLimit = Math.max(1, Math.min(limit, MAX_SIBLING_LIST_KEYS)); + long deadlineNanos = operationDeadlineNanos(); + List entries = new ArrayList<>(); + try { + Iterable> results = buildClient(deadlineNanos).listObjects(ListObjectsArgs.builder() + .bucket(properties.getBucket()) + .prefix(prefix == null ? "" : prefix) + .recursive(true) + .maxKeys(safeLimit) + .build()); + for (Result result : results) { + Item item = result.get(); + entries.add(new String[]{item.objectName(), + item.lastModified() == null ? "" : item.lastModified().toString()}); + } + } catch (Exception ex) { + throw new IllegalStateException("transient storage list failed prefix=" + prefix + + " err=" + ex.getMessage(), ex); + } + entries.sort((left, right) -> right[1].compareTo(left[1])); + List keys = new ArrayList<>(entries.size()); + for (String[] entry : entries) { + keys.add(entry[0]); + } + return keys; + } + public void deleteObject(String objectKey) { deleteObject(objectKey, true, operationDeadlineNanos()); } diff --git a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TransientPayloadStorageService.java b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TransientPayloadStorageService.java index a728e81a..b4fa9c6a 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TransientPayloadStorageService.java +++ b/backend-java/src/main/java/com/nanri/aiimage/modules/task/service/TransientPayloadStorageService.java @@ -30,6 +30,7 @@ import java.util.Locale; import java.util.Objects; import java.util.Set; import java.util.UUID; +import java.util.regex.Pattern; import java.util.zip.GZIPInputStream; import java.util.zip.GZIPOutputStream; @@ -150,8 +151,20 @@ public class TransientPayloadStorageService { return decodeStoredPayloadBytes(readLocalPayloadBytes(stripLocalInstanceId(localKey))); } if (pointer.startsWith(RUSTFS_POINTER_PREFIX)) { - return decodeStoredPayloadBytes( - rustfsObjectStorageService.readObjectBytes(pointer.substring(RUSTFS_POINTER_PREFIX.length()))); + String objectKey = pointer.substring(RUSTFS_POINTER_PREFIX.length()); + try { + return decodeStoredPayloadBytes(rustfsObjectStorageService.readObjectBytes(objectKey)); + } catch (RuntimeException readException) { + String sibling = findVersionedChunkSibling(objectKey, readException); + if (sibling == null) { + throw readException; + } + // 2026-09-17 线上任务 28459:行指向的普通 key 被误删,但同一分片槽位的版本化对象还在。 + // 读出它即可让任务按已有数据出结果,不必整单失败。 + log.warn("[transient-payload] chunk 载荷对象不存在,回退同槽位版本化对象 pointer={} sibling={}", + pointer, sibling); + return decodeStoredPayloadBytes(rustfsObjectStorageService.readObjectBytes(sibling)); + } } if (pointer.startsWith(OSS_POINTER_PREFIX)) { return decodeStoredPayloadBytes( @@ -164,6 +177,53 @@ public class TransientPayloadStorageService { return value; } + /** chunk 载荷槽位的 entryKey 形态:{@code chunk-}(版本化写入则形如 {@code chunk--})。 */ + private static final Pattern CHUNK_ENTRY_KEY_PATTERN = Pattern.compile("chunk-\\d+"); + + /** 兄弟对象查找上限:只为找回同槽位对象,不需要列全。 */ + private static final int MAX_SIBLING_LOOKUP_KEYS = 50; + + /** + * 指针对象已不存在时,尝试找回同一分片槽位的版本化兄弟对象。 + * + *

两个条件同时满足才兜底,避免读到无关对象或掩盖真实故障: + *

    + *
  1. 末段 entryKey 是 {@code chunk-} 形态——只有这种槽位才有「版本化兄弟」语义;
  2. + *
  3. 失败原因是对象确实不存在(NoSuchKey)——权限/网络类失败照旧上抛以便重试。
  4. + *
+ */ + private String findVersionedChunkSibling(String objectKey, Throwable cause) { + if (!isObjectMissing(cause)) { + return null; + } + int slash = objectKey.lastIndexOf('/'); + String directory = slash < 0 ? "" : objectKey.substring(0, slash + 1); + String fileName = slash < 0 ? objectKey : objectKey.substring(slash + 1); + if (!fileName.endsWith(".json")) { + return null; + } + String entryKey = fileName.substring(0, fileName.length() - ".json".length()); + if (!CHUNK_ENTRY_KEY_PATTERN.matcher(entryKey).matches()) { + return null; + } + List candidates = rustfsObjectStorageService.listObjectKeysNewestFirst( + directory + entryKey + "-", MAX_SIBLING_LOOKUP_KEYS); + return candidates.isEmpty() ? null : candidates.getFirst(); + } + + /** 对象确已不存在:RustFS 返回 NoSuchKey,message 为 "The specified key does not exist."。 */ + private static boolean isObjectMissing(Throwable error) { + Throwable cursor = error; + while (cursor != null) { + String message = cursor.getMessage(); + if (message != null && message.contains("does not exist")) { + return true; + } + cursor = cursor.getCause(); + } + return false; + } + public void deletePayloadIfPresent(String value) { String pointer = extractPointer(value); if (pointer == null) { @@ -209,8 +269,10 @@ public class TransientPayloadStorageService { * *

判断口径: *

    - *
  • {@code biz_task_chunk.payload_json} 命中 > 1 行(> 1 表示除了 caller 视角下 - * 自己即将释放的那一行之外,至少还有别的 chunk 行也指向同一对象)→ 视为仍被引用。
  • + *
  • {@code biz_task_chunk.payload_json} 命中任意行(≥ 1)→ 视为仍被引用。 + * 曾用 {@code > 1} 作判据,等于放行「恰好还有 1 行引用」的情况,会把对方仍在用的对象 + * 删掉(2026-09-17 线上任务 28459:合并成功后指针未落库 + 旧对象被删 → 该分片永久读不到)。 + * 物理删除本就约定在 DB 行删除之后执行,故调用方正常路径下引用数必然为 0。
  • *
  • {@code biz_task_scope_state.parsed_payload_json}/{@code state_json} 命中 > 0 行 → * 视为仍被引用(这两个字段不是 caller 自身行的常见持有者,命中即非自我引用)。
  • *
@@ -248,7 +310,7 @@ public class TransientPayloadStorageService { // 解析不出 taskId 时按原口径全局查,行为与改造前一致。 Long pointerTaskId = extractTaskId(pointer); Long chunkCount = referencedChunkCount(pointerTaskId, values); - if (chunkCount != null && chunkCount > 1L) { + if (chunkCount != null && chunkCount > 0L) { return true; } Long scopeStateCount = referencedScopeStateCount(pointerTaskId, values); diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/file/service/object/RustfsObjectStorageServiceTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/file/service/object/RustfsObjectStorageServiceTest.java index d068f47b..6fee00e4 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/file/service/object/RustfsObjectStorageServiceTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/file/service/object/RustfsObjectStorageServiceTest.java @@ -297,16 +297,19 @@ class RustfsObjectStorageServiceTest { RustfsObjectStorageService service = new RustfsObjectStorageService( properties, emptyProvider(), provider(retryService), () -> client); + // 补偿删除只对版本化 key(本次写入独占)生效,故这里用 versioned key 覆盖该路径 + String firstKey = "task/a-11111111-2222-3333-4444-555555555555.json"; + String secondKey = "task/b-66666666-7777-8888-9999-000000000000.json"; IllegalStateException firstFailure = assertThrows(IllegalStateException.class, - () -> service.uploadText("task/a.json", "{}", true)); + () -> service.uploadText(firstKey, "{}", true)); IllegalStateException secondFailure = assertThrows(IllegalStateException.class, - () -> service.uploadText("task/b.json", "{}", true)); + () -> service.uploadText(secondKey, "{}", true)); assertTrue(firstFailure.getMessage().contains("not visible")); assertTrue(secondFailure.getMessage().contains("not visible")); verify(client, times(5)).statObject(any(StatObjectArgs.class)); - verify(retryService).enqueue(eq("task/a.json"), same(firstFailure)); - verify(retryService).enqueue(eq("task/b.json"), same(secondFailure)); + verify(retryService).enqueue(eq(firstKey), same(firstFailure)); + verify(retryService).enqueue(eq(secondKey), same(secondFailure)); IllegalStateException rejected = assertThrows(IllegalStateException.class, () -> service.uploadText("task/c.json", "{}", false)); assertTrue(rejected.getMessage().contains("cooldown active")); @@ -360,6 +363,9 @@ class RustfsObjectStorageServiceTest { @Test void completedUploadIsEnqueuedWhenDeadlineExpiresAfterPut() throws Exception { + // 只有版本化 key(本次写入独占)才允许补偿删除;共享 key 的守卫见 + // RustfsUploadCompensationGuardTest#deterministicKeyDerivedFromShapeDoesNotEnqueueCompensation + String versionedKey = "task/chunk-1-0ffc254b-afad-4279-aab2-e85e3ff955e9.json"; TransientStorageProperties properties = configuredProperties(); properties.setOperationTimeoutSeconds(1); RustfsDeleteRetryService retryService = mock(RustfsDeleteRetryService.class); @@ -372,11 +378,11 @@ class RustfsObjectStorageServiceTest { properties, emptyProvider(), provider(retryService), () -> client); IllegalStateException failure = assertThrows(IllegalStateException.class, - () -> service.uploadText("task/a.json", "{}", false)); + () -> service.uploadText(versionedKey, "{}", false)); assertTrue(failure.getMessage().contains("operation timeout")); verify(client).putObject(any(PutObjectArgs.class)); - verify(retryService).enqueue(eq("task/a.json"), same(failure)); + verify(retryService).enqueue(eq(versionedKey), same(failure)); } private static TransientStorageProperties configuredProperties() { diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/file/service/object/RustfsUploadCompensationGuardTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/file/service/object/RustfsUploadCompensationGuardTest.java new file mode 100644 index 00000000..0b6ff14d --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/file/service/object/RustfsUploadCompensationGuardTest.java @@ -0,0 +1,217 @@ +package com.nanri.aiimage.modules.file.service.object; + +import com.nanri.aiimage.config.TransientStorageProperties; +import io.minio.MinioClient; +import io.minio.PutObjectArgs; +import io.minio.StatObjectArgs; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.ObjectProvider; + +import java.io.IOException; + +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.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * 上传失败补偿删除的准入守卫(2026-09-17 线上任务 28459 事故)。 + * + *

事故形态:客户端重传复用**确定性 key**({@code chunk-462.json})。第一次上传 put 已成功、 + * 但随后的 {@code verifyObjectVisible} 抖动 → 旧代码无条件 + * {@code enqueueDeleteRetry(objectKey)};而删除重试队列对该 key 不再做任何引用校验 + * ({@code RustfsDeleteRetryService} → {@code deleteObjectFromRetry} → {@code removeObject}), + * 于是把重传成功后 DB 行仍指向的对象删掉,下游读 chunk 直接 404,整个任务组装失败。 + * + *

口径:只有**本次写入独占的对象**(版本化 key,末段以 UUID 结尾)才允许补偿删除; + * 确定性 key 会被重写复用,删它就可能删掉别的行正在引用的对象。 + */ +class RustfsUploadCompensationGuardTest { + + // ===== 正常路径 ===== + + /** 唯一 key(版本化)+ 验证失败 + 显式允许 → 补偿删除入队。 */ + @Test + void uniqueKeyVerifyFailureEnqueuesCompensation() throws Exception { + RustfsDeleteRetryService retryService = mock(RustfsDeleteRetryService.class); + RustfsObjectStorageService service = new RustfsObjectStorageService( + configuredProperties(), emptyProvider(), provider(retryService), RustfsUploadCompensationGuardTest::putOkButStatFails); + + assertThrows(IllegalStateException.class, + () -> service.uploadBytes("task/chunk-462-0ffc254b-afad-4279-aab2-e85e3ff955e9.json", + "{}".getBytes(), true, true)); + + verify(retryService).enqueue(eq("task/chunk-462-0ffc254b-afad-4279-aab2-e85e3ff955e9.json"), + any(Throwable.class)); + } + + /** 三参重载按 key 形态判定:版本化 key 自动允许补偿删除。 */ + @Test + void versionedKeyDerivedFromShapeEnqueuesCompensation() throws Exception { + RustfsDeleteRetryService retryService = mock(RustfsDeleteRetryService.class); + RustfsObjectStorageService service = new RustfsObjectStorageService( + configuredProperties(), emptyProvider(), provider(retryService), RustfsUploadCompensationGuardTest::putOkButStatFails); + + assertThrows(IllegalStateException.class, + () -> service.uploadBytes("task/chunk-77-d03bf3b0-abbb-469d-9bb8-a7c01c309477.json", + "{}".getBytes(), true)); + + verify(retryService).enqueue(anyString(), any(Throwable.class)); + } + + // ===== 事故根因:共享 key 必须不删 ===== + + /** 确定性(共享)key + 显式禁止 → 不入队(异常仍照旧抛出,失败归属不变)。 */ + @Test + void sharedKeyVerifyFailureDoesNotEnqueueCompensation() throws Exception { + RustfsDeleteRetryService retryService = mock(RustfsDeleteRetryService.class); + RustfsObjectStorageService service = new RustfsObjectStorageService( + configuredProperties(), emptyProvider(), provider(retryService), RustfsUploadCompensationGuardTest::putOkButStatFails); + + assertThrows(IllegalStateException.class, + () -> service.uploadBytes("task/chunk-462.json", "{}".getBytes(), true, false)); + + verify(retryService, never()).enqueue(anyString(), any(Throwable.class)); + } + + /** 三参重载按 key 形态判定:确定性 key 自动禁止补偿删除(本次事故的直接修复点)。 */ + @Test + void deterministicKeyDerivedFromShapeDoesNotEnqueueCompensation() throws Exception { + RustfsDeleteRetryService retryService = mock(RustfsDeleteRetryService.class); + RustfsObjectStorageService service = new RustfsObjectStorageService( + configuredProperties(), emptyProvider(), provider(retryService), RustfsUploadCompensationGuardTest::putOkButStatFails); + + assertThrows(IllegalStateException.class, + () -> service.uploadBytes("task/chunk-462.json", "{}".getBytes(), true)); + + verify(retryService, never()).enqueue(anyString(), any(Throwable.class)); + } + + /** verify=false 上传即成功,不得产生任何补偿删除。 */ + @Test + void verifyDisabledNeverEnqueuesCompensation() throws Exception { + RustfsDeleteRetryService retryService = mock(RustfsDeleteRetryService.class); + RustfsObjectStorageService service = new RustfsObjectStorageService( + configuredProperties(), emptyProvider(), provider(retryService), RustfsUploadCompensationGuardTest::putOkButStatFails); + + String key = service.uploadBytes("task/chunk-462-uuid-x.json", "{}".getBytes(), false, true); + + assertEquals("task/chunk-462-uuid-x.json", key); + verify(retryService, never()).enqueue(anyString(), any(Throwable.class)); + } + + // ===== 异常路径 ===== + + /** put 从未成功过:对象本就不存在,不该产生删除动作(与开关无关)。 */ + @Test + void putNeverCompletedDoesNotEnqueueCompensation() throws Exception { + RustfsDeleteRetryService retryService = mock(RustfsDeleteRetryService.class); + MinioClient client = mock(MinioClient.class); + when(client.putObject(any(PutObjectArgs.class))).thenThrow(new IOException("put failed")); + TransientStorageProperties properties = configuredProperties(); + properties.setUploadMaxRetries(1); + RustfsObjectStorageService service = new RustfsObjectStorageService( + properties, emptyProvider(), provider(retryService), () -> client); + + assertThrows(RuntimeException.class, + () -> service.uploadBytes("task/chunk-462-uuid-x.json", "{}".getBytes(), true, true)); + + verify(retryService, never()).enqueue(anyString(), any(Throwable.class)); + } + + /** 补偿删除入队本身失败:不得把上传失败改写成入队异常,原异常必须继续上抛。 */ + @Test + void enqueueFailureDoesNotMaskUploadFailure() throws Exception { + RustfsDeleteRetryService retryService = mock(RustfsDeleteRetryService.class); + org.mockito.Mockito.doThrow(new IllegalStateException("queue down")) + .when(retryService).enqueue(anyString(), any(Throwable.class)); + RustfsObjectStorageService service = new RustfsObjectStorageService( + configuredProperties(), emptyProvider(), provider(retryService), RustfsUploadCompensationGuardTest::putOkButStatFails); + + IllegalStateException ex = assertThrows(IllegalStateException.class, + () -> service.uploadBytes("task/chunk-9-11111111-2222-3333-4444-555555555555.json", + "{}".getBytes(), true, true)); + + assertFalse(ex.getMessage() == null || ex.getMessage().contains("queue down"), + "入队异常不应成为对外错误;实际=" + ex.getMessage()); + } + + // ===== key 形态判定的边界 ===== + + @Test + void versionedKeyShapeIsRecognized() { + assertTrue(RustfsObjectStorageService.isVersionedObjectKey( + "task-chunk/appearance_patent/28459/hash/chunk-462-0ffc254b-afad-4279-aab2-e85e3ff955e9.json")); + assertTrue(RustfsObjectStorageService.isVersionedObjectKey( + "task-chunk/x/1/hash/0ffc254b-afad-4279-aab2-e85e3ff955e9.json"), + "整个 entryKey 就是 UUID 时也应识别为唯一 key"); + assertTrue(RustfsObjectStorageService.isVersionedObjectKey( + "task-chunk/x/1/hash/chunk-1-0ffc254b-afad-4279-aab2-e85e3ff955e9"), + "无 .json 后缀但末段是 UUID,同样视为唯一"); + } + + @Test + void deterministicKeyShapeIsNotRecognized() { + assertFalse(RustfsObjectStorageService.isVersionedObjectKey("task-chunk/x/1/hash/chunk-462.json")); + assertFalse(RustfsObjectStorageService.isVersionedObjectKey("task-parsed/x/1/hash/latest.json")); + assertFalse(RustfsObjectStorageService.isVersionedObjectKey("task-result-payload/x/1/hash/28459.json"), + "提交号这类 key 会被重写复用,不能当成唯一 key"); + assertFalse(RustfsObjectStorageService.isVersionedObjectKey( + "task-chunk/x/1/0ffc254b-afad-4279-aab2-e85e3ff955e9/chunk-462.json"), + "UUID 出现在中间段不算唯一——重写复用的判据只看末段"); + } + + @Test + void malformedKeyShapeDoesNotThrowAndIsTreatedAsShared() { + assertFalse(RustfsObjectStorageService.isVersionedObjectKey(null)); + assertFalse(RustfsObjectStorageService.isVersionedObjectKey("")); + assertFalse(RustfsObjectStorageService.isVersionedObjectKey(" ")); + assertFalse(RustfsObjectStorageService.isVersionedObjectKey("chunk-462.json")); + assertFalse(RustfsObjectStorageService.isVersionedObjectKey("-afad-4279-aab2-e85e3ff955e9.json".substring(0, 20))); + } + + // ===== 辅助 ===== + + private static MinioClient putOkButStatFails() { + MinioClient client = mock(MinioClient.class); + try { + when(client.putObject(any(PutObjectArgs.class))).thenReturn(null); + when(client.statObject(any(StatObjectArgs.class))).thenThrow(new IOException("stat 抖动")); + } catch (Exception ex) { + throw new IllegalStateException(ex); + } + return client; + } + + private static TransientStorageProperties configuredProperties() { + TransientStorageProperties properties = new TransientStorageProperties(); + properties.setEndpoint("http://127.0.0.1:9000"); + properties.setBucket("bucket"); + properties.setAccessKeyId("ak"); + properties.setAccessKeySecret("sk"); + properties.setBaseRetryDelayMillis(0); + properties.setRetryJitterMillis(0); + return properties; + } + + private static ObjectProvider provider(T value) { + @SuppressWarnings("unchecked") + ObjectProvider provider = mock(ObjectProvider.class); + when(provider.getIfAvailable()).thenReturn(value); + return provider; + } + + private static ObjectProvider emptyProvider() { + @SuppressWarnings("unchecked") + ObjectProvider provider = mock(ObjectProvider.class); + when(provider.getIfAvailable()).thenReturn(null); + return provider; + } +} diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskPayloadReferenceCheckTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskPayloadReferenceCheckTest.java index 40642e49..fb91f5bd 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskPayloadReferenceCheckTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TaskPayloadReferenceCheckTest.java @@ -33,8 +33,13 @@ import static org.mockito.Mockito.when; /** * task-150:清理前引用检查契约(plan 09)。 * payload 被 biz_task_chunk / biz_task_scope_state 引用则不清理; - * chunk 仅剩自身一行(count=1)不算共享引用;查询异常保守保留; + * chunk 命中任意行(含仅剩 1 行)都不算「可安全清理」;查询异常保守保留; * 检查只读无副作用;候选值批量反查。 + * + *

2026-09-17 线上任务 28459 之后收紧:原契约把 count==1 当作「调用方自己那行」放行删除, + * 但调用方无法证明那一行就是自己(如合并已把自己那行指向新对象、另一行仍指向旧对象时 count 恰为 1), + * 删除即造成该分片永久读不到 → 整单失败。物理删除本就约定在 DB 行删除之后执行,正常路径引用数必然为 0; + * 宁可留下孤儿对象(有保留期清理兜底),也不删掉可能仍被引用的对象。 */ @ExtendWith(MockitoExtension.class) class TaskPayloadReferenceCheckTest { @@ -100,13 +105,14 @@ class TaskPayloadReferenceCheckTest { verify(rustfsObjectStorageService, never()).deleteObject(any()); } + /** count==1 无法证明那一行就是调用方自己(可能正是别的行仍在用)→ 保守不删。 */ @Test - void singleChunkRowIsOwnRowNotSharedReference() { + void singleChunkRowBlocksDeleteBecauseOwnershipCannotBeProven() { when(taskChunkMapper.selectCount(any())).thenReturn(1L); service.deletePayloadIfPresent(RUSTFS_VALUE); - verify(rustfsObjectStorageService).deleteObject("payload-key"); + verify(rustfsObjectStorageService, never()).deleteObject(any()); } @Test diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TransientPayloadChunkSiblingFallbackTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TransientPayloadChunkSiblingFallbackTest.java new file mode 100644 index 00000000..45054715 --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TransientPayloadChunkSiblingFallbackTest.java @@ -0,0 +1,230 @@ +package com.nanri.aiimage.modules.task.service; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.nanri.aiimage.config.InstanceMetadata; +import com.nanri.aiimage.config.StorageProperties; +import com.nanri.aiimage.config.TransientStorageProperties; +import com.nanri.aiimage.modules.file.service.object.RustfsObjectStorageService; +import com.nanri.aiimage.modules.file.service.oss.OssStorageService; +import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; +import com.nanri.aiimage.modules.task.mapper.TaskScopeStateMapper; +import org.junit.jupiter.api.Test; + +import java.nio.charset.StandardCharsets; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * chunk 载荷对象缺失时的「同槽位版本化兄弟对象」读兜底(2026-09-17 线上任务 28459)。 + * + *

事故后 DB 行指向的普通 key({@code chunk-462.json})已被误删,但同槽位的版本化对象 + * ({@code chunk-462-.json})还在。兜底让读路径自动改用兄弟对象,任务即可按已有数据出结果, + * 而不是整单 FAILED——这是防止「删错对象 → 任务永久失败」的最后一道网。 + */ +class TransientPayloadChunkSiblingFallbackTest { + + private static final String DIR = "task-chunk/appearance_patent/28459/" + + "2248d39710545b47b7c7035fc60c4e33924a16174abf65dd1c5a230918e6f61e/"; + private static final String CHUNK_KEY = DIR + "chunk-462.json"; + private static final String CHUNK_POINTER = "rustfs:" + CHUNK_KEY; + private static final String SIBLING_KEY = DIR + "chunk-462-0ffc254b-afad-4279-aab2-e85e3ff955e9.json"; + + // ===== 正常路径 ===== + + /** 对象在:直接读,不得触发列对象(正常读路径不能因为兜底变慢)。 */ + @Test + void presentObjectIsReadWithoutListing() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + when(rustfs.readObjectBytes(CHUNK_KEY)).thenReturn(bytes("[\"原始\"]")); + + String payload = service(rustfs).resolvePayload(CHUNK_POINTER, "read chunk failed"); + + assertEquals("[\"原始\"]", payload); + verify(rustfs, never()).listObjectKeysNewestFirst(anyString(), anyInt()); + } + + /** 对象缺失 + 存在兄弟对象 → 读兄弟对象(本次事故的自动恢复路径)。 */ + @Test + void missingChunkObjectFallsBackToVersionedSibling() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + when(rustfs.readObjectBytes(CHUNK_KEY)).thenThrow(missingKey()); + when(rustfs.listObjectKeysNewestFirst(DIR + "chunk-462-", 50)).thenReturn(List.of(SIBLING_KEY)); + when(rustfs.readObjectBytes(SIBLING_KEY)).thenReturn(bytes("[\"合并后\"]")); + + String payload = service(rustfs).resolvePayload(CHUNK_POINTER, "read chunk failed"); + + assertEquals("[\"合并后\"]", payload); + verify(rustfs).listObjectKeysNewestFirst(DIR + "chunk-462-", 50); + } + + /** 多个兄弟对象:取列表首个(实现约定按最后修改时间倒序,最新在前)。 */ + @Test + void newestSiblingIsPreferred() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + String newest = DIR + "chunk-462-ffffffff-ffff-ffff-ffff-ffffffffffff.json"; + when(rustfs.readObjectBytes(CHUNK_KEY)).thenThrow(missingKey()); + when(rustfs.listObjectKeysNewestFirst(anyString(), anyInt())).thenReturn(List.of(newest, SIBLING_KEY)); + when(rustfs.readObjectBytes(newest)).thenReturn(bytes("[\"最新\"]")); + + assertEquals("[\"最新\"]", service(rustfs).resolvePayload(CHUNK_POINTER, "read chunk failed")); + verify(rustfs).readObjectBytes(newest); + verify(rustfs, never()).readObjectBytes(SIBLING_KEY); + } + + /** 指针是 JSON 编码形式(DB 里的常见写法)时兜底同样生效。 */ + @Test + void jsonQuotedPointerStillFallsBack() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + when(rustfs.readObjectBytes(CHUNK_KEY)).thenThrow(missingKey()); + when(rustfs.listObjectKeysNewestFirst(anyString(), anyInt())).thenReturn(List.of(SIBLING_KEY)); + when(rustfs.readObjectBytes(SIBLING_KEY)).thenReturn(bytes("[\"带引号\"]")); + + assertEquals("[\"带引号\"]", + service(rustfs).resolvePayload('"' + CHUNK_POINTER + '"', "read chunk failed")); + } + + /** 非指针内容(直接就是载荷 JSON)原样返回,不触碰对象存储。 */ + @Test + void rawContentIsReturnedAsIs() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + + assertEquals("[1,2,3]", service(rustfs).resolvePayload("[1,2,3]", "read chunk failed")); + verify(rustfs, never()).readObjectBytes(anyString()); + } + + // ===== 边界:不该兜底的场景 ===== + + /** 非 chunk 槽位(如 row-1.json)缺失 → 不兜底,照旧抛错(避免误配无关对象)。 */ + @Test + void nonChunkSlotDoesNotFallBack() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + String key = "task-result-item/similar-asin/9/hash/row-1.json"; + when(rustfs.readObjectBytes(key)).thenThrow(missingKey()); + + assertThrows(IllegalStateException.class, + () -> service(rustfs).resolvePayload("rustfs:" + key, "read item failed")); + + verify(rustfs, never()).listObjectKeysNewestFirst(anyString(), anyInt()); + } + + /** parsed 槽位(latest.json)缺失 → 不兜底(该槽位没有版本化兄弟语义)。 */ + @Test + void parsedSlotDoesNotFallBack() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + String key = "task-parsed/appearance_patent/28459/hash/latest.json"; + when(rustfs.readObjectBytes(key)).thenThrow(missingKey()); + + assertThrows(IllegalStateException.class, + () -> service(rustfs).resolvePayload("rustfs:" + key, "read parsed failed")); + + verify(rustfs, never()).listObjectKeysNewestFirst(anyString(), anyInt()); + } + + /** 非缺失类失败(权限/网络)→ 不兜底,照旧上抛以便重试。 */ + @Test + void nonMissingFailureDoesNotFallBack() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + when(rustfs.readObjectBytes(CHUNK_KEY)) + .thenThrow(new IllegalStateException("Access Denied")); + + assertThrows(IllegalStateException.class, + () -> service(rustfs).resolvePayload(CHUNK_POINTER, "read chunk failed")); + + verify(rustfs, never()).listObjectKeysNewestFirst(anyString(), anyInt()); + } + + // ===== 异常路径 ===== + + /** 没有兄弟对象 → 抛错,且错误信息必须保留原 errorMessage 与原指针(便于定位)。 */ + @Test + void noSiblingRethrowsWithOriginalPointer() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + when(rustfs.readObjectBytes(CHUNK_KEY)).thenThrow(missingKey()); + when(rustfs.listObjectKeysNewestFirst(anyString(), anyInt())).thenReturn(List.of()); + + IllegalStateException ex = assertThrows(IllegalStateException.class, + () -> service(rustfs).resolvePayload(CHUNK_POINTER, "read chunk failed")); + + assertTrue(ex.getMessage().contains("read chunk failed"), "错误信息=" + ex.getMessage()); + assertTrue(ex.getMessage().contains("chunk-462.json"), "错误信息=" + ex.getMessage()); + } + + /** 兄弟对象读取也失败 → 异常上抛,不得静默返回空内容。 */ + @Test + void siblingReadFailurePropagates() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + when(rustfs.readObjectBytes(CHUNK_KEY)).thenThrow(missingKey()); + when(rustfs.listObjectKeysNewestFirst(anyString(), anyInt())).thenReturn(List.of(SIBLING_KEY)); + when(rustfs.readObjectBytes(SIBLING_KEY)).thenThrow(new IllegalStateException("read timeout")); + + IllegalStateException ex = assertThrows(IllegalStateException.class, + () -> service(rustfs).resolvePayload(CHUNK_POINTER, "read chunk failed")); + + assertTrue(ex.getMessage().contains("read chunk failed"), "错误信息=" + ex.getMessage()); + } + + /** 列对象本身失败(桶不可用)→ 仍按原错误抛出,错误信息不含兄弟查找细节。 */ + @Test + void listFailureDoesNotMaskOriginalError() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + when(rustfs.readObjectBytes(CHUNK_KEY)).thenThrow(missingKey()); + when(rustfs.listObjectKeysNewestFirst(anyString(), anyInt())) + .thenThrow(new IllegalStateException("bucket unavailable")); + + IllegalStateException ex = assertThrows(IllegalStateException.class, + () -> service(rustfs).resolvePayload(CHUNK_POINTER, "read chunk failed")); + + assertTrue(ex.getMessage().contains("read chunk failed"), "错误信息=" + ex.getMessage()); + } + + /** local 指针不走 rustfs 兄弟兜底。 */ + @Test + void localPointerDoesNotUseChunkSiblingFallback() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + + assertThrows(IllegalStateException.class, + () -> service(rustfs).resolvePayload("local:i/server-121/" + CHUNK_KEY, "read chunk failed")); + + verify(rustfs, never()).listObjectKeysNewestFirst(anyString(), anyInt()); + } + + // ===== 辅助 ===== + + private static IllegalStateException missingKey() { + return new IllegalStateException("The specified key does not exist."); + } + + private static byte[] bytes(String value) { + return value.getBytes(StandardCharsets.UTF_8); + } + + private TransientPayloadStorageService service(RustfsObjectStorageService rustfs) { + TransientStorageProperties transientProperties = new TransientStorageProperties(); + transientProperties.setEnabled(true); + transientProperties.setEndpoint("http://127.0.0.1:9000"); + transientProperties.setBucket("json-server"); + transientProperties.setAccessKeyId("ak"); + transientProperties.setAccessKeySecret("sk"); + StorageProperties storageProperties = new StorageProperties(); + storageProperties.setLocalTempDir(System.getProperty("java.io.tmpdir")); + return new TransientPayloadStorageService( + transientProperties, + storageProperties, + rustfs, + mock(OssStorageService.class), + new ObjectMapper(), + new InstanceMetadata("server-121"), + mock(TaskChunkMapper.class), + mock(TaskScopeStateMapper.class)); + } +} diff --git a/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TransientPayloadDeleteReferenceGuardTest.java b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TransientPayloadDeleteReferenceGuardTest.java new file mode 100644 index 00000000..90bd2d1e --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/modules/task/service/TransientPayloadDeleteReferenceGuardTest.java @@ -0,0 +1,180 @@ +package com.nanri.aiimage.modules.task.service; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.nanri.aiimage.config.InstanceMetadata; +import com.nanri.aiimage.config.StorageProperties; +import com.nanri.aiimage.config.TransientStorageProperties; +import com.nanri.aiimage.modules.file.service.oss.OssStorageService; +import com.nanri.aiimage.modules.file.service.object.RustfsObjectStorageService; +import com.nanri.aiimage.modules.task.mapper.TaskChunkMapper; +import com.nanri.aiimage.modules.task.mapper.TaskScopeStateMapper; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * 临时载荷删除前的引用守卫(2026-09-17 线上任务 28459 事故的第二条同族缺陷)。 + * + *

守卫本意是「还有别的行引用同一 pointer 就不物理删」。原实现用 + * {@code chunkCount > 1} 作判据,等价于「恰好还有 1 行引用时可以删」—— + * 这在 caller 已经把自己那行挪走、但**另一行**仍指向旧对象的场景下,会把 + * 对方仍在用的对象删掉。docs 里 CollectDataDeleteConsistencyTest 已明确口径是 + * 「跨任务共享指针仍引用 → 跳过删除」,故判据收紧为「命中任意一行即不删」。 + */ +class TransientPayloadDeleteReferenceGuardTest { + + private static final String RUSTFS_CHUNK_POINTER = + "rustfs:task-chunk/appearance_patent/28459/2248d39710545b47b7c7035fc60c4e33924a16174abf65dd1c5a230918e6f61e/chunk-462.json"; + + // ===== 正常路径 ===== + + /** 没有任何行引用 → 允许物理删除。 */ + @Test + void zeroReferencesAllowsPhysicalDelete() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + TaskChunkMapper chunkMapper = mock(TaskChunkMapper.class); + TaskScopeStateMapper scopeMapper = mock(TaskScopeStateMapper.class); + when(chunkMapper.selectCount(any())).thenReturn(0L); + when(scopeMapper.selectCount(any())).thenReturn(0L); + + service(rustfs, chunkMapper, scopeMapper).deletePayloadIfPresent(RUSTFS_CHUNK_POINTER); + + verify(rustfs).deleteObject(eq("task-chunk/appearance_patent/28459/2248d39710545b47b7c7035fc60c4e33924a16174abf65dd1c5a230918e6f61e/chunk-462.json")); + } + + /** 多条行引用 → 不删(原有行为必须保持)。 */ + @Test + void multipleReferencesBlockDelete() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + TaskChunkMapper chunkMapper = mock(TaskChunkMapper.class); + when(chunkMapper.selectCount(any())).thenReturn(3L); + + service(rustfs, chunkMapper, mock(TaskScopeStateMapper.class)).deletePayloadIfPresent(RUSTFS_CHUNK_POINTER); + + verify(rustfs, never()).deleteObject(anyString()); + } + + // ===== 本次修复点 ===== + + /** 恰好 1 行仍引用 → 必须不删(原实现会删掉这唯一一行仍在用的对象)。 */ + @Test + void singleReferencingChunkRowBlocksDelete() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + TaskChunkMapper chunkMapper = mock(TaskChunkMapper.class); + when(chunkMapper.selectCount(any())).thenReturn(1L); + + service(rustfs, chunkMapper, mock(TaskScopeStateMapper.class)).deletePayloadIfPresent(RUSTFS_CHUNK_POINTER); + + verify(rustfs, never()).deleteObject(anyString()); + } + + /** scope_state 单行引用同样阻断删除(该字段命中即非自我引用)。 */ + @Test + void singleScopeStateReferenceBlocksDelete() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + TaskChunkMapper chunkMapper = mock(TaskChunkMapper.class); + TaskScopeStateMapper scopeMapper = mock(TaskScopeStateMapper.class); + when(chunkMapper.selectCount(any())).thenReturn(0L); + when(scopeMapper.selectCount(any())).thenReturn(1L); + + service(rustfs, chunkMapper, scopeMapper).deletePayloadIfPresent(RUSTFS_CHUNK_POINTER); + + verify(rustfs, never()).deleteObject(anyString()); + } + + // ===== 边界 ===== + + /** 传入 JSON 编码(带引号)的指针——DB 里常见写法,必须能解析并同样受守卫约束。 */ + @Test + void jsonQuotedPointerIsParsedAndStillGuarded() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + TaskChunkMapper chunkMapper = mock(TaskChunkMapper.class); + when(chunkMapper.selectCount(any())).thenReturn(1L); + + service(rustfs, chunkMapper, mock(TaskScopeStateMapper.class)) + .deletePayloadIfPresent('"' + RUSTFS_CHUNK_POINTER + '"'); + + verify(rustfs, never()).deleteObject(anyString()); + } + + /** 非指针值原样传入 → 不解析、不查库、不删。 */ + @Test + void nonPointerValueIsIgnoredWithoutLookup() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + TaskChunkMapper chunkMapper = mock(TaskChunkMapper.class); + + TransientPayloadStorageService service = service(rustfs, chunkMapper, mock(TaskScopeStateMapper.class)); + assertDoesNotThrow(() -> service.deletePayloadIfPresent("{\"rows\":[1,2,3]}")); + + verify(chunkMapper, never()).selectCount(any()); + verify(rustfs, never()).deleteObject(anyString()); + } + + /** 空值不得抛异常、不得触发任何删除。 */ + @Test + void nullOrBlankValueIsIgnored() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + TaskChunkMapper chunkMapper = mock(TaskChunkMapper.class); + TransientPayloadStorageService service = service(rustfs, chunkMapper, mock(TaskScopeStateMapper.class)); + + assertDoesNotThrow(() -> service.deletePayloadIfPresent(null)); + assertDoesNotThrow(() -> service.deletePayloadIfPresent(" ")); + + verify(chunkMapper, never()).selectCount(any()); + verify(rustfs, never()).deleteObject(anyString()); + } + + // ===== 异常路径 ===== + + /** 反查引用的查询抛异常 → 保守不删(原实现即如此,必须保持)。 */ + @Test + void referenceQueryFailureIsConservativeAndSkipsDelete() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + TaskChunkMapper chunkMapper = mock(TaskChunkMapper.class); + when(chunkMapper.selectCount(any())).thenThrow(new IllegalStateException("db down")); + + assertDoesNotThrow(() -> + service(rustfs, chunkMapper, mock(TaskScopeStateMapper.class)) + .deletePayloadIfPresent(RUSTFS_CHUNK_POINTER)); + + verify(rustfs, never()).deleteObject(anyString()); + } + + /** 跨实例的 local 指针:不物理删(留给持有实例的清理任务)。 */ + @Test + void crossInstanceLocalPointerIsNotDeleted() { + RustfsObjectStorageService rustfs = mock(RustfsObjectStorageService.class); + TaskChunkMapper chunkMapper = mock(TaskChunkMapper.class); + when(chunkMapper.selectCount(any())).thenReturn(0L); + + service(rustfs, chunkMapper, mock(TaskScopeStateMapper.class)) + .deletePayloadIfPresent("local:i/server-110/task-chunk/x/1/hash/chunk-1.json"); + + verify(rustfs, never()).deleteObject(anyString()); + } + + private TransientPayloadStorageService service(RustfsObjectStorageService rustfs, + TaskChunkMapper chunkMapper, + TaskScopeStateMapper scopeMapper) { + TransientStorageProperties transientProperties = new TransientStorageProperties(); + transientProperties.setEnabled(true); + StorageProperties storageProperties = new StorageProperties(); + storageProperties.setLocalTempDir(System.getProperty("java.io.tmpdir")); + return new TransientPayloadStorageService( + transientProperties, + storageProperties, + rustfs, + mock(OssStorageService.class), + new ObjectMapper(), + new InstanceMetadata("server-121"), + chunkMapper, + scopeMapper); + } +}