use crate::builder::bigquery::Query;
use crate::error::QueryError;
use crate::query::client_builder::ClientBuilder;
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>,
}
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();
}
if let Some(retry_policy) = builder.config.retry_policy {
job_service_builder = job_service_builder.with_retry_policy(retry_policy);
}
if let Some(backoff_policy) = builder.config.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);
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(),
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(())
}
}