From 8f74be00db3af8931f19232f298e4729d157c0c1 Mon Sep 17 00:00:00 2001 From: Sitao Wang Date: Thu, 16 Jul 2026 22:59:30 -0400 Subject: [PATCH 1/5] Make gRPC client futures Send by replacing async closures with owned futures in the retry helper --- components/spider-client/src/grpc/job.rs | 54 ++++++++++--------- .../spider-client/src/grpc/resource_group.rs | 20 +++---- components/spider-utils/src/grpc/retry.rs | 54 ++++++++++--------- 3 files changed, 69 insertions(+), 59 deletions(-) diff --git a/components/spider-client/src/grpc/job.rs b/components/spider-client/src/grpc/job.rs index e45cc5537..e3fb25f8d 100644 --- a/components/spider-client/src/grpc/job.rs +++ b/components/spider-client/src/grpc/job.rs @@ -90,11 +90,11 @@ impl JobOrchestrationClient { compressed_serialized_task_graph, compressed_serialized_inputs, }; - let response = call_with_retry(self.retry_config, async || { - self.connection_pool - .get_client() - .register_job(request.clone()) - .await + let pool = self.connection_pool.clone(); + let response = call_with_retry(self.retry_config, move || { + let mut client = pool.get_client(); + let request = request.clone(); + async move { client.register_job(request).await } }) .await .map_err(|status| job_status_to_error(&status))? @@ -119,8 +119,11 @@ impl JobOrchestrationClient { let request = storage::JobIdRequest { job_id: job_id.get(), }; - let response = call_with_retry(self.retry_config, async || { - self.connection_pool.get_client().start_job(request).await + let pool = self.connection_pool.clone(); + let response = call_with_retry(self.retry_config, move || { + let mut client = pool.get_client(); + let request = request; + async move { client.start_job(request).await } }) .await .map_err(|status| job_status_to_error(&status))? @@ -145,8 +148,11 @@ impl JobOrchestrationClient { let request = storage::JobIdRequest { job_id: job_id.get(), }; - let response = call_with_retry(self.retry_config, async || { - self.connection_pool.get_client().cancel_job(request).await + let pool = self.connection_pool.clone(); + let response = call_with_retry(self.retry_config, move || { + let mut client = pool.get_client(); + let request = request; + async move { client.cancel_job(request).await } }) .await .map_err(|status| job_status_to_error(&status))? @@ -171,11 +177,11 @@ impl JobOrchestrationClient { let request = storage::JobIdRequest { job_id: job_id.get(), }; - let response = call_with_retry(self.retry_config, async || { - self.connection_pool - .get_client() - .get_job_state(request) - .await + let pool = self.connection_pool.clone(); + let response = call_with_retry(self.retry_config, move || { + let mut client = pool.get_client(); + let request = request; + async move { client.get_job_state(request).await } }) .await .map_err(|status| job_status_to_error(&status))? @@ -202,11 +208,11 @@ impl JobOrchestrationClient { let request = storage::JobIdRequest { job_id: job_id.get(), }; - let response = call_with_retry(self.retry_config, async || { - self.connection_pool - .get_client() - .get_job_outputs(request) - .await + let pool = self.connection_pool.clone(); + let response = call_with_retry(self.retry_config, move || { + let mut client = pool.get_client(); + let request = request; + async move { client.get_job_outputs(request).await } }) .await .map_err(|status| job_status_to_error(&status))? @@ -231,11 +237,11 @@ impl JobOrchestrationClient { let request = storage::JobIdRequest { job_id: job_id.get(), }; - let response = call_with_retry(self.retry_config, async || { - self.connection_pool - .get_client() - .get_job_error(request) - .await + let pool = self.connection_pool.clone(); + let response = call_with_retry(self.retry_config, move || { + let mut client = pool.get_client(); + let request = request; + async move { client.get_job_error(request).await } }) .await .map_err(|status| job_status_to_error(&status))? diff --git a/components/spider-client/src/grpc/resource_group.rs b/components/spider-client/src/grpc/resource_group.rs index 2aaad9b03..44d410aca 100644 --- a/components/spider-client/src/grpc/resource_group.rs +++ b/components/spider-client/src/grpc/resource_group.rs @@ -73,11 +73,11 @@ impl ResourceGroupManagementClient { external_resource_group_id, password, }; - let response = call_with_retry(self.retry_config, async || { - self.connection_pool - .get_client() - .add_resource_group(request.clone()) - .await + let pool = self.connection_pool.clone(); + let response = call_with_retry(self.retry_config, move || { + let mut client = pool.get_client(); + let request = request.clone(); + async move { client.add_resource_group(request).await } }) .await .map_err(|status| resource_group_status_to_error(&status))? @@ -107,11 +107,11 @@ impl ResourceGroupManagementClient { resource_group_id: resource_group_id.get(), password, }; - call_with_retry(self.retry_config, async || { - self.connection_pool - .get_client() - .verify_resource_group(request.clone()) - .await + let pool = self.connection_pool.clone(); + call_with_retry(self.retry_config, move || { + let mut client = pool.get_client(); + let request = request.clone(); + async move { client.verify_resource_group(request).await } }) .await .map_err(|status| resource_group_status_to_error(&status))?; diff --git a/components/spider-utils/src/grpc/retry.rs b/components/spider-utils/src/grpc/retry.rs index b2fc0d4db..450505627 100644 --- a/components/spider-utils/src/grpc/retry.rs +++ b/components/spider-utils/src/grpc/retry.rs @@ -1,6 +1,7 @@ //! An async retry helper for transient gRPC call failures. use std::error::Error; +use std::future::Future; use std::time::Duration; use rand::Rng; @@ -43,7 +44,8 @@ impl Default for RetryConfig { /// /// * `ResponseType` - The success value produced by `grpc_call`. /// * `ErrorType` - The error produced by `grpc_call`. -/// * `GrpcCall` - The async closure performing the gRPC call. +/// * `GrpcCall` - The closure performing the gRPC call. +/// * `FutureType` - The `Send` future returned by `grpc_call`. /// * `RetriableCheck` - Classifies an error as retriable or not. /// /// # Returns @@ -60,7 +62,8 @@ impl Default for RetryConfig { pub async fn execute_with_retry< ResponseType, ErrorType, - GrpcCall: AsyncFnMut() -> Result, + GrpcCall: FnMut() -> FutureType, + FutureType: Future> + Send, RetriableCheck: Fn(&ErrorType) -> bool, >( max_retries: usize, @@ -88,7 +91,8 @@ pub async fn execute_with_retry< /// # Type Parameters /// /// * `ResponseType` - The success value produced by `grpc_call`. -/// * `GrpcCall` - The async closure performing the gRPC round-trip. +/// * `GrpcCall` - The closure performing the gRPC round-trip. +/// * `FutureType` - The `Send` future returned by `grpc_call`. /// /// # Returns /// @@ -102,7 +106,8 @@ pub async fn execute_with_retry< /// budget is exhausted. pub async fn call_with_retry< ResponseType, - GrpcCall: AsyncFnMut() -> Result, + GrpcCall: FnMut() -> FutureType, + FutureType: Future> + Send, >( retry_config: RetryConfig, grpc_call: GrpcCall, @@ -164,7 +169,8 @@ fn backoff(retry: usize, max_backoff: Duration) -> Duration { #[cfg(test)] mod tests { - use std::cell::Cell; + use std::sync::atomic::AtomicUsize; + use std::sync::atomic::Ordering; use std::time::Duration; use tonic::Code; @@ -180,12 +186,12 @@ mod tests { #[tokio::test] async fn succeeds_on_first_attempt() { - let calls = Cell::new(0usize); + let calls = AtomicUsize::new(0); let result: Result = execute_with_retry( 3, TEST_MAX_BACKOFF, async || { - calls.set(calls.get() + 1); + calls.fetch_add(1, Ordering::Relaxed); Ok(42) }, |_error| true, @@ -193,18 +199,17 @@ mod tests { .await; assert_eq!(result, Ok(42)); - assert_eq!(calls.get(), 1); + assert_eq!(calls.load(Ordering::Relaxed), 1); } #[tokio::test] async fn succeeds_after_retriable_failures() { - let calls = Cell::new(0usize); + let calls = AtomicUsize::new(0); let result: Result = execute_with_retry( 5, TEST_MAX_BACKOFF, async || { - let attempt = calls.get(); - calls.set(attempt + 1); + let attempt = calls.fetch_add(1, Ordering::Relaxed); if attempt < 2 { Err(-1) } else { Ok(7) } }, |_error| true, @@ -212,17 +217,17 @@ mod tests { .await; assert_eq!(result, Ok(7)); - assert_eq!(calls.get(), 3); + assert_eq!(calls.load(Ordering::Relaxed), 3); } #[tokio::test] async fn non_retriable_error_returns_immediately() { - let calls = Cell::new(0usize); + let calls = AtomicUsize::new(0); let result: Result = execute_with_retry( 3, TEST_MAX_BACKOFF, async || { - calls.set(calls.get() + 1); + calls.fetch_add(1, Ordering::Relaxed); Err(99) }, |error| *error != 99, @@ -230,18 +235,18 @@ mod tests { .await; assert_eq!(result, Err(99)); - assert_eq!(calls.get(), 1); + assert_eq!(calls.load(Ordering::Relaxed), 1); } #[tokio::test] async fn retries_are_exhausted() { let max_retries = 4usize; - let calls = Cell::new(0usize); + let calls = AtomicUsize::new(0); let result: Result = execute_with_retry( max_retries, TEST_MAX_BACKOFF, async || { - calls.set(calls.get() + 1); + calls.fetch_add(1, Ordering::Relaxed); Err(-7) }, |_error| true, @@ -249,7 +254,7 @@ mod tests { .await; assert_eq!(result, Err(-7)); - assert_eq!(calls.get(), max_retries + 1); + assert_eq!(calls.load(Ordering::Relaxed), max_retries + 1); } #[tokio::test] @@ -258,10 +263,9 @@ mod tests { max_retries: 5, max_backoff: TEST_MAX_BACKOFF, }; - let calls = Cell::new(0usize); + let calls = AtomicUsize::new(0); let result: Result = call_with_retry(config, async || { - let attempt = calls.get(); - calls.set(attempt + 1); + let attempt = calls.fetch_add(1, Ordering::Relaxed); if attempt < 2 { Err(Status::unavailable("connection lost")) } else { @@ -274,7 +278,7 @@ mod tests { result.expect("call_with_retry should succeed after retriable failures"), 11 ); - assert_eq!(calls.get(), 3); + assert_eq!(calls.load(Ordering::Relaxed), 3); } #[tokio::test] @@ -283,9 +287,9 @@ mod tests { max_retries: 5, max_backoff: TEST_MAX_BACKOFF, }; - let calls = Cell::new(0usize); + let calls = AtomicUsize::new(0); let result: Result = call_with_retry(config, async || { - calls.set(calls.get() + 1); + calls.fetch_add(1, Ordering::Relaxed); Err(Status::not_found("missing")) }) .await; @@ -296,7 +300,7 @@ mod tests { .code(), Code::NotFound ); - assert_eq!(calls.get(), 1); + assert_eq!(calls.load(Ordering::Relaxed), 1); } #[test] From 5bf13bd831320ea67e96f3ca6cba90faa3a0cbe3 Mon Sep 17 00:00:00 2001 From: Sitao Wang Date: Thu, 16 Jul 2026 23:13:25 -0400 Subject: [PATCH 2/5] Address comment --- components/spider-client/src/grpc/job.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/spider-client/src/grpc/job.rs b/components/spider-client/src/grpc/job.rs index e3fb25f8d..3d60b4226 100644 --- a/components/spider-client/src/grpc/job.rs +++ b/components/spider-client/src/grpc/job.rs @@ -121,7 +121,7 @@ impl JobOrchestrationClient { }; let pool = self.connection_pool.clone(); let response = call_with_retry(self.retry_config, move || { - let mut client = pool.get_client(); + let mut client = pool.get_client().clone(); let request = request; async move { client.start_job(request).await } }) From 2102eacb24b30f5138e72f5d198b87ca7185acd0 Mon Sep 17 00:00:00 2001 From: Sitao Wang Date: Thu, 16 Jul 2026 23:27:39 -0400 Subject: [PATCH 3/5] Fix lint --- components/spider-client/src/grpc/job.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/spider-client/src/grpc/job.rs b/components/spider-client/src/grpc/job.rs index 3d60b4226..e3fb25f8d 100644 --- a/components/spider-client/src/grpc/job.rs +++ b/components/spider-client/src/grpc/job.rs @@ -121,7 +121,7 @@ impl JobOrchestrationClient { }; let pool = self.connection_pool.clone(); let response = call_with_retry(self.retry_config, move || { - let mut client = pool.get_client().clone(); + let mut client = pool.get_client(); let request = request; async move { client.start_job(request).await } }) From f6e42f120eaf446f2d802204e2d937e5f8fd9fa8 Mon Sep 17 00:00:00 2001 From: LinZhihao-723 Date: Sat, 18 Jul 2026 17:24:58 -0400 Subject: [PATCH 4/5] Move request construction inside the closure. --- components/spider-client/src/grpc/job.rs | 46 ++++++++----------- .../spider-client/src/grpc/resource_group.rs | 18 ++++---- 2 files changed, 28 insertions(+), 36 deletions(-) diff --git a/components/spider-client/src/grpc/job.rs b/components/spider-client/src/grpc/job.rs index e3fb25f8d..471f793b5 100644 --- a/components/spider-client/src/grpc/job.rs +++ b/components/spider-client/src/grpc/job.rs @@ -85,15 +85,14 @@ impl JobOrchestrationClient { .to_zstd_compressed_json() .map_err(|error| ClientError::Serialization(error.to_string()))?; let compressed_serialized_inputs = serialize_inputs(inputs)?; - let request = storage::RegisterJobRequest { - resource_group_id: resource_group_id.get(), - compressed_serialized_task_graph, - compressed_serialized_inputs, - }; let pool = self.connection_pool.clone(); let response = call_with_retry(self.retry_config, move || { let mut client = pool.get_client(); - let request = request.clone(); + let request = storage::RegisterJobRequest { + resource_group_id: resource_group_id.get(), + compressed_serialized_task_graph: compressed_serialized_task_graph.clone(), + compressed_serialized_inputs: compressed_serialized_inputs.clone(), + }; async move { client.register_job(request).await } }) .await @@ -116,13 +115,12 @@ impl JobOrchestrationClient { /// * Forwards [`JobOrchestrationServiceClient::start_job`]'s status on failure. /// * Forwards [`job_state_response_to_result`]'s return values on failure. pub async fn start_job(&self, job_id: JobId) -> Result { - let request = storage::JobIdRequest { - job_id: job_id.get(), - }; let pool = self.connection_pool.clone(); let response = call_with_retry(self.retry_config, move || { let mut client = pool.get_client(); - let request = request; + let request = storage::JobIdRequest { + job_id: job_id.get(), + }; async move { client.start_job(request).await } }) .await @@ -145,13 +143,12 @@ impl JobOrchestrationClient { /// * Forwards [`JobOrchestrationServiceClient::cancel_job`]'s status on failure. /// * Forwards [`job_state_response_to_result`]'s return values on failure. pub async fn cancel_job(&self, job_id: JobId) -> Result { - let request = storage::JobIdRequest { - job_id: job_id.get(), - }; let pool = self.connection_pool.clone(); let response = call_with_retry(self.retry_config, move || { let mut client = pool.get_client(); - let request = request; + let request = storage::JobIdRequest { + job_id: job_id.get(), + }; async move { client.cancel_job(request).await } }) .await @@ -174,13 +171,12 @@ impl JobOrchestrationClient { /// * Forwards [`JobOrchestrationServiceClient::get_job_state`]'s status on failure. /// * Forwards [`job_state_response_to_result`]'s return values on failure. pub async fn get_job_state(&self, job_id: JobId) -> Result { - let request = storage::JobIdRequest { - job_id: job_id.get(), - }; let pool = self.connection_pool.clone(); let response = call_with_retry(self.retry_config, move || { let mut client = pool.get_client(); - let request = request; + let request = storage::JobIdRequest { + job_id: job_id.get(), + }; async move { client.get_job_state(request).await } }) .await @@ -205,13 +201,12 @@ impl JobOrchestrationClient { /// [`ClientError::Deserialization`]. /// * Forwards [`JobOrchestrationServiceClient::get_job_outputs`]'s status on failure. pub async fn get_job_outputs(&self, job_id: JobId) -> Result, ClientError> { - let request = storage::JobIdRequest { - job_id: job_id.get(), - }; let pool = self.connection_pool.clone(); let response = call_with_retry(self.retry_config, move || { let mut client = pool.get_client(); - let request = request; + let request = storage::JobIdRequest { + job_id: job_id.get(), + }; async move { client.get_job_outputs(request).await } }) .await @@ -234,13 +229,12 @@ impl JobOrchestrationClient { /// /// * Forwards [`JobOrchestrationServiceClient::get_job_error`]'s status on failure. pub async fn get_job_error(&self, job_id: JobId) -> Result { - let request = storage::JobIdRequest { - job_id: job_id.get(), - }; let pool = self.connection_pool.clone(); let response = call_with_retry(self.retry_config, move || { let mut client = pool.get_client(); - let request = request; + let request = storage::JobIdRequest { + job_id: job_id.get(), + }; async move { client.get_job_error(request).await } }) .await diff --git a/components/spider-client/src/grpc/resource_group.rs b/components/spider-client/src/grpc/resource_group.rs index 44d410aca..7159e676e 100644 --- a/components/spider-client/src/grpc/resource_group.rs +++ b/components/spider-client/src/grpc/resource_group.rs @@ -69,14 +69,13 @@ impl ResourceGroupManagementClient { external_resource_group_id: String, password: Vec, ) -> Result { - let request = storage::AddResourceGroupRequest { - external_resource_group_id, - password, - }; let pool = self.connection_pool.clone(); let response = call_with_retry(self.retry_config, move || { let mut client = pool.get_client(); - let request = request.clone(); + let request = storage::AddResourceGroupRequest { + external_resource_group_id: external_resource_group_id.clone(), + password: password.clone(), + }; async move { client.add_resource_group(request).await } }) .await @@ -103,14 +102,13 @@ impl ResourceGroupManagementClient { resource_group_id: ResourceGroupId, password: Vec, ) -> Result<(), ClientError> { - let request = storage::VerifyResourceGroupRequest { - resource_group_id: resource_group_id.get(), - password, - }; let pool = self.connection_pool.clone(); call_with_retry(self.retry_config, move || { let mut client = pool.get_client(); - let request = request.clone(); + let request = storage::VerifyResourceGroupRequest { + resource_group_id: resource_group_id.get(), + password: password.clone(), + }; async move { client.verify_resource_group(request).await } }) .await From 6bdefadb6cfdd9cb87886609043921c4499f8e68 Mon Sep 17 00:00:00 2001 From: LinZhihao-723 Date: Sat, 18 Jul 2026 17:42:12 -0400 Subject: [PATCH 5/5] Add compile-time Send+Sync assertion. --- components/spider-client/src/client.rs | 29 ++++++++++++++++++++++++++ 1 file changed, 29 insertions(+) diff --git a/components/spider-client/src/client.rs b/components/spider-client/src/client.rs index 3879db990..2a66ccb0c 100644 --- a/components/spider-client/src/client.rs +++ b/components/spider-client/src/client.rs @@ -290,3 +290,32 @@ impl SpiderClientBuilder { } const DEFAULT_POOL_SIZE: NonZeroUsize = NonZeroUsize::new(8).unwrap(); + +/// Compile-time assertion that the public client handles are `Send + Sync`. +const _: () = { + const fn assert_send_sync() {} + assert_send_sync::(); + assert_send_sync::(); +}; + +/// Compile-time assertion that every public [`SpiderClient`] async method returns a `Send` future. +/// +/// This function is never called; it exists only to type-check the `Send` bound on each returned +/// future. The arguments are taken as parameters (never constructed) so no real values are needed. +#[expect(dead_code, reason = "compile-time-only `Send` assertion; never called")] +fn assert_client_futures_send( + client: &SpiderClient, + resource_group_id: ResourceGroupId, + job_id: JobId, + task_graph: &TaskGraph, +) { + const fn assert_send(_: &FutureType) {} + assert_send(&client.submit_job(resource_group_id, task_graph, Vec::new())); + assert_send(&client.start_job(job_id)); + assert_send(&client.cancel_job(job_id)); + assert_send(&client.get_job_state(job_id)); + assert_send(&client.get_job_outputs(job_id)); + assert_send(&client.get_job_error(job_id)); + assert_send(&client.add_resource_group(String::new(), Vec::new())); + assert_send(&client.verify_resource_group(resource_group_id, Vec::new())); +}