diff --git a/Cargo.lock b/Cargo.lock index 3d264e31d3..dddfbf8be9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2021,6 +2021,7 @@ dependencies = [ "thiserror", "time", "tokio", + "uuid", ] [[package]] diff --git a/src/bigquery/Cargo.toml b/src/bigquery/Cargo.toml index e6c923a5bf..c275a9b3bc 100644 --- a/src/bigquery/Cargo.toml +++ b/src/bigquery/Cargo.toml @@ -41,6 +41,7 @@ wkt.workspace = true gaxi = { workspace = true, features = ["_internal-common", "_internal-grpc-client", "_internal-http-client"] } tokio = { workspace = true, features = ["time"] } rust_decimal.workspace = true +uuid.workspace = true [dev-dependencies] anyhow.workspace = true diff --git a/src/bigquery/src/query/execution.rs b/src/bigquery/src/query/execution.rs index c18e3cc3d4..527b42f3a1 100644 --- a/src/bigquery/src/query/execution.rs +++ b/src/bigquery/src/query/execution.rs @@ -16,6 +16,7 @@ use crate::error::QueryError; use crate::query::{Query, Result}; use google_cloud_bigquery_v2::client::JobService; use google_cloud_bigquery_v2::model::{InsertJobRequest, PostQueryRequest}; +use google_cloud_gax::options::RequestOptionsBuilder as _; use std::sync::Arc; pub(crate) struct PostQueryExecutor { @@ -40,6 +41,9 @@ impl PostQueryExecutor { let res = self .job_service .query() + // requests to jobs.query are idempotent because every request + // carries a generated request_id. + .with_idempotency(true) .with_request(self.request) .send() .await?; @@ -97,6 +101,9 @@ impl InsertJobExecutor { .job_service .insert_job() .with_request(self.request) + // jobs.insert is idempotent because every request + // carries a generated UUID job_id. + .with_idempotency(true) .send() .await?; diff --git a/src/bigquery/src/query/run_query.rs b/src/bigquery/src/query/run_query.rs index f2479e06b2..b43c6899fb 100644 --- a/src/bigquery/src/query/run_query.rs +++ b/src/bigquery/src/query/run_query.rs @@ -19,9 +19,13 @@ use crate::query::{Query, Result}; use google_cloud_bigquery_v2::client::JobService; use google_cloud_bigquery_v2::model::query_request::JobCreationMode; use google_cloud_bigquery_v2::model::{ - InsertJobRequest, Job, JobConfiguration, PostQueryRequest, QueryRequest, + InsertJobRequest, Job, JobConfiguration, JobReference, PostQueryRequest, QueryRequest, }; use std::sync::Arc; +use uuid::Uuid; + +pub(crate) const JOB_ID_PREFIX: &str = "job_"; +pub(crate) const QUERY_REQUEST_ID_PREFIX: &str = "req_"; /// A unified request builder for configuring and running a SQL query. /// It automatically routes to either `jobs.query` (fast path) or `jobs.insert` (job path) @@ -67,8 +71,11 @@ impl RunQuery { if self.request.force_job_path() { // Route to jobs.insert + let job_ref = generate_job_reference(&project_id, &self.request.location); let job_config: JobConfiguration = self.request.into(); - let job = Job::new().set_configuration(job_config); + let job = Job::new() + .set_configuration(job_config) + .set_job_reference(job_ref); let req = InsertJobRequest::new() .set_job(job) .set_project_id(project_id); @@ -77,12 +84,15 @@ impl RunQuery { .execute() .await } else { + let query_request_id = generate_prefixed_id(QUERY_REQUEST_ID_PREFIX); // Route to jobs.query let query_request: QueryRequest = self.request.into(); - let query_request = query_request.set_format_options( - google_cloud_bigquery_v2::model::DataFormatOptions::new() - .set_use_int64_timestamp(true), - ); + let query_request = query_request + .set_format_options( + google_cloud_bigquery_v2::model::DataFormatOptions::new() + .set_use_int64_timestamp(true), + ) + .set_request_id(query_request_id); let req = PostQueryRequest::new() .set_project_id(project_id) .set_query_request(query_request); @@ -94,6 +104,23 @@ impl RunQuery { } } +fn generate_job_reference(project_id: &str, location: &str) -> JobReference { + let job_id = generate_prefixed_id(JOB_ID_PREFIX); + let mut job_ref = JobReference::new() + .set_project_id(project_id.to_string()) + .set_job_id(job_id); + + if !location.is_empty() { + job_ref = job_ref.set_location(location.to_string()); + } + + job_ref +} + +fn generate_prefixed_id(prefix: &str) -> String { + format!("{prefix}{}", Uuid::new_v4().simple()) +} + include!("../generated/run_query_builder.rs"); #[cfg(test)] @@ -105,6 +132,7 @@ mod tests { Job, JobConfiguration, JobReference, JobStatus, QueryRequest, QueryResponse, }; use google_cloud_gax::response::Response; + use uuid::Uuid; type TestResult = anyhow::Result<()>; @@ -141,10 +169,38 @@ mod tests { assert!(matches!(res, Err(QueryError::MissingProjectId))); } + #[test] + fn test_generate_prefixed_id() { + let job_id = generate_prefixed_id(JOB_ID_PREFIX); + assert!(job_id.starts_with(JOB_ID_PREFIX), "{job_id:?}"); + assert!( + Uuid::parse_str(&job_id[JOB_ID_PREFIX.len()..]).is_ok(), + "{job_id:?}" + ); + + let req_id = generate_prefixed_id(QUERY_REQUEST_ID_PREFIX); + assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX), "{req_id:?}"); + assert!(req_id.len() <= 36, "{req_id:?}"); // bigquery limits request id to 36 characters + assert!( + Uuid::parse_str(&req_id[QUERY_REQUEST_ID_PREFIX.len()..]).is_ok(), + "{req_id:?}" + ); + } + + #[test] + fn test_generate_job_reference() { + let job_ref = generate_job_reference("my-project", "us-central1"); + assert_eq!(job_ref.project_id, "my-project"); + assert!(job_ref.job_id.starts_with(JOB_ID_PREFIX), "{job_ref:?}"); + assert_eq!(job_ref.location.as_deref(), Some("us-central1")); + } + #[tokio::test] async fn test_run_jobs_insert() -> TestResult { let mut mock = MockJobService::new(); - mock.expect_insert_job().returning(|_, _| { + mock.expect_insert_job().returning(|req, _| { + let job_ref = req.job.as_ref().unwrap().job_reference.as_ref().unwrap(); + assert!(job_ref.job_id.starts_with(JOB_ID_PREFIX), "{job_ref:?}"); let job_ref = JobReference::new() .set_job_id("test-job") .set_project_id("my-project"); @@ -167,8 +223,12 @@ mod tests { #[tokio::test] async fn test_run_jobs_query() -> TestResult { let mut mock = MockJobService::new(); - mock.expect_query() - .returning(move |_, _| Ok(Response::from(QueryResponse::new()))); + mock.expect_query().returning(move |req, _| { + let req_id = &req.query_request.as_ref().unwrap().request_id; + assert!(req_id.starts_with(QUERY_REQUEST_ID_PREFIX), "{req_id:?}"); + assert!(req_id.len() <= 36, "{req_id:?}"); // bigquery limits request id to 36 characters + Ok(Response::from(QueryResponse::new())) + }); let job_service = create_job_service(mock); let run_query = RunQuery::new(job_service, "SELECT 1".to_string()).with_project_id("my-project");