use crate::builder::bigquery::Query;
use crate::error::QueryError;
use crate::query::client_builder::ClientBuilder;
use crate::query::execution::check_job_status;
use crate::query::{Query as QueryHandle, Result as QueryResult};
use google_cloud_bigquery_v2::client::JobService;
use google_cloud_bigquery_v2::model::JobReference;
use google_cloud_gax::client_builder::Result as BuilderResult;
use std::sync::Arc;
#[derive(Clone, Debug)]
pub struct BigQuery {
job_service: Arc<JobService>,
project_id: Option<String>,
}
pub(super) mod info {
pub(crate) const NAME: &str = env!("CARGO_PKG_NAME");
pub(crate) const VERSION: &str = env!("CARGO_PKG_VERSION");
}
impl BigQuery {
pub fn builder() -> ClientBuilder {
ClientBuilder::new()
}
pub(crate) async fn new(builder: ClientBuilder) -> BuilderResult<Self> {
let mut job_service_builder = JobService::builder();
if let Some(creds) = builder.config.cred {
job_service_builder = job_service_builder.with_credentials(creds);
}
if let Some(endpoint) = builder.config.endpoint {
job_service_builder = job_service_builder.with_endpoint(endpoint);
}
if let Some(universe_domain) = builder.config.universe_domain {
job_service_builder = job_service_builder.with_universe_domain(universe_domain);
}
if builder.config.tracing {
job_service_builder = job_service_builder.with_tracing();
}
let retry_policy = builder
.config
.retry_policy
.unwrap_or_else(crate::query::retry_policy::default_retry_policy);
job_service_builder = job_service_builder.with_retry_policy(retry_policy);
let backoff_policy = builder
.config
.backoff_policy
.unwrap_or_else(crate::query::retry_policy::default_backoff_policy);
job_service_builder = job_service_builder.with_backoff_policy(backoff_policy);
job_service_builder =
job_service_builder.with_retry_throttler(builder.config.retry_throttler);
job_service_builder =
job_service_builder.with_extension(gaxi::api_header::XGoogApiClient {
name: info::NAME,
version: info::VERSION,
library_type: gaxi::api_header::GCCL,
});
let job_service = Arc::new(job_service_builder.build().await?);
Ok(BigQuery {
job_service,
project_id: builder.project_id,
})
}
pub fn query<S: Into<String>>(&self, sql: S) -> Query {
let builder = Query::new(self.job_service.clone(), sql.into());
self.project_id
.as_deref()
.into_iter()
.fold(builder, |builder, project_id| {
builder.with_project_id(project_id)
})
}
pub async fn attach_job(&self, mut job_ref: JobReference) -> QueryResult<QueryHandle> {
if job_ref.project_id.is_empty()
&& let Some(proj) = &self.project_id
{
job_ref.project_id = proj.clone();
}
let req = self
.job_service
.get_job()
.set_job_id(job_ref.job_id.clone())
.set_project_id(job_ref.project_id.clone());
let req = job_ref
.location
.clone()
.into_iter()
.fold(req, |req, location| req.set_location(location));
let job = req.send().await?;
let is_query = job
.configuration
.as_ref()
.and_then(|c| c.query.as_ref())
.is_some();
if !is_query {
return Err(QueryError::UnsupportedJobType);
}
Ok(QueryHandle::from_job(
self.job_service.clone(),
check_job_status(job)?,
None,
None,
))
}
}
#[cfg(test)]
mod tests {
use super::BigQuery;
use crate::error::QueryError;
use crate::query::tests::{MockJobService, create_job_service};
use google_cloud_auth::credentials::anonymous::Builder as Anonymous;
use google_cloud_bigquery_v2::client::JobService;
use google_cloud_bigquery_v2::model::{
Job, JobConfiguration, JobConfigurationQuery, JobReference,
};
use google_cloud_gax::response::Response;
use std::sync::Arc;
impl BigQuery {
fn from_job_service(job_service: Arc<JobService>, project_id: Option<String>) -> Self {
Self {
job_service,
project_id,
}
}
}
#[tokio::test]
async fn test_bigquery_builder() -> anyhow::Result<()> {
let client = BigQuery::builder()
.with_credentials(Anonymous::new().build())
.build()
.await?;
assert!(client.project_id.is_none());
Ok(())
}
#[tokio::test]
async fn test_bigquery_builder_with_project_id() -> anyhow::Result<()> {
let client = BigQuery::builder()
.with_project_id("test-proj")
.with_credentials(Anonymous::new().build())
.build()
.await?;
assert_eq!(client.project_id.as_deref(), Some("test-proj"));
Ok(())
}
#[tokio::test]
async fn test_bigquery_query_inherits_project_id() -> anyhow::Result<()> {
let client = BigQuery::builder()
.with_project_id("test-proj")
.with_credentials(Anonymous::new().build())
.build()
.await?;
let query_builder = client.query("SELECT 1");
assert_eq!(query_builder.project_id.as_deref(), Some("test-proj"));
Ok(())
}
#[tokio::test]
async fn test_bigquery_query_without_project_id() -> anyhow::Result<()> {
let client = BigQuery::builder()
.with_credentials(Anonymous::new().build())
.build()
.await?;
let query_builder = client.query("SELECT 1");
assert!(query_builder.project_id.is_none());
Ok(())
}
#[tokio::test]
async fn test_bigquery_attach_job() -> anyhow::Result<()> {
let mut mock = MockJobService::new();
mock.expect_get_job().returning(|req, _| {
assert_eq!(req.project_id, "test-proj");
assert_eq!(req.job_id, "job_123");
let job = Job::new()
.set_job_reference(
JobReference::new()
.set_project_id("test-proj")
.set_job_id("job_123"),
)
.set_configuration(
JobConfiguration::new()
.set_query(JobConfigurationQuery::new().set_query("SELECT 1")),
);
Ok(Response::from(job))
});
let client = BigQuery::from_job_service(create_job_service(mock), None);
let job_ref = JobReference::new()
.set_project_id("test-proj")
.set_job_id("job_123");
let query = client.attach_job(job_ref).await?;
let job_ref = query
.metadata()
.job_reference
.as_ref()
.expect("job_reference should be set");
assert_eq!(job_ref.project_id, "test-proj");
assert_eq!(job_ref.job_id, "job_123");
Ok(())
}
#[tokio::test]
async fn test_bigquery_attach_job_inherits_project_id() -> anyhow::Result<()> {
let mut mock = MockJobService::new();
mock.expect_get_job().returning(|req, _| {
assert_eq!(req.project_id, "client-proj");
assert_eq!(req.job_id, "job_456");
let job = Job::new()
.set_job_reference(
JobReference::new()
.set_project_id("client-proj")
.set_job_id("job_456"),
)
.set_configuration(
JobConfiguration::new()
.set_query(JobConfigurationQuery::new().set_query("SELECT 1")),
);
Ok(Response::from(job))
});
let client =
BigQuery::from_job_service(create_job_service(mock), Some("client-proj".to_string()));
let job_ref = JobReference::new().set_job_id("job_456");
let query = client.attach_job(job_ref).await?;
let job_ref = query
.metadata()
.job_reference
.as_ref()
.expect("job_reference should be set");
assert_eq!(job_ref.project_id, "client-proj");
assert_eq!(job_ref.job_id, "job_456");
Ok(())
}
#[tokio::test]
async fn test_bigquery_attach_job_missing_project_id() -> anyhow::Result<()> {
let client = BigQuery::builder()
.with_credentials(Anonymous::new().build())
.build()
.await?;
let job_ref = JobReference::new().set_job_id("job_789");
let err = client
.attach_job(job_ref)
.await
.expect_err("should return an error when project_id is missing");
assert!(
matches!(&err, QueryError::Rpc { source } if source.is_binding()),
"expected Binding error for missing project ID, got {err:?}"
);
Ok(())
}
#[tokio::test]
async fn test_bigquery_attach_job_empty_job_id() -> anyhow::Result<()> {
let client = BigQuery::builder()
.with_project_id("client-proj")
.with_credentials(Anonymous::new().build())
.build()
.await?;
let job_ref = JobReference::new();
let err = client
.attach_job(job_ref)
.await
.expect_err("should return an error when job_id is empty");
assert!(
matches!(&err, QueryError::Rpc { source } if source.is_binding()),
"expected Binding error for empty job ID, got {err:?}"
);
Ok(())
}
#[tokio::test]
async fn test_bigquery_attach_job_unsupported_job_type() -> anyhow::Result<()> {
let mut mock = MockJobService::new();
mock.expect_get_job().returning(|_, _| {
let job = Job::new().set_configuration(JobConfiguration::new());
Ok(Response::from(job))
});
let client =
BigQuery::from_job_service(create_job_service(mock), Some("client-proj".to_string()));
let job_ref = JobReference::new().set_job_id("job_extract");
let err = client
.attach_job(job_ref)
.await
.expect_err("should return an error for non-query job");
assert!(
matches!(&err, QueryError::UnsupportedJobType),
"expected UnsupportedJobType, got {err:?}"
);
Ok(())
}
#[tokio::test]
async fn test_bigquery_attach_job_failed_job() -> anyhow::Result<()> {
use google_cloud_bigquery_v2::model::{ErrorProto, JobConfigurationQuery, JobStatus};
let mut mock = MockJobService::new();
mock.expect_get_job().returning(|_, _| {
let err_proto = ErrorProto::new()
.set_reason("invalidQuery")
.set_message("Syntax error");
let job = Job::new()
.set_configuration(
JobConfiguration::new()
.set_query(JobConfigurationQuery::new().set_query("SELECT * FROM")),
)
.set_status(
JobStatus::new()
.set_state("DONE")
.set_error_result(err_proto.clone())
.set_errors(vec![err_proto]),
);
Ok(Response::from(job))
});
let client =
BigQuery::from_job_service(create_job_service(mock), Some("client-proj".to_string()));
let job_ref = JobReference::new().set_job_id("job_failed");
let err = client
.attach_job(job_ref)
.await
.expect_err("should return an error for failed query job");
assert!(
matches!(&err, QueryError::JobFailed { errors } if errors.len() == 1 && errors[0].reason == "invalidQuery"),
"expected JobFailed, got {err:?}"
);
Ok(())
}
#[tokio::test]
async fn test_bigquery_calls_send_veneer_header_not_gapic() -> anyhow::Result<()> {
use httptest::{Expectation, Server, all_of, matchers::*, responders::*};
use serde_json::json;
let server = Server::run();
server.expect(
Expectation::matching(all_of![
request::method_path("GET", "/bigquery/v2/projects/test-proj/jobs/job_123"),
request::headers(contains((
"x-goog-api-client",
matches(format!("gccl/{}", env!("CARGO_PKG_VERSION"))),
))),
not(request::headers(contains((
"x-goog-api-client",
matches("gapic/"),
)))),
])
.respond_with(json_encoded(json!({
"jobReference": {
"projectId": "test-proj",
"jobId": "job_123"
},
"configuration": {
"query": {
"query": "SELECT 1"
}
},
"status": {
"state": "DONE"
}
}))),
);
let client = BigQuery::builder()
.with_endpoint(server.url_str(""))
.with_credentials(Anonymous::new().build())
.with_project_id("test-proj")
.build()
.await?;
let job_ref = JobReference::new().set_job_id("job_123");
let query = client.attach_job(job_ref).await?;
let metadata = query.metadata();
assert_eq!(
metadata.job_reference.as_ref().map(|j| j.job_id.as_str()),
Some("job_123")
);
Ok(())
}
}