From e3508611ee081e3284db1f7c2d790d0dbca8a5ad Mon Sep 17 00:00:00 2001 From: Darren Bolduc Date: Fri, 31 Jul 2026 21:14:48 +0000 Subject: [PATCH 1/2] impl(bigquery): validate table format for writes --- Cargo.lock | 1 + src/bigquery-write/Cargo.toml | 1 + .../src/arrow/writer_builder.rs | 70 +++++++++++++++++-- 3 files changed, 66 insertions(+), 6 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 9ceb88c055..3ef42c4d40 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2197,6 +2197,7 @@ dependencies = [ "serde", "serde_json", "serde_with", + "test-case", "thiserror", "tokio", "tokio-stream", diff --git a/src/bigquery-write/Cargo.toml b/src/bigquery-write/Cargo.toml index b9f5caad21..cbb51afa7c 100644 --- a/src/bigquery-write/Cargo.toml +++ b/src/bigquery-write/Cargo.toml @@ -55,6 +55,7 @@ gaxi = { workspace = true, features = ["_internal-common" [dev-dependencies] anyhow.workspace = true +test-case.workspace = true bigquery-write-grpc-mock = { path = "grpc-mock" } google-cloud-auth.workspace = true google-cloud-bigquery-write = { path = ".", features = ["default-rustls-provider"] } diff --git a/src/bigquery-write/src/arrow/writer_builder.rs b/src/bigquery-write/src/arrow/writer_builder.rs index e0c88abf1e..14cbce3fee 100644 --- a/src/bigquery-write/src/arrow/writer_builder.rs +++ b/src/bigquery-write/src/arrow/writer_builder.rs @@ -12,10 +12,13 @@ // See the License for the specific language governing permissions and // limitations under the License. -use crate::Result; use crate::arrow::DefaultWriter; use crate::model::ArrowSchema; use crate::transport::Transport; +use crate::{Error, Result}; +use gaxi::path_parameter::{PathMismatchBuilder, try_match}; +use gaxi::routing_parameter::Segment; +use google_cloud_gax::error::binding::BindingError; use std::sync::Arc; /// A builder to create a stream writer @@ -31,31 +34,86 @@ impl WriterBuilder { Self { inner, schema } } - // TODO(#6224) - add an example showing the format of `table` /// Create a writer for the [default stream] for the given table. /// /// [default stream]: https://docs.cloud.google.com/bigquery/docs/write-api#default_stream pub fn default>(self, table: T) -> Result { - // TODO(#6249) - validate table resource format - let mut write_stream = table.into(); + let table = table.into(); + validate_table(table.as_str())?; + let mut write_stream = table; write_stream.push_str("/streams/_default"); Ok(DefaultWriter::new(self.inner, write_stream, self.schema)) } } +fn validate_table(table: &str) -> Result<()> { + try_match( + Some(table), + &[ + Segment::Literal("projects/"), + Segment::SingleWildcard, + Segment::Literal("/datasets/"), + Segment::SingleWildcard, + Segment::Literal("/tables/"), + Segment::SingleWildcard, + ], + ) + .ok_or_else(|| { + let builder = PathMismatchBuilder::default().maybe_add( + Some(table), + &[ + Segment::Literal("projects/"), + Segment::SingleWildcard, + Segment::Literal("/datasets/"), + Segment::SingleWildcard, + Segment::Literal("/tables/"), + Segment::SingleWildcard, + ], + "table", + "projects/*/datasets/*/tables/*", + ); + Error::binding(BindingError { + paths: vec![builder.build()], + }) + }) + .map(|_| ()) +} + #[cfg(test)] mod tests { use super::*; use crate::transport::tests::test_transport; + use test_case::test_case; #[tokio::test] async fn default() -> anyhow::Result<()> { let transport = Arc::new(test_transport("http://ignored:1".to_string()).await?); let schema = ArrowSchema::new().set_serialized_schema("test"); let builder = WriterBuilder::new(transport, schema.clone()); - let writer = builder.default("projects/p/tables/t")?; - assert_eq!(writer.write_stream, "projects/p/tables/t/streams/_default"); + let writer = builder.default("projects/p/datasets/d/tables/t")?; + assert_eq!( + writer.write_stream, + "projects/p/datasets/d/tables/t/streams/_default" + ); assert_eq!(writer.schema, schema); Ok(()) } + + #[test_case("projects/p")] + #[test_case("projects/p/tables/t")] + #[test_case("projects/p/datasets/d/tables/")] + #[test_case("projects/p/instances/i/tables/t")] + #[test_case("projects/p/datasets/d/tables/t/streams")] + #[test_case("projects/p/datasets/d/tables/t/streams/_default")] + #[tokio::test] + async fn bad_table_format(table: &str) -> anyhow::Result<()> { + let transport = Arc::new(test_transport("http://ignored:1".to_string()).await?); + let schema = ArrowSchema::new().set_serialized_schema("test"); + let builder = WriterBuilder::new(transport, schema.clone()); + let err = builder + .default(table) + .expect_err("should fail locally on bad format"); + assert!(err.is_binding(), "{err:?}"); + Ok(()) + } } From eba8ad016e8b66701ae6cd75b7b1485c421513c2 Mon Sep 17 00:00:00 2001 From: Darren Bolduc Date: Fri, 31 Jul 2026 21:21:11 +0000 Subject: [PATCH 2/2] clean up; ty bot --- .../src/arrow/writer_builder.rs | 49 ++++++++----------- 1 file changed, 20 insertions(+), 29 deletions(-) diff --git a/src/bigquery-write/src/arrow/writer_builder.rs b/src/bigquery-write/src/arrow/writer_builder.rs index 14cbce3fee..78dd687763 100644 --- a/src/bigquery-write/src/arrow/writer_builder.rs +++ b/src/bigquery-write/src/arrow/writer_builder.rs @@ -47,36 +47,27 @@ impl WriterBuilder { } fn validate_table(table: &str) -> Result<()> { - try_match( - Some(table), - &[ - Segment::Literal("projects/"), - Segment::SingleWildcard, - Segment::Literal("/datasets/"), - Segment::SingleWildcard, - Segment::Literal("/tables/"), - Segment::SingleWildcard, - ], - ) - .ok_or_else(|| { - let builder = PathMismatchBuilder::default().maybe_add( - Some(table), - &[ - Segment::Literal("projects/"), - Segment::SingleWildcard, - Segment::Literal("/datasets/"), - Segment::SingleWildcard, - Segment::Literal("/tables/"), - Segment::SingleWildcard, - ], - "table", - "projects/*/datasets/*/tables/*", - ); - Error::binding(BindingError { - paths: vec![builder.build()], + let segments = &[ + Segment::Literal("projects/"), + Segment::SingleWildcard, + Segment::Literal("/datasets/"), + Segment::SingleWildcard, + Segment::Literal("/tables/"), + Segment::SingleWildcard, + ]; + try_match(Some(table), segments) + .ok_or_else(|| { + let builder = PathMismatchBuilder::default().maybe_add( + Some(table), + segments, + "table", + "projects/*/datasets/*/tables/*", + ); + Error::binding(BindingError { + paths: vec![builder.build()], + }) }) - }) - .map(|_| ()) + .map(|_| ()) } #[cfg(test)]