Skip to main content

google_cloud_bigquery/query/
retry_policy.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//! Defines a retry policy for BigQuery.
16
17use crate::error::QueryError;
18use google_cloud_bigquery_v2::model::ErrorProto;
19use google_cloud_gax::backoff_policy::BackoffPolicy;
20use google_cloud_gax::backoff_policy::BackoffPolicyArg;
21use google_cloud_gax::error::Error as GaxError;
22use google_cloud_gax::error::rpc::{Code, StatusDetails};
23use google_cloud_gax::exponential_backoff::ExponentialBackoffBuilder;
24use google_cloud_gax::retry_policy::{RetryPolicy, RetryPolicyExt};
25use google_cloud_gax::retry_result::RetryResult;
26use google_cloud_gax::retry_state::RetryState;
27use std::sync::Arc;
28use std::time::Duration;
29
30/// Follows the RPC retry strategy recommended by the BigQuery guides on
31/// [error handling].
32///
33/// ```
34/// # async fn sample() -> anyhow::Result<()> {
35/// # use google_cloud_bigquery::client::BigQuery;
36/// # use google_cloud_bigquery::query::retry_policy::RetryableErrors;
37/// # use google_cloud_gax::retry_policy::RetryPolicyExt;
38/// let policy = RetryableErrors.with_attempt_limit(6);
39/// let client = BigQuery::builder()
40///     .with_retry_policy(policy)
41///     .build()
42///     .await?;
43/// # Ok(())
44/// # }
45/// ```
46///
47/// This policy must be decorated to limit the duration of the retry loop or
48/// the number of attempts.
49///
50/// [error handling]: https://cloud.google.com/bigquery/docs/error-messages
51#[derive(Clone, Debug)]
52pub struct RetryableErrors;
53
54impl RetryPolicy for RetryableErrors {
55    fn on_error(&self, state: &RetryState, error: GaxError) -> RetryResult {
56        if error.is_transient_and_before_rpc() {
57            return RetryResult::Continue(error);
58        }
59        if !state.idempotent {
60            return RetryResult::Permanent(error);
61        }
62        if error.is_io() || error.is_timeout() {
63            return RetryResult::Continue(error);
64        }
65        if error.is_transport() && error.http_status_code().is_none() {
66            return RetryResult::Continue(error);
67        }
68        if let Some(429 | 500 | 502 | 503 | 504) = error.http_status_code() {
69            return RetryResult::Continue(error);
70        }
71        if let Some(status) = error.status() {
72            return match status.code {
73                Code::Aborted
74                | Code::DeadlineExceeded
75                | Code::Internal
76                | Code::ResourceExhausted
77                | Code::Unavailable
78                | Code::Unknown => RetryResult::Continue(error),
79                _ => RetryResult::Permanent(error),
80            };
81        }
82        RetryResult::Permanent(error)
83    }
84}
85
86pub(crate) fn default_retry_policy() -> Arc<dyn RetryPolicy> {
87    Arc::new(RetryableErrors.with_attempt_limit(6))
88}
89
90pub(crate) fn default_backoff_policy() -> Arc<dyn BackoffPolicy> {
91    Arc::new(
92        ExponentialBackoffBuilder::default()
93            .with_initial_delay(Duration::from_secs(1))
94            .with_maximum_delay(Duration::from_secs(32))
95            .with_scaling(2.0)
96            .build()
97            .expect("valid backoff configuration"),
98    )
99}
100
101/// The result of evaluating a BigQuery job error against a [`JobRetryPolicy`].
102#[derive(Debug)]
103pub(crate) enum JobRetryResult {
104    Continue(Duration, #[allow(dead_code)] QueryError),
105    Exhausted(QueryError),
106    Permanent(QueryError),
107}
108
109impl JobRetryResult {
110    #[allow(dead_code)]
111    pub(crate) fn is_continue(&self) -> bool {
112        matches!(self, Self::Continue(_, _))
113    }
114
115    #[allow(dead_code)]
116    pub(crate) fn is_exhausted(&self) -> bool {
117        matches!(self, Self::Exhausted(_))
118    }
119
120    #[allow(dead_code)]
121    pub(crate) fn is_permanent(&self) -> bool {
122        matches!(self, Self::Permanent(_))
123    }
124}
125
126/// A policy trait for handling BigQuery job execution retries.
127///
128/// Note: This trait is kept crate-internal for now. It is a copy/adaptation of GAX's
129/// `RetryPolicy` trait specifically tailored for BigQuery job-level re-issuance.
130/// In the future, we plan to refine this trait and discuss how to make it public
131/// so customers can provide their own custom job retry policies.
132pub(crate) trait JobRetryPolicy<S = RetryState>: Send + Sync + std::fmt::Debug {
133    fn on_error(&self, state: &S, error: QueryError) -> JobRetryResult;
134}
135
136#[derive(Clone, Debug)]
137pub(crate) struct RetryableJobErrors {
138    attempt_limit: u32,
139    backoff: Arc<dyn BackoffPolicy>,
140}
141
142impl Default for RetryableJobErrors {
143    fn default() -> Self {
144        Self {
145            attempt_limit: 3,
146            backoff: default_backoff_policy(),
147        }
148    }
149}
150
151impl RetryableJobErrors {
152    #[allow(dead_code)]
153    pub fn with_attempt_limit(mut self, attempt_limit: u32) -> Self {
154        self.attempt_limit = attempt_limit;
155        self
156    }
157
158    #[allow(dead_code)]
159    pub fn with_backoff_policy<V: Into<BackoffPolicyArg>>(mut self, v: V) -> Self {
160        self.backoff = v.into().into();
161        self
162    }
163}
164
165impl JobRetryPolicy for RetryableJobErrors {
166    fn on_error(&self, state: &RetryState, error: QueryError) -> JobRetryResult {
167        if !is_query_error_retryable(&error) {
168            return JobRetryResult::Permanent(error);
169        }
170        if state.attempt_count >= self.attempt_limit {
171            return JobRetryResult::Exhausted(error);
172        }
173        let delay = self.backoff.on_failure(state);
174        JobRetryResult::Continue(delay, error)
175    }
176}
177
178pub(crate) fn default_job_retry_policy() -> Arc<dyn JobRetryPolicy> {
179    Arc::new(RetryableJobErrors::default())
180}
181
182pub(crate) fn is_query_error_retryable(err: &QueryError) -> bool {
183    match err {
184        QueryError::JobFailed { errors } => is_retryable_errors(errors),
185        QueryError::Rpc { source } => is_rpc_error_retryable(source),
186        _ => false,
187    }
188}
189
190pub(crate) fn is_rpc_error_retryable(error: &GaxError) -> bool {
191    if let Some(status) = error.status() {
192        for detail in &status.details {
193            if let StatusDetails::ErrorInfo(info) = detail
194                && is_retryable_error_reason(&info.reason)
195            {
196                return true;
197            }
198        }
199    }
200    false
201}
202
203pub(crate) fn is_retryable_errors(errors: &[ErrorProto]) -> bool {
204    !errors.is_empty() && errors.iter().all(|e| is_retryable_error_reason(&e.reason))
205}
206
207pub(crate) fn is_retryable_error_reason(reason: &str) -> bool {
208    matches!(
209        reason,
210        "backendError"
211            | "jobBackendError"
212            | "rateLimitExceeded"
213            | "jobRateLimitExceeded"
214            | "internalError"
215            | "jobInternalError"
216    )
217}
218
219/// Returns true if `error` reports a conflict with a resource that already
220/// exists, such as `409 Already Exists: Job my-project:US.job_1234567890`.
221pub(crate) fn is_duplicate_job_error(error: &QueryError) -> bool {
222    let QueryError::Rpc { source } = error else {
223        return false;
224    };
225    source.http_status_code() == Some(409)
226        || source
227            .status()
228            .is_some_and(|s| s.code == Code::AlreadyExists)
229}
230
231#[cfg(test)]
232mod tests {
233    use super::*;
234    use crate::query::tests::create_test_backoff_policy;
235    use google_cloud_bigquery_v2::model::ErrorProto;
236    use google_cloud_gax::error::CredentialsError;
237    use google_cloud_gax::error::rpc::{Code, Status};
238    use google_cloud_gax::retry_state::RetryState;
239    use google_cloud_rpc::model::ErrorInfo;
240    use http::HeaderMap;
241    use test_case::test_case;
242
243    #[test_case("backendError", true)]
244    #[test_case("jobBackendError", true)]
245    #[test_case("rateLimitExceeded", true)]
246    #[test_case("jobRateLimitExceeded", true)]
247    #[test_case("internalError", true)]
248    #[test_case("jobInternalError", true)]
249    #[test_case("invalidQuery", false)]
250    #[test_case("notFound", false)]
251    fn test_is_retryable_error_reason(reason: &str, expected: bool) {
252        assert_eq!(is_retryable_error_reason(reason), expected);
253    }
254
255    #[test]
256    fn test_is_retryable_errors() {
257        assert!(!is_retryable_errors(&[]));
258
259        let non_retryable = vec![ErrorProto::new().set_reason("invalidQuery")];
260        assert!(!is_retryable_errors(&non_retryable));
261
262        let mixed = vec![
263            ErrorProto::new().set_reason("invalidQuery"),
264            ErrorProto::new().set_reason("backendError"),
265        ];
266        assert!(!is_retryable_errors(&mixed));
267
268        let retryable = vec![
269            ErrorProto::new().set_reason("rateLimitExceeded"),
270            ErrorProto::new().set_reason("backendError"),
271        ];
272        assert!(is_retryable_errors(&retryable));
273    }
274
275    #[test]
276    fn test_retryable_errors_on_error() {
277        let p = RetryableErrors;
278        let idempotent = RetryState::new(true);
279        let non_idempotent = RetryState::new(false);
280
281        let retryable_codes = [
282            Code::Aborted,
283            Code::DeadlineExceeded,
284            Code::Internal,
285            Code::ResourceExhausted,
286            Code::Unavailable,
287            Code::Unknown,
288        ];
289        for code in retryable_codes {
290            let err = || GaxError::service(Status::default().set_code(code));
291            assert!(p.on_error(&idempotent, err()).is_continue(), "{code:?}");
292            assert!(
293                p.on_error(&non_idempotent, err()).is_permanent(),
294                "{code:?}"
295            );
296        }
297
298        let permanent_codes = [
299            Code::NotFound,
300            Code::PermissionDenied,
301            Code::InvalidArgument,
302        ];
303        for code in permanent_codes {
304            let err = || GaxError::service(Status::default().set_code(code));
305            assert!(p.on_error(&idempotent, err()).is_permanent(), "{code:?}");
306            assert!(
307                p.on_error(&non_idempotent, err()).is_permanent(),
308                "{code:?}"
309            );
310        }
311
312        let retryable_http = [429, 500, 502, 503, 504];
313        for code in retryable_http {
314            let err = || GaxError::http(code, HeaderMap::new(), bytes::Bytes::new());
315            assert!(p.on_error(&idempotent, err()).is_continue(), "HTTP {code}");
316            assert!(
317                p.on_error(&non_idempotent, err()).is_permanent(),
318                "HTTP {code}"
319            );
320        }
321
322        let permanent_http = [400, 404, 408, 409, 501];
323        for code in permanent_http {
324            let err = || GaxError::http(code, HeaderMap::new(), bytes::Bytes::new());
325            assert!(p.on_error(&idempotent, err()).is_permanent(), "HTTP {code}");
326            assert!(
327                p.on_error(&non_idempotent, err()).is_permanent(),
328                "HTTP {code}"
329            );
330        }
331
332        let io = || GaxError::io("connection reset");
333        assert!(p.on_error(&idempotent, io()).is_continue());
334        assert!(p.on_error(&non_idempotent, io()).is_permanent());
335
336        let timeout = || GaxError::timeout("deadline");
337        assert!(p.on_error(&idempotent, timeout()).is_continue());
338        assert!(p.on_error(&non_idempotent, timeout()).is_permanent());
339
340        // Failures before the request is sent cannot have reached the service,
341        // so they are retried either way.
342        let before_rpc =
343            || GaxError::authentication(CredentialsError::from_msg(true, "token refresh failed"));
344        assert!(p.on_error(&idempotent, before_rpc()).is_continue());
345        assert!(p.on_error(&non_idempotent, before_rpc()).is_continue());
346    }
347
348    #[test]
349    fn test_is_duplicate_job_error() {
350        let rpc = |source| QueryError::Rpc { source };
351        let status = |code| GaxError::service(Status::default().set_code(code));
352        let http = |code| GaxError::http(code, HeaderMap::new(), bytes::Bytes::new());
353
354        assert!(is_duplicate_job_error(&rpc(http(409))));
355        assert!(is_duplicate_job_error(&rpc(status(Code::AlreadyExists))));
356
357        assert!(!is_duplicate_job_error(&rpc(status(Code::Aborted))));
358        assert!(!is_duplicate_job_error(&rpc(http(500))));
359        assert!(!is_duplicate_job_error(&QueryError::JobFailed {
360            errors: vec![ErrorProto::new().set_reason("duplicate")],
361        }));
362    }
363
364    #[test]
365    fn test_duplicate_job_error_is_not_retryable() {
366        const BQ_DUPLICATE_PAYLOAD: &[u8] = br#"{
367  "error": {
368    "code": 409,
369    "message": "Already Exists: Job my-project:US.job_1234567890",
370    "errors": [
371      {
372        "message": "Already Exists: Job my-project:US.job_1234567890",
373        "domain": "global",
374        "reason": "duplicate"
375      }
376    ],
377    "status": "ALREADY_EXISTS"
378  }
379}"#;
380        let status = Status::try_from(&bytes::Bytes::from_static(BQ_DUPLICATE_PAYLOAD))
381            .expect("should deserialize BigQuery REST error");
382        let err = QueryError::Rpc {
383            source: GaxError::service(status),
384        };
385        assert!(is_duplicate_job_error(&err), "{err:?}");
386        assert!(!is_query_error_retryable(&err), "{err:?}");
387    }
388
389    #[test]
390    fn test_is_rpc_error_retryable() {
391        use google_cloud_rpc::model::ErrorInfo;
392
393        // BigQuery REST API JSON error with backendError
394        const BQ_REST_PAYLOAD: &[u8] = br#"{
395  "error": {
396    "code": 400,
397    "message": "The job encountered an error during execution. Retrying the job may solve the problem.",
398    "errors": [
399      {
400        "message": "The job encountered an error during execution. Retrying the job may solve the problem.",
401        "domain": "global",
402        "reason": "backendError"
403      }
404    ],
405    "status": "INVALID_ARGUMENT"
406  }
407}"#;
408        let status = Status::try_from(&bytes::Bytes::from_static(BQ_REST_PAYLOAD))
409            .expect("should deserialize BigQuery REST error");
410        let err = GaxError::service(status);
411        assert!(is_rpc_error_retryable(&err));
412
413        // ErrorInfo detail with retryable reason
414        let status = Status::default()
415            .set_code(Code::InvalidArgument)
416            .set_message("Error occurred")
417            .set_details(vec![StatusDetails::ErrorInfo(
418                ErrorInfo::new().set_reason("backendError"),
419            )]);
420        let err = GaxError::service(status);
421        assert!(is_rpc_error_retryable(&err));
422
423        // ErrorInfo detail with non-retryable reason
424        let status = Status::default()
425            .set_code(Code::InvalidArgument)
426            .set_message("Error occurred")
427            .set_details(vec![StatusDetails::ErrorInfo(
428                ErrorInfo::new().set_reason("invalidQuery"),
429            )]);
430        let err = GaxError::service(status);
431        assert!(!is_rpc_error_retryable(&err));
432
433        // Status without ErrorInfo details
434        let status = Status::default()
435            .set_code(Code::InvalidArgument)
436            .set_message("Syntax error: Unexpected identifier");
437        let err = GaxError::service(status);
438        assert!(!is_rpc_error_retryable(&err));
439    }
440
441    #[test]
442    fn test_job_retryable_errors() {
443        let policy = RetryableJobErrors::default();
444        let state = RetryState::default();
445
446        let retryable_err = QueryError::JobFailed {
447            errors: vec![ErrorProto::new().set_reason("backendError")],
448        };
449        assert!(policy.on_error(&state, retryable_err).is_continue());
450
451        let permanent_err = QueryError::JobFailed {
452            errors: vec![ErrorProto::new().set_reason("invalidQuery")],
453        };
454        assert!(policy.on_error(&state, permanent_err).is_permanent());
455
456        let rpc_retryable_err = QueryError::Rpc {
457            source: GaxError::service(
458                Status::default()
459                    .set_code(Code::InvalidArgument)
460                    .set_details(vec![StatusDetails::ErrorInfo(
461                        ErrorInfo::new().set_reason("backendError"),
462                    )]),
463            ),
464        };
465        assert!(policy.on_error(&state, rpc_retryable_err).is_continue());
466
467        let rpc_permanent_err = QueryError::Rpc {
468            source: GaxError::service(
469                Status::default()
470                    .set_code(Code::InvalidArgument)
471                    .set_message("Syntax error: Unexpected identifier"),
472            ),
473        };
474        assert!(policy.on_error(&state, rpc_permanent_err).is_permanent());
475    }
476
477    #[test]
478    fn test_job_attempt_limit() {
479        let policy = default_job_retry_policy(); // default attempt_limit is 3
480        let retryable_err = || QueryError::JobFailed {
481            errors: vec![ErrorProto::new().set_reason("backendError")],
482        };
483
484        let mut state = RetryState::default();
485        state.attempt_count = 1;
486        assert!(policy.on_error(&state, retryable_err()).is_continue());
487
488        state.attempt_count = 2;
489        assert!(policy.on_error(&state, retryable_err()).is_continue());
490
491        state.attempt_count = 3;
492        assert!(policy.on_error(&state, retryable_err()).is_exhausted());
493    }
494
495    #[test]
496    fn test_job_backoff_policy() {
497        let mut backoff = create_test_backoff_policy();
498        backoff
499            .expect_on_failure()
500            .return_const(Duration::from_secs(5));
501
502        let policy = RetryableJobErrors::default().with_backoff_policy(backoff);
503        let retryable_err = QueryError::JobFailed {
504            errors: vec![ErrorProto::new().set_reason("backendError")],
505        };
506
507        let state = RetryState::default();
508        if let JobRetryResult::Continue(delay, _) = policy.on_error(&state, retryable_err) {
509            assert_eq!(delay, Duration::from_secs(5));
510        } else {
511            panic!("expected Continue with 5s delay");
512        }
513    }
514
515    #[test]
516    fn test_default_retry_policy_is_bounded() {
517        let policy = default_retry_policy();
518        let err = || GaxError::service(Status::default().set_code(Code::Unavailable));
519
520        for attempt in 1u32..6u32 {
521            let state = RetryState::new(true).set_attempt_count(attempt);
522            assert!(
523                policy.on_error(&state, err()).is_continue(),
524                "attempt {attempt} should continue"
525            );
526            assert_eq!(policy.remaining_time(&state), None);
527        }
528
529        let state = RetryState::new(true).set_attempt_count(6u32);
530        assert!(
531            policy.on_error(&state, err()).is_exhausted(),
532            "attempt 6 should be exhausted"
533        );
534    }
535}