google_cloud_bigquery/query/
builder.rs1use 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#[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 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 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 pub async fn send(self) -> Result<QueryHandle> {
143 Box::pin(RetryContext::new(self).execute()).await
145 }
146
147 pub async fn until_done(self) -> Result<CompleteQuery> {
178 if self.request.dry_run {
179 return Err(QueryError::DryRun);
180 }
181 Box::pin(async move { self.send().await?.until_done().await }).await
183 }
184}
185
186pub(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
202pub(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::response::Response;
229
230 const BIGQUERY_REQ_ID_LIMIT: usize = 36;
232
233 type TestResult = anyhow::Result<()>;
234
235 #[test]
236 fn test_new() {
237 let job_service = create_job_service(MockJobService::new());
238 let sql = "SELECT 1".to_string();
239 let query_builder = Query::new(job_service, sql.clone());
240 assert_eq!(query_builder.request.query, sql);
241 assert_eq!(
242 query_builder.request.use_legacy_sql,
243 Some(wkt::BoolValue::from(false))
244 );
245 assert_eq!(
246 query_builder.request.job_creation_mode,
247 JobCreationMode::JobCreationOptional
248 );
249 assert_eq!(query_builder.project_id, None);
250 }
251
252 #[test]
253 fn test_with_project_id() {
254 let job_service = create_job_service(MockJobService::new());
255 let query_builder =
256 Query::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");
257 assert_eq!(query_builder.project_id.unwrap(), "my-project");
258 }
259
260 #[tokio::test]
261 async fn test_run_missing_project_id() -> anyhow::Result<()> {
262 let client = BigQuery::builder()
263 .with_credentials(Anonymous::new().build())
264 .build()
265 .await?;
266 let query_builder = client.query("SELECT 1");
267 let err = query_builder
268 .send()
269 .await
270 .expect_err("should return an error when project_id is missing");
271 assert!(
272 matches!(&err, QueryError::Rpc { source } if source.is_binding()),
273 "expected Binding error for missing project ID, got {err:?}"
274 );
275 Ok(())
276 }
277
278 #[tokio::test]
279 async fn test_run_missing_project_id_force_job_path() -> anyhow::Result<()> {
280 let client = BigQuery::builder()
281 .with_credentials(Anonymous::new().build())
282 .build()
283 .await?;
284 let query_builder = client.query("SELECT 1").set_allow_large_results(true);
285 let err = query_builder
286 .send()
287 .await
288 .expect_err("should return an error when project_id is missing");
289 assert!(
290 matches!(&err, QueryError::Rpc { source } if source.is_binding()),
291 "expected Binding error for missing project ID on job path, got {err:?}"
292 );
293 Ok(())
294 }
295
296 #[tokio::test]
297 async fn test_run_until_done_missing_project_id() -> anyhow::Result<()> {
298 let client = BigQuery::builder()
299 .with_credentials(Anonymous::new().build())
300 .build()
301 .await?;
302 let query_builder = client.query("SELECT 1");
303 let err = query_builder
304 .until_done()
305 .await
306 .expect_err("should return an error when project_id is missing");
307 assert!(
308 matches!(&err, QueryError::Rpc { source } if source.is_binding()),
309 "expected Binding error for missing project ID, got {err:?}"
310 );
311 Ok(())
312 }
313
314 #[test]
315 fn test_generate_prefixed_id() {
316 let job_id = generate_prefixed_id(JOB_ID_PREFIX);
317 assert!(job_id.starts_with(JOB_ID_PREFIX), "{job_id:?}");
318 assert!(
319 Uuid::parse_str(&job_id[JOB_ID_PREFIX.len()..]).is_ok(),
320 "{job_id:?}"
321 );
322
323 let req_id = generate_prefixed_id(QUERY_REQUEST_ID_PREFIX);
324 assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX), "{req_id:?}");
325 assert!(req_id.len() <= BIGQUERY_REQ_ID_LIMIT, "{req_id:?}");
326 assert!(
327 Uuid::parse_str(&req_id[QUERY_REQUEST_ID_PREFIX.len()..]).is_ok(),
328 "{req_id:?}"
329 );
330 }
331
332 #[test]
333 fn test_generate_job_reference() {
334 let job_ref = generate_job_reference("my-project", "us-central1");
335 assert_eq!(job_ref.project_id, "my-project");
336 assert!(job_ref.job_id.starts_with(JOB_ID_PREFIX), "{job_ref:?}");
337 assert_eq!(job_ref.location.as_deref(), Some("us-central1"));
338 }
339
340 #[tokio::test]
341 async fn test_run_jobs_insert() -> TestResult {
342 let mut mock = MockJobService::new();
343 mock.expect_insert_job().returning(|req, _| {
344 let job_ref = req.job.as_ref().unwrap().job_reference.as_ref().unwrap();
345 assert!(job_ref.job_id.starts_with(JOB_ID_PREFIX), "{job_ref:?}");
346 let job_ref = JobReference::new()
347 .set_job_id("test-job")
348 .set_project_id("my-project");
349 let job = Job::new()
350 .set_job_reference(job_ref)
351 .set_status(JobStatus::new().set_state("DONE"));
352 Ok(google_cloud_gax::response::Response::from(job))
353 });
354 mock.expect_query().never();
355
356 let job_service = create_job_service(mock);
357
358 let query_builder = Query::new(job_service, "SELECT 1".to_string())
359 .with_project_id("my-project")
360 .set_allow_large_results(true);
361 let query = query_builder.send().await?;
362 assert!(query.completed, "{query:?}");
363
364 Ok(())
365 }
366
367 #[tokio::test]
368 async fn test_run_jobs_query() -> TestResult {
369 let mut mock = MockJobService::new();
370 mock.expect_query().returning(move |req, _| {
371 let req_id = &req.query_request.as_ref().unwrap().request_id;
372 assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX), "{req_id:?}");
373 assert!(req_id.len() <= BIGQUERY_REQ_ID_LIMIT, "{req_id:?}");
374 Ok(Response::from(
375 QueryResponse::new().set_query_id("some_query_id"),
376 ))
377 });
378 mock.expect_insert_job().never();
379
380 let job_service = create_job_service(mock);
381 let query_builder =
382 Query::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");
383 let query = query_builder.send().await?;
384 assert!(!query.completed, "{query:?}");
385 assert_eq!(query.metadata.query_id, "some_query_id");
386
387 Ok(())
388 }
389
390 #[test]
391 fn test_force_job_path() {
392 let job_service = create_job_service(MockJobService::new());
393 let mut query_builder = Query::new(job_service, "SELECT 1".to_string());
394 assert!(!query_builder.request.force_job_path());
395
396 query_builder = query_builder.set_allow_large_results(true);
398 assert!(query_builder.request.force_job_path());
399 }
400
401 #[test]
402 fn test_request_conversions() {
403 let req = QueryRequest::default()
404 .set_query("SELECT 1".to_string())
405 .set_dry_run(true)
406 .set_use_legacy_sql(true);
407
408 let query_request: JobsQueryRequest = req.clone().into();
409 assert_eq!(query_request.query, "SELECT 1");
410 assert!(query_request.dry_run);
411 assert_eq!(
412 query_request.use_legacy_sql,
413 Some(wkt::BoolValue::from(true))
414 );
415
416 let job_config: JobConfiguration = req.into();
417 let job_query = job_config.query.as_ref().unwrap();
418 assert_eq!(job_query.query, "SELECT 1");
419 assert_eq!(job_query.use_legacy_sql, Some(wkt::BoolValue::from(true)));
420 }
421
422 #[tokio::test(start_paused = true)]
423 async fn test_run_reissue_on_retryable_job_failed() -> TestResult {
424 let mut mock = MockJobService::new();
425 let mut seq = mockall::Sequence::new();
426 let first_request_id = std::sync::Arc::new(std::sync::Mutex::new(String::new()));
427 let first_request_id_clone = first_request_id.clone();
428
429 mock.expect_query()
430 .in_sequence(&mut seq)
431 .times(1)
432 .returning(move |req, _| {
433 let req_id = req.query_request.as_ref().unwrap().request_id.clone();
434 assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX));
435 *first_request_id_clone.lock().unwrap() = req_id;
436 let err_proto = ErrorProto::new()
437 .set_reason("backendError")
438 .set_message("temporary server issue");
439 Ok(Response::from(
440 QueryResponse::new().set_errors(vec![err_proto]),
441 ))
442 });
443
444 mock.expect_query()
445 .in_sequence(&mut seq)
446 .times(1)
447 .returning(move |req, _| {
448 let req_id = req.query_request.as_ref().unwrap().request_id.clone();
449 assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX));
450 assert_ne!(
451 req_id,
452 *first_request_id.lock().unwrap(),
453 "reissued query must generate a fresh request_id to avoid 409 duplicate ID conflicts"
454 );
455 Ok(Response::from(
456 QueryResponse::new().set_query_id("q_success"),
457 ))
458 });
459
460 let job_service = create_job_service(mock);
461 let query_builder =
462 Query::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");
463 let query = query_builder.send().await?;
464 assert_eq!(query.metadata.query_id, "q_success");
465
466 Ok(())
467 }
468
469 #[tokio::test]
470 async fn test_run_jobs_query_with_max_results() -> TestResult {
471 let mut mock = MockJobService::new();
472 mock.expect_query().returning(move |req, _| {
473 assert_eq!(
474 req.query_request.as_ref().and_then(|r| r.max_results),
475 Some(100)
476 );
477 Ok(Response::from(QueryResponse::new()))
478 });
479 let job_service = create_job_service(mock);
480 let query_builder = Query::new(job_service, "SELECT 1".to_string())
481 .with_project_id("my-project")
482 .set_max_results(100_u32);
483 let query = query_builder.send().await?;
484 assert_eq!(query.max_results, Some(100));
485
486 Ok(())
487 }
488
489 #[tokio::test]
490 async fn test_run_jobs_insert_with_max_results() -> TestResult {
491 let mut mock = MockJobService::new();
492 mock.expect_insert_job().returning(|_, _| {
493 let job_ref = JobReference::new()
494 .set_job_id("test-job")
495 .set_project_id("my-project");
496 let job = Job::new()
497 .set_job_reference(job_ref)
498 .set_status(JobStatus::new().set_state("DONE"));
499 Ok(google_cloud_gax::response::Response::from(job))
500 });
501 let job_service = create_job_service(mock);
502
503 let query_builder = Query::new(job_service, "SELECT 1".to_string())
504 .with_project_id("my-project")
505 .set_allow_large_results(true)
506 .set_max_results(50_u32);
507 let query = query_builder.send().await?;
508 assert_eq!(query.max_results, Some(50));
509
510 Ok(())
511 }
512
513 #[tokio::test]
514 async fn test_until_done_jobs_query() -> TestResult {
515 let mut mock = MockJobService::new();
516 mock.expect_query().returning(move |_, _| {
517 Ok(Response::from(
518 QueryResponse::new()
519 .set_job_complete(true)
520 .set_query_id("some_query_id"),
521 ))
522 });
523 let job_service = create_job_service(mock);
524 let query = Query::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");
525 let complete = query.until_done().await?;
526 assert_eq!(complete.metadata().query_id, "some_query_id");
527
528 Ok(())
529 }
530
531 #[tokio::test]
532 async fn test_until_done_dry_run_returns_error() -> TestResult {
533 let mut mock = MockJobService::new();
534 mock.expect_insert_job().never();
535 mock.expect_query().never();
536 let job_service = create_job_service(mock);
537 let query = Query::new(job_service, "SELECT 1".to_string())
538 .with_project_id("my-project")
539 .set_dry_run(true);
540 let err = query.until_done().await.unwrap_err();
541 assert!(
542 matches!(err, QueryError::DryRun),
543 "expected DryRun error, got {err:?}"
544 );
545
546 Ok(())
547 }
548}