Skip to main content

google_cloud_bigquery/query/
builder.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::QueryRequest;
17use crate::query::execution::RetryContext;
18use crate::query::retry_policy::{JobRetryPolicy, default_job_retry_policy};
19use crate::query::{CompleteQuery, Query as QueryHandle, Result};
20use google_cloud_bigquery_v2::client::JobService;
21use google_cloud_bigquery_v2::model::JobReference;
22use google_cloud_bigquery_v2::model::query_request::JobCreationMode;
23use std::sync::Arc;
24use uuid::Uuid;
25
26pub(crate) const JOB_ID_PREFIX: &str = "job_";
27pub(crate) const QUERY_REQUEST_ID_PREFIX: &str = "req_";
28
29/// A builder for configuring and executing a SQL query.
30///
31/// [`BigQuery::query()`](crate::client::BigQuery::query) returns a [`Query`] builder.
32///
33/// Use this builder to configure query parameters, dataset defaults, geographic location,
34/// result limits, and caching before executing the query with [`send()`](Query::send)
35/// or [`until_done()`](Query::until_done).
36///
37/// # Example
38///
39/// ```
40/// # use google_cloud_bigquery::client::BigQuery;
41/// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
42/// let mut rows = client
43///     .query("SELECT name FROM `bigquery-public-data.usa_names.usa_1910_2013` WHERE state = 'TX' LIMIT 100")
44///     .set_location("US")
45///     .set_page_size(50_u32)
46///     .until_done()
47///     .await?
48///     .read();
49///
50/// while let Some(row) = rows.next().await.transpose()? {
51///     let name: String = row.get("name")?;
52///     println!("Name: {name}");
53/// }
54/// # Ok(())
55/// # }
56/// ```
57#[derive(Clone, Debug)]
58pub struct Query {
59    pub(crate) job_service: Arc<JobService>,
60    pub(crate) request: QueryRequest,
61    pub(crate) project_id: Option<String>,
62    pub(crate) job_retry_policy: Arc<dyn JobRetryPolicy>,
63}
64
65impl Query {
66    /// Creates a new `Query` builder for the given SQL query.
67    pub(crate) fn new(job_service: Arc<JobService>, sql: String) -> Self {
68        Self {
69            job_service,
70            request: QueryRequest::default()
71                .set_query(sql)
72                .set_use_legacy_sql(wkt::BoolValue::from(false))
73                .set_job_creation_mode(JobCreationMode::JobCreationOptional),
74            project_id: None,
75            job_retry_policy: default_job_retry_policy(),
76        }
77    }
78
79    /// Sets the project ID to override the default client project ID for query
80    /// execution and billing.
81    ///
82    /// If you configured a default project ID on the
83    /// [`BigQuery`][crate::client::BigQuery] client via
84    /// [`ClientBuilder::with_project_id`](crate::builder::bigquery::ClientBuilder::with_project_id),
85    /// the query inherits it automatically. Calling this method overrides the
86    /// project ID for this specific query.
87    ///
88    /// You must specify a project ID either on the client or on the query
89    /// builder before calling [`send()`](Query::send).
90    ///
91    /// # Example
92    ///
93    /// ```
94    /// # use google_cloud_bigquery::client::BigQuery;
95    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
96    /// let query_handle = client
97    ///     .query("SELECT 1 AS count")
98    ///     .with_project_id("my-project-id")
99    ///     .send()
100    ///     .await?;
101    /// # Ok(())
102    /// # }
103    /// ```
104    pub fn with_project_id<S: Into<String>>(mut self, project_id: S) -> Self {
105        self.project_id = Some(project_id.into());
106        self
107    }
108
109    /// Executes the SQL query.
110    ///
111    /// Returns a [`Query`](QueryHandle) handle representing the query execution.
112    /// You can call [`until_done()`](QueryHandle::until_done) on the returned handle
113    /// to wait for results, or inspect the initial [`metadata()`](QueryHandle::metadata).
114    ///
115    /// You must configure a target project ID on either the
116    /// [`BigQuery`][crate::client::BigQuery] client via
117    /// [`ClientBuilder::with_project_id`](crate::builder::bigquery::ClientBuilder::with_project_id)
118    /// or on this builder via [`with_project_id()`](Query::with_project_id).
119    ///
120    /// # Example
121    ///
122    /// ```
123    /// # use google_cloud_bigquery::client::BigQuery;
124    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
125    /// let query_handle = client
126    ///     .query("SELECT CURRENT_TIMESTAMP() AS now")
127    ///     .send()
128    ///     .await?;
129    ///
130    /// let completed_query = query_handle
131    ///     .until_done()
132    ///     .await?;
133    ///
134    /// let mut rows = completed_query.read();
135    /// if let Some(row) = rows.next().await.transpose()? {
136    ///     let now: String = row.get("now")?;
137    ///     println!("Current time: {now}");
138    /// }
139    /// # Ok(())
140    /// # }
141    /// ```
142    pub async fn send(self) -> Result<QueryHandle> {
143        // Box heavy RPC call future to avoid large stack frames.
144        Box::pin(RetryContext::new(self).execute()).await
145    }
146
147    /// Sends the query execution request and waits until execution completes.
148    ///
149    /// This is a convenience method equivalent to calling
150    /// [`.send().await?.until_done().await`](Query::send).
151    ///
152    /// # Errors
153    ///
154    /// Returns [`QueryError::DryRun`](crate::error::QueryError::DryRun) if the query was configured as a dry run.
155    ///
156    /// Returns an error if a remote service or network failure happens during
157    /// execution or polling, or if the BigQuery job fails due to runtime execution errors.
158    ///
159    /// # Example
160    ///
161    /// ```
162    /// # use google_cloud_bigquery::client::BigQuery;
163    /// # async fn sample(client: BigQuery) -> anyhow::Result<()> {
164    /// let completed_query = client
165    ///     .query("SELECT CURRENT_TIMESTAMP() AS now")
166    ///     .until_done()
167    ///     .await?;
168    ///
169    /// let mut rows = completed_query.read();
170    /// if let Some(row) = rows.next().await.transpose()? {
171    ///     let now: String = row.get("now")?;
172    ///     println!("Current time: {now}");
173    /// }
174    /// # Ok(())
175    /// # }
176    /// ```
177    pub async fn until_done(self) -> Result<CompleteQuery> {
178        if self.request.dry_run {
179            return Err(QueryError::DryRun);
180        }
181        // Box heavy RPC call future to avoid large stack frames.
182        Box::pin(async move { self.send().await?.until_done().await }).await
183    }
184}
185
186// Create a job reference with a generated job ID.
187//
188// BigQuery does not strictly define a format for job IDs, just a limit in size.
189// See https://docs.cloud.google.com/bigquery/docs/reference/rest/v2/JobReference
190pub(crate) fn generate_job_reference(project_id: &str, location: &str) -> JobReference {
191    let job_id = generate_prefixed_id(JOB_ID_PREFIX);
192    let mut job_ref = JobReference::new()
193        .set_project_id(project_id.to_string())
194        .set_job_id(job_id);
195
196    if !location.is_empty() {
197        job_ref = job_ref.set_location(location.to_string());
198    }
199    job_ref
200}
201
202// Create a random ID with the given prefix and a UUID.
203//
204// BigQuery does not strictly define a format for request and job IDs, just a limit in size.
205// However, request IDs are more restrictive and have a limit of 36 characters.
206// UUID v4 simple format is used to reduce length.
207//
208// See https://docs.cloud.google.com/bigquery/docs/reference/rest/v2/jobs/query#QueryRequest for request ID
209// and https://docs.cloud.google.com/bigquery/docs/reference/rest/v2/JobReference for job ID.
210pub(crate) fn generate_prefixed_id(prefix: &str) -> String {
211    format!("{prefix}{}", Uuid::new_v4().simple())
212}
213
214include!("../generated/builder.rs");
215
216#[cfg(test)]
217mod tests {
218    use super::*;
219    use crate::client::BigQuery;
220    use crate::error::QueryError;
221    use crate::query::tests::{MockJobService, create_job_service};
222    use google_cloud_auth::credentials::anonymous::Builder as Anonymous;
223    use google_cloud_bigquery_v2::model::query_request::JobCreationMode;
224    use google_cloud_bigquery_v2::model::{
225        ErrorProto, Job, JobConfiguration, JobReference, JobStatus,
226        QueryRequest as JobsQueryRequest, QueryResponse,
227    };
228    use google_cloud_gax::error::Error as GaxError;
229    use google_cloud_gax::error::rpc::Status;
230    use google_cloud_gax::response::Response;
231
232    // bigquery limits request id to 36 characters
233    const BIGQUERY_REQ_ID_LIMIT: usize = 36;
234
235    type TestResult = anyhow::Result<()>;
236
237    #[test]
238    fn test_new() {
239        let job_service = create_job_service(MockJobService::new());
240        let sql = "SELECT 1".to_string();
241        let query_builder = Query::new(job_service, sql.clone());
242        assert_eq!(query_builder.request.query, sql);
243        assert_eq!(
244            query_builder.request.use_legacy_sql,
245            Some(wkt::BoolValue::from(false))
246        );
247        assert_eq!(
248            query_builder.request.job_creation_mode,
249            JobCreationMode::JobCreationOptional
250        );
251        assert_eq!(query_builder.project_id, None);
252    }
253
254    #[test]
255    fn test_with_project_id() {
256        let job_service = create_job_service(MockJobService::new());
257        let query_builder =
258            Query::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");
259        assert_eq!(query_builder.project_id.unwrap(), "my-project");
260    }
261
262    #[tokio::test]
263    async fn test_run_missing_project_id() -> anyhow::Result<()> {
264        let client = BigQuery::builder()
265            .with_credentials(Anonymous::new().build())
266            .build()
267            .await?;
268        let query_builder = client.query("SELECT 1");
269        let err = query_builder
270            .send()
271            .await
272            .expect_err("should return an error when project_id is missing");
273        assert!(
274            matches!(&err, QueryError::Rpc { source } if source.is_binding()),
275            "expected Binding error for missing project ID, got {err:?}"
276        );
277        Ok(())
278    }
279
280    #[tokio::test]
281    async fn test_run_missing_project_id_force_job_path() -> anyhow::Result<()> {
282        let client = BigQuery::builder()
283            .with_credentials(Anonymous::new().build())
284            .build()
285            .await?;
286        let query_builder = client.query("SELECT 1").set_allow_large_results(true);
287        let err = query_builder
288            .send()
289            .await
290            .expect_err("should return an error when project_id is missing");
291        assert!(
292            matches!(&err, QueryError::Rpc { source } if source.is_binding()),
293            "expected Binding error for missing project ID on job path, got {err:?}"
294        );
295        Ok(())
296    }
297
298    #[tokio::test]
299    async fn test_run_until_done_missing_project_id() -> anyhow::Result<()> {
300        let client = BigQuery::builder()
301            .with_credentials(Anonymous::new().build())
302            .build()
303            .await?;
304        let query_builder = client.query("SELECT 1");
305        let err = query_builder
306            .until_done()
307            .await
308            .expect_err("should return an error when project_id is missing");
309        assert!(
310            matches!(&err, QueryError::Rpc { source } if source.is_binding()),
311            "expected Binding error for missing project ID, got {err:?}"
312        );
313        Ok(())
314    }
315
316    #[test]
317    fn test_generate_prefixed_id() {
318        let job_id = generate_prefixed_id(JOB_ID_PREFIX);
319        assert!(job_id.starts_with(JOB_ID_PREFIX), "{job_id:?}");
320        assert!(
321            Uuid::parse_str(&job_id[JOB_ID_PREFIX.len()..]).is_ok(),
322            "{job_id:?}"
323        );
324
325        let req_id = generate_prefixed_id(QUERY_REQUEST_ID_PREFIX);
326        assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX), "{req_id:?}");
327        assert!(req_id.len() <= BIGQUERY_REQ_ID_LIMIT, "{req_id:?}");
328        assert!(
329            Uuid::parse_str(&req_id[QUERY_REQUEST_ID_PREFIX.len()..]).is_ok(),
330            "{req_id:?}"
331        );
332    }
333
334    #[test]
335    fn test_generate_job_reference() {
336        let job_ref = generate_job_reference("my-project", "us-central1");
337        assert_eq!(job_ref.project_id, "my-project");
338        assert!(job_ref.job_id.starts_with(JOB_ID_PREFIX), "{job_ref:?}");
339        assert_eq!(job_ref.location.as_deref(), Some("us-central1"));
340    }
341
342    #[tokio::test]
343    async fn test_run_jobs_insert() -> TestResult {
344        let mut mock = MockJobService::new();
345        mock.expect_insert_job().returning(|req, _| {
346            let job_ref = req.job.as_ref().unwrap().job_reference.as_ref().unwrap();
347            assert!(job_ref.job_id.starts_with(JOB_ID_PREFIX), "{job_ref:?}");
348            let job_ref = JobReference::new()
349                .set_job_id("test-job")
350                .set_project_id("my-project");
351            let job = Job::new()
352                .set_job_reference(job_ref)
353                .set_status(JobStatus::new().set_state("DONE"));
354            Ok(google_cloud_gax::response::Response::from(job))
355        });
356        mock.expect_query().never();
357
358        let job_service = create_job_service(mock);
359
360        let query_builder = Query::new(job_service, "SELECT 1".to_string())
361            .with_project_id("my-project")
362            .set_allow_large_results(true);
363        let query = query_builder.send().await?;
364        assert!(query.completed, "{query:?}");
365
366        Ok(())
367    }
368
369    #[tokio::test]
370    async fn test_run_jobs_query() -> TestResult {
371        let mut mock = MockJobService::new();
372        mock.expect_query().returning(move |req, _| {
373            let req_id = &req.query_request.as_ref().unwrap().request_id;
374            assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX), "{req_id:?}");
375            assert!(req_id.len() <= BIGQUERY_REQ_ID_LIMIT, "{req_id:?}");
376            Ok(Response::from(
377                QueryResponse::new().set_query_id("some_query_id"),
378            ))
379        });
380        mock.expect_insert_job().never();
381
382        let job_service = create_job_service(mock);
383        let query_builder =
384            Query::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");
385        let query = query_builder.send().await?;
386        assert!(!query.completed, "{query:?}");
387        assert_eq!(query.metadata.query_id, "some_query_id");
388
389        Ok(())
390    }
391
392    #[test]
393    fn test_force_job_path() {
394        let job_service = create_job_service(MockJobService::new());
395        let mut query_builder = Query::new(job_service, "SELECT 1".to_string());
396        assert!(!query_builder.request.force_job_path());
397
398        // setting a jobs.insert exclusive field
399        query_builder = query_builder.set_allow_large_results(true);
400        assert!(query_builder.request.force_job_path());
401    }
402
403    #[test]
404    fn test_request_conversions() {
405        let req = QueryRequest::default()
406            .set_query("SELECT 1".to_string())
407            .set_dry_run(true)
408            .set_use_legacy_sql(true);
409
410        let query_request: JobsQueryRequest = req.clone().into();
411        assert_eq!(query_request.query, "SELECT 1");
412        assert!(query_request.dry_run);
413        assert_eq!(
414            query_request.use_legacy_sql,
415            Some(wkt::BoolValue::from(true))
416        );
417
418        let job_config: JobConfiguration = req.into();
419        let job_query = job_config.query.as_ref().unwrap();
420        assert_eq!(job_query.query, "SELECT 1");
421        assert_eq!(job_query.use_legacy_sql, Some(wkt::BoolValue::from(true)));
422    }
423
424    #[tokio::test(start_paused = true)]
425    async fn test_run_reissue_on_retryable_job_failed() -> TestResult {
426        let mut mock = MockJobService::new();
427        let mut seq = mockall::Sequence::new();
428        let first_request_id = std::sync::Arc::new(std::sync::Mutex::new(String::new()));
429        let first_request_id_clone = first_request_id.clone();
430
431        mock.expect_query()
432            .in_sequence(&mut seq)
433            .times(1)
434            .returning(move |req, _| {
435                let req_id = req.query_request.as_ref().unwrap().request_id.clone();
436                assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX));
437                *first_request_id_clone.lock().unwrap() = req_id;
438                let err_proto = ErrorProto::new()
439                    .set_reason("backendError")
440                    .set_message("temporary server issue");
441                Ok(Response::from(
442                    QueryResponse::new().set_errors(vec![err_proto]),
443                ))
444            });
445
446        mock.expect_query()
447            .in_sequence(&mut seq)
448            .times(1)
449            .returning(move |req, _| {
450                let req_id = req.query_request.as_ref().unwrap().request_id.clone();
451                assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX));
452                assert_ne!(
453                    req_id,
454                    *first_request_id.lock().unwrap(),
455                    "reissued query must generate a fresh request_id to avoid 409 duplicate ID conflicts"
456                );
457                Ok(Response::from(
458                    QueryResponse::new().set_query_id("q_success"),
459                ))
460            });
461
462        let job_service = create_job_service(mock);
463        let query_builder =
464            Query::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");
465        let query = query_builder.send().await?;
466        assert_eq!(query.metadata.query_id, "q_success");
467
468        Ok(())
469    }
470
471    #[tokio::test(start_paused = true)]
472    async fn test_run_reissue_on_retryable_rpc_error() -> TestResult {
473        let mut mock = MockJobService::new();
474        let mut seq = mockall::Sequence::new();
475        let first_request_id = std::sync::Arc::new(std::sync::Mutex::new(String::new()));
476        let first_request_id_clone = first_request_id.clone();
477
478        mock.expect_query()
479            .in_sequence(&mut seq)
480            .times(1)
481            .returning(move |req, _| {
482                let req_id = req.query_request.as_ref().unwrap().request_id.clone();
483                assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX));
484                *first_request_id_clone.lock().unwrap() = req_id;
485                const BQ_REST_PAYLOAD: &[u8] = br#"{
486  "error": {
487    "code": 400,
488    "message": "The job encountered an error during execution. Retrying the job may solve the problem.",
489    "errors": [
490      {
491        "message": "The job encountered an error during execution. Retrying the job may solve the problem.",
492        "domain": "global",
493        "reason": "backendError"
494      }
495    ],
496    "status": "INVALID_ARGUMENT"
497  }
498}"#;
499                let status = Status::try_from(&bytes::Bytes::from_static(BQ_REST_PAYLOAD)).unwrap();
500                Err(GaxError::service(status))
501            });
502
503        mock.expect_query()
504            .in_sequence(&mut seq)
505            .times(1)
506            .returning(move |req, _| {
507                let req_id = req.query_request.as_ref().unwrap().request_id.clone();
508                assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX));
509                assert_ne!(
510                    req_id,
511                    *first_request_id.lock().unwrap(),
512                    "reissued query must generate a fresh request_id to avoid 409 duplicate ID conflicts"
513                );
514                Ok(Response::from(
515                    QueryResponse::new().set_query_id("q_success"),
516                ))
517            });
518
519        let job_service = create_job_service(mock);
520        let query_builder =
521            Query::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");
522        let query = query_builder.send().await?;
523        assert_eq!(query.metadata.query_id, "q_success");
524
525        Ok(())
526    }
527
528    #[tokio::test]
529    async fn test_run_jobs_query_with_page_size() -> TestResult {
530        let mut mock = MockJobService::new();
531        mock.expect_query().returning(move |req, _| {
532            assert_eq!(
533                req.query_request.as_ref().and_then(|r| r.max_results),
534                Some(100)
535            );
536            Ok(Response::from(QueryResponse::new()))
537        });
538        let job_service = create_job_service(mock);
539        let query_builder = Query::new(job_service, "SELECT 1".to_string())
540            .with_project_id("my-project")
541            .set_page_size(100_u32);
542        let query = query_builder.send().await?;
543        assert_eq!(query.page_size, Some(100));
544
545        Ok(())
546    }
547
548    #[tokio::test]
549    async fn test_run_jobs_insert_with_page_size() -> TestResult {
550        let mut mock = MockJobService::new();
551        mock.expect_insert_job().returning(|_, _| {
552            let job_ref = JobReference::new()
553                .set_job_id("test-job")
554                .set_project_id("my-project");
555            let job = Job::new()
556                .set_job_reference(job_ref)
557                .set_status(JobStatus::new().set_state("DONE"));
558            Ok(google_cloud_gax::response::Response::from(job))
559        });
560        let job_service = create_job_service(mock);
561
562        let query_builder = Query::new(job_service, "SELECT 1".to_string())
563            .with_project_id("my-project")
564            .set_allow_large_results(true)
565            .set_page_size(50_u32);
566        let query = query_builder.send().await?;
567        assert_eq!(query.page_size, Some(50));
568
569        Ok(())
570    }
571
572    #[tokio::test]
573    async fn test_until_done_jobs_query() -> TestResult {
574        let mut mock = MockJobService::new();
575        mock.expect_query().returning(move |_, _| {
576            Ok(Response::from(
577                QueryResponse::new()
578                    .set_job_complete(true)
579                    .set_query_id("some_query_id"),
580            ))
581        });
582        let job_service = create_job_service(mock);
583        let query = Query::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");
584        let complete = query.until_done().await?;
585        assert_eq!(complete.metadata().query_id, "some_query_id");
586
587        Ok(())
588    }
589
590    #[tokio::test]
591    async fn test_until_done_dry_run_returns_error() -> TestResult {
592        let mut mock = MockJobService::new();
593        mock.expect_insert_job().never();
594        mock.expect_query().never();
595        let job_service = create_job_service(mock);
596        let query = Query::new(job_service, "SELECT 1".to_string())
597            .with_project_id("my-project")
598            .set_dry_run(true);
599        let err = query.until_done().await.unwrap_err();
600        assert!(
601            matches!(err, QueryError::DryRun),
602            "expected DryRun error, got {err:?}"
603        );
604
605        Ok(())
606    }
607}