diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 7efb32beb..3dc6c0139 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -12,7 +12,7 @@ concurrency: cancel-in-progress: true env: - BUILDER_VERSION: v0.9.96 + BUILDER_VERSION: latest BUILDER_SOURCE: releases BUILDER_HOST: https://d19elf31gohf1l.cloudfront.net PACKAGE_NAME: aws-crt-java diff --git a/crt/aws-c-common b/crt/aws-c-common index 3c69b871d..3a5e638ae 160000 --- a/crt/aws-c-common +++ b/crt/aws-c-common @@ -1 +1 @@ -Subproject commit 3c69b871dfa1815231802febf1bb6899f84cccdb +Subproject commit 3a5e638aee99c9f1d65696f7df8c4e5d3dc5805c diff --git a/crt/aws-c-mqtt b/crt/aws-c-mqtt index 2ef9605ec..e35b9ca3f 160000 --- a/crt/aws-c-mqtt +++ b/crt/aws-c-mqtt @@ -1 +1 @@ -Subproject commit 2ef9605ec9c50bea3f921e08022ddd57eed70901 +Subproject commit e35b9ca3f9fcbf1a972c831c5e79046ce56959d1 diff --git a/crt/aws-c-s3 b/crt/aws-c-s3 index a852faa2d..469cbd020 160000 --- a/crt/aws-c-s3 +++ b/crt/aws-c-s3 @@ -1 +1 @@ -Subproject commit a852faa2df3ab2b31fb4cfd64fd3379a2f4ae22e +Subproject commit 469cbd020db52c329631a614e3b8401f3fda7717 diff --git a/crt/aws-c-sdkutils b/crt/aws-c-sdkutils index cb14fea36..528b9dfff 160000 --- a/crt/aws-c-sdkutils +++ b/crt/aws-c-sdkutils @@ -1 +1 @@ -Subproject commit cb14fea362c82c995eebd34e2e96590ab4e0ed58 +Subproject commit 528b9dfff4a804b334875ecf8a0471f7d1366f24 diff --git a/src/main/java/software/amazon/awssdk/crt/s3/ResumeToken.java b/src/main/java/software/amazon/awssdk/crt/s3/ResumeToken.java index 376e01298..e5da66677 100644 --- a/src/main/java/software/amazon/awssdk/crt/s3/ResumeToken.java +++ b/src/main/java/software/amazon/awssdk/crt/s3/ResumeToken.java @@ -64,6 +64,17 @@ public ResumeToken build() { private long numPartsCompleted; private String uploadId; + /* Download specific fields, populated by native code. */ + private String etag; + private String versionId; + private String s3ObjectLastModified; + private long objectSize; + private long objectRangeStart; + private long objectRangeEnd; + private long continuesDownloadedBytes; + private long totalDownloadedBytes; + private long fileLastModifiedEpochNs; + public ResumeToken(PutResumeTokenBuilder builder) { this.nativeType = S3MetaRequestOptions.MetaRequestType.PUT_OBJECT.getNativeValue(); this.partSize = builder.partSize; @@ -122,4 +133,122 @@ public String getUploadId() { return uploadId; } + + /****** + * Download Specific fields. + ******/ + + private void validateDownloadToken(String field) { + if (getType() != S3MetaRequestOptions.MetaRequestType.GET_OBJECT) { + throw new IllegalArgumentException( + "ResumeToken - " + field + " is only defined for Get Object Resume tokens"); + } + } + + /** + * ETag of the S3 object being downloaded, captured from the first response. + * May be null/empty if the download was paused before the first response arrived. + * + * @return etag of the object + */ + public String getEtag() { + validateDownloadToken("etag"); + return etag; + } + + /** + * Version ID of the S3 object being downloaded. + * Optional: null/empty when the bucket is unversioned or the version id was not captured. + * + * @return version id of the object + */ + public String getVersionId() { + validateDownloadToken("version id"); + return versionId; + } + + /** + * Last-Modified of the S3 object being downloaded, in HTTP-date format + * (RFC 9110 5.6.7, "Wed, 09 Oct 2024 22:28:00 GMT"), captured from the first + * response. The exact string in the response header. + * Optional: null/empty if the value was not captured before the pause. + * + * @return Last-Modified of the object as an HTTP-date string + */ + public String getS3ObjectLastModified() { + validateDownloadToken("s3 object last modified"); + return s3ObjectLastModified; + } + + /** + * Total size of the S3 object being downloaded, regardless of any Range header + * on the request (for a ranged download this is larger than the range being + * fetched). 0 if the download was paused before the object size was discovered. + * + * @return total object size in bytes + */ + public long getObjectSize() { + validateDownloadToken("object size"); + return objectSize; + } + + /** + * Absolute byte offset in the object where the download's range starts. + * 0 for a download without a Range header. + * + * @return range start offset in bytes + */ + public long getObjectRangeStart() { + validateDownloadToken("object range start"); + return objectRangeStart; + } + + /** + * Absolute byte offset in the object where the download's range ends (inclusive). + * For a download without a Range header this is object size - 1. + * 0 if the download was paused before the object size was discovered. + * + * @return range end offset in bytes (inclusive) + */ + public long getObjectRangeEnd() { + validateDownloadToken("object range end"); + return objectRangeEnd; + } + + /** + * Number of bytes downloaded continuously from the start of the range, with no + * gaps. Everything before this offset (relative to the object range start) has + * been downloaded. + * + * @return continuously downloaded bytes + */ + public long getContinuesDownloadedBytes() { + validateDownloadToken("continues downloaded bytes"); + return continuesDownloadedBytes; + } + + /** + * Total number of bytes downloaded before the pause. May be greater than + * {@link #getContinuesDownloadedBytes()} when parts completed out of order, + * leaving gaps. Equals it when delivery was strictly in order. + * + * @return total downloaded bytes + */ + public long getTotalDownloadedBytes() { + validateDownloadToken("total downloaded bytes"); + return totalDownloadedBytes; + } + + /** + * Last-modified time of the local receive file (nanoseconds since the Unix + * epoch), captured after the file handle was closed during the pause. + * Only set when the download was writing to a file; 0 for downloads that + * deliver via body callback or when the timestamp could not be queried. + * + * @return local receive file last-modified time in epoch nanoseconds + */ + public long getFileLastModifiedEpochNs() { + validateDownloadToken("file last modified epoch ns"); + return fileLastModifiedEpochNs; + } } diff --git a/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequest.java b/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequest.java index 7824bac78..7ea426459 100644 --- a/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequest.java +++ b/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequest.java @@ -5,7 +5,9 @@ package software.amazon.awssdk.crt.s3; import java.util.concurrent.CompletableFuture; +import software.amazon.awssdk.crt.CRT; import software.amazon.awssdk.crt.CrtResource; +import software.amazon.awssdk.crt.CrtRuntimeException; public class S3MetaRequest extends CrtResource { @@ -21,6 +23,23 @@ private void onShutdownComplete() { this.shutdownComplete.complete(null); } + /** + * Called from native when an async pause completes. The resume token is null when no + * resumable state was captured, which is not an error. + * + * @param future the future to complete with the pause result + * @param errorCode 0 on success, otherwise the CRT error code that made the pause fail + * @param resumeToken the resume token, or null if no resumable state was captured + */ + private static void onPauseComplete( + CompletableFuture future, int errorCode, ResumeToken resumeToken) { + if (errorCode != CRT.AWS_CRT_SUCCESS) { + future.completeExceptionally(new CrtRuntimeException(errorCode)); + return; + } + future.complete(resumeToken); + } + /** * Determines whether a resource releases its dependencies at the same time the * native handle is released or if it waits. Resources that wait are responsible @@ -75,6 +94,38 @@ public ResumeToken pause() { return s3MetaRequestPause(getNativeHandle()); } + /** + * Asynchronously pause the meta request. Works for both uploads (PUT) and downloads (GET). + * The returned future completes once all in-flight work has finished (in-flight parts for + * uploads, file writes for downloads) and the resume token is ready. + *

