Skip to main content

google_cloud_bigquery/query/
query_handle.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
15use 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/// A handle representing a running or completed SQL query execution.
32///
33/// [`Query::send()`](crate::builder::bigquery::Query::send) returns a [`Query`].
34///
35/// To obtain the final result set, call [`until_done()`](Query::until_done),
36/// which waits for the query execution to complete and returns a
37/// [`CompleteQuery`].
38///
39/// # Example
40///
41/// ```
42/// # use google_cloud_bigquery::client::BigQuery;
43/// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
44/// let query_handle = client
45///     .query("SELECT 42 AS answer")
46///     .send()
47///     .await?;
48///
49/// // Poll until execution completes.
50/// let completed = query_handle.until_done().await?;
51/// # Ok(())
52/// # }
53/// ```
54#[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    /// Returns the initial metadata from query execution.
106    ///
107    /// Depending on how the query was executed, the metadata contains:
108    /// - [`job_reference`][QueryMetadata::job_reference]: The reference to the BigQuery job, if one was created.
109    /// - [`query_id`][QueryMetadata::query_id]: The unique ID of the query if executed without creating a job.
110    /// - [`job_creation_reason`][QueryMetadata::job_creation_reason]: The reason why a job was created (controlled via `set_job_creation_mode`).
111    /// - [`job_complete`][QueryMetadata::job_complete]: Whether the query completed immediately without requiring polling.
112    ///
113    /// To wait for the query to finish and retrieve full results and final
114    /// metadata, call [`until_done`][Self::until_done].
115    pub fn metadata(&self) -> &QueryMetadata {
116        &self.metadata
117    }
118
119    /// Build a request to fetch full [Job] execution metadata from the service for this query.
120    ///
121    /// > Returns `None` if the query was executed without creating a job
122    /// > (for example, when using [`JobCreationMode::JobCreationOptional`],
123    /// > or when the query was a dry run).
124    ///
125    /// [Job]: https://docs.cloud.google.com/bigquery/docs/reference/rest/v2/Job
126    /// [`JobCreationMode::JobCreationOptional`]: google_cloud_bigquery_v2::model::query_request::JobCreationMode::JobCreationOptional
127    ///
128    /// # Example
129    ///
130    /// ```
131    /// # use google_cloud_bigquery::client::BigQuery;
132    /// # use google_cloud_bigquery_v2::model::query_request::JobCreationMode;
133    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
134    /// let query = client
135    ///     .query("SELECT 1")
136    ///     .set_job_creation_mode(JobCreationMode::JobCreationRequired) // Forces a job to be created
137    ///     .send()
138    ///     .await?;
139    ///
140    /// match query.get_job() {
141    ///     Some(req) => {
142    ///         let job_info = req.send().await?;
143    ///         println!("Executed by user: {}", job_info.user_email);
144    ///     }
145    ///     None => {
146    ///         println!(
147    ///             "Query was run without creating a job. Query ID: {}",
148    ///             query.metadata().query_id
149    ///         );
150    ///     }
151    /// }
152    /// # Ok(())
153    /// # }
154    /// ```
155    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    /// Waits for query execution to complete.
168    ///
169    /// If the query completed immediately, this method returns a
170    /// [`CompleteQuery`] without making additional
171    /// network calls. Otherwise, it polls the service until the query finishes.
172    ///
173    /// # Errors
174    ///
175    /// Returns [`QueryError::DryRun`] if the query was configured as a dry run.
176    ///
177    /// Returns an error if a remote service or network failure happens during
178    /// polling, or if the BigQuery job fails due to runtime execution errors.
179    ///
180    /// # Example
181    ///
182    /// ```
183    /// # use google_cloud_bigquery::client::BigQuery;
184    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
185    /// let query_handle = client
186    ///     .query("SELECT 1 + 1 AS result")
187    ///     .send()
188    ///     .await?;
189    ///
190    /// let complete = query_handle.until_done().await?;
191    /// # Ok(())
192    /// # }
193    /// ```
194    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/// A handle representing a successfully completed query ready for reading
258/// results.
259///
260/// [`Query::until_done()`] returns a [`CompleteQuery`].
261///
262/// This handle provides access to cached execution metadata, schema
263/// definitions, and a row iterator via [`read()`](CompleteQuery::read).
264///
265/// # Example
266///
267/// ```
268/// # use google_cloud_bigquery::client::BigQuery;
269/// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
270/// let complete = client
271///     .query("SELECT 'done' AS status")
272///     .until_done()
273///     .await?;
274///
275/// println!("Cache hit: {:?}", complete.metadata().cache_hit);
276/// # Ok(())
277/// # }
278/// ```
279#[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        // DDL/DML queries have no schema.
300        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        // DDL/DML queries have no schema.
327        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    /// Read the result set of the query.
346    ///
347    /// # Example
348    ///
349    /// ```
350    /// # use google_cloud_bigquery::client::BigQuery;
351    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
352    /// let mut rows = client
353    ///     .query("SELECT 100 AS score")
354    ///     .until_done()
355    ///     .await?
356    ///     .read();
357    ///
358    /// while let Some(row) = rows.next().await.transpose()? {
359    ///     let score: i64 = row.get("score");
360    ///     println!("Score: {score}");
361    /// }
362    /// # Ok(())
363    /// # }
364    /// ```
365    pub fn read(self) -> RowIterator {
366        RowIterator::new(self)
367    }
368
369    /// Returns a reference to the cached summary metadata for this query.
370    ///
371    /// The returned [`CompleteQueryMetadata`] contains
372    /// summary statistics such as total rows, schema details, cache hit
373    /// indicators, and estimated bytes processed without making additional
374    /// RPCs.
375    ///
376    /// # Example
377    ///
378    /// ```
379    /// # use google_cloud_bigquery::client::BigQuery;
380    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
381    /// let completed = client
382    ///     .query("SELECT 'metadata_check'")
383    ///     .until_done()
384    ///     .await?;
385    ///
386    /// let meta = completed.metadata();
387    /// println!("Total rows: {:?}", meta.total_rows);
388    /// # Ok(())
389    /// # }
390    /// ```
391    pub fn metadata(&self) -> &CompleteQueryMetadata {
392        &self.metadata
393    }
394
395    /// Build a request to fetch full [Job] execution metadata from the service for this query.
396    ///
397    /// > Returns `None` if the query was executed without creating a job
398    /// > (for example, when using [`JobCreationMode::JobCreationOptional`],
399    /// > or when the query was a dry run).
400    ///
401    /// [Job]: https://docs.cloud.google.com/bigquery/docs/reference/rest/v2/Job
402    /// [`JobCreationMode::JobCreationOptional`]: google_cloud_bigquery_v2::model::query_request::JobCreationMode::JobCreationOptional
403    ///
404    /// # Example
405    ///
406    /// ```
407    /// # use google_cloud_bigquery::client::BigQuery;
408    /// # use google_cloud_bigquery_v2::model::query_request::JobCreationMode;
409    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
410    /// let completed = client
411    ///     .query("SELECT 1")
412    ///     .set_job_creation_mode(JobCreationMode::JobCreationRequired) // Forces a job to be created
413    ///     .until_done()
414    ///     .await?;
415    ///
416    /// match completed.get_job() {
417    ///     Some(req) => {
418    ///         let job_info = req.send().await?;
419    ///         println!("Executed by user: {}", job_info.user_email);
420    ///     }
421    ///     None => {
422    ///         println!(
423    ///             "Query was run without creating a job. Query ID: {}",
424    ///             completed.metadata().query_id
425    ///         );
426    ///     }
427    /// }
428    /// # Ok(())
429    /// # }
430    /// ```
431    pub fn get_job(&self) -> Option<GetJob> {
432        build_get_job(&self.job_service, self.job_ref.as_ref()?)
433    }
434}
435
436/// Builds a `jobs.get` request from a job reference, or `None` if the
437/// reference cannot identify a job. Dry-run queries return a job reference
438/// without a job ID.
439fn 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
458/// Helper function to poll getQueryResults until a job finishes.
459pub(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        // TODO(#5592): limit retry attempts or add cancellation mechanism
493        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        // First poll on initial_job_id fails with retryable error (attempt 0 -> 1)
832        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        // Reissue succeeds and returns reissued_job_id
845        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        // Second poll on reissued_job_id fails with retryable error, but attempt limit (1) is exhausted!
864        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        // Dry-runs return a reference that has a location and project id, but no job id.
940        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}