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 8e5452b5..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; @@ -46,8 +50,12 @@ public class S3StorageService implements ObjectStorageService { .connectionAcquisitionTimeout(properties.getConnectionAcquisitionTimeout()); this.s3Client = buildS3Client(httpClientBuilder); this.s3Presigner = buildPresigner(); - log.info("Initialized S3 storage client for bucket '{}' (bucket verification is deferred until first storage access)", - properties.getBucket()); + if (properties.isAutoCreateBucket()) { + log.info("Initialized S3 storage client for bucket '{}' (bucket auto-creation is deferred until first write)", + properties.getBucket()); + } else { + log.info("Initialized S3 storage client for bucket '{}'", properties.getBucket()); + } } protected S3Client buildS3Client(ApacheHttpClient.Builder httpClientBuilder) { @@ -92,27 +100,79 @@ public class S3StorageService implements ObjectStorageService { return; } try { - s3Client.headBucket(HeadBucketRequest.builder().bucket(properties.getBucket()).build()); - } catch (NoSuchBucketException e) { log.info("Bucket '{}' does not exist, creating...", properties.getBucket()); s3Client.createBucket(CreateBucketRequest.builder().bucket(properties.getBucket()).build()); + } catch (BucketAlreadyExistsException | BucketAlreadyOwnedByYouException e) { + log.debug("Bucket '{}' was created concurrently, continuing", properties.getBucket()); } bucketPrepared = true; } } - @Override public void putObject(String key, InputStream data, long size, String contentType) { + 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()) + .key(key) + .contentType(contentType) + .contentLength(size) + .build(), + requestBody); + } + + private Path stagePutObjectBody(InputStream data) throws IOException { + Path stagedBody = Files.createTempFile("skillhub-s3-upload-", ".tmp"); try { - ensureBucketPrepared(); - s3Client.putObject(PutObjectRequest.builder().bucket(properties.getBucket()).key(key).contentType(contentType).contentLength(size).build(), RequestBody.fromInputStream(data, size)); + 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, RequestBody.fromFile(stagedBody), size, contentType); + bucketPrepared = true; + } catch (NoSuchBucketException e) { + ensureBucketPrepared(); + 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); } } @Override public InputStream getObject(String key) { try { - ensureBucketPrepared(); return s3Client.getObject(GetObjectRequest.builder().bucket(properties.getBucket()).key(key).build()); } catch (RuntimeException e) { throw new StorageAccessException("getObject", key, e); @@ -121,7 +181,6 @@ public class S3StorageService implements ObjectStorageService { @Override public void deleteObject(String key) { try { - ensureBucketPrepared(); s3Client.deleteObject(DeleteObjectRequest.builder().bucket(properties.getBucket()).key(key).build()); } catch (RuntimeException e) { throw new StorageAccessException("deleteObject", key, e); @@ -131,7 +190,6 @@ public class S3StorageService implements ObjectStorageService { @Override public void deleteObjects(List keys) { if (keys.isEmpty()) return; try { - ensureBucketPrepared(); List ids = keys.stream().map(k -> ObjectIdentifier.builder().key(k).build()).toList(); s3Client.deleteObjects(DeleteObjectsRequest.builder().bucket(properties.getBucket()).delete(Delete.builder().objects(ids).build()).build()); } catch (RuntimeException e) { @@ -141,7 +199,6 @@ public class S3StorageService implements ObjectStorageService { @Override public boolean exists(String key) { try { - ensureBucketPrepared(); s3Client.headObject(HeadObjectRequest.builder().bucket(properties.getBucket()).key(key).build()); return true; } @@ -151,7 +208,6 @@ public class S3StorageService implements ObjectStorageService { @Override public ObjectMetadata getMetadata(String key) { try { - ensureBucketPrepared(); HeadObjectResponse resp = s3Client.headObject(HeadObjectRequest.builder().bucket(properties.getBucket()).key(key).build()); return new ObjectMetadata(resp.contentLength(), resp.contentType(), resp.lastModified()); } catch (RuntimeException e) { 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 0fa1887c..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 @@ -5,6 +5,7 @@ import software.amazon.awssdk.services.s3.model.GetObjectRequest; import software.amazon.awssdk.core.sync.RequestBody; import software.amazon.awssdk.http.apache.ApacheHttpClient; import software.amazon.awssdk.services.s3.S3Client; +import software.amazon.awssdk.services.s3.model.BucketAlreadyOwnedByYouException; import software.amazon.awssdk.services.s3.model.CreateBucketRequest; import software.amazon.awssdk.services.s3.model.CreateBucketResponse; import software.amazon.awssdk.services.s3.model.HeadBucketRequest; @@ -15,11 +16,15 @@ 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; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; @@ -75,16 +80,78 @@ class S3StorageServiceTest { verify(client).putObject(any(PutObjectRequest.class), any(RequestBody.class)); } + @Test + void putObjectShouldWriteDirectlyWhenBucketAlreadyExists() { + S3Client client = mock(S3Client.class); + S3Presigner presigner = mock(S3Presigner.class); + when(client.putObject(any(PutObjectRequest.class), any(RequestBody.class))) + .thenReturn(PutObjectResponse.builder().eTag("etag").build()); + TestableS3StorageService service = new TestableS3StorageService(properties(true), client, presigner); + + service.init(); + byte[] content = "hello".getBytes(StandardCharsets.UTF_8); + service.putObject("packages/demo.zip", new ByteArrayInputStream(content), content.length, "application/zip"); + + verify(client, never()).headBucket(any(HeadBucketRequest.class)); + verify(client, never()).createBucket(any(CreateBucketRequest.class)); + verify(client).putObject(any(PutObjectRequest.class), any(RequestBody.class)); + } + + @Test + void getObjectShouldNotCreateBucketWhenAutoCreateIsEnabled() { + S3Client client = mock(S3Client.class); + S3Presigner presigner = mock(S3Presigner.class); + when(client.getObject(any(GetObjectRequest.class))) + .thenThrow(NoSuchBucketException.builder().message("missing").build()); + TestableS3StorageService service = new TestableS3StorageService(properties(true), client, presigner); + + service.init(); + assertThatThrownBy(() -> service.getObject("packages/demo.zip")) + .isInstanceOf(StorageAccessException.class); + + verify(client, never()).headBucket(any(HeadBucketRequest.class)); + verify(client, never()).createBucket(any(CreateBucketRequest.class)); + verify(client).getObject(any(GetObjectRequest.class)); + } + + @Test + 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))) + .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); + + service.init(); + byte[] content = "hello".getBytes(StandardCharsets.UTF_8); + service.putObject("packages/demo-1.zip", new ByteArrayInputStream(content), content.length, "application/zip"); + + 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 void putObjectShouldCreateBucketOnlyOnceWhenAutoCreateIsEnabled() { S3Client client = mock(S3Client.class); S3Presigner presigner = mock(S3Presigner.class); - doThrow(NoSuchBucketException.builder().message("missing").build()) - .when(client).headBucket(any(HeadBucketRequest.class)); + when(client.putObject(any(PutObjectRequest.class), any(RequestBody.class))) + .thenThrow(NoSuchBucketException.builder().message("missing").build()) + .thenReturn(PutObjectResponse.builder().eTag("etag-1").build()) + .thenReturn(PutObjectResponse.builder().eTag("etag-2").build()); when(client.createBucket(any(CreateBucketRequest.class))) .thenReturn(CreateBucketResponse.builder().build()); - when(client.putObject(any(PutObjectRequest.class), any(RequestBody.class))) - .thenReturn(PutObjectResponse.builder().eTag("etag").build()); TestableS3StorageService service = new TestableS3StorageService(properties(true), client, presigner); service.init(); @@ -92,7 +159,26 @@ class S3StorageServiceTest { service.putObject("packages/demo-1.zip", new ByteArrayInputStream(content), content.length, "application/zip"); service.putObject("packages/demo-2.zip", new ByteArrayInputStream(content), content.length, "application/zip"); - verify(client, times(1)).headBucket(any(HeadBucketRequest.class)); + verify(client, never()).headBucket(any(HeadBucketRequest.class)); + verify(client, times(1)).createBucket(any(CreateBucketRequest.class)); + verify(client, times(3)).putObject(any(PutObjectRequest.class), any(RequestBody.class)); + } + + @Test + void putObjectShouldRetryWhenBucketWasCreatedConcurrently() { + S3Client client = mock(S3Client.class); + S3Presigner presigner = mock(S3Presigner.class); + when(client.putObject(any(PutObjectRequest.class), any(RequestBody.class))) + .thenThrow(NoSuchBucketException.builder().message("missing").build()) + .thenReturn(PutObjectResponse.builder().eTag("etag").build()); + doThrow(BucketAlreadyOwnedByYouException.builder().message("exists").build()) + .when(client).createBucket(any(CreateBucketRequest.class)); + TestableS3StorageService service = new TestableS3StorageService(properties(true), client, presigner); + + service.init(); + byte[] content = "hello".getBytes(StandardCharsets.UTF_8); + service.putObject("packages/demo.zip", new ByteArrayInputStream(content), content.length, "application/zip"); + verify(client, times(1)).createBucket(any(CreateBucketRequest.class)); verify(client, times(2)).putObject(any(PutObjectRequest.class), any(RequestBody.class)); } @@ -104,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()) {