From ac75b3620b77c08b59a57c75d1ee8ba2af0f6cdc Mon Sep 17 00:00:00 2001 From: Junbo Wang Date: Tue, 1 Sep 2026 18:18:59 +0800 Subject: [PATCH 1/2] [rust] Preserve highest_field_id during Schema JSON deserialization --- .../crates/fluss/src/metadata/json_serde.rs | 52 +++++++++++++++++++ fluss-rust/crates/fluss/src/metadata/table.rs | 17 +++++- 2 files changed, 68 insertions(+), 1 deletion(-) diff --git a/fluss-rust/crates/fluss/src/metadata/json_serde.rs b/fluss-rust/crates/fluss/src/metadata/json_serde.rs index c504660e92b..7d1ea255e7a 100644 --- a/fluss-rust/crates/fluss/src/metadata/json_serde.rs +++ b/fluss-rust/crates/fluss/src/metadata/json_serde.rs @@ -543,6 +543,23 @@ impl JsonSerde for Schema { } } + if let Some(highest_field_id_node) = node.get(Self::HIGHEST_FIELD_ID) { + let highest_field_id = + highest_field_id_node + .as_i64() + .ok_or_else(|| Error::JsonSerdeError { + message: format!("{} must be an integer", Self::HIGHEST_FIELD_ID), + })?; + let highest_field_id = + i32::try_from(highest_field_id).map_err(|_| Error::JsonSerdeError { + message: format!( + "{} value {highest_field_id} does not fit in i32", + Self::HIGHEST_FIELD_ID + ), + })?; + schema_builder = schema_builder.highest_field_id(highest_field_id); + } + schema_builder.build() } } @@ -984,6 +1001,41 @@ mod tests { assert_eq!(round_tripped, original); } + #[test] + fn schema_deserialization_preserves_highest_field_id_above_maximum() { + let json = json!({ + "version": 1, + "columns": [ + {"name": "id", "data_type": {"type": "INTEGER"}, "id": 0}, + {"name": "name", "data_type": {"type": "STRING"}, "id": 1} + ], + "highest_field_id": 10 + }); + + let schema = Schema::deserialize_json(&json).unwrap(); + assert_eq!(schema.highest_field_id(), 10); + } + + #[test] + fn schema_deserialization_rejects_highest_field_id_below_maximum() { + let json = json!({ + "version": 1, + "columns": [ + {"name": "id", "data_type": {"type": "INTEGER"}, "id": 0}, + {"name": "name", "data_type": {"type": "STRING"}, "id": 1} + ], + "highest_field_id": 0 + }); + + let error = Schema::deserialize_json(&json).unwrap_err(); + assert!( + error + .to_string() + .contains("must be greater than or equal to the maximum field id (1)"), + "{error}" + ); + } + #[test] fn schema_rejects_duplicate_nested_field_ids() { let nested = DataTypes::row(vec![ diff --git a/fluss-rust/crates/fluss/src/metadata/table.rs b/fluss-rust/crates/fluss/src/metadata/table.rs index c03a60065b2..dc6a86f2653 100644 --- a/fluss-rust/crates/fluss/src/metadata/table.rs +++ b/fluss-rust/crates/fluss/src/metadata/table.rs @@ -253,6 +253,7 @@ pub struct SchemaBuilder { columns: Vec, primary_key: Option, auto_increment_col_names: Vec, + highest_field_id: Option, } impl SchemaBuilder { @@ -330,9 +331,23 @@ impl SchemaBuilder { Ok(self) } + pub(crate) fn highest_field_id(mut self, highest_field_id: i32) -> Self { + self.highest_field_id = Some(highest_field_id); + self + } + pub fn build(&self) -> Result { let columns = Self::normalize_columns(&self.columns, self.primary_key.as_ref())?; - let (columns_with_ids, highest_field_id) = Self::assign_all_field_ids(columns)?; + let (columns_with_ids, maximum_field_id) = Self::assign_all_field_ids(columns)?; + let highest_field_id = self.highest_field_id.unwrap_or(maximum_field_id); + + if !columns_with_ids.is_empty() && highest_field_id < maximum_field_id { + return Err(IllegalArgument { + message: format!( + "Highest field id ({highest_field_id}) must be greater than or equal to the maximum field id ({maximum_field_id})" + ), + }); + } if !self.auto_increment_col_names.is_empty() && self.primary_key.is_none() { return Err(IllegalArgument { From a9b6077cf71f44fc3e490c8b4b730e477cebf838 Mon Sep 17 00:00:00 2001 From: Junbo Wang Date: Tue, 1 Sep 2026 19:54:51 +0800 Subject: [PATCH 2/2] [rust] Clamp deserialized highest_field_id to schema maximum --- .../crates/fluss/src/metadata/json_serde.rs | 19 +++++++++++-------- fluss-rust/crates/fluss/src/metadata/table.rs | 13 ++++--------- 2 files changed, 15 insertions(+), 17 deletions(-) diff --git a/fluss-rust/crates/fluss/src/metadata/json_serde.rs b/fluss-rust/crates/fluss/src/metadata/json_serde.rs index 7d1ea255e7a..cc636a4ad6b 100644 --- a/fluss-rust/crates/fluss/src/metadata/json_serde.rs +++ b/fluss-rust/crates/fluss/src/metadata/json_serde.rs @@ -1017,7 +1017,7 @@ mod tests { } #[test] - fn schema_deserialization_rejects_highest_field_id_below_maximum() { + fn schema_deserialization_clamps_highest_field_id_to_computed_minimum() { let json = json!({ "version": 1, "columns": [ @@ -1027,13 +1027,16 @@ mod tests { "highest_field_id": 0 }); - let error = Schema::deserialize_json(&json).unwrap_err(); - assert!( - error - .to_string() - .contains("must be greater than or equal to the maximum field id (1)"), - "{error}" - ); + let schema = Schema::deserialize_json(&json).unwrap(); + assert_eq!(schema.highest_field_id(), 1); + + let empty_json = json!({ + "version": 1, + "columns": [], + "highest_field_id": -3 + }); + let empty_schema = Schema::deserialize_json(&empty_json).unwrap(); + assert_eq!(empty_schema.highest_field_id(), -1); } #[test] diff --git a/fluss-rust/crates/fluss/src/metadata/table.rs b/fluss-rust/crates/fluss/src/metadata/table.rs index dc6a86f2653..405f54872fe 100644 --- a/fluss-rust/crates/fluss/src/metadata/table.rs +++ b/fluss-rust/crates/fluss/src/metadata/table.rs @@ -339,15 +339,10 @@ impl SchemaBuilder { pub fn build(&self) -> Result { let columns = Self::normalize_columns(&self.columns, self.primary_key.as_ref())?; let (columns_with_ids, maximum_field_id) = Self::assign_all_field_ids(columns)?; - let highest_field_id = self.highest_field_id.unwrap_or(maximum_field_id); - - if !columns_with_ids.is_empty() && highest_field_id < maximum_field_id { - return Err(IllegalArgument { - message: format!( - "Highest field id ({highest_field_id}) must be greater than or equal to the maximum field id ({maximum_field_id})" - ), - }); - } + let highest_field_id = self + .highest_field_id + .unwrap_or(maximum_field_id) + .max(maximum_field_id); if !self.auto_increment_col_names.is_empty() && self.primary_key.is_none() { return Err(IllegalArgument {