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 84b88b98406..c0a4f11b47c 100644 --- a/fluss-rust/bindings/cpp/src/types.rs +++ b/fluss-rust/bindings/cpp/src/types.rs @@ -456,13 +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. -/// Used by both AppendWriter and UpsertWriter. -pub fn resolve_row_types( - row: &fcore::row::GenericRow<'_>, +/// 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<'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 @@ -471,17 +483,98 @@ pub fn resolve_row_types( 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 /// 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 +602,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 +757,371 @@ pub fn core_scan_batches_to_ffi( batches: ffi_batches, }) } + +#[cfg(test)] +mod tests { + 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; + + /// 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(resolved)) => { + assert_eq!(input.as_ptr(), resolved.as_ptr()); + } + (Datum::Blob(input), Datum::Blob(resolved)) => { + assert_eq!(input.as_ptr(), resolved.as_ptr()); + } + _ => 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. + 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, 0).unwrap(); + assert_no_copy_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), 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_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![ + DataField::new("string", DataTypes::string(), None), + DataField::new("bytes", DataTypes::bytes(), None), + ])), + ) + .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("12.34"))), + Datum::String(Cow::Owned(String::from("kept"))), + ], + }))], + }; + 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_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] + 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), 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}"); + } +}