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::error::Error as GaxError;
229 use google_cloud_gax::error::rpc::Status;
230 use google_cloud_gax::response::Response;
231
232 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 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}