diff --git a/src/generated/cloud/bigquery/v2/src/operation.rs b/src/generated/cloud/bigquery/v2/src/operation.rs index 37f522e736..c1d6b8c3e7 100644 --- a/src/generated/cloud/bigquery/v2/src/operation.rs +++ b/src/generated/cloud/bigquery/v2/src/operation.rs @@ -104,14 +104,14 @@ pub enum JobPollerError { #[derive(Debug)] pub struct JobPoller { policy: JobRetryPolicy, - builder: InsertJob, + builder: Box, } impl JobPoller { pub(crate) fn new(builder: InsertJob) -> Self { Self { policy: JobRetryPolicy::default(), - builder, + builder: Box::new(builder), } } @@ -128,38 +128,44 @@ impl JobPoller { } /// Polls the job until it is done, returning the final Job status. - pub async fn until_done(self) -> Result { - let mut attempts = 0_u32; - let mut builder = self.builder; - let backoff = self.policy.backoff; - let start_time = std::time::Instant::now(); - - loop { - // NOTE: the client library intercepts errors and retries internally - // according to the policies set on `builder`. - let job_result = builder.clone().poller().until_done().await?; - attempts += 1; - - if let Some(status) = &job_result.status - && let Some(err) = &status.error_result - { - if is_retryable_job_error(&err.reason) - && attempts < self.policy.job_level_attempt_limit + pub fn until_done( + self, + ) -> std::pin::Pin> + Send>> + { + Box::pin(async move { + let mut attempts = 0_u32; + let mut builder = self.builder; + let backoff = self.policy.backoff; + let start_time = std::time::Instant::now(); + + loop { + // NOTE: the client library intercepts errors and retries internally + // according to the policies set on `builder`. + let poller = Box::new(builder.clone().poller()); + let job_result = poller.until_done().await?; + attempts += 1; + + if let Some(status) = &job_result.status + && let Some(err) = &status.error_result { - let retry_job = prepare_job_for_retry(job_result); - builder = builder.set_job(retry_job); - - let retry_state = RetryState::new(true) - .set_start(start_time) - .set_attempt_count(attempts); - let delay = backoff.on_failure(&retry_state); - tokio::time::sleep(delay).await; - continue; + if is_retryable_job_error(&err.reason) + && attempts < self.policy.job_level_attempt_limit + { + let retry_job = prepare_job_for_retry(job_result); + *builder = builder.set_job(retry_job); + + let retry_state = RetryState::new(true) + .set_start(start_time) + .set_attempt_count(attempts); + let delay = backoff.on_failure(&retry_state); + tokio::time::sleep(delay).await; + continue; + } + return Err(JobPollerError::ErrorProto(err.clone())); } - return Err(JobPollerError::ErrorProto(err.clone())); + return Ok(job_result); } - return Ok(job_result); - } + }) } } diff --git a/src/lro/src/internal/discovery.rs b/src/lro/src/internal/discovery.rs index b14d865d82..9d3ede6c9e 100644 --- a/src/lro/src/internal/discovery.rs +++ b/src/lro/src/internal/discovery.rs @@ -172,7 +172,7 @@ where None } async fn until_done(self) -> Result { - crate::until_done(self).await + Box::pin(crate::until_done(self)).await } #[cfg(feature = "unstable-stream")]