diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index f537de192..9dffb6523 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -824,10 +824,7 @@ impl ReplicationPool { let cancel_token = CancellationToken::new(); resyncer.register_cancel_token(&opts, cancel_token.clone()).await; tokio::spawn(async move { - resyncer - .clone() - .resync_bucket(cancel_token, storage, false, opts.clone()) - .await; + Box::pin(resyncer.clone().resync_bucket(cancel_token, storage, false, opts.clone())).await; resyncer.clear_cancel_token(&opts).await; }); @@ -914,7 +911,7 @@ impl ReplicationPool { }; tokio::spawn(async move { resync.register_cancel_token(&opts, ctx.clone()).await; - resync.clone().resync_bucket(ctx, storage, true, opts.clone()).await; + Box::pin(resync.clone().resync_bucket(ctx, storage, true, opts.clone())).await; resync.clear_cancel_token(&opts).await; }); } diff --git a/crates/ecstore/src/tier/warm_backend_gcs.rs b/crates/ecstore/src/tier/warm_backend_gcs.rs index 87ab63139..f2d19b898 100644 --- a/crates/ecstore/src/tier/warm_backend_gcs.rs +++ b/crates/ecstore/src/tier/warm_backend_gcs.rs @@ -107,11 +107,12 @@ impl WarmBackend for WarmBackendGCS { ReaderImpl::Body(content_body) => content_body.to_vec(), ReaderImpl::ObjectBody(mut content_body) => content_body.read_all().await?, }; - let Ok(res) = self - .client - .write_object(&self.bucket, &self.get_dest(object), Bytes::from(d)) - .send_buffered() - .await + let Ok(res) = Box::pin( + self.client + .write_object(&self.bucket, &self.get_dest(object), Bytes::from(d)) + .send_buffered(), + ) + .await else { return Err(std::io::Error::other("write_object error")); }; diff --git a/rustfs/src/admin/router.rs b/rustfs/src/admin/router.rs index a6493f503..754a2e809 100644 --- a/rustfs/src/admin/router.rs +++ b/rustfs/src/admin/router.rs @@ -2186,7 +2186,7 @@ async fn handle_misc_extension_request(req: &mut S3Request, route: &MiscEx MiscExtRoute::ObjectLambda { bucket, object } => { let get_req = build_object_lambda_get_request(req, bucket, object)?; let usecase = DefaultObjectUsecase::from_global(); - let get_resp = usecase.execute_get_object(get_req).await?; + let get_resp = Box::pin(usecase.execute_get_object(get_req)).await?; invoke_object_lambda_target(req, bucket, object, get_resp).await } MiscExtRoute::ListenNotification { bucket } => { @@ -2386,7 +2386,7 @@ where return handle_replication_extension_request(&mut req, &ext_req).await; } if let Some(ext_req) = parse_misc_extension_request(&req.method, &req.uri) { - return handle_misc_extension_request(&mut req, &ext_req).await; + return Box::pin(handle_misc_extension_request(&mut req, &ext_req)).await; } // Console requests should be handled by console router first (including OPTIONS) diff --git a/rustfs/src/app/lifecycle_transition_api_test.rs b/rustfs/src/app/lifecycle_transition_api_test.rs index ee6524ef9..7c8ec80b0 100644 --- a/rustfs/src/app/lifecycle_transition_api_test.rs +++ b/rustfs/src/app/lifecycle_transition_api_test.rs @@ -373,8 +373,7 @@ async fn put_and_copy_object_transition_immediately_via_usecases() { .build() .unwrap(); - usecase - .execute_put_object(&fs, build_request(put_input, Method::PUT)) + Box::pin(usecase.execute_put_object(&fs, build_request(put_input, Method::PUT))) .await .expect("Failed to put object through usecase"); @@ -410,8 +409,7 @@ async fn put_and_copy_object_transition_immediately_via_usecases() { .build() .unwrap(); - usecase - .execute_copy_object(build_request(copy_input, Method::PUT)) + Box::pin(usecase.execute_copy_object(build_request(copy_input, Method::PUT))) .await .expect("Failed to copy object through usecase"); @@ -468,8 +466,7 @@ async fn complete_multipart_upload_transitions_immediately_via_usecase() { .build() .unwrap(); - usecase - .execute_complete_multipart_upload(build_request(complete_input, Method::POST)) + Box::pin(usecase.execute_complete_multipart_upload(build_request(complete_input, Method::POST))) .await .expect("Failed to complete multipart upload through usecase"); @@ -526,8 +523,7 @@ async fn delete_transitioned_object_removes_remote_tier_copy_via_usecase() { ); insert_header(&mut req.headers, SUFFIX_FORCE_DELETE, "true"); - usecase - .execute_delete_object(req) + Box::pin(usecase.execute_delete_object(req)) .await .expect("Failed to delete object through usecase"); diff --git a/rustfs/src/app/multipart_usecase.rs b/rustfs/src/app/multipart_usecase.rs index a4219825b..9f9f9f041 100644 --- a/rustfs/src/app/multipart_usecase.rs +++ b/rustfs/src/app/multipart_usecase.rs @@ -1258,7 +1258,9 @@ mod tests { .unwrap(); let req = build_request(input, Method::POST); - let err = make_usecase().execute_complete_multipart_upload(req).await.unwrap_err(); + let err = Box::pin(make_usecase().execute_complete_multipart_upload(req)) + .await + .unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidPart); } @@ -1285,7 +1287,9 @@ mod tests { .unwrap(); let req = build_request(input, Method::POST); - let err = make_usecase().execute_complete_multipart_upload(req).await.unwrap_err(); + let err = Box::pin(make_usecase().execute_complete_multipart_upload(req)) + .await + .unwrap_err(); assert_ne!(err.code(), &S3ErrorCode::InvalidPartOrder); } @@ -1312,7 +1316,9 @@ mod tests { .unwrap(); let req = build_request(input, Method::POST); - let err = make_usecase().execute_complete_multipart_upload(req).await.unwrap_err(); + let err = Box::pin(make_usecase().execute_complete_multipart_upload(req)) + .await + .unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidPartOrder); } @@ -1362,7 +1368,9 @@ mod tests { let mut req = build_request(input, Method::POST); req.headers.insert(header_name, HeaderValue::from_str(header_value).unwrap()); - let err = make_usecase().execute_complete_multipart_upload(req).await.unwrap_err(); + let err = Box::pin(make_usecase().execute_complete_multipart_upload(req)) + .await + .unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest, "header {header_name} should be rejected"); } } @@ -1471,7 +1479,7 @@ mod tests { .unwrap(); let req = build_request(input, Method::PUT); - let err = make_usecase().execute_upload_part_copy(req).await.unwrap_err(); + let err = Box::pin(make_usecase().execute_upload_part_copy(req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InternalError); } diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index db45284ad..2cf304e7e 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -1417,7 +1417,7 @@ impl DefaultObjectUsecase { return Err(s3_error!(InvalidStorageClass)); } if is_put_object_extract_requested(&req.headers) { - return self.execute_put_object_extract(req).await; + return Box::pin(self.execute_put_object_extract(req)).await; } let input = std::mem::take(&mut req.input); @@ -4416,7 +4416,7 @@ mod tests { let usecase = DefaultObjectUsecase::without_context(); let fs = FS::new(); - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); + let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::NotImplemented); } @@ -4435,7 +4435,7 @@ mod tests { let usecase = DefaultObjectUsecase::without_context(); let fs = FS::new(); - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); + let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::NotImplemented); } @@ -4454,7 +4454,7 @@ mod tests { let usecase = DefaultObjectUsecase::without_context(); let fs = FS::new(); - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); + let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidStorageClass); } @@ -4474,7 +4474,7 @@ mod tests { let usecase = DefaultObjectUsecase::without_context(); let fs = FS::new(); - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); + let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::NotImplemented); } @@ -4494,7 +4494,7 @@ mod tests { let usecase = DefaultObjectUsecase::without_context(); let fs = FS::new(); - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); + let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::NotImplemented); } @@ -4511,7 +4511,7 @@ mod tests { let usecase = DefaultObjectUsecase::without_context(); let fs = FS::new(); - let err = usecase.execute_put_object(&fs, req).await.unwrap_err(); + let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidStorageClass); } @@ -4527,7 +4527,7 @@ mod tests { let req = build_request(input, Method::GET); let usecase = DefaultObjectUsecase::without_context(); - let err = usecase.execute_get_object(req).await.unwrap_err(); + let err = Box::pin(usecase.execute_get_object(req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } @@ -4547,7 +4547,7 @@ mod tests { let req = build_request(input, Method::PUT); let usecase = DefaultObjectUsecase::without_context(); - let err = usecase.execute_copy_object(req).await.unwrap_err(); + let err = Box::pin(usecase.execute_copy_object(req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); } @@ -4562,7 +4562,7 @@ mod tests { let req = build_request(input, Method::DELETE); let usecase = DefaultObjectUsecase::without_context(); - let err = usecase.execute_delete_object(req).await.unwrap_err(); + let err = Box::pin(usecase.execute_delete_object(req)).await.unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidArgument); } diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index f6dc2b60d..072bedd92 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -228,14 +228,14 @@ impl S3 for FS { req: S3Request, ) -> S3Result> { let usecase = DefaultMultipartUsecase::from_global(); - usecase.execute_complete_multipart_upload(req).await + Box::pin(usecase.execute_complete_multipart_upload(req)).await } /// Copy an object from one location to another #[instrument(level = "debug", skip(self, req))] async fn copy_object(&self, req: S3Request) -> S3Result> { let usecase = DefaultObjectUsecase::from_global(); - usecase.execute_copy_object(req).await + Box::pin(usecase.execute_copy_object(req)).await } #[instrument( @@ -345,7 +345,7 @@ impl S3 for FS { #[instrument(level = "debug", skip(self, req))] async fn delete_object(&self, req: S3Request) -> S3Result> { let usecase = DefaultObjectUsecase::from_global(); - usecase.execute_delete_object(req).await + Box::pin(usecase.execute_delete_object(req)).await } #[instrument(level = "debug", skip(self))] @@ -611,7 +611,7 @@ impl S3 for FS { )] async fn get_object(&self, req: S3Request) -> S3Result> { let usecase = DefaultObjectUsecase::from_global(); - usecase.execute_get_object(req).await + Box::pin(usecase.execute_get_object(req)).await } async fn get_object_acl(&self, req: S3Request) -> S3Result> { @@ -1102,7 +1102,7 @@ impl S3 for FS { #[instrument(level = "debug", skip(self, req))] async fn put_object(&self, req: S3Request) -> S3Result> { let usecase = DefaultObjectUsecase::from_global(); - usecase.execute_put_object(self, req).await + Box::pin(usecase.execute_put_object(self, req)).await } async fn put_object_acl(&self, req: S3Request) -> S3Result> { @@ -1447,6 +1447,6 @@ impl S3 for FS { async fn upload_part_copy(&self, req: S3Request) -> S3Result> { record_s3_op(S3Operation::UploadPartCopy, &req.input.bucket); let usecase = DefaultMultipartUsecase::from_global(); - usecase.execute_upload_part_copy(req).await + Box::pin(usecase.execute_upload_part_copy(req)).await } }