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, Schema};
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) max_results: 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 max_results: 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: QueryMetadata::from(initial_job),
81 retry_context,
82 max_results,
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 max_results: 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,
101 max_results,
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 max_results,
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 max_results,
215 ));
216 }
217
218 let job_ref = metadata
219 .job_reference
220 .as_ref()
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(10))
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 max_results,
237 ));
238 }
239 Err(err) => {
240 let Some(retry_ctx) = retry_context else {
241 return Err(err);
242 };
243 match retry_ctx.on_error(err) {
244 JobRetryResult::Continue(delay, _) => {
245 self = retry_ctx.reissue(delay).await?;
246 }
247 JobRetryResult::Permanent(e) | JobRetryResult::Exhausted(e) => {
248 return Err(e);
249 }
250 }
251 }
252 }
253 }
254 }
255}
256
257#[derive(Clone, Debug)]
280pub struct CompleteQuery {
281 pub(crate) job_service: Arc<JobService>,
282 pub(crate) job_ref: Option<JobReference>,
283 pub(crate) cached_rows: VecDeque<wkt::Struct>,
284 pub(crate) schema: Arc<Schema>,
285 pub(crate) page_token: Option<String>,
286 pub(crate) metadata: CompleteQueryMetadata,
287 pub(crate) max_results: 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 max_results: Option<u32>,
296 ) -> Self {
297 let cached_rows = VecDeque::from(std::mem::take(&mut res.rows));
298 let metadata = CompleteQueryMetadata::from(res);
299 let schema = metadata.schema.clone().unwrap_or_default();
301 let schema = Arc::new(Schema::new(schema));
302 let page_token = if metadata.page_token.is_empty() {
303 None
304 } else {
305 Some(metadata.page_token.clone())
306 };
307 Self {
308 job_service,
309 job_ref: Some(job_ref.clone()),
310 cached_rows,
311 page_token,
312 schema,
313 metadata,
314 max_results,
315 }
316 }
317
318 pub(crate) fn from_query_metadata(
319 job_service: Arc<JobService>,
320 metadata: QueryMetadata,
321 cached_rows: VecDeque<wkt::Struct>,
322 max_results: Option<u32>,
323 ) -> Self {
324 let job_ref = metadata.job_reference.clone();
325 let metadata = CompleteQueryMetadata::from(metadata);
326 let schema = metadata.schema.clone().unwrap_or_default();
328 let schema = Arc::new(Schema::new(schema));
329 let page_token = if metadata.page_token.is_empty() {
330 None
331 } else {
332 Some(metadata.page_token.clone())
333 };
334 Self {
335 job_service,
336 job_ref,
337 cached_rows,
338 page_token,
339 schema,
340 metadata,
341 max_results,
342 }
343 }
344
345 pub fn read(self) -> RowIterator {
366 RowIterator::new(self)
367 }
368
369 pub fn metadata(&self) -> &CompleteQueryMetadata {
392 &self.metadata
393 }
394
395 pub fn get_job(&self) -> Option<GetJob> {
432 build_get_job(&self.job_service, self.job_ref.as_ref()?)
433 }
434}
435
436fn build_get_job(job_service: &JobService, job_ref: &JobReference) -> Option<GetJob> {
440 if job_ref.job_id.is_empty() {
441 return None;
442 }
443
444 let req = job_service
445 .get_job()
446 .set_job_id(job_ref.job_id.clone())
447 .set_project_id(job_ref.project_id.clone());
448
449 let req = job_ref
450 .location
451 .clone()
452 .into_iter()
453 .fold(req, |req, location| req.set_location(location));
454
455 Some(req)
456}
457
458pub(crate) async fn poll_query_results(
460 job_service: &JobService,
461 job_ref: &JobReference,
462 backoff_policy: Arc<dyn PollingBackoffPolicy>,
463) -> Result<GetQueryResultsResponse> {
464 let mut state = PollingState::default();
465
466 loop {
467 let mut req = GetQueryResultsRequest::new()
468 .set_max_results(0u32)
469 .set_project_id(job_ref.project_id.clone())
470 .set_job_id(job_ref.job_id.clone());
471 if let Some(location) = job_ref.location.clone() {
472 req = req.set_location(location);
473 }
474
475 let res = job_service
476 .get_query_results()
477 .with_request(req)
478 .send()
479 .await?;
480
481 if !res.errors.is_empty() {
482 return Err(QueryError::JobFailed { errors: res.errors });
483 }
484
485 let completed = res.job_complete.unwrap_or(false);
486 if completed {
487 return Ok(res);
488 }
489
490 let delay = backoff_policy.wait_period(&state);
491 tokio::time::sleep(delay).await;
492 state.attempt_count += 1;
494 }
495}
496
497#[cfg(test)]
498mod tests {
499 use super::*;
500 use crate::query::builder::{QUERY_REQUEST_ID_PREFIX, Query as QueryBuilder};
501 use crate::query::retry_policy::RetryableJobErrors;
502 use crate::query::tests::{MockJobService, create_job_service, create_test_backoff_policy};
503 use google_cloud_bigquery_v2::model::{
504 ErrorProto, GetQueryResultsResponse, Job, JobConfiguration, JobReference, QueryResponse,
505 TableFieldSchema, TableSchema,
506 };
507 use google_cloud_gax::error::Error as GaxError;
508 use google_cloud_gax::error::rpc::{Code, Status};
509 use google_cloud_gax::response::Response;
510 use std::time::Duration;
511
512 type TestResult = anyhow::Result<()>;
513
514 impl CompleteQuery {
515 pub(crate) fn from_query_response(
516 job_service: Arc<JobService>,
517 mut query_res: QueryResponse,
518 max_results: Option<u32>,
519 ) -> Self {
520 let cached_rows = std::mem::take(&mut query_res.rows).into();
521 let metadata = QueryMetadata::from(query_res);
522 Self::from_query_metadata(job_service, metadata, cached_rows, max_results)
523 }
524 }
525
526 #[tokio::test]
527 async fn test_query_until_done_already_completed() -> TestResult {
528 let job_service = create_job_service(MockJobService::new());
529 let job_ref = JobReference::new()
530 .set_project_id("some_project")
531 .set_job_id("some_job_id");
532 let query_res = QueryResponse::new()
533 .set_job_complete(true)
534 .set_job_reference(job_ref.clone())
535 .set_schema(TableSchema::new())
536 .set_page_token("some_page_token")
537 .set_rows([wkt::Struct::new()])
538 .set_cache_hit(true);
539
540 let query = Query::from_query_response(job_service, query_res, None, None);
541
542 let completed = query.until_done().await?;
543 assert_eq!(completed.job_ref.as_ref().unwrap().job_id, "some_job_id");
544 assert_eq!(completed.page_token, Some("some_page_token".to_string()));
545 assert_eq!(completed.cached_rows.len(), 1);
546
547 let metadata = completed.metadata();
548 assert_eq!(metadata.cache_hit, Some(true));
549 assert_eq!(metadata.job_complete, Some(true));
550 assert_eq!(metadata.page_token, "some_page_token".to_string());
551
552 Ok(())
553 }
554
555 #[tokio::test]
556 async fn test_query_until_done_preserves_max_results() -> TestResult {
557 let job_service = create_job_service(MockJobService::new());
558 let job_ref = JobReference::new()
559 .set_project_id("some_project")
560 .set_job_id("some_job_id");
561 let query_res = QueryResponse::new()
562 .set_job_complete(true)
563 .set_job_reference(job_ref.clone())
564 .set_schema(TableSchema::new());
565
566 let query = Query::from_query_response(job_service, query_res, None, Some(42));
567
568 let completed = query.until_done().await?;
569 assert_eq!(completed.max_results, Some(42));
570
571 Ok(())
572 }
573
574 #[tokio::test]
575 async fn test_query_until_done_polls_success() -> TestResult {
576 let mut mock = MockJobService::new();
577 mock.expect_get_query_results()
578 .returning(|req, _| {
579 assert_eq!(req.project_id, "some_project");
580 assert_eq!(req.job_id, "some_job_id");
581 assert_eq!(req.max_results, Some(0));
582 assert_eq!(req.location, "us-central1");
583 let res = GetQueryResultsResponse::new()
584 .set_job_complete(true)
585 .set_job_reference(JobReference::new().set_job_id(req.job_id))
586 .set_schema(TableSchema::new())
587 .set_page_token("")
588 .set_rows(vec![wkt::Struct::new(), wkt::Struct::new()])
589 .set_cache_hit(false);
590 Ok(Response::from(res))
591 })
592 .times(1);
593 let job_service = create_job_service(mock);
594 let job_ref = JobReference::new()
595 .set_project_id("some_project")
596 .set_job_id("some_job_id")
597 .set_location("us-central1");
598 let query_res = QueryResponse::new()
599 .set_job_complete(false)
600 .set_job_reference(job_ref);
601
602 let query = Query::from_query_response(job_service, query_res, None, None);
603
604 let completed = query.until_done().await?;
605 assert_eq!(completed.job_ref.as_ref().unwrap().job_id, "some_job_id");
606 assert_eq!(completed.page_token, None);
607 assert_eq!(completed.cached_rows.len(), 2);
608
609 let metadata = completed.metadata();
610 assert_eq!(metadata.cache_hit, Some(false));
611 assert_eq!(metadata.job_complete, Some(true));
612 assert_eq!(metadata.page_token, "".to_string());
613
614 Ok(())
615 }
616
617 #[tokio::test(start_paused = true)]
618 async fn test_poll_query_results_loops_until_complete() -> TestResult {
619 let mut mock = MockJobService::new();
620 let mut backoff_policy = create_test_backoff_policy();
621 backoff_policy
622 .expect_wait_period()
623 .times(2)
624 .return_const(Duration::from_millis(1));
625
626 let mut seq = mockall::Sequence::new();
627
628 mock.expect_get_query_results()
629 .in_sequence(&mut seq)
630 .times(2)
631 .returning(|_, _| {
632 Ok(Response::from(
633 GetQueryResultsResponse::new().set_job_complete(false),
634 ))
635 });
636
637 mock.expect_get_query_results()
638 .in_sequence(&mut seq)
639 .times(1)
640 .returning(|_, _| {
641 Ok(Response::from(
642 GetQueryResultsResponse::new().set_job_complete(true),
643 ))
644 });
645
646 let job_service = create_job_service(mock);
647 let job_ref = JobReference::new()
648 .set_project_id("some_project")
649 .set_job_id("some_job_id");
650
651 let res = poll_query_results(&job_service, &job_ref, Arc::new(backoff_policy)).await?;
652
653 assert!(res.job_complete.unwrap(), "{res:?}");
654
655 Ok(())
656 }
657
658 #[tokio::test]
659 async fn test_query_until_done_job_failed_error() -> TestResult {
660 let mut mock = MockJobService::new();
661 mock.expect_get_query_results().returning(|req, _| {
662 assert_eq!(req.project_id, "some_project");
663 assert_eq!(req.job_id, "some_job_id");
664 assert_eq!(req.max_results, Some(0));
665 let err_proto = ErrorProto::new()
666 .set_reason("invalidQuery")
667 .set_message("Syntax error");
668 let res = GetQueryResultsResponse::new().set_errors(vec![err_proto]);
669 Ok(Response::from(res))
670 });
671 let job_service = create_job_service(mock);
672 let job_ref = JobReference::new()
673 .set_project_id("some_project")
674 .set_job_id("some_job_id");
675 let query_res = QueryResponse::new()
676 .set_job_complete(false)
677 .set_job_reference(job_ref);
678
679 let query = Query::from_query_response(job_service, query_res, None, None);
680
681 let err = query.until_done().await.unwrap_err();
682 let errors = match err {
683 QueryError::JobFailed { errors } => errors,
684 _ => panic!("expected QueryError::JobFailed, got {err:?}"),
685 };
686 assert_eq!(
687 errors,
688 [ErrorProto::new()
689 .set_reason("invalidQuery")
690 .set_message("Syntax error")]
691 );
692
693 Ok(())
694 }
695
696 #[tokio::test]
697 async fn test_query_until_done_rpc_error() -> TestResult {
698 let mut mock = MockJobService::new();
699 mock.expect_get_query_results().returning(|req, _| {
700 assert_eq!(req.project_id, "some_project");
701 assert_eq!(req.job_id, "some_job_id");
702 assert_eq!(req.max_results, Some(0));
703 let status = Status::default()
704 .set_code(Code::InvalidArgument)
705 .set_message("simulated bad request");
706 Err(GaxError::service(status))
707 });
708 let job_service = create_job_service(mock);
709 let job_ref = JobReference::new()
710 .set_project_id("some_project")
711 .set_job_id("some_job_id");
712 let query_res = QueryResponse::new()
713 .set_job_complete(false)
714 .set_job_reference(job_ref);
715
716 let query = Query::from_query_response(job_service, query_res, None, None);
717
718 let err = query.until_done().await.unwrap_err();
719 let source = match err {
720 QueryError::Rpc { source } => source,
721 _ => panic!("expected QueryError::Rpc, got {err:?}"),
722 };
723 assert_eq!(source.status().unwrap().code, Code::InvalidArgument);
724
725 Ok(())
726 }
727
728 #[tokio::test]
729 async fn test_complete_query_read() -> TestResult {
730 let job_service = create_job_service(MockJobService::new());
731 let job_ref = JobReference::new()
732 .set_project_id("some_project")
733 .set_job_id("some_job_id");
734 let schema = TableSchema::new().set_fields([TableFieldSchema::new()
735 .set_name("name")
736 .set_type("STRING")
737 .set_mode("NULLABLE")]);
738 let row = serde_json::Map::from_iter([(
739 "f".to_string(),
740 serde_json::json!([{ "v": "test_name" }]),
741 )]);
742 let query_res = QueryResponse::new()
743 .set_job_complete(true)
744 .set_job_reference(job_ref)
745 .set_schema(schema)
746 .set_rows(vec![row]);
747
748 let complete_query = CompleteQuery::from_query_response(job_service, query_res, None);
749
750 let mut iter = complete_query.read();
751 let row = iter.next().await.expect("should return first row")?;
752 assert_eq!(row.get::<String, _>("name"), "test_name");
753 assert!(iter.next().await.is_none(), "{iter:?}");
754
755 Ok(())
756 }
757
758 #[tokio::test(start_paused = true)]
759 async fn test_query_until_done_reissue_on_retryable_job_failed() -> TestResult {
760 let mut mock = MockJobService::new();
761 let mut seq = mockall::Sequence::new();
762
763 mock.expect_get_query_results()
764 .in_sequence(&mut seq)
765 .times(1)
766 .returning(|req, _| {
767 assert_eq!(req.job_id, "initial_job_id");
768 let err_proto = ErrorProto::new()
769 .set_reason("backendError")
770 .set_message("temporary server issue");
771 let res = GetQueryResultsResponse::new().set_errors(vec![err_proto]);
772 Ok(Response::from(res))
773 });
774
775 mock.expect_query()
776 .in_sequence(&mut seq)
777 .times(1)
778 .returning(|req, _| {
779 let req_id = &req.query_request.as_ref().unwrap().request_id;
780 assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX));
781 assert!(uuid::Uuid::parse_str(&req_id[QUERY_REQUEST_ID_PREFIX.len()..]).is_ok());
782 let new_job_ref = JobReference::new()
783 .set_project_id("some_project")
784 .set_job_id("reissued_job_id");
785 Ok(Response::from(
786 QueryResponse::new()
787 .set_job_complete(false)
788 .set_job_reference(new_job_ref)
789 .set_schema(TableSchema::new()),
790 ))
791 });
792
793 mock.expect_get_query_results()
794 .in_sequence(&mut seq)
795 .times(1)
796 .returning(|req, _| {
797 assert_eq!(req.job_id, "reissued_job_id");
798 let res = GetQueryResultsResponse::new()
799 .set_job_complete(true)
800 .set_schema(TableSchema::new());
801 Ok(Response::from(res))
802 });
803
804 let job_service = create_job_service(mock);
805 let job_ref = JobReference::new()
806 .set_project_id("some_project")
807 .set_job_id("initial_job_id");
808 let query_res = QueryResponse::new()
809 .set_job_complete(false)
810 .set_job_reference(job_ref);
811
812 let query_builder = QueryBuilder::new(job_service.clone(), "SELECT 1".to_string())
813 .with_project_id("some_project");
814 let retry_context = Some(RetryContext::new(query_builder));
815
816 let query = Query::from_query_response(job_service, query_res, retry_context, None);
817
818 let completed = query.until_done().await?;
819 assert_eq!(
820 completed.job_ref.as_ref().unwrap().job_id,
821 "reissued_job_id"
822 );
823 Ok(())
824 }
825
826 #[tokio::test(start_paused = true)]
827 async fn test_query_until_done_reissue_retry_exhausted() -> TestResult {
828 let mut mock = MockJobService::new();
829 let mut seq = mockall::Sequence::new();
830
831 mock.expect_get_query_results()
833 .in_sequence(&mut seq)
834 .times(1)
835 .returning(|req, _| {
836 assert_eq!(req.job_id, "initial_job_id");
837 let err_proto = ErrorProto::new()
838 .set_reason("backendError")
839 .set_message("first temporary server issue");
840 let res = GetQueryResultsResponse::new().set_errors(vec![err_proto]);
841 Ok(Response::from(res))
842 });
843
844 mock.expect_query()
846 .in_sequence(&mut seq)
847 .times(1)
848 .returning(|req, _| {
849 let req_id = &req.query_request.as_ref().unwrap().request_id;
850 assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX));
851 assert!(uuid::Uuid::parse_str(&req_id[QUERY_REQUEST_ID_PREFIX.len()..]).is_ok());
852 let new_job_ref = JobReference::new()
853 .set_project_id("some_project")
854 .set_job_id("reissued_job_id");
855 Ok(Response::from(
856 QueryResponse::new()
857 .set_job_complete(false)
858 .set_job_reference(new_job_ref)
859 .set_schema(TableSchema::new()),
860 ))
861 });
862
863 mock.expect_get_query_results()
865 .in_sequence(&mut seq)
866 .times(1)
867 .returning(|req, _| {
868 assert_eq!(req.job_id, "reissued_job_id");
869 let err_proto = ErrorProto::new()
870 .set_reason("backendError")
871 .set_message("second temporary server issue");
872 let res = GetQueryResultsResponse::new().set_errors(vec![err_proto]);
873 Ok(Response::from(res))
874 });
875
876 let job_service = create_job_service(mock);
877 let job_ref = JobReference::new()
878 .set_project_id("some_project")
879 .set_job_id("initial_job_id");
880
881 let query_res = QueryResponse::new()
882 .set_job_complete(false)
883 .set_job_reference(job_ref)
884 .set_schema(TableSchema::new());
885 let mut query_builder = QueryBuilder::new(job_service.clone(), "SELECT 1".to_string())
886 .with_project_id("some_project");
887 query_builder.job_retry_policy =
888 Arc::new(RetryableJobErrors::default().with_attempt_limit(1));
889 let retry_context = Some(RetryContext::new(query_builder));
890
891 let query = Query::from_query_response(job_service, query_res, retry_context, None);
892
893 let err = query.until_done().await.unwrap_err();
894 let errors = match err {
895 QueryError::JobFailed { errors } => errors,
896 _ => panic!("expected QueryError::JobFailed, got {err:?}"),
897 };
898 assert_eq!(errors.len(), 1);
899 assert_eq!(errors[0].reason, "backendError");
900 assert_eq!(errors[0].message, "second temporary server issue");
901
902 Ok(())
903 }
904
905 #[tokio::test]
906 async fn test_query_get_job_success() -> TestResult {
907 let mut mock = MockJobService::new();
908 mock.expect_get_job().returning(|req, _| {
909 assert_eq!(req.project_id, "some_project");
910 assert_eq!(req.job_id, "some_job_id");
911 assert_eq!(req.location, "us-central1");
912 let res = Job::new()
913 .set_job_reference(JobReference::new().set_job_id(req.job_id))
914 .set_user_email("test@example.com");
915 Ok(Response::from(res))
916 });
917 let job_service = create_job_service(mock);
918 let job_ref = JobReference::new()
919 .set_project_id("some_project")
920 .set_job_id("some_job_id")
921 .set_location("us-central1");
922 let query_res = QueryResponse::new()
923 .set_schema(TableSchema::new())
924 .set_job_reference(job_ref);
925
926 let query = Query::from_query_response(job_service.clone(), query_res.clone(), None, None);
927 let job = query.get_job().unwrap().send().await?;
928 assert_eq!(job.user_email, "test@example.com");
929
930 let complete_query = CompleteQuery::from_query_response(job_service, query_res, None);
931 let job = complete_query.get_job().unwrap().send().await?;
932 assert_eq!(job.user_email, "test@example.com");
933 Ok(())
934 }
935
936 #[tokio::test]
937 async fn test_query_get_job_empty_job_id() -> TestResult {
938 let job_service = create_job_service(MockJobService::new());
939 let job_ref = JobReference::new()
941 .set_location("us-central1")
942 .set_project_id("some-project");
943 let query_res = QueryResponse::new()
944 .set_schema(TableSchema::new())
945 .set_job_reference(job_ref);
946
947 let query = Query::from_query_response(job_service.clone(), query_res.clone(), None, None);
948 assert!(query.get_job().is_none(), "{query:?}");
949
950 let complete_query = CompleteQuery::from_query_response(job_service, query_res, None);
951 assert!(complete_query.get_job().is_none(), "{complete_query:?}");
952 Ok(())
953 }
954
955 #[tokio::test]
956 async fn test_query_get_job_rpc_error() -> TestResult {
957 let mut mock = MockJobService::new();
958 mock.expect_get_job().returning(|req, _| {
959 assert_eq!(req.project_id, "some_project");
960 assert_eq!(req.job_id, "some_job_id");
961 let status = Status::default()
962 .set_code(Code::NotFound)
963 .set_message("job not found");
964 Err(GaxError::service(status))
965 });
966 let job_service = create_job_service(mock);
967 let job_ref = JobReference::new()
968 .set_project_id("some_project")
969 .set_job_id("some_job_id");
970 let query_res = QueryResponse::new()
971 .set_schema(TableSchema::new())
972 .set_job_reference(job_ref);
973
974 let query = Query::from_query_response(job_service.clone(), query_res.clone(), None, None);
975 let req = query.get_job().unwrap();
976 let err = req.send().await.unwrap_err();
977 assert_eq!(err.status().unwrap().code, Code::NotFound);
978
979 let complete_query = CompleteQuery::from_query_response(job_service, query_res, None);
980 let req = complete_query.get_job().unwrap();
981 let err = req.send().await.unwrap_err();
982 assert_eq!(err.status().unwrap().code, Code::NotFound);
983 Ok(())
984 }
985
986 #[tokio::test]
987 async fn test_query_until_done_dry_run_job_returns_error() -> TestResult {
988 let job_service = create_job_service(MockJobService::new());
989 let job_ref = JobReference::new()
990 .set_project_id("some_project")
991 .set_location("US");
992 let job = Job::new()
993 .set_job_reference(job_ref)
994 .set_configuration(JobConfiguration::new().set_dry_run(true));
995
996 let query = Query::from_job(job_service, job, None, None);
997 let err = query.until_done().await.unwrap_err();
998 assert!(
999 matches!(err, QueryError::DryRun),
1000 "expected DryRun error, got {err:?}"
1001 );
1002 Ok(())
1003 }
1004}