fix(storage): stage retryable s3 upload bodies

This commit is contained in:
dongmucat 2026-04-21 13:09:20 +08:00
parent 37aa366233
commit 1e9d22b528
2 changed files with 58 additions and 5 deletions

View file

@ -18,8 +18,12 @@ import software.amazon.awssdk.services.s3.presigner.model.GetObjectPresignReques
import software.amazon.awssdk.services.s3.presigner.model.PresignedGetObjectRequest;
import java.io.InputStream;
import java.io.IOException;
import java.net.URI;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.time.Duration;
import java.util.List;
@ -106,6 +110,10 @@ public class S3StorageService implements ObjectStorageService {
}
private void putObjectInternal(String key, InputStream data, long size, String contentType) {
putObjectInternal(key, RequestBody.fromInputStream(data, size), size, contentType);
}
private void putObjectInternal(String key, RequestBody requestBody, long size, String contentType) {
s3Client.putObject(
PutObjectRequest.builder()
.bucket(properties.getBucket())
@ -113,25 +121,53 @@ public class S3StorageService implements ObjectStorageService {
.contentType(contentType)
.contentLength(size)
.build(),
RequestBody.fromInputStream(data, size));
requestBody);
}
private Path stagePutObjectBody(InputStream data) throws IOException {
Path stagedBody = Files.createTempFile("skillhub-s3-upload-", ".tmp");
try {
Files.copy(data, stagedBody, StandardCopyOption.REPLACE_EXISTING);
return stagedBody;
} catch (IOException e) {
Files.deleteIfExists(stagedBody);
throw e;
}
}
private void deleteStagedBody(Path stagedBody) {
if (stagedBody == null) {
return;
}
try {
Files.deleteIfExists(stagedBody);
} catch (IOException e) {
log.warn("Failed to clean up staged S3 upload body {}", stagedBody, e);
}
}
@Override public void putObject(String key, InputStream data, long size, String contentType) {
Path stagedBody = null;
try {
if (!properties.isAutoCreateBucket() || bucketPrepared) {
putObjectInternal(key, data, size, contentType);
return;
}
stagedBody = stagePutObjectBody(data);
try {
putObjectInternal(key, data, size, contentType);
putObjectInternal(key, RequestBody.fromFile(stagedBody), size, contentType);
bucketPrepared = true;
} catch (NoSuchBucketException e) {
ensureBucketPrepared();
putObjectInternal(key, data, size, contentType);
putObjectInternal(key, RequestBody.fromFile(stagedBody), size, contentType);
}
} catch (IOException e) {
throw new StorageAccessException("putObject", key, e);
} catch (RuntimeException e) {
throw new StorageAccessException("putObject", key, e);
} finally {
deleteStagedBody(stagedBody);
}
}

View file

@ -16,9 +16,12 @@ import software.amazon.awssdk.services.s3.presigner.S3Presigner;
import software.amazon.awssdk.services.s3.presigner.model.GetObjectPresignRequest;
import java.io.ByteArrayInputStream;
import java.io.InputStream;
import java.net.URI;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@ -115,9 +118,16 @@ class S3StorageServiceTest {
void putObjectShouldCreateBucketAndRetryWhenMissing() {
S3Client client = mock(S3Client.class);
S3Presigner presigner = mock(S3Presigner.class);
List<String> uploadedBodies = new ArrayList<>();
when(client.putObject(any(PutObjectRequest.class), any(RequestBody.class)))
.thenThrow(NoSuchBucketException.builder().message("missing").build())
.thenReturn(PutObjectResponse.builder().eTag("etag").build());
.thenAnswer(invocation -> {
uploadedBodies.add(readBody(invocation.getArgument(1)));
throw NoSuchBucketException.builder().message("missing").build();
})
.thenAnswer(invocation -> {
uploadedBodies.add(readBody(invocation.getArgument(1)));
return PutObjectResponse.builder().eTag("etag").build();
});
when(client.createBucket(any(CreateBucketRequest.class)))
.thenReturn(CreateBucketResponse.builder().build());
TestableS3StorageService service = new TestableS3StorageService(properties(true), client, presigner);
@ -129,6 +139,7 @@ class S3StorageServiceTest {
verify(client, never()).headBucket(any(HeadBucketRequest.class));
verify(client, times(1)).createBucket(any(CreateBucketRequest.class));
verify(client, times(2)).putObject(any(PutObjectRequest.class), any(RequestBody.class));
assertThat(uploadedBodies).containsExactly("hello", "hello");
}
@Test
@ -179,6 +190,12 @@ class S3StorageServiceTest {
return properties;
}
private String readBody(RequestBody body) throws Exception {
try (InputStream inputStream = body.contentStreamProvider().newStream()) {
return new String(inputStream.readAllBytes(), StandardCharsets.UTF_8);
}
}
private URI presignGetObjectUrl(boolean forcePathStyle) {
S3StorageService storageService = new S3StorageService(createProperties(forcePathStyle));
try (var presigner = storageService.buildPresigner()) {