google_cloud_bigquery/query/
retry_policy.rs1use 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#[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#[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
126pub(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
219pub(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 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 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 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 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 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(); 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}