Skip to main content

google_cloud_bigquery_v2/
job_poller.rs

1// Copyright 2026 Google LLC
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     https://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Helpers to await the completion of an [InsertJob] request
16
17use crate::builder::job_service::InsertJob;
18use crate::model::Job;
19use google_cloud_gax::backoff_policy::BackoffPolicy;
20use google_cloud_gax::error::Error as GaxError;
21use google_cloud_gax::error::rpc::{Code, Status};
22use google_cloud_gax::exponential_backoff::ExponentialBackoff;
23use google_cloud_gax::retry_state::RetryState;
24use google_cloud_lro::Poller;
25
26impl google_cloud_lro::internal::DiscoveryOperation for Job {
27    fn name(&self) -> Option<&String> {
28        self.job_reference.as_ref().map(|r| &r.job_id)
29    }
30
31    fn done(&self) -> bool {
32        self.status
33            .as_ref()
34            .map(|s| s.state == "DONE")
35            .unwrap_or(false)
36    }
37
38    fn error(&self) -> Option<Status> {
39        self.status.as_ref().and_then(|s| {
40            s.error_result.as_ref().map(|e| {
41                Status::default()
42                    .set_code(Code::Unknown)
43                    .set_message(e.message.clone())
44            })
45        })
46    }
47}
48
49/// Determines if a BigQuery job failure reason is transient and eligible for
50/// job-level retry.
51///
52/// Returns `true` for retryable reasons (`jobBackendError`,
53/// `jobInternalError`, `jobRateLimitExceeded`, `tableUnavailable`) per
54/// BigQuery error handling specification.
55#[allow(dead_code)]
56pub(crate) fn is_retryable_job_error(reason: &str) -> bool {
57    matches!(
58        reason,
59        "jobBackendError" | "jobInternalError" | "jobRateLimitExceeded" | "tableUnavailable"
60    )
61}
62
63/// Prepares a `Job` instance for retry by assigning a new synthetic job ID
64/// and clearing existing execution status.
65///
66/// To preserve idempotency and avoid job execution collisions, each job-level
67/// retry must use a unique job ID while retaining original reference details
68/// (project ID, location) and configuration settings.
69#[allow(dead_code)]
70pub(crate) fn prepare_job_for_retry(mut job: Job) -> Job {
71    job.job_reference.get_or_insert_default().job_id = uuid::Uuid::new_v4().to_string();
72    job.status = None;
73    job
74}
75
76/// Configuration policy for BigQuery job-level retries.
77#[derive(Debug)]
78pub(crate) struct JobRetryPolicy {
79    /// Maximum number of general job-level attempts for retryable job errors.
80    pub job_level_attempt_limit: u32,
81    /// Backoff strategy between retry attempts.
82    pub backoff: ExponentialBackoff,
83}
84
85impl Default for JobRetryPolicy {
86    fn default() -> Self {
87        Self {
88            job_level_attempt_limit: 3,
89            backoff: ExponentialBackoff::default(),
90        }
91    }
92}
93
94/// Errors returned by the JobPoller.
95#[derive(Debug, thiserror::Error)]
96#[non_exhaustive]
97pub enum JobPollerError {
98    /// An error occurred during the RPC or LRO polling.
99    #[error(transparent)]
100    Rpc(#[from] GaxError),
101    /// The job completed, but the BigQuery service reported an internal error.
102    #[error("BigQuery job failed ({}): {}", .error_result.reason, .error_result.message)]
103    JobFailed {
104        /// Final error result of the job.
105        error_result: Box<crate::model::ErrorProto>,
106        /// Errors and warnings encountered during the running of the job.
107        errors: Vec<crate::model::ErrorProto>,
108    },
109}
110
111/// A poller that monitors the status of an inserted BigQuery job and handles retries.
112#[derive(Debug)]
113pub struct JobPoller {
114    policy: JobRetryPolicy,
115    // Because the builder holds a `Job` which is >7kB, we ought to store it
116    // on the heap.
117    //
118    // This is important across `await` points where stack variables
119    // are captured in the async fn's state machine. This bloats the size of
120    // the returned `Future`, which can potentially overflow the stack.
121    //
122    // See #6391.
123    builder: Box<InsertJob>,
124}
125
126impl JobPoller {
127    pub(crate) fn new(builder: InsertJob) -> Self {
128        Self {
129            policy: JobRetryPolicy::default(),
130            builder: Box::new(builder),
131        }
132    }
133
134    /// Sets the maximum number of job-level attempts.
135    pub fn with_attempt_limit(mut self, limit: u32) -> Self {
136        self.policy.job_level_attempt_limit = limit;
137        self
138    }
139
140    /// Sets the exponential backoff policy for job-level retries.
141    pub fn with_job_retry_backoff(mut self, backoff: ExponentialBackoff) -> Self {
142        self.policy.backoff = backoff;
143        self
144    }
145
146    /// Polls the job until it is done, returning the final Job status.
147    pub async fn until_done(self) -> Result<Job, JobPollerError> {
148        let mut attempts = 0_u32;
149        let mut builder = self.builder;
150        let backoff = self.policy.backoff;
151        let start_time = std::time::Instant::now();
152
153        loop {
154            // NOTE: the client library intercepts errors and retries internally
155            // according to the policies set on `builder`.
156            {
157                let poller = (*builder).clone().poller();
158                let job = poller.until_done().await?;
159                let Some(status) = &job.status else {
160                    return Ok(job);
161                };
162                let Some(err) = &status.error_result else {
163                    return Ok(job);
164                };
165
166                attempts += 1;
167                if !is_retryable_job_error(&err.reason)
168                    || attempts >= self.policy.job_level_attempt_limit
169                {
170                    return Err(JobPollerError::JobFailed {
171                        error_result: Box::new(err.clone()),
172                        errors: status.errors.clone(),
173                    });
174                }
175
176                let job = prepare_job_for_retry(job);
177                *builder = (*builder).set_job(job);
178
179                // We use a block so that `job` (~7kB) is not allocated on the
180                // stack across the sleep `await` point.
181            }
182
183            let retry_state = RetryState::new(true)
184                .set_start(start_time)
185                .set_attempt_count(attempts);
186            let delay = backoff.on_failure(&retry_state);
187            tokio::time::sleep(delay).await;
188        }
189    }
190}
191
192impl InsertJob {
193    /// Returns a `JobPoller`, which can retry on [job-level errors].
194    ///
195    /// If the job fails with an internal error, the `JobPoller` will retry the
196    /// `InsertJob` operation. Note that the client library will supply a
197    /// synthetic job ID for any retries.
198    ///
199    /// ```no_run
200    /// # async fn example(builder: google_cloud_bigquery_v2::builder::job_service::InsertJob) -> Result<(), Box<dyn std::error::Error>> {
201    /// let job = builder.into_job_poller().until_done().await?;
202    /// # Ok(())
203    /// # }
204    /// ```
205    ///
206    /// [job-level errors]: https://docs.cloud.google.com/bigquery/docs/error-messages#errortable
207    pub fn into_job_poller(self) -> JobPoller {
208        JobPoller::new(self)
209    }
210}
211
212#[cfg(test)]
213mod tests {
214    use super::*;
215    use crate::model::{
216        ErrorProto, JobConfiguration, JobConfigurationQuery, JobReference, JobStatus,
217    };
218    use google_cloud_lro::internal::DiscoveryOperation;
219
220    #[test]
221    fn name_none() {
222        let job = Job::default();
223        assert_eq!(job.name(), None);
224    }
225
226    #[test]
227    fn name_some() {
228        let job = Job::new().set_job_reference(JobReference::new().set_job_id("test-id"));
229        assert_eq!(job.name().map(|s| s.as_str()), Some("test-id"));
230    }
231
232    #[test]
233    fn done_none() {
234        let job = Job::default();
235        assert!(!job.done());
236    }
237
238    #[test]
239    fn done_false() {
240        let job = Job::new().set_status(JobStatus::new().set_state("RUNNING"));
241        assert!(!job.done());
242    }
243
244    #[test]
245    fn done_true() {
246        let job = Job::new().set_status(JobStatus::new().set_state("DONE"));
247        assert!(job.done());
248    }
249
250    #[test]
251    fn error_none() {
252        let job = Job::default();
253        assert!(job.error().is_none());
254
255        let job_no_error = Job::new().set_status(JobStatus::new().set_state("DONE"));
256        assert!(job_no_error.error().is_none());
257    }
258
259    #[test]
260    fn error_some() {
261        let job = Job::new()
262            .set_status(JobStatus::new().set_error_result(ErrorProto::new().set_message("failed")));
263        let err = job.error().expect("should have error");
264        assert_eq!(err.code, Code::Unknown);
265        assert_eq!(err.message, "failed");
266    }
267
268    #[test]
269    fn retryable_job_errors() {
270        assert!(is_retryable_job_error("jobBackendError"));
271        assert!(is_retryable_job_error("jobInternalError"));
272        assert!(is_retryable_job_error("jobRateLimitExceeded"));
273        assert!(is_retryable_job_error("tableUnavailable"));
274
275        assert!(!is_retryable_job_error("invalidQuery"));
276        assert!(!is_retryable_job_error("accessDenied"));
277        assert!(!is_retryable_job_error("notFound"));
278        assert!(!is_retryable_job_error("backendError"));
279        assert!(!is_retryable_job_error(""));
280    }
281
282    #[test]
283    fn job_retry_policy_defaults() {
284        let policy = JobRetryPolicy::default();
285        assert_eq!(policy.job_level_attempt_limit, 3);
286    }
287
288    #[test]
289    fn prepare_job_for_retry_generates_new_id_and_resets_status() {
290        let original_job = Job::new()
291            .set_job_reference(
292                JobReference::new()
293                    .set_project_id("test-project")
294                    .set_job_id("original-job-id")
295                    .set_location("US"),
296            )
297            .set_status(
298                JobStatus::new().set_state("DONE").set_error_result(
299                    ErrorProto::new()
300                        .set_reason("jobBackendError")
301                        .set_message("backend failed"),
302                ),
303            );
304
305        let retried_job = prepare_job_for_retry(original_job);
306
307        assert!(retried_job.status.is_none());
308
309        let ref_data = retried_job
310            .job_reference
311            .expect("should have job reference");
312        assert_eq!(ref_data.project_id, "test-project");
313        assert_eq!(ref_data.location.as_deref(), Some("US"));
314        assert_ne!(ref_data.job_id, "original-job-id");
315        assert!(uuid::Uuid::parse_str(&ref_data.job_id).is_ok());
316    }
317
318    #[test]
319    fn prepare_job_for_retry_handles_none_job_reference() {
320        let original_job = Job::new().set_status(JobStatus::new().set_state("DONE"));
321
322        let retried_job = prepare_job_for_retry(original_job);
323        assert!(retried_job.status.is_none());
324
325        let ref_data = retried_job
326            .job_reference
327            .expect("should create job reference when missing");
328        assert!(uuid::Uuid::parse_str(&ref_data.job_id).is_ok());
329    }
330
331    #[test]
332    fn prepare_job_for_retry_preserves_job_configuration_and_metadata() {
333        let original_job = Job::new()
334            .set_job_reference(
335                JobReference::new()
336                    .set_project_id("my-project")
337                    .set_job_id("initial-id")
338                    .set_location("EU"),
339            )
340            .set_configuration(
341                JobConfiguration::new()
342                    .set_query(JobConfigurationQuery::new().set_query("SELECT 42"))
343                    .set_labels([("env".to_string(), "test".to_string())]),
344            )
345            .set_user_email("user@example.com")
346            .set_status(
347                JobStatus::new().set_state("DONE").set_error_result(
348                    ErrorProto::new()
349                        .set_reason("jobInternalError")
350                        .set_message("internal error"),
351                ),
352            );
353
354        let retried = prepare_job_for_retry(original_job);
355
356        // Status must be reset to None for retry submission
357        assert!(retried.status.is_none());
358
359        // Configuration and user_email must be preserved
360        assert_eq!(
361            retried
362                .configuration
363                .as_ref()
364                .and_then(|c| c.query.as_ref())
365                .map(|q| q.query.as_str()),
366            Some("SELECT 42")
367        );
368        assert_eq!(
369            retried
370                .configuration
371                .as_ref()
372                .and_then(|c| c.labels.get("env").map(|s| s.as_str())),
373            Some("test")
374        );
375        assert_eq!(retried.user_email.as_str(), "user@example.com");
376
377        // JobReference metadata preserved, but job_id replaced with a new valid UUID
378        let ref_data = retried.job_reference.expect("must have reference");
379        assert_eq!(ref_data.project_id, "my-project");
380        assert_eq!(ref_data.location.as_deref(), Some("EU"));
381        assert_ne!(ref_data.job_id, "initial-id");
382        assert!(uuid::Uuid::parse_str(&ref_data.job_id).is_ok());
383    }
384
385    #[test]
386    fn custom_retry_policy_builder() {
387        let mut policy = JobRetryPolicy::default();
388        assert_eq!(policy.job_level_attempt_limit, 3);
389
390        policy.job_level_attempt_limit = 5;
391        assert_eq!(policy.job_level_attempt_limit, 5);
392    }
393
394    #[test]
395    fn job_poller_error_job_failed_display() {
396        let err_proto = ErrorProto::new()
397            .set_reason("invalidQuery")
398            .set_message("syntax error");
399        let sub_error = ErrorProto::new()
400            .set_reason("invalid")
401            .set_message("detailed error");
402
403        let poller_err = JobPollerError::JobFailed {
404            error_result: Box::new(err_proto),
405            errors: vec![sub_error],
406        };
407
408        let JobPollerError::JobFailed {
409            error_result,
410            errors,
411        } = &poller_err
412        else {
413            panic!("expected JobPollerError::JobFailed");
414        };
415        assert_eq!(error_result.reason, "invalidQuery");
416        assert_eq!(errors.len(), 1);
417        assert_eq!(errors[0].reason, "invalid");
418        assert_eq!(
419            poller_err.to_string(),
420            "BigQuery job failed (invalidQuery): syntax error"
421        );
422    }
423}