From 3c29fa15a8e9e3d24eec147b4b4e50ff33f9bf18 Mon Sep 17 00:00:00 2001 From: naivedogger Date: Tue, 8 Sep 2026 11:50:00 +0800 Subject: [PATCH 1/2] [c++] Avoid redundant copies during row type resolution Borrow unchanged STRING and BYTES values from the input row while preserving type conversions and validation. Add regression tests for borrowed storage, mixed conversions, nested rows, and invalid values. Fixes apache/fluss#4239 --- fluss-rust/bindings/cpp/src/types.rs | 168 +++++++++++++++++++++++++-- 1 file changed, 159 insertions(+), 9 deletions(-) diff --git a/fluss-rust/bindings/cpp/src/types.rs b/fluss-rust/bindings/cpp/src/types.rs index 84b88b98406..f3d31657193 100644 --- a/fluss-rust/bindings/cpp/src/types.rs +++ b/fluss-rust/bindings/cpp/src/types.rs @@ -457,11 +457,12 @@ pub fn core_database_info_to_ffi(info: &fcore::metadata::DatabaseInfo) -> ffi::F /// Resolve types in a GenericRow using schema metadata. /// Narrows Int32 → Int8/Int16, parses decimal strings, etc. -/// Used by both AppendWriter and UpsertWriter. -pub fn resolve_row_types( - row: &fcore::row::GenericRow<'_>, +/// Unchanged STRING and BYTES values borrow the input row's storage. +/// Used by append, upsert, delete, lookup, and prefix lookup. +pub fn resolve_row_types<'a>( + row: &'a fcore::row::GenericRow<'_>, schema: Option<&fcore::metadata::Schema>, -) -> Result> { +) -> Result> { let mut out = fcore::row::GenericRow::new(row.values.len()); for (idx, datum) in row.values.iter().enumerate() { @@ -477,11 +478,11 @@ pub fn resolve_row_types( /// Resolve a single datum against its (optional) target column type, recursing /// into nested ROW values. Narrows Int32 → Int8/Int16, parses decimal strings, /// and leaves already-typed ARRAY/MAP binaries (built by the writers) untouched. -fn resolve_datum( - datum: &fcore::row::Datum<'_>, +fn resolve_datum<'a>( + datum: &'a fcore::row::Datum<'_>, target: Option<&fcore::metadata::DataType>, idx: usize, -) -> Result> { +) -> Result> { Ok(match datum { Datum::Null => Datum::Null, Datum::Bool(v) => Datum::Bool(*v), @@ -509,9 +510,9 @@ fn resolve_datum( .map_err(|e| anyhow!("Column {idx}: {e}"))?; Datum::Decimal(decimal) } - _ => Datum::String(Cow::Owned(cow.to_string())), + _ => Datum::String(Cow::Borrowed(cow.as_ref())), }, - Datum::Blob(cow) => Datum::Blob(Cow::Owned(cow.to_vec())), + Datum::Blob(cow) => Datum::Blob(Cow::Borrowed(cow.as_ref())), Datum::Decimal(d) => Datum::Decimal(d.clone()), Datum::Date(d) => Datum::Date(*d), Datum::Time(t) => Datum::Time(*t), @@ -664,3 +665,152 @@ pub fn core_scan_batches_to_ffi( batches: ffi_batches, }) } + +#[cfg(test)] +mod tests { + use super::resolve_row_types; + use fluss::metadata::{DataField, DataType, DataTypes, DecimalType, RowType, Schema}; + use fluss::row::{Datum, Decimal, GenericRow}; + use std::borrow::Cow; + + fn assert_borrowed_values(input: &[Datum<'_>], resolved: &[Datum<'_>]) { + assert_eq!(resolved, input); + for (input, resolved) in input.iter().zip(resolved) { + match (input, resolved) { + (Datum::String(input), Datum::String(Cow::Borrowed(resolved))) => { + assert_eq!(input.as_ptr(), resolved.as_ptr()); + } + (Datum::Blob(input), Datum::Blob(Cow::Borrowed(resolved))) => { + assert_eq!(input.as_ptr(), resolved.as_ptr()); + } + _ => panic!("expected borrowed STRING or BYTES, got {resolved:?}"), + } + } + } + + #[test] + fn test_resolve_row_types_borrows_strings_and_bytes() { + // Include both setter-owned values and values already borrowing external storage. + let string = String::from("borrowed string"); + let bytes = vec![0, 1, 255]; + let row = GenericRow { + values: vec![ + Datum::String(Cow::Owned(String::from("owned string"))), + Datum::Blob(Cow::Owned(vec![2, 3, 255])), + Datum::String(Cow::Borrowed(&string)), + Datum::Blob(Cow::Borrowed(&bytes)), + Datum::String(Cow::Owned(String::new())), + Datum::Blob(Cow::Owned(Vec::new())), + ], + }; + let schema = Schema::builder() + .column("owned_string", DataTypes::string()) + .column("owned_bytes", DataTypes::bytes()) + .column("borrowed_string", DataTypes::string()) + .column("borrowed_bytes", DataTypes::bytes()) + .column("empty_string", DataTypes::string()) + .column("empty_bytes", DataTypes::bytes()) + .build() + .unwrap(); + + for schema in [Some(&schema), None] { + let resolved = resolve_row_types(&row, schema).unwrap(); + assert_borrowed_values(&row.values, &resolved.values); + } + } + + #[test] + fn test_resolve_row_types_preserves_conversions() { + let schema = Schema::builder() + .column("tiny", DataTypes::tinyint()) + .column("small", DataTypes::smallint()) + .column( + "decimal", + DataType::Decimal(DecimalType::new(5, 2).unwrap()), + ) + .column("string", DataTypes::string()) + .column("bytes", DataTypes::bytes()) + .column("nullable", DataTypes::string()) + .build() + .unwrap(); + let row = GenericRow { + values: vec![ + Datum::Int32(127), + Datum::Int32(-32768), + Datum::String(Cow::Owned(String::from("123.45"))), + Datum::String(Cow::Owned(String::from("unchanged"))), + Datum::Blob(Cow::Owned(vec![1, 2, 3])), + Datum::Null, + ], + }; + let resolved = resolve_row_types(&row, Some(&schema)).unwrap(); + assert_eq!(resolved.values[0], Datum::Int8(127)); + assert_eq!(resolved.values[1], Datum::Int16(-32768)); + assert_eq!( + resolved.values[2], + Datum::Decimal(Decimal::from_unscaled_long(12345, 5, 2).unwrap()) + ); + assert_borrowed_values(&row.values[3..5], &resolved.values[3..5]); + assert_eq!(resolved.values[5], Datum::Null); + } + + #[test] + fn test_resolve_row_types_borrows_nested_values() { + let schema = Schema::builder() + .column( + "nested", + DataType::Row(RowType::new(vec![ + DataField::new("string", DataTypes::string(), None), + DataField::new("bytes", DataTypes::bytes(), None), + ])), + ) + .build() + .unwrap(); + let row = GenericRow { + values: vec![Datum::Row(Box::new(GenericRow { + values: vec![ + Datum::String(Cow::Owned(String::from("nested string"))), + Datum::Blob(Cow::Owned(vec![1, 2, 3])), + ], + }))], + }; + let resolved = resolve_row_types(&row, Some(&schema)).unwrap(); + let (Datum::Row(input), Datum::Row(nested)) = (&row.values[0], &resolved.values[0]) else { + panic!("expected nested rows"); + }; + assert_borrowed_values(&input.values, &nested.values); + } + + #[test] + fn test_resolve_row_types_preserves_validation() { + let cases = [ + (DataTypes::tinyint(), Datum::Int32(128), "overflows TinyInt"), + ( + DataTypes::smallint(), + Datum::Int32(32768), + "overflows SmallInt", + ), + ( + DataType::Decimal(DecimalType::new(5, 2).unwrap()), + Datum::String(Cow::Borrowed("invalid")), + "invalid decimal string", + ), + ( + DataType::Decimal(DecimalType::new(5, 2).unwrap()), + Datum::String(Cow::Borrowed("1234.56")), + "Decimal precision overflow", + ), + ]; + for (data_type, datum, expected_error) in cases { + let schema = Schema::builder() + .column("value", data_type) + .build() + .unwrap(); + let row = GenericRow { + values: vec![datum], + }; + let error = resolve_row_types(&row, Some(&schema)).unwrap_err(); + assert!(error.to_string().contains(expected_error), "{error}"); + } + } +} From 47191bbb8206d2ee1b2e50906d94ab244dc93132 Mon Sep 17 00:00:00 2001 From: naivedogger Date: Thu, 17 Sep 2026 20:03:54 +0800 Subject: [PATCH 2/2] [c++] Skip row rebuild when no type conversion is needed resolve_row_types now returns the input row borrowed when no column needs converting, building a second row only for actual conversions or for padding short upsert/delete rows to full schema width. The lookup and prefix-lookup paths compact and resolve in one pass via the new resolve_dense_row_types, which likewise skips the rebuild for already-dense rows. Addresses the follow-up from the review of #4246. --- fluss-rust/bindings/cpp/src/lib.rs | 110 ++++----- fluss-rust/bindings/cpp/src/types.rs | 355 +++++++++++++++++++++++++-- 2 files changed, 377 insertions(+), 88 deletions(-) diff --git a/fluss-rust/bindings/cpp/src/lib.rs b/fluss-rust/bindings/cpp/src/lib.rs index 7d1224b28c7..a9fc19ac9d9 100644 --- a/fluss-rust/bindings/cpp/src/lib.rs +++ b/fluss-rust/bindings/cpp/src/lib.rs @@ -2155,12 +2155,12 @@ unsafe fn delete_append_writer(writer: *mut AppendWriter) { impl AppendWriter { fn append(&mut self, row: &GenericRowInner) -> ffi::FfiPtrResult { let schema = self.table_info.get_schema(); - let generic_row = match types::resolve_row_types(&row.row, Some(schema)) { + let generic_row = match types::resolve_row_types(&row.row, Some(schema), 0) { Ok(r) => r, Err(e) => return client_err_ptr(e.to_string()), }; - let result_future = match self.inner.append(&generic_row) { + let result_future = match self.inner.append(generic_row.as_ref()) { Ok(f) => f, Err(e) => return err_ptr_from_core(&e), }; @@ -2243,25 +2243,17 @@ unsafe fn delete_upsert_writer(writer: *mut UpsertWriter) { } impl UpsertWriter { - /// Pad row with Null to full schema width. - /// This allows callers to only set the fields they care about. - fn pad_row<'a>(&self, mut row: fcore::row::GenericRow<'a>) -> fcore::row::GenericRow<'a> { - let num_columns = self.table_info.get_schema().columns().len(); - if row.values.len() < num_columns { - row.values.resize(num_columns, fcore::row::Datum::Null); - } - row - } - fn upsert(&mut self, row: &GenericRowInner) -> ffi::FfiPtrResult { let schema = self.table_info.get_schema(); - let generic_row = match types::resolve_row_types(&row.row, Some(schema)) { - Ok(r) => r, - Err(e) => return client_err_ptr(e.to_string()), - }; - let generic_row = self.pad_row(generic_row); + // Resolve types and pad to full schema width, so callers may set only + // the fields they care about. + let generic_row = + match types::resolve_row_types(&row.row, Some(schema), schema.columns().len()) { + Ok(r) => r, + Err(e) => return client_err_ptr(e.to_string()), + }; - let result_future = match self.inner.upsert(&generic_row) { + let result_future = match self.inner.upsert(generic_row.as_ref()) { Ok(f) => f, Err(e) => return err_ptr_from_core(&e), }; @@ -2274,13 +2266,15 @@ impl UpsertWriter { fn delete_row(&mut self, row: &GenericRowInner) -> ffi::FfiPtrResult { let schema = self.table_info.get_schema(); - let generic_row = match types::resolve_row_types(&row.row, Some(schema)) { - Ok(r) => r, - Err(e) => return client_err_ptr(e.to_string()), - }; - let generic_row = self.pad_row(generic_row); + // Resolve types and pad to full schema width, so callers may set only + // the fields they care about. + let generic_row = + match types::resolve_row_types(&row.row, Some(schema), schema.columns().len()) { + Ok(r) => r, + Err(e) => return client_err_ptr(e.to_string()), + }; - let result_future = match self.inner.delete(&generic_row) { + let result_future = match self.inner.delete(generic_row.as_ref()) { Ok(f) => f, Err(e) => return err_ptr_from_core(&e), }; @@ -2311,35 +2305,24 @@ unsafe fn delete_lookuper(lookuper: *mut Lookuper) { } impl Lookuper { - /// Build a dense PK-only row from a (possibly sparse) input row. - /// The user may set PK values at their full schema positions (e.g. [0, 2]) - /// via name-based Set(). We compact them into [0, 1, …] to match - /// the lookup_row_type the core KeyEncoder expects. - fn dense_pk_row<'a>(&self, mut row: fcore::row::GenericRow<'a>) -> fcore::row::GenericRow<'a> { - let pk_indices = self.table_info.get_schema().primary_key_indexes(); - let mut dense = fcore::row::GenericRow::new(pk_indices.len()); - for (dense_idx, &schema_idx) in pk_indices.iter().enumerate() { - if schema_idx < row.values.len() { - dense.values[dense_idx] = - std::mem::replace(&mut row.values[schema_idx], fcore::row::Datum::Null); - } - } - dense - } - fn lookup(&mut self, pk_row: &GenericRowInner) -> Box { let schema = self.table_info.get_schema(); - let generic_row = match types::resolve_row_types(&pk_row.row, Some(schema)) { - Ok(r) => self.dense_pk_row(r), - Err(e) => { - return Box::new(LookupResultInner::from_error( - CLIENT_ERROR_CODE, - e.to_string(), - )); - } - }; + // Compact PK values (set at their full schema positions, e.g. [0, 2]) + // into the dense PK-only row the core KeyEncoder expects. Skips the + // rebuild when the row is already dense and needs no conversion. + let pk_indices = schema.primary_key_indexes(); + let generic_row = + match types::resolve_dense_row_types(&pk_row.row, Some(schema), &pk_indices) { + Ok(r) => r, + Err(e) => { + return Box::new(LookupResultInner::from_error( + CLIENT_ERROR_CODE, + e.to_string(), + )); + } + }; - let lookup_result = match RUNTIME.block_on(self.inner.lookup(&generic_row)) { + let lookup_result = match RUNTIME.block_on(self.inner.lookup(generic_row.as_ref())) { Ok(r) => r, Err(e) => { let ffi_err = err_from_core_error(&e); @@ -2391,23 +2374,18 @@ unsafe fn delete_prefix_lookuper(lookuper: *mut PrefixLookuper) { } impl PrefixLookuper { - /// Compact a sparse input row (prefix columns set at their schema positions) - /// into the dense, lookup-column-ordered row the core prefix encoder expects. - fn dense_prefix_row<'a>(&self, mut row: GenericRow<'a>) -> GenericRow<'a> { - let mut dense = GenericRow::new(self.lookup_column_indices.len()); - for (dense_idx, &schema_idx) in self.lookup_column_indices.iter().enumerate() { - if schema_idx < row.values.len() { - dense.values[dense_idx] = - std::mem::replace(&mut row.values[schema_idx], Datum::Null); - } - } - dense - } - fn prefix_lookup(&mut self, prefix_row: &GenericRowInner) -> Box { let schema = self.table_info.get_schema(); - let generic_row = match types::resolve_row_types(&prefix_row.row, Some(schema)) { - Ok(r) => self.dense_prefix_row(r), + // Compact prefix values (set at their full schema positions) into the + // dense, lookup-column-ordered row the core prefix encoder expects. + // Skips the rebuild when the row is already dense and needs no + // conversion. + let generic_row = match types::resolve_dense_row_types( + &prefix_row.row, + Some(schema), + &self.lookup_column_indices, + ) { + Ok(r) => r, Err(e) => { return Box::new(PrefixLookupResultInner::from_error( CLIENT_ERROR_CODE, @@ -2416,7 +2394,7 @@ impl PrefixLookuper { } }; - let lookup_result = match RUNTIME.block_on(self.inner.lookup(&generic_row)) { + let lookup_result = match RUNTIME.block_on(self.inner.lookup(generic_row.as_ref())) { Ok(r) => r, Err(e) => { let ffi_err = err_from_core_error(&e); diff --git a/fluss-rust/bindings/cpp/src/types.rs b/fluss-rust/bindings/cpp/src/types.rs index f3d31657193..c0a4f11b47c 100644 --- a/fluss-rust/bindings/cpp/src/types.rs +++ b/fluss-rust/bindings/cpp/src/types.rs @@ -456,14 +456,25 @@ pub fn core_database_info_to_ffi(info: &fcore::metadata::DatabaseInfo) -> ffi::F } /// Resolve types in a GenericRow using schema metadata. -/// Narrows Int32 → Int8/Int16, parses decimal strings, etc. -/// Unchanged STRING and BYTES values borrow the input row's storage. +/// Narrows Int32 → Int8/Int16, parses decimal strings, etc., and pads rows +/// shorter than `min_width` with trailing Nulls (the upsert/delete writers +/// require full schema width; pass 0 to keep the row's own width). +/// +/// When no column needs converting and the row is already `min_width` wide, +/// the input row is returned as-is (borrowed) and no second row is built. +/// Otherwise a new row is built, in which unchanged STRING and BYTES values +/// borrow the input row's storage. /// Used by append, upsert, delete, lookup, and prefix lookup. pub fn resolve_row_types<'a>( - row: &'a fcore::row::GenericRow<'_>, + row: &'a fcore::row::GenericRow<'a>, schema: Option<&fcore::metadata::Schema>, -) -> Result> { - let mut out = fcore::row::GenericRow::new(row.values.len()); + min_width: usize, +) -> Result>> { + if row.values.len() >= min_width && !row_needs_resolution(row, schema) { + return Ok(Cow::Borrowed(row)); + } + + let mut out = fcore::row::GenericRow::new(row.values.len().max(min_width)); for (idx, datum) in row.values.iter().enumerate() { let target = schema @@ -472,7 +483,88 @@ pub fn resolve_row_types<'a>( out.set_field(idx, resolve_datum(datum, target, idx)?); } - Ok(out) + Ok(Cow::Owned(out)) +} + +/// Resolve the columns at `indices` (schema positions) of a possibly sparse +/// input row into a dense row in the given order, resolving each value +/// against its column type; positions beyond the input row's width become +/// Null. Used by lookup (primary-key positions) and prefix lookup to compact +/// values set at their full schema positions into the dense row the core key +/// encoders expect. +/// +/// When the indices are exactly [0, 1, …] and the row is already that wide, +/// this is plain [`resolve_row_types`] and may return the input row borrowed. +pub fn resolve_dense_row_types<'a>( + row: &'a fcore::row::GenericRow<'a>, + schema: Option<&fcore::metadata::Schema>, + indices: &[usize], +) -> Result>> { + // The row is already dense: plain resolution (possibly borrowed) suffices. + if row.values.len() == indices.len() + && indices + .iter() + .enumerate() + .all(|(dense_idx, &schema_idx)| schema_idx == dense_idx) + { + return resolve_row_types(row, schema, 0); + } + + let mut dense = fcore::row::GenericRow::new(indices.len()); + for (dense_idx, &schema_idx) in indices.iter().enumerate() { + let target = schema + .and_then(|s| s.columns().get(schema_idx)) + .map(|c| c.data_type()); + let resolved = match row.values.get(schema_idx) { + Some(datum) => resolve_datum(datum, target, schema_idx)?, + None => Datum::Null, + }; + dense.set_field(dense_idx, resolved); + } + Ok(Cow::Owned(dense)) +} + +/// Whether any field of the row would be changed by `resolve_row_types`: +/// an Int32 targeted at a narrower integer type, or a String targeted at a +/// Decimal column (recursively through nested rows). +fn row_needs_resolution( + row: &fcore::row::GenericRow<'_>, + schema: Option<&fcore::metadata::Schema>, +) -> bool { + row.values.iter().enumerate().any(|(idx, datum)| { + let target = schema + .and_then(|s| s.columns().get(idx)) + .map(|c| c.data_type()); + datum_needs_resolution(datum, target) + }) +} + +/// Whether `resolve_datum` would change this datum. Mirrors the conversion +/// branches of `resolve_datum`; every other datum/target combination is +/// passed through unchanged, so resolution can be skipped for them. +fn datum_needs_resolution( + datum: &fcore::row::Datum<'_>, + target: Option<&fcore::metadata::DataType>, +) -> bool { + match datum { + Datum::Int32(_) => matches!( + target, + Some(fcore::metadata::DataType::TinyInt(_)) + | Some(fcore::metadata::DataType::SmallInt(_)) + ), + Datum::String(_) => matches!(target, Some(fcore::metadata::DataType::Decimal(_))), + Datum::Row(nested) => { + let field_types = match target { + Some(fcore::metadata::DataType::Row(rt)) => Some(rt.fields()), + _ => None, + }; + nested.values.iter().enumerate().any(|(i, d)| { + let field_type = field_types.and_then(|f| f.get(i)).map(|f| f.data_type()); + datum_needs_resolution(d, field_type) + }) + } + _ => false, + } } /// Resolve a single datum against its (optional) target column type, recursing @@ -668,26 +760,64 @@ pub fn core_scan_batches_to_ffi( #[cfg(test)] mod tests { - use super::resolve_row_types; + use super::{resolve_dense_row_types, resolve_row_types}; use fluss::metadata::{DataField, DataType, DataTypes, DecimalType, RowType, Schema}; use fluss::row::{Datum, Decimal, GenericRow}; use std::borrow::Cow; - fn assert_borrowed_values(input: &[Datum<'_>], resolved: &[Datum<'_>]) { + /// Asserts that STRING/BYTES values were resolved without copying: either + /// the input row was returned as-is (fast path) or the resolved values + /// borrow the input row's storage. + fn assert_no_copy_values(input: &[Datum<'_>], resolved: &[Datum<'_>]) { assert_eq!(resolved, input); for (input, resolved) in input.iter().zip(resolved) { match (input, resolved) { - (Datum::String(input), Datum::String(Cow::Borrowed(resolved))) => { + (Datum::String(input), Datum::String(resolved)) => { assert_eq!(input.as_ptr(), resolved.as_ptr()); } - (Datum::Blob(input), Datum::Blob(Cow::Borrowed(resolved))) => { + (Datum::Blob(input), Datum::Blob(resolved)) => { assert_eq!(input.as_ptr(), resolved.as_ptr()); } - _ => panic!("expected borrowed STRING or BYTES, got {resolved:?}"), + _ => panic!("expected STRING or BYTES values, got {resolved:?}"), } } } + #[test] + fn test_resolve_row_types_returns_input_row_when_no_conversion() { + // No datum needs converting, whatever the target types are. + let row = GenericRow { + values: vec![ + Datum::Int32(7), // Int column: pass-through + Datum::Int8(1), // already narrow + Datum::Null, + Datum::String(Cow::Owned(String::from("s"))), + ], + }; + let schema = Schema::builder() + .column("i", DataTypes::int()) + .column("t", DataTypes::tinyint()) + .column("n", DataTypes::string()) + .column("s", DataTypes::string()) + .build() + .unwrap(); + + for schema in [Some(&schema), None] { + let resolved = resolve_row_types(&row, schema, 0).unwrap(); + assert!(matches!(resolved, Cow::Borrowed(_))); + assert!(std::ptr::eq(resolved.as_ref(), &row)); + } + + // A row that needs conversion is rebuilt instead. + let schema = Schema::builder() + .column("t", DataTypes::tinyint()) + .build() + .unwrap(); + let resolved = resolve_row_types(&row, Some(&schema), 0).unwrap(); + assert!(matches!(resolved, Cow::Owned(_))); + assert_eq!(resolved.values[0], Datum::Int8(7)); + } + #[test] fn test_resolve_row_types_borrows_strings_and_bytes() { // Include both setter-owned values and values already borrowing external storage. @@ -714,8 +844,8 @@ mod tests { .unwrap(); for schema in [Some(&schema), None] { - let resolved = resolve_row_types(&row, schema).unwrap(); - assert_borrowed_values(&row.values, &resolved.values); + let resolved = resolve_row_types(&row, schema, 0).unwrap(); + assert_no_copy_values(&row.values, &resolved.values); } } @@ -743,20 +873,52 @@ mod tests { Datum::Null, ], }; - let resolved = resolve_row_types(&row, Some(&schema)).unwrap(); + let resolved = resolve_row_types(&row, Some(&schema), 0).unwrap(); + assert!(matches!(resolved, Cow::Owned(_))); assert_eq!(resolved.values[0], Datum::Int8(127)); assert_eq!(resolved.values[1], Datum::Int16(-32768)); assert_eq!( resolved.values[2], Datum::Decimal(Decimal::from_unscaled_long(12345, 5, 2).unwrap()) ); - assert_borrowed_values(&row.values[3..5], &resolved.values[3..5]); + assert_no_copy_values(&row.values[3..5], &resolved.values[3..5]); assert_eq!(resolved.values[5], Datum::Null); } + #[test] + fn test_resolve_row_types_pads_short_rows_to_min_width() { + let row = GenericRow { + values: vec![ + Datum::String(Cow::Owned(String::from("a"))), + Datum::Int32(1), + ], + }; + let schema = Schema::builder() + .column("a", DataTypes::string()) + .column("b", DataTypes::int()) + .column("c", DataTypes::string()) + .column("d", DataTypes::string()) + .build() + .unwrap(); + + // Wide enough: the input row is returned as-is. + let resolved = resolve_row_types(&row, Some(&schema), 2).unwrap(); + assert!(matches!(resolved, Cow::Borrowed(_))); + + // Short row: rebuilt at min_width with trailing Nulls, values borrowed. + let resolved = resolve_row_types(&row, Some(&schema), 4).unwrap(); + assert!(matches!(resolved, Cow::Owned(_))); + assert_eq!(resolved.values.len(), 4); + assert_no_copy_values(&row.values[..1], &resolved.values[..1]); + assert_eq!(resolved.values[1], Datum::Int32(1)); + assert_eq!(resolved.values[2], Datum::Null); + assert_eq!(resolved.values[3], Datum::Null); + } + #[test] fn test_resolve_row_types_borrows_nested_values() { let schema = Schema::builder() + .column("tiny", DataTypes::tinyint()) .column( "nested", DataType::Row(RowType::new(vec![ @@ -766,19 +928,62 @@ mod tests { ) .build() .unwrap(); + // The outer row needs a conversion, so the nested row is rebuilt too. + let row = GenericRow { + values: vec![ + Datum::Int32(1), + Datum::Row(Box::new(GenericRow { + values: vec![ + Datum::String(Cow::Owned(String::from("nested string"))), + Datum::Blob(Cow::Owned(vec![1, 2, 3])), + ], + })), + ], + }; + let resolved = resolve_row_types(&row, Some(&schema), 0).unwrap(); + assert_eq!(resolved.values[0], Datum::Int8(1)); + let (Datum::Row(input), Datum::Row(nested)) = (&row.values[1], &resolved.values[1]) else { + panic!("expected nested rows"); + }; + assert_no_copy_values(&input.values, &nested.values); + } + + #[test] + fn test_resolve_row_types_nested_conversion_forces_rebuild() { + // Only the nested field needs converting; the fast path must not + // skip it. + let schema = Schema::builder() + .column( + "nested", + DataType::Row(RowType::new(vec![ + DataField::new( + "decimal", + DataType::Decimal(DecimalType::new(5, 2).unwrap()), + None, + ), + DataField::new("string", DataTypes::string(), None), + ])), + ) + .build() + .unwrap(); let row = GenericRow { values: vec![Datum::Row(Box::new(GenericRow { values: vec![ - Datum::String(Cow::Owned(String::from("nested string"))), - Datum::Blob(Cow::Owned(vec![1, 2, 3])), + Datum::String(Cow::Owned(String::from("12.34"))), + Datum::String(Cow::Owned(String::from("kept"))), ], }))], }; - let resolved = resolve_row_types(&row, Some(&schema)).unwrap(); - let (Datum::Row(input), Datum::Row(nested)) = (&row.values[0], &resolved.values[0]) else { - panic!("expected nested rows"); + let resolved = resolve_row_types(&row, Some(&schema), 0).unwrap(); + assert!(matches!(resolved, Cow::Owned(_))); + let Datum::Row(nested) = &resolved.values[0] else { + panic!("expected nested row"); }; - assert_borrowed_values(&input.values, &nested.values); + assert_eq!( + nested.values[0], + Datum::Decimal(Decimal::from_unscaled_long(1234, 5, 2).unwrap()) + ); + assert_eq!(nested.values[1], Datum::String(Cow::Borrowed("kept"))); } #[test] @@ -809,8 +1014,114 @@ mod tests { let row = GenericRow { values: vec![datum], }; - let error = resolve_row_types(&row, Some(&schema)).unwrap_err(); + let error = resolve_row_types(&row, Some(&schema), 0).unwrap_err(); assert!(error.to_string().contains(expected_error), "{error}"); } } + + #[test] + fn test_resolve_dense_row_types_identity_returns_input_row() { + // PK columns at schema positions [0, 1] and a 2-wide row: already dense. + let row = GenericRow { + values: vec![ + Datum::Int32(1), + Datum::String(Cow::Owned(String::from("pk"))), + ], + }; + let schema = Schema::builder() + .column("a", DataTypes::int()) + .column("b", DataTypes::string()) + .build() + .unwrap(); + + let resolved = resolve_dense_row_types(&row, Some(&schema), &[0, 1]).unwrap(); + assert!(matches!(resolved, Cow::Borrowed(_))); + assert!(std::ptr::eq(resolved.as_ref(), &row)); + } + + #[test] + fn test_resolve_dense_row_types_compacts_and_converts() { + // PK columns sit at schema positions [0, 2]; values are set at their + // full schema positions in a wider row. + let schema = Schema::builder() + .column("a", DataTypes::tinyint()) + .column("filler", DataTypes::string()) + .column("dec", DataType::Decimal(DecimalType::new(5, 2).unwrap())) + .build() + .unwrap(); + let row = GenericRow { + values: vec![ + Datum::Int32(7), // a: TinyInt, narrow + Datum::String(Cow::Owned(String::from("not a pk"))), // skipped + Datum::String(Cow::Owned(String::from("12.34"))), // dec: parse + ], + }; + + let resolved = resolve_dense_row_types(&row, Some(&schema), &[0, 2]).unwrap(); + assert!(matches!(resolved, Cow::Owned(_))); + assert_eq!(resolved.values.len(), 2); + assert_eq!(resolved.values[0], Datum::Int8(7)); + assert_eq!( + resolved.values[1], + Datum::Decimal(Decimal::from_unscaled_long(1234, 5, 2).unwrap()) + ); + } + + #[test] + fn test_resolve_dense_row_types_borrows_unchanged_values() { + let string = String::from("borrowed"); + let bytes = vec![9, 9]; + let schema = Schema::builder() + .column("a", DataTypes::string()) + .column("b", DataTypes::bytes()) + .column("c", DataTypes::string()) + .build() + .unwrap(); + let row = GenericRow { + values: vec![ + Datum::String(Cow::Borrowed(&string)), + Datum::Blob(Cow::Borrowed(&bytes)), + Datum::Null, // not part of the dense projection + ], + }; + + let resolved = resolve_dense_row_types(&row, Some(&schema), &[0, 1]).unwrap(); + assert_eq!(resolved.values.len(), 2); + assert_no_copy_values(&row.values[..2], &resolved.values); + } + + #[test] + fn test_resolve_dense_row_types_out_of_range_is_null() { + let schema = Schema::builder() + .column("a", DataTypes::string()) + .column("b", DataTypes::string()) + .build() + .unwrap(); + // Only the first PK value is set; position 1 is beyond the row width. + let row = GenericRow { + values: vec![Datum::String(Cow::Owned(String::from("only")))], + }; + + let resolved = resolve_dense_row_types(&row, Some(&schema), &[0, 1]).unwrap(); + assert!(matches!(resolved, Cow::Owned(_))); + assert_eq!(resolved.values.len(), 2); + assert_eq!(resolved.values[0], Datum::String(Cow::Borrowed("only"))); + assert_eq!(resolved.values[1], Datum::Null); + } + + #[test] + fn test_resolve_dense_row_types_preserves_validation() { + let schema = Schema::builder() + .column("a", DataTypes::string()) + .column("dec", DataType::Decimal(DecimalType::new(5, 2).unwrap())) + .build() + .unwrap(); + let row = GenericRow { + values: vec![Datum::Null, Datum::String(Cow::Owned(String::from("oops")))], + }; + + // Errors report the schema position, as full-row resolution does. + let error = resolve_dense_row_types(&row, Some(&schema), &[1]).unwrap_err(); + assert!(error.to_string().contains("Column 1"), "{error}"); + } }