google_cloud_bigquery/query/
client.rs1use crate::builder::bigquery::Query;
16use crate::error::QueryError;
17use crate::query::client_builder::ClientBuilder;
18use crate::query::execution::check_job_status;
19use crate::query::{Query as QueryHandle, Result as QueryResult};
20use google_cloud_bigquery_v2::client::JobService;
21use google_cloud_bigquery_v2::model::JobReference;
22use google_cloud_gax::client_builder::Result as BuilderResult;
23use std::sync::Arc;
24
25#[derive(Clone, Debug)]
66pub struct BigQuery {
67 job_service: Arc<JobService>,
68 project_id: Option<String>,
69}
70
71pub(super) mod info {
72 pub(crate) const NAME: &str = env!("CARGO_PKG_NAME");
73 pub(crate) const VERSION: &str = env!("CARGO_PKG_VERSION");
74}
75
76impl BigQuery {
77 pub fn builder() -> ClientBuilder {
90 ClientBuilder::new()
91 }
92
93 pub(crate) async fn new(builder: ClientBuilder) -> BuilderResult<Self> {
94 let mut job_service_builder = JobService::builder();
95 if let Some(creds) = builder.config.cred {
96 job_service_builder = job_service_builder.with_credentials(creds);
97 }
98 if let Some(endpoint) = builder.config.endpoint {
99 job_service_builder = job_service_builder.with_endpoint(endpoint);
100 }
101 if let Some(universe_domain) = builder.config.universe_domain {
102 job_service_builder = job_service_builder.with_universe_domain(universe_domain);
103 }
104 if builder.config.tracing {
105 job_service_builder = job_service_builder.with_tracing();
106 }
107 let retry_policy = builder
108 .config
109 .retry_policy
110 .unwrap_or_else(crate::query::retry_policy::default_retry_policy);
111 job_service_builder = job_service_builder.with_retry_policy(retry_policy);
112
113 let backoff_policy = builder
114 .config
115 .backoff_policy
116 .unwrap_or_else(crate::query::retry_policy::default_backoff_policy);
117 job_service_builder = job_service_builder.with_backoff_policy(backoff_policy);
118 job_service_builder =
119 job_service_builder.with_retry_throttler(builder.config.retry_throttler);
120
121 job_service_builder =
122 job_service_builder.with_extension(gaxi::api_header::XGoogApiClient {
123 name: info::NAME,
124 version: info::VERSION,
125 library_type: gaxi::api_header::GCCL,
126 });
127
128 let job_service = Arc::new(job_service_builder.build().await?);
129
130 Ok(BigQuery {
131 job_service,
132 project_id: builder.project_id,
133 })
134 }
135
136 pub fn query<S: Into<String>>(&self, sql: S) -> Query {
174 let builder = Query::new(self.job_service.clone(), sql.into());
175 self.project_id
176 .as_deref()
177 .into_iter()
178 .fold(builder, |builder, project_id| {
179 builder.with_project_id(project_id)
180 })
181 }
182
183 pub async fn attach_job(&self, mut job_ref: JobReference) -> QueryResult<QueryHandle> {
211 if job_ref.project_id.is_empty()
212 && let Some(proj) = &self.project_id
213 {
214 job_ref.project_id = proj.clone();
215 }
216
217 let req = self
218 .job_service
219 .get_job()
220 .set_job_id(job_ref.job_id.clone())
221 .set_project_id(job_ref.project_id.clone());
222
223 let req = job_ref
224 .location
225 .clone()
226 .into_iter()
227 .fold(req, |req, location| req.set_location(location));
228
229 let job = req.send().await?;
230
231 let is_query = job
232 .configuration
233 .as_ref()
234 .and_then(|c| c.query.as_ref())
235 .is_some();
236 if !is_query {
237 return Err(QueryError::UnsupportedJobType);
238 }
239
240 Ok(QueryHandle::from_job(
241 self.job_service.clone(),
242 check_job_status(job)?,
243 None,
244 None,
245 ))
246 }
247}
248
249#[cfg(test)]
250mod tests {
251 use super::BigQuery;
252 use crate::error::QueryError;
253 use crate::query::tests::{MockJobService, create_job_service};
254 use google_cloud_auth::credentials::anonymous::Builder as Anonymous;
255 use google_cloud_bigquery_v2::client::JobService;
256 use google_cloud_bigquery_v2::model::{
257 Job, JobConfiguration, JobConfigurationQuery, JobReference,
258 };
259 use google_cloud_gax::response::Response;
260 use std::sync::Arc;
261
262 impl BigQuery {
263 fn from_job_service(job_service: Arc<JobService>, project_id: Option<String>) -> Self {
264 Self {
265 job_service,
266 project_id,
267 }
268 }
269 }
270
271 #[tokio::test]
272 async fn test_bigquery_builder() -> anyhow::Result<()> {
273 let client = BigQuery::builder()
274 .with_credentials(Anonymous::new().build())
275 .build()
276 .await?;
277 assert!(client.project_id.is_none());
278 Ok(())
279 }
280
281 #[tokio::test]
282 async fn test_bigquery_builder_with_project_id() -> anyhow::Result<()> {
283 let client = BigQuery::builder()
284 .with_project_id("test-proj")
285 .with_credentials(Anonymous::new().build())
286 .build()
287 .await?;
288 assert_eq!(client.project_id.as_deref(), Some("test-proj"));
289 Ok(())
290 }
291
292 #[tokio::test]
293 async fn test_bigquery_query_inherits_project_id() -> anyhow::Result<()> {
294 let client = BigQuery::builder()
295 .with_project_id("test-proj")
296 .with_credentials(Anonymous::new().build())
297 .build()
298 .await?;
299 let query_builder = client.query("SELECT 1");
300 assert_eq!(query_builder.project_id.as_deref(), Some("test-proj"));
301 Ok(())
302 }
303
304 #[tokio::test]
305 async fn test_bigquery_query_without_project_id() -> anyhow::Result<()> {
306 let client = BigQuery::builder()
307 .with_credentials(Anonymous::new().build())
308 .build()
309 .await?;
310 let query_builder = client.query("SELECT 1");
311 assert!(query_builder.project_id.is_none());
312 Ok(())
313 }
314
315 #[tokio::test]
316 async fn test_bigquery_attach_job() -> anyhow::Result<()> {
317 let mut mock = MockJobService::new();
318 mock.expect_get_job().returning(|req, _| {
319 assert_eq!(req.project_id, "test-proj");
320 assert_eq!(req.job_id, "job_123");
321 let job = Job::new()
322 .set_job_reference(
323 JobReference::new()
324 .set_project_id("test-proj")
325 .set_job_id("job_123"),
326 )
327 .set_configuration(
328 JobConfiguration::new()
329 .set_query(JobConfigurationQuery::new().set_query("SELECT 1")),
330 );
331 Ok(Response::from(job))
332 });
333 let client = BigQuery::from_job_service(create_job_service(mock), None);
334 let job_ref = JobReference::new()
335 .set_project_id("test-proj")
336 .set_job_id("job_123");
337 let query = client.attach_job(job_ref).await?;
338 let job_ref = query
339 .metadata()
340 .job_reference
341 .as_ref()
342 .expect("job_reference should be set");
343 assert_eq!(job_ref.project_id, "test-proj");
344 assert_eq!(job_ref.job_id, "job_123");
345 Ok(())
346 }
347
348 #[tokio::test]
349 async fn test_bigquery_attach_job_inherits_project_id() -> anyhow::Result<()> {
350 let mut mock = MockJobService::new();
351 mock.expect_get_job().returning(|req, _| {
352 assert_eq!(req.project_id, "client-proj");
353 assert_eq!(req.job_id, "job_456");
354 let job = Job::new()
355 .set_job_reference(
356 JobReference::new()
357 .set_project_id("client-proj")
358 .set_job_id("job_456"),
359 )
360 .set_configuration(
361 JobConfiguration::new()
362 .set_query(JobConfigurationQuery::new().set_query("SELECT 1")),
363 );
364 Ok(Response::from(job))
365 });
366 let client =
367 BigQuery::from_job_service(create_job_service(mock), Some("client-proj".to_string()));
368 let job_ref = JobReference::new().set_job_id("job_456");
369 let query = client.attach_job(job_ref).await?;
370 let job_ref = query
371 .metadata()
372 .job_reference
373 .as_ref()
374 .expect("job_reference should be set");
375 assert_eq!(job_ref.project_id, "client-proj");
376 assert_eq!(job_ref.job_id, "job_456");
377 Ok(())
378 }
379
380 #[tokio::test]
381 async fn test_bigquery_attach_job_missing_project_id() -> anyhow::Result<()> {
382 let client = BigQuery::builder()
383 .with_credentials(Anonymous::new().build())
384 .build()
385 .await?;
386 let job_ref = JobReference::new().set_job_id("job_789");
387 let err = client
388 .attach_job(job_ref)
389 .await
390 .expect_err("should return an error when project_id is missing");
391 assert!(
392 matches!(&err, QueryError::Rpc { source } if source.is_binding()),
393 "expected Binding error for missing project ID, got {err:?}"
394 );
395 Ok(())
396 }
397
398 #[tokio::test]
399 async fn test_bigquery_attach_job_empty_job_id() -> anyhow::Result<()> {
400 let client = BigQuery::builder()
401 .with_project_id("client-proj")
402 .with_credentials(Anonymous::new().build())
403 .build()
404 .await?;
405 let job_ref = JobReference::new();
406 let err = client
407 .attach_job(job_ref)
408 .await
409 .expect_err("should return an error when job_id is empty");
410 assert!(
411 matches!(&err, QueryError::Rpc { source } if source.is_binding()),
412 "expected Binding error for empty job ID, got {err:?}"
413 );
414 Ok(())
415 }
416
417 #[tokio::test]
418 async fn test_bigquery_attach_job_unsupported_job_type() -> anyhow::Result<()> {
419 let mut mock = MockJobService::new();
420 mock.expect_get_job().returning(|_, _| {
421 let job = Job::new().set_configuration(JobConfiguration::new());
422 Ok(Response::from(job))
423 });
424 let client =
425 BigQuery::from_job_service(create_job_service(mock), Some("client-proj".to_string()));
426 let job_ref = JobReference::new().set_job_id("job_extract");
427 let err = client
428 .attach_job(job_ref)
429 .await
430 .expect_err("should return an error for non-query job");
431 assert!(
432 matches!(&err, QueryError::UnsupportedJobType),
433 "expected UnsupportedJobType, got {err:?}"
434 );
435 Ok(())
436 }
437
438 #[tokio::test]
439 async fn test_bigquery_attach_job_failed_job() -> anyhow::Result<()> {
440 use google_cloud_bigquery_v2::model::{ErrorProto, JobConfigurationQuery, JobStatus};
441
442 let mut mock = MockJobService::new();
443 mock.expect_get_job().returning(|_, _| {
444 let err_proto = ErrorProto::new()
445 .set_reason("invalidQuery")
446 .set_message("Syntax error");
447 let job = Job::new()
448 .set_configuration(
449 JobConfiguration::new()
450 .set_query(JobConfigurationQuery::new().set_query("SELECT * FROM")),
451 )
452 .set_status(
453 JobStatus::new()
454 .set_state("DONE")
455 .set_error_result(err_proto.clone())
456 .set_errors(vec![err_proto]),
457 );
458 Ok(Response::from(job))
459 });
460 let client =
461 BigQuery::from_job_service(create_job_service(mock), Some("client-proj".to_string()));
462 let job_ref = JobReference::new().set_job_id("job_failed");
463 let err = client
464 .attach_job(job_ref)
465 .await
466 .expect_err("should return an error for failed query job");
467 assert!(
468 matches!(&err, QueryError::JobFailed { errors } if errors.len() == 1 && errors[0].reason == "invalidQuery"),
469 "expected JobFailed, got {err:?}"
470 );
471 Ok(())
472 }
473
474 #[tokio::test]
475 async fn test_bigquery_calls_send_veneer_header_not_gapic() -> anyhow::Result<()> {
476 use httptest::{Expectation, Server, all_of, matchers::*, responders::*};
477 use serde_json::json;
478
479 let server = Server::run();
480 server.expect(
481 Expectation::matching(all_of![
482 request::method_path("GET", "/bigquery/v2/projects/test-proj/jobs/job_123"),
483 request::headers(contains((
484 "x-goog-api-client",
485 matches(format!("gccl/{}", env!("CARGO_PKG_VERSION"))),
486 ))),
487 not(request::headers(contains((
488 "x-goog-api-client",
489 matches("gapic/"),
490 )))),
491 ])
492 .respond_with(json_encoded(json!({
493 "jobReference": {
494 "projectId": "test-proj",
495 "jobId": "job_123"
496 },
497 "configuration": {
498 "query": {
499 "query": "SELECT 1"
500 }
501 },
502 "status": {
503 "state": "DONE"
504 }
505 }))),
506 );
507
508 let client = BigQuery::builder()
509 .with_endpoint(server.url_str(""))
510 .with_credentials(Anonymous::new().build())
511 .with_project_id("test-proj")
512 .build()
513 .await?;
514
515 let job_ref = JobReference::new().set_job_id("job_123");
516 let query = client.attach_job(job_ref).await?;
517 let metadata = query.metadata();
518 assert_eq!(
519 metadata.job_reference.as_ref().map(|j| j.job_id.as_str()),
520 Some("job_123")
521 );
522
523 Ok(())
524 }
525}