1use crate::error::QueryError;
16use crate::generated::{CompleteQueryMetadata, QueryMetadata};
17use crate::query::execution::RetryContext;
18use crate::query::retry_policy::JobRetryResult;
19use crate::query::{Result, RowIterator};
20use google_cloud_bigquery_v2::builder::job_service::GetJob;
21use google_cloud_bigquery_v2::client::JobService;
22use google_cloud_bigquery_v2::model::{
23 GetQueryResultsRequest, GetQueryResultsResponse, Job, JobReference, QueryResponse,
24};
25use google_cloud_gax::exponential_backoff::ExponentialBackoffBuilder;
26use google_cloud_gax::polling_backoff_policy::PollingBackoffPolicy;
27use google_cloud_gax::polling_state::PollingState;
28use std::collections::VecDeque;
29use std::sync::Arc;
30
31#[derive(Clone, Debug)]
55pub struct Query {
56 pub(crate) job_service: Arc<JobService>,
57 pub(crate) completed: bool,
58 pub(crate) metadata: QueryMetadata,
59 pub(crate) cached_rows: Option<VecDeque<wkt::Struct>>,
60 pub(crate) page_size: Option<u32>,
61 pub(crate) retry_context: Option<RetryContext>,
62}
63
64impl Query {
65 pub(crate) fn from_job(
66 job_service: Arc<JobService>,
67 initial_job: Job,
68 retry_context: Option<&RetryContext>,
69 page_size: Option<u32>,
70 ) -> Self {
71 let completed = initial_job
72 .status
73 .as_ref()
74 .map(|s| s.state == "DONE")
75 .unwrap_or(false);
76 Self {
77 job_service,
78 completed,
79 cached_rows: None,
80 metadata: build_query_metadata_from_job(initial_job),
81 retry_context: retry_context.filter(|_| !completed).cloned(),
82 page_size,
83 }
84 }
85
86 pub(crate) fn from_query_response(
87 job_service: Arc<JobService>,
88 mut query_response: QueryResponse,
89 retry_context: Option<&RetryContext>,
90 page_size: Option<u32>,
91 ) -> Self {
92 let completed = query_response.job_complete.unwrap_or(false);
93 let cached_rows = VecDeque::from(std::mem::take(&mut query_response.rows));
94 let metadata = QueryMetadata::from(query_response);
95 Self {
96 job_service,
97 completed,
98 cached_rows: Some(cached_rows),
99 metadata,
100 retry_context: retry_context.filter(|_| !completed).cloned(),
101 page_size,
102 }
103 }
104
105 pub fn metadata(&self) -> &QueryMetadata {
116 &self.metadata
117 }
118
119 pub fn get_job(&self) -> Option<GetJob> {
156 build_get_job(&self.job_service, self.metadata.job_reference.as_ref()?)
157 }
158
159 pub(crate) fn is_dry_run(&self) -> bool {
160 self.metadata
161 .configuration
162 .as_ref()
163 .and_then(|c| c.dry_run)
164 .unwrap_or(false)
165 }
166
167 pub async fn until_done(mut self) -> Result<CompleteQuery> {
195 if self.is_dry_run() {
196 return Err(QueryError::DryRun);
197 }
198
199 loop {
200 let Query {
201 job_service,
202 completed,
203 metadata,
204 cached_rows,
205 page_size,
206 retry_context,
207 } = self;
208
209 if completed && let Some(cached_rows) = cached_rows {
210 return Ok(CompleteQuery::from_query_metadata(
211 job_service,
212 metadata,
213 cached_rows,
214 page_size,
215 ));
216 }
217
218 let job_ref = metadata
219 .job_reference
220 .clone()
221 .expect("query job should have job reference at this point");
222
223 let backoff_policy = Arc::new(
224 ExponentialBackoffBuilder::default()
225 .with_initial_delay(std::time::Duration::from_secs(1))
226 .build()
227 .expect("valid backoff configuration"),
228 );
229
230 match poll_query_results(&job_service, &job_ref, backoff_policy).await {
231 Ok(res) => {
232 return Ok(CompleteQuery::from_get_query_results_response(
233 job_service,
234 &job_ref,
235 res,
236 metadata,
237 page_size,
238 ));
239 }
240 Err(err) => {
241 let Some(retry_ctx) = retry_context else {
242 return Err(err);
243 };
244 match retry_ctx.on_error(err) {
245 JobRetryResult::Continue(delay, _) => {
246 self = retry_ctx.reissue(delay).await?;
247 }
248 JobRetryResult::Permanent(e) | JobRetryResult::Exhausted(e) => {
249 return Err(e);
250 }
251 }
252 }
253 }
254 }
255 }
256}
257
258#[derive(Clone, Debug)]
281pub struct CompleteQuery {
282 pub(crate) job_service: Arc<JobService>,
283 pub(crate) job_ref: Option<JobReference>,
284 pub(crate) cached_rows: VecDeque<wkt::Struct>,
285 pub(crate) page_token: Option<String>,
286 pub(crate) metadata: CompleteQueryMetadata,
287 pub(crate) page_size: Option<u32>,
288}
289
290impl CompleteQuery {
291 pub(crate) fn from_get_query_results_response(
292 job_service: Arc<JobService>,
293 job_ref: &JobReference,
294 mut res: GetQueryResultsResponse,
295 initial_metadata: QueryMetadata,
296 page_size: Option<u32>,
297 ) -> Self {
298 let cached_rows = VecDeque::from(std::mem::take(&mut res.rows));
299 let metadata =
300 build_complete_query_metadata_from_get_query_results(initial_metadata, res, job_ref);
301 Self::from_complete_metadata(
302 job_service,
303 Some(job_ref.clone()),
304 metadata,
305 cached_rows,
306 page_size,
307 )
308 }
309
310 pub(crate) fn from_query_metadata(
311 job_service: Arc<JobService>,
312 metadata: QueryMetadata,
313 cached_rows: VecDeque<wkt::Struct>,
314 page_size: Option<u32>,
315 ) -> Self {
316 let job_ref = metadata.job_reference.clone();
317 let metadata = CompleteQueryMetadata::from(metadata);
318 Self::from_complete_metadata(job_service, job_ref, metadata, cached_rows, page_size)
319 }
320
321 pub(crate) fn from_complete_metadata(
322 job_service: Arc<JobService>,
323 job_ref: Option<JobReference>,
324 metadata: CompleteQueryMetadata,
325 cached_rows: VecDeque<wkt::Struct>,
326 page_size: Option<u32>,
327 ) -> Self {
328 let page_token = if metadata.page_token.is_empty() {
329 None
330 } else {
331 Some(metadata.page_token.clone())
332 };
333 Self {
334 job_service,
335 job_ref,
336 cached_rows,
337 page_token,
338 metadata,
339 page_size,
340 }
341 }
342
343 pub fn read(self) -> RowIterator {
364 RowIterator::new(self)
365 }
366
367 pub fn metadata(&self) -> &CompleteQueryMetadata {
390 &self.metadata
391 }
392
393 pub fn get_job(&self) -> Option<GetJob> {
430 build_get_job(&self.job_service, self.job_ref.as_ref()?)
431 }
432}
433
434pub(crate) fn build_get_job(job_service: &JobService, job_ref: &JobReference) -> Option<GetJob> {
438 if job_ref.job_id.is_empty() {
439 return None;
440 }
441
442 let req = job_service
443 .get_job()
444 .set_job_id(job_ref.job_id.clone())
445 .set_project_id(job_ref.project_id.clone());
446
447 let req = job_ref
448 .location
449 .clone()
450 .into_iter()
451 .fold(req, |req, location| req.set_location(location));
452
453 Some(req)
454}
455
456pub(crate) async fn poll_query_results(
458 job_service: &JobService,
459 job_ref: &JobReference,
460 backoff_policy: Arc<dyn PollingBackoffPolicy>,
461) -> Result<GetQueryResultsResponse> {
462 let mut state = PollingState::default();
463
464 loop {
465 let mut req = GetQueryResultsRequest::new()
466 .set_max_results(0u32)
467 .set_project_id(job_ref.project_id.clone())
468 .set_job_id(job_ref.job_id.clone());
469 if let Some(location) = job_ref.location.clone() {
470 req = req.set_location(location);
471 }
472
473 let res = job_service
474 .get_query_results()
475 .with_request(req)
476 .send()
477 .await?;
478
479 if !res.errors.is_empty() {
480 return Err(QueryError::JobFailed { errors: res.errors });
481 }
482
483 let completed = res.job_complete.unwrap_or(false);
484 if completed {
485 return Ok(res);
486 }
487
488 let delay = backoff_policy.wait_period(&state);
489 tokio::time::sleep(delay).await;
490 state.attempt_count += 1;
492 }
493}
494
495fn build_query_metadata_from_job(resp: Job) -> QueryMetadata {
501 let mut metadata = QueryMetadata::from(resp);
502
503 let query_stats = metadata.statistics.as_ref().and_then(|s| s.query.as_ref());
504
505 metadata.job_complete = metadata.status.as_ref().map(|s| s.state == "DONE");
506 metadata.errors = metadata
507 .status
508 .as_ref()
509 .map(|s| s.errors.clone())
510 .unwrap_or_default();
511 metadata.schema = query_stats.and_then(|q| q.schema.clone());
512 metadata.total_bytes_processed =
513 query_stats
514 .and_then(|q| q.total_bytes_processed)
515 .or_else(|| {
516 metadata
517 .statistics
518 .as_ref()
519 .and_then(|s| s.total_bytes_processed)
520 });
521 metadata.total_bytes_billed = query_stats.and_then(|q| q.total_bytes_billed);
522 metadata.total_slot_ms = query_stats
523 .and_then(|q| q.total_slot_ms)
524 .or_else(|| metadata.statistics.as_ref().and_then(|s| s.total_slot_ms));
525 metadata.cache_hit = query_stats.and_then(|q| q.cache_hit);
526 metadata.num_dml_affected_rows = query_stats.and_then(|q| q.num_dml_affected_rows);
527 metadata.dml_stats = query_stats.and_then(|q| q.dml_stats.clone());
528 metadata.statement_type = query_stats
529 .map(|q| q.statement_type.clone())
530 .unwrap_or_default();
531 metadata.session_info = metadata
532 .statistics
533 .as_ref()
534 .and_then(|s| s.session_info.clone());
535
536 metadata.creation_time = metadata
537 .statistics
538 .as_ref()
539 .and_then(|s| (s.creation_time > 0).then_some(s.creation_time));
540 metadata.start_time = metadata
541 .statistics
542 .as_ref()
543 .and_then(|s| (s.start_time > 0).then_some(s.start_time));
544 metadata.end_time = metadata
545 .statistics
546 .as_ref()
547 .and_then(|s| (s.end_time > 0).then_some(s.end_time));
548
549 if metadata.location.is_empty() {
550 metadata.location = metadata
551 .job_reference
552 .as_ref()
553 .and_then(|r| r.location.clone())
554 .unwrap_or_default();
555 }
556
557 metadata
558}
559
560fn build_complete_query_metadata_from_get_query_results(
563 initial_metadata: QueryMetadata,
564 res: GetQueryResultsResponse,
565 job_ref: &JobReference,
566) -> CompleteQueryMetadata {
567 let mut metadata = CompleteQueryMetadata::from(initial_metadata);
568 if res.schema.is_some() {
569 metadata.schema = res.schema;
570 }
571 if res.total_rows.is_some() {
572 metadata.total_rows = res.total_rows;
573 }
574 if !res.page_token.is_empty() {
575 metadata.page_token = res.page_token;
576 }
577 if res.total_bytes_processed.is_some() {
578 metadata.total_bytes_processed = res.total_bytes_processed;
579 }
580 if res.job_complete.is_some() {
581 metadata.job_complete = res.job_complete;
582 }
583 if !res.errors.is_empty() {
584 metadata.errors = res.errors;
585 }
586 if res.cache_hit.is_some() {
587 metadata.cache_hit = res.cache_hit;
588 }
589 if res.num_dml_affected_rows.is_some() {
590 metadata.num_dml_affected_rows = res.num_dml_affected_rows;
591 }
592 if !res.etag.is_empty() {
593 metadata.etag = res.etag;
594 }
595 if res.job_reference.is_some() {
596 metadata.job_reference = res.job_reference;
597 }
598 if metadata.location.is_empty()
599 && let Some(loc) = job_ref.location.clone()
600 {
601 metadata.location = loc;
602 }
603
604 metadata
605}
606
607#[cfg(test)]
608mod tests {
609 use super::*;
610 use crate::query::builder::{QUERY_REQUEST_ID_PREFIX, Query as QueryBuilder};
611 use crate::query::retry_policy::RetryableJobErrors;
612 use crate::query::tests::{MockJobService, create_job_service, create_test_backoff_policy};
613 use google_cloud_bigquery_v2::model::{
614 ErrorProto, GetQueryResultsResponse, Job, JobConfiguration, JobReference, QueryResponse,
615 TableFieldSchema, TableSchema,
616 };
617 use google_cloud_gax::error::Error as GaxError;
618 use google_cloud_gax::error::rpc::{Code, Status};
619 use google_cloud_gax::response::Response;
620 use std::time::Duration;
621
622 type TestResult = anyhow::Result<()>;
623
624 impl CompleteQuery {
625 pub(crate) fn from_query_response(
626 job_service: Arc<JobService>,
627 mut query_res: QueryResponse,
628 page_size: Option<u32>,
629 ) -> Self {
630 let cached_rows = std::mem::take(&mut query_res.rows).into();
631 let metadata = QueryMetadata::from(query_res);
632 Self::from_query_metadata(job_service, metadata, cached_rows, page_size)
633 }
634 }
635
636 #[tokio::test]
637 async fn test_query_until_done_already_completed() -> TestResult {
638 let job_service = create_job_service(MockJobService::new());
639 let job_ref = JobReference::new()
640 .set_project_id("some_project")
641 .set_job_id("some_job_id");
642 let query_res = QueryResponse::new()
643 .set_job_complete(true)
644 .set_job_reference(job_ref.clone())
645 .set_schema(TableSchema::new())
646 .set_page_token("some_page_token")
647 .set_rows([wkt::Struct::new()])
648 .set_cache_hit(true);
649
650 let query = Query::from_query_response(job_service, query_res, None, None);
651
652 let completed = query.until_done().await?;
653 assert_eq!(completed.job_ref.as_ref().unwrap().job_id, "some_job_id");
654 assert_eq!(completed.page_token, Some("some_page_token".to_string()));
655 assert_eq!(completed.cached_rows.len(), 1);
656
657 let metadata = completed.metadata();
658 assert_eq!(metadata.cache_hit, Some(true));
659 assert_eq!(metadata.job_complete, Some(true));
660 assert_eq!(metadata.page_token, "some_page_token".to_string());
661
662 Ok(())
663 }
664
665 #[tokio::test]
666 async fn test_query_until_done_preserves_page_size() -> TestResult {
667 let job_service = create_job_service(MockJobService::new());
668 let job_ref = JobReference::new()
669 .set_project_id("some_project")
670 .set_job_id("some_job_id");
671 let query_res = QueryResponse::new()
672 .set_job_complete(true)
673 .set_job_reference(job_ref.clone())
674 .set_schema(TableSchema::new());
675
676 let query = Query::from_query_response(job_service, query_res, None, Some(42));
677
678 let completed = query.until_done().await?;
679 assert_eq!(completed.page_size, Some(42));
680
681 Ok(())
682 }
683
684 #[tokio::test]
685 async fn test_query_until_done_polls_success() -> TestResult {
686 let mut mock = MockJobService::new();
687 mock.expect_get_query_results()
688 .returning(|req, _| {
689 assert_eq!(req.project_id, "some_project");
690 assert_eq!(req.job_id, "some_job_id");
691 assert_eq!(req.max_results, Some(0));
692 assert_eq!(req.location, "us-central1");
693 let res = GetQueryResultsResponse::new()
694 .set_job_complete(true)
695 .set_job_reference(JobReference::new().set_job_id(req.job_id))
696 .set_schema(TableSchema::new())
697 .set_page_token("some_page_token")
698 .set_rows(vec![wkt::Struct::new(), wkt::Struct::new()])
699 .set_cache_hit(false);
700 Ok(Response::from(res))
701 })
702 .times(1);
703 let job_service = create_job_service(mock);
704 let job_ref = JobReference::new()
705 .set_project_id("some_project")
706 .set_job_id("some_job_id")
707 .set_location("us-central1");
708 let query_res = QueryResponse::new()
709 .set_job_complete(false)
710 .set_job_reference(job_ref);
711
712 let query = Query::from_query_response(job_service, query_res, None, None);
713
714 let completed = query.until_done().await?;
715 assert_eq!(completed.job_ref.as_ref().unwrap().job_id, "some_job_id");
716 assert_eq!(completed.page_token, Some("some_page_token".to_string()));
717 assert_eq!(completed.cached_rows.len(), 2);
718
719 let metadata = completed.metadata();
720 assert_eq!(metadata.cache_hit, Some(false));
721 assert_eq!(metadata.job_complete, Some(true));
722 assert_eq!(metadata.page_token, "some_page_token");
723
724 Ok(())
725 }
726
727 #[tokio::test(start_paused = true)]
728 async fn test_poll_query_results_loops_until_complete() -> TestResult {
729 let mut mock = MockJobService::new();
730 let mut backoff_policy = create_test_backoff_policy();
731 backoff_policy
732 .expect_wait_period()
733 .times(2)
734 .return_const(Duration::from_millis(1));
735
736 let mut seq = mockall::Sequence::new();
737
738 mock.expect_get_query_results()
739 .in_sequence(&mut seq)
740 .times(2)
741 .returning(|_, _| {
742 Ok(Response::from(
743 GetQueryResultsResponse::new().set_job_complete(false),
744 ))
745 });
746
747 mock.expect_get_query_results()
748 .in_sequence(&mut seq)
749 .times(1)
750 .returning(|_, _| {
751 Ok(Response::from(
752 GetQueryResultsResponse::new().set_job_complete(true),
753 ))
754 });
755
756 let job_service = create_job_service(mock);
757 let job_ref = JobReference::new()
758 .set_project_id("some_project")
759 .set_job_id("some_job_id");
760
761 let res = poll_query_results(&job_service, &job_ref, Arc::new(backoff_policy)).await?;
762
763 assert!(res.job_complete.unwrap(), "{res:?}");
764
765 Ok(())
766 }
767
768 #[tokio::test]
769 async fn test_query_until_done_job_failed_error() -> TestResult {
770 let mut mock = MockJobService::new();
771 mock.expect_get_query_results().returning(|req, _| {
772 assert_eq!(req.project_id, "some_project");
773 assert_eq!(req.job_id, "some_job_id");
774 assert_eq!(req.max_results, Some(0));
775 let err_proto = ErrorProto::new()
776 .set_reason("invalidQuery")
777 .set_message("Syntax error");
778 let res = GetQueryResultsResponse::new().set_errors(vec![err_proto]);
779 Ok(Response::from(res))
780 });
781 let job_service = create_job_service(mock);
782 let job_ref = JobReference::new()
783 .set_project_id("some_project")
784 .set_job_id("some_job_id");
785 let query_res = QueryResponse::new()
786 .set_job_complete(false)
787 .set_job_reference(job_ref);
788
789 let query = Query::from_query_response(job_service, query_res, None, None);
790
791 let err = query.until_done().await.unwrap_err();
792 let errors = match err {
793 QueryError::JobFailed { errors } => errors,
794 _ => panic!("expected QueryError::JobFailed, got {err:?}"),
795 };
796 assert_eq!(
797 errors,
798 [ErrorProto::new()
799 .set_reason("invalidQuery")
800 .set_message("Syntax error")]
801 );
802
803 Ok(())
804 }
805
806 #[tokio::test]
807 async fn test_query_until_done_rpc_error() -> TestResult {
808 let mut mock = MockJobService::new();
809 mock.expect_get_query_results().returning(|req, _| {
810 assert_eq!(req.project_id, "some_project");
811 assert_eq!(req.job_id, "some_job_id");
812 assert_eq!(req.max_results, Some(0));
813 let status = Status::default()
814 .set_code(Code::InvalidArgument)
815 .set_message("simulated bad request");
816 Err(GaxError::service(status))
817 });
818 let job_service = create_job_service(mock);
819 let job_ref = JobReference::new()
820 .set_project_id("some_project")
821 .set_job_id("some_job_id");
822 let query_res = QueryResponse::new()
823 .set_job_complete(false)
824 .set_job_reference(job_ref);
825
826 let query = Query::from_query_response(job_service, query_res, None, None);
827
828 let err = query.until_done().await.unwrap_err();
829 let source = match err {
830 QueryError::Rpc { source } => source,
831 _ => panic!("expected QueryError::Rpc, got {err:?}"),
832 };
833 assert_eq!(source.status().unwrap().code, Code::InvalidArgument);
834
835 Ok(())
836 }
837
838 #[tokio::test]
839 async fn test_complete_query_read() -> TestResult {
840 let job_service = create_job_service(MockJobService::new());
841 let job_ref = JobReference::new()
842 .set_project_id("some_project")
843 .set_job_id("some_job_id");
844 let schema = TableSchema::new().set_fields([TableFieldSchema::new()
845 .set_name("name")
846 .set_type("STRING")
847 .set_mode("NULLABLE")]);
848 let row = serde_json::Map::from_iter([(
849 "f".to_string(),
850 serde_json::json!([{ "v": "test_name" }]),
851 )]);
852 let query_res = QueryResponse::new()
853 .set_job_complete(true)
854 .set_job_reference(job_ref)
855 .set_schema(schema)
856 .set_rows(vec![row]);
857
858 let complete_query = CompleteQuery::from_query_response(job_service, query_res, None);
859
860 let mut iter = complete_query.read();
861 let row = iter.next().await.expect("should return first row")?;
862 assert_eq!(row.get::<String, _>("name")?, "test_name");
863 assert!(iter.next().await.is_none(), "{iter:?}");
864
865 Ok(())
866 }
867
868 #[tokio::test(start_paused = true)]
869 async fn test_query_until_done_reissue_on_retryable_job_failed() -> TestResult {
870 let mut mock = MockJobService::new();
871 let mut seq = mockall::Sequence::new();
872
873 mock.expect_get_query_results()
874 .in_sequence(&mut seq)
875 .times(1)
876 .returning(|req, _| {
877 assert_eq!(req.job_id, "initial_job_id");
878 let err_proto = ErrorProto::new()
879 .set_reason("backendError")
880 .set_message("temporary server issue");
881 let res = GetQueryResultsResponse::new().set_errors(vec![err_proto]);
882 Ok(Response::from(res))
883 });
884
885 mock.expect_query()
886 .in_sequence(&mut seq)
887 .times(1)
888 .returning(|req, _| {
889 let req_id = &req.query_request.as_ref().unwrap().request_id;
890 assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX));
891 assert!(uuid::Uuid::parse_str(&req_id[QUERY_REQUEST_ID_PREFIX.len()..]).is_ok());
892 let new_job_ref = JobReference::new()
893 .set_project_id("some_project")
894 .set_job_id("reissued_job_id");
895 Ok(Response::from(
896 QueryResponse::new()
897 .set_job_complete(false)
898 .set_job_reference(new_job_ref)
899 .set_schema(TableSchema::new()),
900 ))
901 });
902
903 mock.expect_get_query_results()
904 .in_sequence(&mut seq)
905 .times(1)
906 .returning(|req, _| {
907 assert_eq!(req.job_id, "reissued_job_id");
908 let res = GetQueryResultsResponse::new()
909 .set_job_complete(true)
910 .set_schema(TableSchema::new());
911 Ok(Response::from(res))
912 });
913
914 let job_service = create_job_service(mock);
915 let job_ref = JobReference::new()
916 .set_project_id("some_project")
917 .set_job_id("initial_job_id");
918 let query_res = QueryResponse::new()
919 .set_job_complete(false)
920 .set_job_reference(job_ref);
921
922 let query_builder = QueryBuilder::new(job_service.clone(), "SELECT 1".to_string())
923 .with_project_id("some_project");
924 let mut retry_context = RetryContext::new(query_builder);
925 retry_context.state.attempt_count = 1;
926
927 let query = Query::from_query_response(job_service, query_res, Some(&retry_context), None);
928
929 let completed = query.until_done().await?;
930 assert_eq!(
931 completed.job_ref.as_ref().unwrap().job_id,
932 "reissued_job_id"
933 );
934 Ok(())
935 }
936
937 #[tokio::test(start_paused = true)]
938 async fn test_query_until_done_reissue_retry_exhausted() -> TestResult {
939 let mut mock = MockJobService::new();
940 let mut seq = mockall::Sequence::new();
941
942 mock.expect_get_query_results()
944 .in_sequence(&mut seq)
945 .times(1)
946 .returning(|req, _| {
947 assert_eq!(req.job_id, "initial_job_id");
948 let err_proto = ErrorProto::new()
949 .set_reason("backendError")
950 .set_message("first temporary server issue");
951 let res = GetQueryResultsResponse::new().set_errors(vec![err_proto]);
952 Ok(Response::from(res))
953 });
954
955 mock.expect_query()
957 .in_sequence(&mut seq)
958 .times(1)
959 .returning(|req, _| {
960 let req_id = &req.query_request.as_ref().unwrap().request_id;
961 assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX));
962 assert!(uuid::Uuid::parse_str(&req_id[QUERY_REQUEST_ID_PREFIX.len()..]).is_ok());
963 let new_job_ref = JobReference::new()
964 .set_project_id("some_project")
965 .set_job_id("reissued_job_id");
966 Ok(Response::from(
967 QueryResponse::new()
968 .set_job_complete(false)
969 .set_job_reference(new_job_ref)
970 .set_schema(TableSchema::new()),
971 ))
972 });
973
974 mock.expect_get_query_results()
976 .in_sequence(&mut seq)
977 .times(1)
978 .returning(|req, _| {
979 assert_eq!(req.job_id, "reissued_job_id");
980 let err_proto = ErrorProto::new()
981 .set_reason("backendError")
982 .set_message("second temporary server issue");
983 let res = GetQueryResultsResponse::new().set_errors(vec![err_proto]);
984 Ok(Response::from(res))
985 });
986
987 let job_service = create_job_service(mock);
988 let job_ref = JobReference::new()
989 .set_project_id("some_project")
990 .set_job_id("initial_job_id");
991
992 let query_res = QueryResponse::new()
993 .set_job_complete(false)
994 .set_job_reference(job_ref)
995 .set_schema(TableSchema::new());
996 let mut query_builder = QueryBuilder::new(job_service.clone(), "SELECT 1".to_string())
997 .with_project_id("some_project");
998 query_builder.job_retry_policy =
999 Arc::new(RetryableJobErrors::default().with_attempt_limit(2));
1000 let mut retry_context = RetryContext::new(query_builder);
1001 retry_context.state.attempt_count = 1;
1002
1003 let query = Query::from_query_response(job_service, query_res, Some(&retry_context), None);
1004
1005 let err = query.until_done().await.unwrap_err();
1006 let errors = match err {
1007 QueryError::JobFailed { errors } => errors,
1008 _ => panic!("expected QueryError::JobFailed, got {err:?}"),
1009 };
1010 assert_eq!(errors.len(), 1);
1011 assert_eq!(errors[0].reason, "backendError");
1012 assert_eq!(errors[0].message, "second temporary server issue");
1013
1014 Ok(())
1015 }
1016
1017 #[tokio::test]
1018 async fn test_query_get_job_success() -> TestResult {
1019 let mut mock = MockJobService::new();
1020 mock.expect_get_job().returning(|req, _| {
1021 assert_eq!(req.project_id, "some_project");
1022 assert_eq!(req.job_id, "some_job_id");
1023 assert_eq!(req.location, "us-central1");
1024 let res = Job::new()
1025 .set_job_reference(JobReference::new().set_job_id(req.job_id))
1026 .set_user_email("test@example.com");
1027 Ok(Response::from(res))
1028 });
1029 let job_service = create_job_service(mock);
1030 let job_ref = JobReference::new()
1031 .set_project_id("some_project")
1032 .set_job_id("some_job_id")
1033 .set_location("us-central1");
1034 let query_res = QueryResponse::new()
1035 .set_schema(TableSchema::new())
1036 .set_job_reference(job_ref);
1037
1038 let query = Query::from_query_response(job_service.clone(), query_res.clone(), None, None);
1039 let job = query.get_job().unwrap().send().await?;
1040 assert_eq!(job.user_email, "test@example.com");
1041
1042 let complete_query = CompleteQuery::from_query_response(job_service, query_res, None);
1043 let job = complete_query.get_job().unwrap().send().await?;
1044 assert_eq!(job.user_email, "test@example.com");
1045 Ok(())
1046 }
1047
1048 #[tokio::test]
1049 async fn test_query_get_job_empty_job_id() -> TestResult {
1050 let job_service = create_job_service(MockJobService::new());
1051 let job_ref = JobReference::new()
1053 .set_location("us-central1")
1054 .set_project_id("some-project");
1055 let query_res = QueryResponse::new()
1056 .set_schema(TableSchema::new())
1057 .set_job_reference(job_ref);
1058
1059 let query = Query::from_query_response(job_service.clone(), query_res.clone(), None, None);
1060 assert!(query.get_job().is_none(), "{query:?}");
1061
1062 let complete_query = CompleteQuery::from_query_response(job_service, query_res, None);
1063 assert!(complete_query.get_job().is_none(), "{complete_query:?}");
1064 Ok(())
1065 }
1066
1067 #[tokio::test]
1068 async fn test_query_get_job_rpc_error() -> TestResult {
1069 let mut mock = MockJobService::new();
1070 mock.expect_get_job().returning(|req, _| {
1071 assert_eq!(req.project_id, "some_project");
1072 assert_eq!(req.job_id, "some_job_id");
1073 let status = Status::default()
1074 .set_code(Code::NotFound)
1075 .set_message("job not found");
1076 Err(GaxError::service(status))
1077 });
1078 let job_service = create_job_service(mock);
1079 let job_ref = JobReference::new()
1080 .set_project_id("some_project")
1081 .set_job_id("some_job_id");
1082 let query_res = QueryResponse::new()
1083 .set_schema(TableSchema::new())
1084 .set_job_reference(job_ref);
1085
1086 let query = Query::from_query_response(job_service.clone(), query_res.clone(), None, None);
1087 let req = query.get_job().unwrap();
1088 let err = req.send().await.unwrap_err();
1089 assert_eq!(err.status().unwrap().code, Code::NotFound);
1090
1091 let complete_query = CompleteQuery::from_query_response(job_service, query_res, None);
1092 let req = complete_query.get_job().unwrap();
1093 let err = req.send().await.unwrap_err();
1094 assert_eq!(err.status().unwrap().code, Code::NotFound);
1095 Ok(())
1096 }
1097
1098 #[tokio::test]
1099 async fn test_query_until_done_dry_run_job_returns_error() -> TestResult {
1100 let job_service = create_job_service(MockJobService::new());
1101 let job_ref = JobReference::new()
1102 .set_project_id("some_project")
1103 .set_location("US");
1104 let job = Job::new()
1105 .set_job_reference(job_ref)
1106 .set_configuration(JobConfiguration::new().set_dry_run(true));
1107
1108 let query = Query::from_job(job_service, job, None, None);
1109 let err = query.until_done().await.unwrap_err();
1110 assert!(
1111 matches!(err, QueryError::DryRun),
1112 "expected DryRun error, got {err:?}"
1113 );
1114 Ok(())
1115 }
1116
1117 #[tokio::test(start_paused = true)]
1118 async fn test_query_until_done_initial_poll_delay() -> TestResult {
1119 let mut mock = MockJobService::new();
1120 let mut seq = mockall::Sequence::new();
1121 mock.expect_get_query_results()
1122 .times(1)
1123 .in_sequence(&mut seq)
1124 .returning(|_, _| {
1125 Ok(Response::from(
1126 GetQueryResultsResponse::new().set_job_complete(false),
1127 ))
1128 });
1129 mock.expect_get_query_results()
1130 .times(1)
1131 .in_sequence(&mut seq)
1132 .returning(|_, _| {
1133 Ok(Response::from(
1134 GetQueryResultsResponse::new().set_job_complete(true),
1135 ))
1136 });
1137
1138 let job_service = create_job_service(mock);
1139 let job_ref = JobReference::new()
1140 .set_project_id("some_project")
1141 .set_job_id("some_job_id");
1142 let query_res = QueryResponse::new()
1143 .set_job_complete(false)
1144 .set_job_reference(job_ref);
1145
1146 let start = tokio::time::Instant::now();
1147 let query = Query::from_query_response(job_service, query_res, None, None);
1148 let _completed = query.until_done().await?;
1149 let elapsed = start.elapsed();
1150
1151 assert!(
1152 elapsed <= Duration::from_secs(1),
1153 "expected initial poll delay <= 1s, got {elapsed:?}"
1154 );
1155 Ok(())
1156 }
1157
1158 #[test]
1159 fn test_build_query_metadata_from_job() {
1160 use google_cloud_bigquery_v2::model::{
1161 JobConfigurationQuery, JobCreationReason, JobStatistics, JobStatistics2, JobStatus,
1162 };
1163
1164 let job =
1165 Job::new()
1166 .set_id("rust-sdk-testing:US.job_123")
1167 .set_kind("bigquery#job")
1168 .set_etag("etag123")
1169 .set_self_link("https://www.googleapis.com/bigquery/v2/...")
1170 .set_user_email("user@example.com")
1171 .set_principal_subject("user:user@example.com")
1172 .set_job_creation_reason(JobCreationReason::new().set_code(
1173 google_cloud_bigquery_v2::model::job_creation_reason::Code::Requested,
1174 ))
1175 .set_job_reference(
1176 JobReference::new()
1177 .set_project_id("rust-sdk-testing")
1178 .set_job_id("job_123")
1179 .set_location("US"),
1180 )
1181 .set_configuration(
1182 JobConfiguration::new().set_query(
1183 JobConfigurationQuery::new()
1184 .set_query("SELECT 1 AS one")
1185 .set_use_legacy_sql(false),
1186 ),
1187 )
1188 .set_status(JobStatus::new().set_state("DONE"))
1189 .set_statistics(
1190 JobStatistics::new()
1191 .set_creation_time(1790020718574i64)
1192 .set_start_time(1790020718592i64)
1193 .set_end_time(1790020718847i64)
1194 .set_total_bytes_processed(0i64)
1195 .set_query(
1196 JobStatistics2::new()
1197 .set_total_bytes_processed(0i64)
1198 .set_total_bytes_billed(0i64)
1199 .set_cache_hit(true)
1200 .set_statement_type("SELECT")
1201 .set_schema(TableSchema::new().set_fields([
1202 TableFieldSchema::new().set_name("one").set_type("INTEGER"),
1203 ])),
1204 ),
1205 );
1206
1207 let metadata = build_query_metadata_from_job(job);
1208
1209 assert_eq!(metadata.id, "rust-sdk-testing:US.job_123");
1210 assert_eq!(metadata.job_complete, Some(true));
1211 assert_eq!(metadata.creation_time, Some(1790020718574));
1212 assert_eq!(metadata.start_time, Some(1790020718592));
1213 assert_eq!(metadata.end_time, Some(1790020718847));
1214 assert_eq!(metadata.cache_hit, Some(true));
1215 assert_eq!(metadata.statement_type, "SELECT");
1216 assert_eq!(metadata.location, "US");
1217 assert_eq!(metadata.total_bytes_processed, Some(0));
1218 assert_eq!(metadata.total_bytes_billed, Some(0));
1219 assert!(metadata.schema.is_some());
1220 }
1221
1222 #[tokio::test]
1223 async fn test_query_until_done_preserves_metadata_fields_from_job() -> TestResult {
1224 use google_cloud_bigquery_v2::model::{
1225 JobConfigurationQuery, JobCreationReason, JobStatistics, JobStatistics2, JobStatus,
1226 };
1227
1228 let mut mock = MockJobService::new();
1229 mock.expect_get_query_results()
1230 .returning(|req, _| {
1231 let res = GetQueryResultsResponse::new()
1232 .set_job_complete(true)
1233 .set_job_reference(JobReference::new().set_job_id(req.job_id))
1234 .set_schema(TableSchema::new())
1235 .set_rows(vec![wkt::Struct::new()])
1236 .set_cache_hit(true);
1237 Ok(Response::from(res))
1238 })
1239 .times(1);
1240
1241 let job_service = create_job_service(mock);
1242 let job =
1243 Job::new()
1244 .set_id("rust-sdk-testing:US.job_123")
1245 .set_kind("bigquery#job")
1246 .set_job_creation_reason(JobCreationReason::new().set_code(
1247 google_cloud_bigquery_v2::model::job_creation_reason::Code::Requested,
1248 ))
1249 .set_job_reference(
1250 JobReference::new()
1251 .set_project_id("rust-sdk-testing")
1252 .set_job_id("job_123")
1253 .set_location("US"),
1254 )
1255 .set_configuration(
1256 JobConfiguration::new().set_query(
1257 JobConfigurationQuery::new()
1258 .set_query("SELECT 1 AS one")
1259 .set_use_legacy_sql(false),
1260 ),
1261 )
1262 .set_status(JobStatus::new().set_state("DONE"))
1263 .set_statistics(
1264 JobStatistics::new()
1265 .set_creation_time(1790020718574i64)
1266 .set_start_time(1790020718592i64)
1267 .set_end_time(1790020718847i64)
1268 .set_total_bytes_processed(0i64)
1269 .set_query(
1270 JobStatistics2::new()
1271 .set_total_bytes_processed(0i64)
1272 .set_total_bytes_billed(0i64)
1273 .set_cache_hit(true)
1274 .set_statement_type("SELECT"),
1275 ),
1276 );
1277
1278 let query = Query::from_job(job_service, job, None, None);
1279 let complete = query.until_done().await?;
1280 let meta = complete.metadata();
1281
1282 assert_eq!(meta.creation_time, Some(1790020718574));
1283 assert_eq!(meta.start_time, Some(1790020718592));
1284 assert_eq!(meta.end_time, Some(1790020718847));
1285 assert_eq!(meta.statement_type, "SELECT");
1286 assert_eq!(meta.location, "US");
1287 assert_eq!(meta.cache_hit, Some(true));
1288
1289 Ok(())
1290 }
1291}