diff --git a/server/skillhub-storage/src/main/java/com/iflytek/skillhub/storage/S3StorageService.java b/server/skillhub-storage/src/main/java/com/iflytek/skillhub/storage/S3StorageService.java index ac965f82..d6a6327f 100644 --- a/server/skillhub-storage/src/main/java/com/iflytek/skillhub/storage/S3StorageService.java +++ b/server/skillhub-storage/src/main/java/com/iflytek/skillhub/storage/S3StorageService.java @@ -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); } } diff --git a/server/skillhub-storage/src/test/java/com/iflytek/skillhub/storage/S3StorageServiceTest.java b/server/skillhub-storage/src/test/java/com/iflytek/skillhub/storage/S3StorageServiceTest.java index be4d6c98..4e243952 100644 --- a/server/skillhub-storage/src/test/java/com/iflytek/skillhub/storage/S3StorageServiceTest.java +++ b/server/skillhub-storage/src/test/java/com/iflytek/skillhub/storage/S3StorageServiceTest.java @@ -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 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()) {