mirror of
https://github.com/iflytek/skillhub.git
synced 2026-10-10 03:27:54 +00:00
Merge pull request #336 from iflytek/fix/s3-bucket-access-check
fix(storage): lazily create missing s3 buckets on upload
This commit is contained in:
commit
60a30190bf
2 changed files with 165 additions and 17 deletions
|
|
@ -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<String> keys) {
|
||||
if (keys.isEmpty()) return;
|
||||
try {
|
||||
ensureBucketPrepared();
|
||||
List<ObjectIdentifier> 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) {
|
||||
|
|
|
|||
|
|
@ -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<String> 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()) {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue