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};
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) 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    /// 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                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/// A handle representing a successfully completed query ready for reading
259/// results.
260///
261/// [`Query::until_done()`] returns a [`CompleteQuery`].
262///
263/// This handle provides access to cached execution metadata, schema
264/// definitions, and a row iterator via [`read()`](CompleteQuery::read).
265///
266/// # Example
267///
268/// ```
269/// # use google_cloud_bigquery::client::BigQuery;
270/// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
271/// let complete = client
272///     .query("SELECT 'done' AS status")
273///     .until_done()
274///     .await?;
275///
276/// println!("Cache hit: {:?}", complete.metadata().cache_hit);
277/// # Ok(())
278/// # }
279/// ```
280#[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    /// Read the result set of the query.
344    ///
345    /// # Example
346    ///
347    /// ```
348    /// # use google_cloud_bigquery::client::BigQuery;
349    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
350    /// let mut rows = client
351    ///     .query("SELECT 100 AS score")
352    ///     .until_done()
353    ///     .await?
354    ///     .read();
355    ///
356    /// while let Some(row) = rows.next().await.transpose()? {
357    ///     let score: i64 = row.get("score")?;
358    ///     println!("Score: {score}");
359    /// }
360    /// # Ok(())
361    /// # }
362    /// ```
363    pub fn read(self) -> RowIterator {
364        RowIterator::new(self)
365    }
366
367    /// Returns a reference to the cached summary metadata for this query.
368    ///
369    /// The returned [`CompleteQueryMetadata`] contains
370    /// summary statistics such as total rows, schema details, cache hit
371    /// indicators, and estimated bytes processed without making additional
372    /// RPCs.
373    ///
374    /// # Example
375    ///
376    /// ```
377    /// # use google_cloud_bigquery::client::BigQuery;
378    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
379    /// let completed = client
380    ///     .query("SELECT 'metadata_check'")
381    ///     .until_done()
382    ///     .await?;
383    ///
384    /// let meta = completed.metadata();
385    /// println!("Total rows: {:?}", meta.total_rows);
386    /// # Ok(())
387    /// # }
388    /// ```
389    pub fn metadata(&self) -> &CompleteQueryMetadata {
390        &self.metadata
391    }
392
393    /// Build a request to fetch full [Job] execution metadata from the service for this query.
394    ///
395    /// > Returns `None` if the query was executed without creating a job
396    /// > (for example, when using [`JobCreationMode::JobCreationOptional`],
397    /// > or when the query was a dry run).
398    ///
399    /// [Job]: https://docs.cloud.google.com/bigquery/docs/reference/rest/v2/Job
400    /// [`JobCreationMode::JobCreationOptional`]: google_cloud_bigquery_v2::model::query_request::JobCreationMode::JobCreationOptional
401    ///
402    /// # Example
403    ///
404    /// ```
405    /// # use google_cloud_bigquery::client::BigQuery;
406    /// # use google_cloud_bigquery_v2::model::query_request::JobCreationMode;
407    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
408    /// let completed = client
409    ///     .query("SELECT 1")
410    ///     .set_job_creation_mode(JobCreationMode::JobCreationRequired) // Forces a job to be created
411    ///     .until_done()
412    ///     .await?;
413    ///
414    /// match completed.get_job() {
415    ///     Some(req) => {
416    ///         let job_info = req.send().await?;
417    ///         println!("Executed by user: {}", job_info.user_email);
418    ///     }
419    ///     None => {
420    ///         println!(
421    ///             "Query was run without creating a job. Query ID: {}",
422    ///             completed.metadata().query_id
423    ///         );
424    ///     }
425    /// }
426    /// # Ok(())
427    /// # }
428    /// ```
429    pub fn get_job(&self) -> Option<GetJob> {
430        build_get_job(&self.job_service, self.job_ref.as_ref()?)
431    }
432}
433
434/// Builds a `jobs.get` request from a job reference, or `None` if the
435/// reference cannot identify a job. Dry-run queries return a job reference
436/// without a job ID.
437pub(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
456/// Helper function to poll getQueryResults until a job finishes.
457pub(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        // TODO(#5592): limit retry attempts or add cancellation mechanism
491        state.attempt_count += 1;
492    }
493}
494
495// Helper function to build QueryMetadata from a Job.
496//
497// The generated code handle fields with same name on the root level, but some data
498// that is returned on jobs.query response are under JobStats for a Query Job when using
499// jobs.insert.
500fn 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
560// Helper function to build CompleteQueryMetadata from GetQueryResultsResponse while
561// preserving metadata from the initial query execution.
562fn 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        // First poll on initial_job_id fails with retryable error (attempt 1)
943        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        // Reissue succeeds and returns reissued_job_id (attempt 2)
956        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        // Second poll on reissued_job_id fails with retryable error, but attempt limit (2) is exhausted!
975        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        // Dry-runs return a reference that has a location and project id, but no job id.
1052        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}