+ * For PutObject resume, input stream should always start at the beginning, + * already uploaded parts will be skipped, but checksums on those will be verified if + * the request specified a checksum algorithm. + *

+ * Note: consuming a download (GET) resume token to resume via meta request options is not + * supported yet. To resume a download, issue a new ranged GET starting at + * {@link ResumeToken#getContinuesDownloadedBytes()} (offset from + * {@link ResumeToken#getObjectRangeStart()}) through the end of the original download. + * + * @return future completed with the resume token once the pause completes. The token may be + * null if the request had not progressed far enough to produce one (equivalent to + * restarting the transfer). Completed exceptionally with a CrtRuntimeException if + * the pause failed. + */ + public CompletableFuture pauseAsync() { + if (isNull()) { + throw new IllegalStateException("S3MetaRequest has been closed."); + } + CompletableFuture future = new CompletableFuture<>(); + try { + s3MetaRequestPauseAsync(getNativeHandle(), future); + } catch (Exception e) { + future.completeExceptionally(e); + } + return future; + } + /** * Increment the flow-control window, so that response data continues downloading. *

@@ -114,5 +165,7 @@ public void incrementReadWindow(long bytes) { private static native ResumeToken s3MetaRequestPause(long s3MetaRequest); + private static native void s3MetaRequestPauseAsync(long s3MetaRequest, CompletableFuture future); + private static native void s3MetaRequestIncrementReadWindow(long s3MetaRequest, long bytes); } diff --git a/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequestResponseHandler.java b/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequestResponseHandler.java index 004dac919..5d194dcff 100644 --- a/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequestResponseHandler.java +++ b/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequestResponseHandler.java @@ -79,4 +79,25 @@ default void onProgress(final S3MetaRequestProgress progress) { */ default void onTelemetry(S3RequestMetrics requestMetrics) { } + + /** + * Invoked with a resume token when the meta request fails unexpectedly. + * Allows persisting state for later resume without re-transferring completed parts. + * Supported for both upload (PUT) and download (GET) meta requests. + * Always invoked exactly once on unexpected failure; not invoked on success or + * explicit pause (use {@link S3MetaRequest#pauseAsync} for that). + * The token is null when no resumable state was captured. + *

+ * WARNING: for a file download with + * {@link S3MetaRequestOptions#withResponseFileDeleteOnFailure} set true, the deletion + * is respected — the partial file is deleted on error, leaving nothing to resume on, + * and this callback fires with a null token. Do not set responseFileDeleteOnFailure + * if you intend to resume from this callback's token. + * + * @param errorCode the CRT error code that caused the meta request to fail + * @param resumeToken resumable state captured at failure time, null if no + * resumable state was captured + */ + default void onErrorResumeToken(final int errorCode, final ResumeToken resumeToken) { + } } diff --git a/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequestResponseHandlerNativeAdapter.java b/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequestResponseHandlerNativeAdapter.java index 59818c1f5..aa2ede089 100644 --- a/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequestResponseHandlerNativeAdapter.java +++ b/src/main/java/software/amazon/awssdk/crt/s3/S3MetaRequestResponseHandlerNativeAdapter.java @@ -36,4 +36,8 @@ void onProgress(final S3MetaRequestProgress progress) { void onTelemetry(final S3RequestMetrics requestMetrics) { responseHandler.onTelemetry(requestMetrics); } + + void onErrorResumeToken(final int errorCode, final ResumeToken resumeToken) { + responseHandler.onErrorResumeToken(errorCode, resumeToken); + } } diff --git a/src/main/resources/META-INF/native-image/software.amazon.awssdk/crt/aws-crt/jni-config.json b/src/main/resources/META-INF/native-image/software.amazon.awssdk/crt/aws-crt/jni-config.json index bd44d99a0..a6197871f 100644 --- a/src/main/resources/META-INF/native-image/software.amazon.awssdk/crt/aws-crt/jni-config.json +++ b/src/main/resources/META-INF/native-image/software.amazon.awssdk/crt/aws-crt/jni-config.json @@ -2124,20 +2124,47 @@ { "name": "software.amazon.awssdk.crt.s3.ResumeToken", "fields": [ + { + "name": "continuesDownloadedBytes" + }, + { + "name": "etag" + }, + { + "name": "fileLastModifiedEpochNs" + }, { "name": "nativeType" }, { "name": "numPartsCompleted" }, + { + "name": "objectRangeEnd" + }, + { + "name": "objectRangeStart" + }, + { + "name": "objectSize" + }, { "name": "partSize" }, + { + "name": "s3ObjectLastModified" + }, + { + "name": "totalDownloadedBytes" + }, { "name": "totalNumParts" }, { "name": "uploadId" + }, + { + "name": "versionId" } ], "methods": [ @@ -2204,6 +2231,14 @@ { "name": "software.amazon.awssdk.crt.s3.S3MetaRequest", "methods": [ + { + "name": "onPauseComplete", + "parameterTypes": [ + "java.util.concurrent.CompletableFuture", + "int", + "software.amazon.awssdk.crt.s3.ResumeToken" + ] + }, { "name": "onShutdownComplete", "parameterTypes": [] @@ -2364,6 +2399,13 @@ { "name": "software.amazon.awssdk.crt.s3.S3MetaRequestResponseHandlerNativeAdapter", "methods": [ + { + "name": "onErrorResumeToken", + "parameterTypes": [ + "int", + "software.amazon.awssdk.crt.s3.ResumeToken" + ] + }, { "name": "onFinished", "parameterTypes": [ diff --git a/src/native/java_class_ids.c b/src/native/java_class_ids.c index 2730dfcfd..f09b6dfcf 100644 --- a/src/native/java_class_ids.c +++ b/src/native/java_class_ids.c @@ -710,9 +710,17 @@ struct java_s3_meta_request_properties s3_meta_request_properties; static void s_cache_s3_meta_request_properties(JNIEnv *env) { jclass cls = (*env)->FindClass(env, "software/amazon/awssdk/crt/s3/S3MetaRequest"); AWS_FATAL_ASSERT(cls); + s3_meta_request_properties.s3_meta_request_class = (*env)->NewGlobalRef(env, cls); s3_meta_request_properties.onShutdownComplete = (*env)->GetMethodID(env, cls, "onShutdownComplete", "()V"); AWS_FATAL_ASSERT(s3_meta_request_properties.onShutdownComplete); + + s3_meta_request_properties.on_pause_complete_method_id = (*env)->GetStaticMethodID( + env, + cls, + "onPauseComplete", + "(Ljava/util/concurrent/CompletableFuture;ILsoftware/amazon/awssdk/crt/s3/ResumeToken;)V"); + AWS_FATAL_ASSERT(s3_meta_request_properties.on_pause_complete_method_id); } struct java_s3_meta_request_response_handler_native_adapter_properties @@ -740,6 +748,10 @@ static void s_cache_s3_meta_request_response_handler_native_adapter_properties(J s3_meta_request_response_handler_native_adapter_properties.onTelemetry = (*env)->GetMethodID(env, cls, "onTelemetry", "(Lsoftware/amazon/awssdk/crt/s3/S3RequestMetrics;)V"); + + s3_meta_request_response_handler_native_adapter_properties.onErrorResumeToken = + (*env)->GetMethodID(env, cls, "onErrorResumeToken", "(ILsoftware/amazon/awssdk/crt/s3/ResumeToken;)V"); + AWS_FATAL_ASSERT(s3_meta_request_response_handler_native_adapter_properties.onErrorResumeToken); } struct java_completable_future_properties completable_future_properties; @@ -1120,6 +1132,33 @@ static void s_cache_s3_meta_request_resume_token(JNIEnv *env) { s3_meta_request_resume_token_properties.upload_id_field_id = (*env)->GetFieldID(env, cls, "uploadId", "Ljava/lang/String;"); AWS_FATAL_ASSERT(s3_meta_request_resume_token_properties.upload_id_field_id); + + /* download specific fields */ + s3_meta_request_resume_token_properties.etag_field_id = (*env)->GetFieldID(env, cls, "etag", "Ljava/lang/String;"); + AWS_FATAL_ASSERT(s3_meta_request_resume_token_properties.etag_field_id); + s3_meta_request_resume_token_properties.version_id_field_id = + (*env)->GetFieldID(env, cls, "versionId", "Ljava/lang/String;"); + AWS_FATAL_ASSERT(s3_meta_request_resume_token_properties.version_id_field_id); + s3_meta_request_resume_token_properties.s3_object_last_modified_field_id = + (*env)->GetFieldID(env, cls, "s3ObjectLastModified", "Ljava/lang/String;"); + AWS_FATAL_ASSERT(s3_meta_request_resume_token_properties.s3_object_last_modified_field_id); + s3_meta_request_resume_token_properties.object_size_field_id = (*env)->GetFieldID(env, cls, "objectSize", "J"); + AWS_FATAL_ASSERT(s3_meta_request_resume_token_properties.object_size_field_id); + s3_meta_request_resume_token_properties.object_range_start_field_id = + (*env)->GetFieldID(env, cls, "objectRangeStart", "J"); + AWS_FATAL_ASSERT(s3_meta_request_resume_token_properties.object_range_start_field_id); + s3_meta_request_resume_token_properties.object_range_end_field_id = + (*env)->GetFieldID(env, cls, "objectRangeEnd", "J"); + AWS_FATAL_ASSERT(s3_meta_request_resume_token_properties.object_range_end_field_id); + s3_meta_request_resume_token_properties.continuous_downloaded_bytes_field_id = + (*env)->GetFieldID(env, cls, "continuesDownloadedBytes", "J"); + AWS_FATAL_ASSERT(s3_meta_request_resume_token_properties.continuous_downloaded_bytes_field_id); + s3_meta_request_resume_token_properties.total_downloaded_bytes_field_id = + (*env)->GetFieldID(env, cls, "totalDownloadedBytes", "J"); + AWS_FATAL_ASSERT(s3_meta_request_resume_token_properties.total_downloaded_bytes_field_id); + s3_meta_request_resume_token_properties.file_last_modified_epoch_ns_field_id = + (*env)->GetFieldID(env, cls, "fileLastModifiedEpochNs", "J"); + AWS_FATAL_ASSERT(s3_meta_request_resume_token_properties.file_last_modified_epoch_ns_field_id); } struct java_aws_mqtt5_connack_packet_properties mqtt5_connack_packet_properties; diff --git a/src/native/java_class_ids.h b/src/native/java_class_ids.h index 848a026ab..43e97cbf2 100644 --- a/src/native/java_class_ids.h +++ b/src/native/java_class_ids.h @@ -319,9 +319,11 @@ struct java_s3_client_properties { }; extern struct java_s3_client_properties s3_client_properties; -/* S3Client */ +/* S3MetaRequest */ struct java_s3_meta_request_properties { + jclass s3_meta_request_class; jmethodID onShutdownComplete; + jmethodID on_pause_complete_method_id; }; extern struct java_s3_meta_request_properties s3_meta_request_properties; @@ -332,6 +334,7 @@ struct java_s3_meta_request_response_handler_native_adapter_properties { jmethodID onResponseHeaders; jmethodID onProgress; jmethodID onTelemetry; + jmethodID onErrorResumeToken; }; extern struct java_s3_meta_request_response_handler_native_adapter_properties s3_meta_request_response_handler_native_adapter_properties; @@ -510,6 +513,16 @@ struct java_aws_s3_meta_request_resume_token { jfieldID total_num_parts_field_id; jfieldID num_parts_completed_field_id; jfieldID upload_id_field_id; + /* download specific fields */ + jfieldID etag_field_id; + jfieldID version_id_field_id; + jfieldID s3_object_last_modified_field_id; + jfieldID object_size_field_id; + jfieldID object_range_start_field_id; + jfieldID object_range_end_field_id; + jfieldID continuous_downloaded_bytes_field_id; + jfieldID total_downloaded_bytes_field_id; + jfieldID file_last_modified_epoch_ns_field_id; }; extern struct java_aws_s3_meta_request_resume_token s3_meta_request_resume_token_properties; diff --git a/src/native/s3_client.c b/src/native/s3_client.c index ccaae2fdc..d9acad534 100644 --- a/src/native/s3_client.c +++ b/src/native/s3_client.c @@ -1257,6 +1257,159 @@ static struct aws_s3_meta_request_resume_token *s_native_resume_token_from_java_ return resume_token; } +/* Create a Java ResumeToken object from a native resume token, populating both the + * common fields and the type-specific (upload vs download) fields. + * Returns NULL (with a pending Java exception) on failure. */ +static jobject s_java_resume_token_from_native_new(JNIEnv *env, struct aws_s3_meta_request_resume_token *resume_token) { + + jobject resume_token_jni = (*env)->NewObject( + env, + s3_meta_request_resume_token_properties.s3_meta_request_resume_token_class, + s3_meta_request_resume_token_properties.s3_meta_request_resume_token_constructor_method_id); + if ((*env)->ExceptionCheck(env) || resume_token_jni == NULL) { + return NULL; + } + + enum aws_s3_meta_request_type type = aws_s3_meta_request_resume_token_type(resume_token); + (*env)->SetIntField(env, resume_token_jni, s3_meta_request_resume_token_properties.native_type_field_id, type); + (*env)->SetLongField( + env, + resume_token_jni, + s3_meta_request_resume_token_properties.part_size_field_id, + (jlong)aws_s3_meta_request_resume_token_part_size(resume_token)); + (*env)->SetLongField( + env, + resume_token_jni, + s3_meta_request_resume_token_properties.total_num_parts_field_id, + (jlong)aws_s3_meta_request_resume_token_total_num_parts(resume_token)); + (*env)->SetLongField( + env, + resume_token_jni, + s3_meta_request_resume_token_properties.num_parts_completed_field_id, + (jlong)aws_s3_meta_request_resume_token_num_parts_completed(resume_token)); + + if (type == AWS_S3_META_REQUEST_TYPE_PUT_OBJECT) { + struct aws_byte_cursor upload_id_cur = aws_s3_meta_request_resume_token_upload_id(resume_token); + jstring upload_id_jni = aws_jni_string_from_cursor(env, &upload_id_cur); + (*env)->SetObjectField( + env, resume_token_jni, s3_meta_request_resume_token_properties.upload_id_field_id, upload_id_jni); + (*env)->DeleteLocalRef(env, upload_id_jni); + } else if (type == AWS_S3_META_REQUEST_TYPE_GET_OBJECT) { + struct aws_byte_cursor etag_cur = aws_s3_meta_request_resume_token_etag(resume_token); + if (etag_cur.len > 0) { + jstring etag_jni = aws_jni_string_from_cursor(env, &etag_cur); + (*env)->SetObjectField( + env, resume_token_jni, s3_meta_request_resume_token_properties.etag_field_id, etag_jni); + (*env)->DeleteLocalRef(env, etag_jni); + } + + struct aws_byte_cursor version_id_cur = aws_s3_meta_request_resume_token_version_id(resume_token); + if (version_id_cur.len > 0) { + jstring version_id_jni = aws_jni_string_from_cursor(env, &version_id_cur); + (*env)->SetObjectField( + env, resume_token_jni, s3_meta_request_resume_token_properties.version_id_field_id, version_id_jni); + (*env)->DeleteLocalRef(env, version_id_jni); + } + + struct aws_byte_cursor last_modified_cur = + aws_s3_meta_request_resume_token_s3_object_last_modified(resume_token); + if (last_modified_cur.len > 0) { + jstring last_modified_jni = aws_jni_string_from_cursor(env, &last_modified_cur); + (*env)->SetObjectField( + env, + resume_token_jni, + s3_meta_request_resume_token_properties.s3_object_last_modified_field_id, + last_modified_jni); + (*env)->DeleteLocalRef(env, last_modified_jni); + } + + (*env)->SetLongField( + env, + resume_token_jni, + s3_meta_request_resume_token_properties.object_size_field_id, + (jlong)aws_s3_meta_request_resume_token_object_size(resume_token)); + (*env)->SetLongField( + env, + resume_token_jni, + s3_meta_request_resume_token_properties.object_range_start_field_id, + (jlong)aws_s3_meta_request_resume_token_object_range_start(resume_token)); + (*env)->SetLongField( + env, + resume_token_jni, + s3_meta_request_resume_token_properties.object_range_end_field_id, + (jlong)aws_s3_meta_request_resume_token_object_range_end(resume_token)); + (*env)->SetLongField( + env, + resume_token_jni, + s3_meta_request_resume_token_properties.continuous_downloaded_bytes_field_id, + (jlong)aws_s3_meta_request_resume_token_continuous_downloaded_bytes(resume_token)); + (*env)->SetLongField( + env, + resume_token_jni, + s3_meta_request_resume_token_properties.total_downloaded_bytes_field_id, + (jlong)aws_s3_meta_request_resume_token_total_downloaded_bytes(resume_token)); + (*env)->SetLongField( + env, + resume_token_jni, + s3_meta_request_resume_token_properties.file_last_modified_epoch_ns_field_id, + (jlong)aws_s3_meta_request_resume_token_file_last_modified_epoch_ns(resume_token)); + } + + return resume_token_jni; +} + +static void s_on_s3_meta_request_error_resume_token_callback( + struct aws_s3_meta_request *meta_request, + struct aws_s3_meta_request_resume_token *resume_token, + int error_code, + void *user_data) { + + struct s3_client_make_meta_request_callback_data *callback_data = + (struct s3_client_make_meta_request_callback_data *)user_data; + + /********** JNI ENV ACQUIRE **********/ + struct aws_jvm_env_context jvm_env_context = aws_jni_acquire_thread_env(callback_data->jvm); + JNIEnv *env = jvm_env_context.env; + if (env == NULL) { + /* If we can't get an environment, then the JVM is probably shutting down. Don't crash. */ + return; + } + + if (callback_data->java_s3_meta_request_response_handler_native_adapter != NULL) { + jobject resume_token_jni = NULL; + if (resume_token != NULL) { + resume_token_jni = s_java_resume_token_from_native_new(env, resume_token); + if (resume_token_jni == NULL && aws_jni_check_and_clear_exception(env)) { + AWS_LOGF_ERROR( + AWS_LS_S3_META_REQUEST, + "id=%p: Ignored Exception from S3MetaRequest.onErrorResumeToken token conversion", + (void *)meta_request); + } + } + + (*env)->CallVoidMethod( + env, + callback_data->java_s3_meta_request_response_handler_native_adapter, + s3_meta_request_response_handler_native_adapter_properties.onErrorResumeToken, + error_code, + resume_token_jni); + + if (aws_jni_check_and_clear_exception(env)) { + AWS_LOGF_ERROR( + AWS_LS_S3_META_REQUEST, + "id=%p: Ignored Exception from S3MetaRequest.onErrorResumeToken callback", + (void *)meta_request); + } + + if (resume_token_jni) { + (*env)->DeleteLocalRef(env, resume_token_jni); + } + } + + aws_jni_release_thread_env(callback_data->jvm, &jvm_env_context); + /********** JNI ENV RELEASE **********/ +} + JNIEXPORT jlong JNICALL Java_software_amazon_awssdk_crt_s3_S3Client_s3ClientMakeMetaRequest( JNIEnv *env, jclass jni_class, @@ -1423,6 +1576,7 @@ JNIEXPORT jlong JNICALL Java_software_amazon_awssdk_crt_s3_S3Client_s3ClientMake .progress_callback = s_on_s3_meta_request_progress_callback, .telemetry_callback = s_on_s3_meta_request_telemetry_callback, .shutdown_callback = s_on_s3_meta_request_shutdown_complete_callback, + .on_error_resume_token = s_on_s3_meta_request_error_resume_token_callback, .endpoint = jni_endpoint != NULL ? &endpoint : NULL, .resume_token = resume_token, .object_size_hint = jni_object_size_hint != NULL ? &object_size_hint : NULL, @@ -1549,49 +1703,116 @@ JNIEXPORT jobject JNICALL Java_software_amazon_awssdk_crt_s3_S3MetaRequest_s3Met jobject resume_token_jni = NULL; if (resume_token != NULL) { - resume_token_jni = (*env)->NewObject( - env, - s3_meta_request_resume_token_properties.s3_meta_request_resume_token_class, - s3_meta_request_resume_token_properties.s3_meta_request_resume_token_constructor_method_id); - if ((*env)->ExceptionCheck(env) || resume_token_jni == NULL) { + resume_token_jni = s_java_resume_token_from_native_new(env, resume_token); + if (resume_token_jni == NULL) { aws_jni_throw_runtime_exception(env, "S3MetaRequest.s3MetaRequestPause: Failed to create ResumeToken."); - goto on_done; } + } + + aws_s3_meta_request_resume_token_release(resume_token); + return resume_token_jni; +} + +struct s3_meta_request_pause_async_callback_data { + struct aws_allocator *allocator; + JavaVM *jvm; + jobject java_future; +}; + +static void s_on_s3_meta_request_pause_async_complete( + struct aws_s3_meta_request *meta_request, + struct aws_s3_meta_request_resume_token *resume_token, + int error_code, + void *user_data) { + + (void)meta_request; + + struct s3_meta_request_pause_async_callback_data *pause_callback_data = + (struct s3_meta_request_pause_async_callback_data *)user_data; + + /********** JNI ENV ACQUIRE **********/ + struct aws_jvm_env_context jvm_env_context = aws_jni_acquire_thread_env(pause_callback_data->jvm); + JNIEnv *env = jvm_env_context.env; + if (env == NULL) { + /* If we can't get an environment, then the JVM is probably shutting down. Don't crash. */ + aws_mem_release(pause_callback_data->allocator, pause_callback_data); + return; + } - enum aws_s3_meta_request_type type = aws_s3_meta_request_resume_token_type(resume_token); - if (type != AWS_S3_META_REQUEST_TYPE_PUT_OBJECT) { - aws_jni_throw_runtime_exception(env, "S3MetaRequest.s3MetaRequestPause: Failed to convert resume token."); - goto on_done; + /* Build the token when there is one; a NULL token is not an error (it means no resumable + * state was captured). Let Java own future completion and exception creation. */ + jobject resume_token_jni = NULL; + if (resume_token != NULL) { + resume_token_jni = s_java_resume_token_from_native_new(env, resume_token); + if (resume_token_jni == NULL && aws_jni_check_and_clear_exception(env)) { + AWS_LOGF_ERROR( + AWS_LS_S3_META_REQUEST, + "id=%p: Ignored Exception from S3MetaRequest.pauseAsync token conversion", + (void *)meta_request); } + } - (*env)->SetIntField(env, resume_token_jni, s3_meta_request_resume_token_properties.native_type_field_id, type); - (*env)->SetLongField( - env, - resume_token_jni, - s3_meta_request_resume_token_properties.part_size_field_id, - aws_s3_meta_request_resume_token_part_size(resume_token)); - (*env)->SetLongField( - env, - resume_token_jni, - s3_meta_request_resume_token_properties.total_num_parts_field_id, - aws_s3_meta_request_resume_token_total_num_parts(resume_token)); - (*env)->SetLongField( - env, - resume_token_jni, - s3_meta_request_resume_token_properties.num_parts_completed_field_id, - aws_s3_meta_request_resume_token_num_parts_completed(resume_token)); + (*env)->CallStaticVoidMethod( + env, + s3_meta_request_properties.s3_meta_request_class, + s3_meta_request_properties.on_pause_complete_method_id, + pause_callback_data->java_future, + error_code, + resume_token_jni); - struct aws_byte_cursor upload_id_cur = aws_s3_meta_request_resume_token_upload_id(resume_token); - jstring upload_id_jni = aws_jni_string_from_cursor(env, &upload_id_cur); - (*env)->SetObjectField( - env, resume_token_jni, s3_meta_request_resume_token_properties.upload_id_field_id, upload_id_jni); + if (aws_jni_check_and_clear_exception(env)) { + AWS_LOGF_ERROR( + AWS_LS_S3_META_REQUEST, + "id=%p: Ignored Exception from S3MetaRequest.onPauseComplete", + (void *)meta_request); + } - (*env)->DeleteLocalRef(env, upload_id_jni); + if (resume_token_jni) { + (*env)->DeleteLocalRef(env, resume_token_jni); } -on_done: - aws_s3_meta_request_resume_token_release(resume_token); - return resume_token_jni; + (*env)->DeleteGlobalRef(env, pause_callback_data->java_future); + + aws_jni_release_thread_env(pause_callback_data->jvm, &jvm_env_context); + /********** JNI ENV RELEASE **********/ + + aws_mem_release(pause_callback_data->allocator, pause_callback_data); +} + +JNIEXPORT void JNICALL Java_software_amazon_awssdk_crt_s3_S3MetaRequest_s3MetaRequestPauseAsync( + JNIEnv *env, + jclass jni_class, + jlong jni_s3_meta_request, + jobject java_future) { + + (void)jni_class; + aws_cache_jni_ids(env); + + struct aws_s3_meta_request *meta_request = (struct aws_s3_meta_request *)jni_s3_meta_request; + if (!meta_request) { + aws_raise_error(AWS_ERROR_INVALID_ARGUMENT); + aws_jni_throw_illegal_argument_exception( + env, "S3MetaRequest.s3MetaRequestPauseAsync: Invalid/null meta request"); + return; + } + + struct aws_allocator *allocator = aws_jni_get_allocator(); + struct s3_meta_request_pause_async_callback_data *pause_callback_data = + aws_mem_calloc(allocator, 1, sizeof(struct s3_meta_request_pause_async_callback_data)); + pause_callback_data->allocator = allocator; + + jint jvmresult = (*env)->GetJavaVM(env, &pause_callback_data->jvm); + AWS_FATAL_ASSERT(jvmresult == 0); + + pause_callback_data->java_future = (*env)->NewGlobalRef(env, java_future); + AWS_FATAL_ASSERT(pause_callback_data->java_future != NULL); + + if (aws_s3_meta_request_pause_async(meta_request, s_on_s3_meta_request_pause_async_complete, pause_callback_data)) { + (*env)->DeleteGlobalRef(env, pause_callback_data->java_future); + aws_mem_release(allocator, pause_callback_data); + aws_jni_throw_runtime_exception(env, "S3MetaRequest.s3MetaRequestPauseAsync: Failed to initiate pause"); + return; + } } JNIEXPORT void JNICALL Java_software_amazon_awssdk_crt_s3_S3MetaRequest_s3MetaRequestIncrementReadWindow( diff --git a/src/test/java/software/amazon/awssdk/crt/test/S3ClientTest.java b/src/test/java/software/amazon/awssdk/crt/test/S3ClientTest.java index 82179d31b..46cb8afae 100644 --- a/src/test/java/software/amazon/awssdk/crt/test/S3ClientTest.java +++ b/src/test/java/software/amazon/awssdk/crt/test/S3ClientTest.java @@ -1262,6 +1262,197 @@ public long getLength() { } } + @Test + public void testS3PutPauseAsyncResume() { + skipIfAndroid(); + skipIfNetworkUnavailable(); + Assume.assumeTrue(hasAwsCredentials()); + + S3ClientOptions clientOptions = new S3ClientOptions() + .withRegion(REGION); + try (S3Client client = createS3Client(clientOptions)) { + CompletableFuture onFinishedFuture = new CompletableFuture<>(); + CompletableFuture onProgressFuture = new CompletableFuture<>(); + S3MetaRequestResponseHandler responseHandler = createTestPutPauseResumeHandler(onFinishedFuture, + onProgressFuture); + + final ByteBuffer payload = ByteBuffer.wrap(createTestPayload(128 * 1024 * 1024)); + HttpRequestBodyStream payloadStream = new HttpRequestBodyStream() { + @Override + public boolean sendRequestBody(ByteBuffer outBuffer) { + ByteBufferUtils.transferData(payload, outBuffer); + return payload.remaining() == 0; + } + + @Override + public boolean resetPosition() { + return true; + } + + @Override + public long getLength() { + return payload.capacity(); + } + }; + + HttpHeader[] headers = { new HttpHeader("Host", ENDPOINT), + new HttpHeader("Content-Length", Integer.valueOf(payload.capacity()).toString()), }; + + HttpRequest httpRequest = new HttpRequest("PUT", + uploadObjectPathInit("/put_object_test_async_pause_128MB"), headers, payloadStream); + + S3MetaRequestOptions metaRequestOptions = new S3MetaRequestOptions() + .withMetaRequestType(MetaRequestType.PUT_OBJECT) + .withChecksumAlgorithm(ChecksumAlgorithm.CRC32) + .withHttpRequest(httpRequest) + .withResponseHandler(responseHandler); + + ResumeToken resumeToken; + try (S3MetaRequest metaRequest = client.makeMetaRequest(metaRequestOptions)) { + onProgressFuture.get(); + + /* Async pause: the future completes once in-flight parts have finished. */ + resumeToken = metaRequest.pauseAsync().get(); + Assert.assertNotNull(resumeToken); + Assert.assertEquals(MetaRequestType.PUT_OBJECT, resumeToken.getType()); + Assert.assertNotNull(resumeToken.getUploadId()); + Assert.assertTrue(resumeToken.getPartSize() > 0); + /* download-only getters must reject upload tokens */ + assertThrows(IllegalArgumentException.class, () -> resumeToken.getObjectSize()); + + Throwable thrown = assertThrows(Throwable.class, + () -> onFinishedFuture.get()); + + Assert.assertEquals("AWS_ERROR_S3_PAUSED", ((CrtRuntimeException) thrown.getCause()).errorName); + } + + final ByteBuffer payloadResume = ByteBuffer.wrap(createTestPayload(128 * 1024 * 1024)); + HttpRequestBodyStream payloadStreamResume = new HttpRequestBodyStream() { + @Override + public boolean sendRequestBody(ByteBuffer outBuffer) { + ByteBufferUtils.transferData(payloadResume, outBuffer); + return payloadResume.remaining() == 0; + } + + @Override + public boolean resetPosition() { + return true; + } + + @Override + public long getLength() { + return payloadResume.capacity(); + } + }; + + HttpHeader[] headersResume = { new HttpHeader("Host", ENDPOINT), + new HttpHeader("Content-Length", Integer.valueOf(payloadResume.capacity()).toString()), }; + + HttpRequest httpRequestResume = new HttpRequest("PUT", + uploadObjectPathInit("/put_object_test_async_pause_128MB"), headersResume, payloadStreamResume); + + CompletableFuture onFinishedFutureResume = new CompletableFuture<>(); + CompletableFuture onProgressFutureResume = new CompletableFuture<>(); + S3MetaRequestResponseHandler responseHandlerResume = createTestPutPauseResumeHandler(onFinishedFutureResume, + onProgressFutureResume); + S3MetaRequestOptions metaRequestOptionsResume = new S3MetaRequestOptions() + .withMetaRequestType(MetaRequestType.PUT_OBJECT) + .withHttpRequest(httpRequestResume) + .withResponseHandler(responseHandlerResume) + .withChecksumAlgorithm(ChecksumAlgorithm.CRC32) + .withResumeToken(new ResumeToken.PutResumeTokenBuilder() + .withPartSize(resumeToken.getPartSize()) + .withTotalNumParts(resumeToken.getTotalNumParts()) + .withNumPartsCompleted(resumeToken.getNumPartsCompleted()) + .withUploadId(resumeToken.getUploadId()) + .build()); + + try (S3MetaRequest metaRequest = client.makeMetaRequest(metaRequestOptionsResume)) { + Integer finish = onFinishedFutureResume.get(); + Assert.assertEquals(Integer.valueOf(0), finish); + } + } catch (InterruptedException | ExecutionException ex) { + Assert.fail(ex.getMessage()); + } + } + + @Test + public void testS3GetPauseAsync() { + skipIfAndroid(); + skipIfNetworkUnavailable(); + Assume.assumeTrue(hasAwsCredentials()); + + final long partSize = 1024 * 1024; + final long objectSize = 10 * 1024 * 1024; + + /* Enable backpressure with a 1MB initial window and never grow it: the download + * stalls once the window is exhausted, so the pause below deterministically lands + * mid-download instead of racing the 10MB transfer to completion. */ + S3ClientOptions clientOptions = new S3ClientOptions() + .withRegion(REGION) + .withPartSize(partSize) + .withReadBackpressureEnabled(true) + .withInitialReadWindowSize(partSize); + try (S3Client client = createS3Client(clientOptions)) { + CompletableFuture onFinishedFuture = new CompletableFuture<>(); + CompletableFuture onBodyFuture = new CompletableFuture<>(); + S3MetaRequestResponseHandler responseHandler = new S3MetaRequestResponseHandler() { + @Override + public int onResponseBody(ByteBuffer bodyBytesIn, long objectRangeStart, long objectRangeEnd) { + onBodyFuture.complete(null); + return 0; + } + + @Override + public void onFinished(S3FinishedResponseContext context) { + Log.log(Log.LogLevel.Info, Log.LogSubject.JavaCrtS3, + "Meta request finished with error code " + context.getErrorCode()); + if (context.getErrorCode() != 0) { + onFinishedFuture.completeExceptionally(new CrtRuntimeException(context.getErrorCode())); + return; + } + onFinishedFuture.complete(Integer.valueOf(context.getErrorCode())); + } + }; + + HttpHeader[] headers = { new HttpHeader("Host", ENDPOINT) }; + HttpRequest httpRequest = new HttpRequest("GET", PRE_EXIST_10MB_PATH, headers, null); + + S3MetaRequestOptions metaRequestOptions = new S3MetaRequestOptions() + .withMetaRequestType(MetaRequestType.GET_OBJECT) + .withHttpRequest(httpRequest) + .withResponseHandler(responseHandler); + + try (S3MetaRequest metaRequest = client.makeMetaRequest(metaRequestOptions)) { + onBodyFuture.get(); + + ResumeToken resumeToken = metaRequest.pauseAsync().get(); + Assert.assertNotNull(resumeToken); + Assert.assertEquals(MetaRequestType.GET_OBJECT, resumeToken.getType()); + Assert.assertEquals(partSize, resumeToken.getPartSize()); + Assert.assertEquals(objectSize, resumeToken.getObjectSize()); + Assert.assertEquals(0, resumeToken.getObjectRangeStart()); + Assert.assertEquals(objectSize - 1, resumeToken.getObjectRangeEnd()); + Assert.assertTrue(resumeToken.getContinuesDownloadedBytes() > 0); + Assert.assertTrue( + resumeToken.getContinuesDownloadedBytes() <= resumeToken.getTotalDownloadedBytes()); + Assert.assertNotNull(resumeToken.getEtag()); + Assert.assertFalse(resumeToken.getEtag().isEmpty()); + /* body-callback download: no local receive file */ + Assert.assertEquals(0, resumeToken.getFileLastModifiedEpochNs()); + /* upload-only getters must reject download tokens */ + assertThrows(IllegalArgumentException.class, () -> resumeToken.getUploadId()); + + Throwable thrown = assertThrows(Throwable.class, + () -> onFinishedFuture.get()); + + Assert.assertEquals("AWS_ERROR_S3_PAUSED", ((CrtRuntimeException) thrown.getCause()).errorName); + } + } catch (InterruptedException | ExecutionException ex) { + Assert.fail(ex.getMessage()); + } + } + private void testS3RoundTripWithChecksumHelper(ChecksumAlgorithm algo, ChecksumLocation location, boolean MPU, boolean provide_full_object_checksum) throws IOException {