google_cloud_bigquery_v2/
job_poller.rs1use 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#[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#[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#[derive(Debug)]
78pub(crate) struct JobRetryPolicy {
79 pub job_level_attempt_limit: u32,
81 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#[derive(Debug, thiserror::Error)]
96#[non_exhaustive]
97pub enum JobPollerError {
98 #[error(transparent)]
100 Rpc(#[from] GaxError),
101 #[error("BigQuery job failed ({}): {}", .error_result.reason, .error_result.message)]
103 JobFailed {
104 error_result: Box<crate::model::ErrorProto>,
106 errors: Vec<crate::model::ErrorProto>,
108 },
109}
110
111#[derive(Debug)]
113pub struct JobPoller {
114 policy: JobRetryPolicy,
115 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 pub fn with_attempt_limit(mut self, limit: u32) -> Self {
136 self.policy.job_level_attempt_limit = limit;
137 self
138 }
139
140 pub fn with_job_retry_backoff(mut self, backoff: ExponentialBackoff) -> Self {
142 self.policy.backoff = backoff;
143 self
144 }
145
146 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 {
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 }
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 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 assert!(retried.status.is_none());
358
359 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 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}