From 68d306e4742352b10e5b4c89861141477e673aea Mon Sep 17 00:00:00 2001 From: CGodiksen <36046286+CGodiksen@users.noreply.github.com> Date: Sun, 7 Sep 2025 18:05:35 +0200 Subject: [PATCH 1/8] Remove implementation of PartialEq for Manager --- crates/modelardb_server/src/manager.rs | 11 +---------- 1 file changed, 1 insertion(+), 10 deletions(-) diff --git a/crates/modelardb_server/src/manager.rs b/crates/modelardb_server/src/manager.rs index 4e9936b01..922440430 100644 --- a/crates/modelardb_server/src/manager.rs +++ b/crates/modelardb_server/src/manager.rs @@ -35,7 +35,7 @@ use crate::context::Context; use crate::error::{ModelarDbServerError, Result}; /// Manages metadata related to the manager and provides functionality for interacting with the manager. -#[derive(Clone, Debug)] +#[derive(Clone, Debug, PartialEq)] pub struct Manager { /// Key received from the manager when registering, used to validate future requests that are /// only allowed to come from the manager. @@ -189,15 +189,6 @@ async fn do_action_and_extract_result( }) } -/// Partial equality is implemented so PartialEq can be derived for [`ClusterMode`](crate::ClusterMode). -/// It cannot be derived for [`Manager`] since both `flight_client` and `table_metadata_manager` -/// does not support equality comparisons. -impl PartialEq for Manager { - fn eq(&self, other: &Self) -> bool { - self.key == other.key - } -} - #[cfg(test)] mod tests { use super::*; From b4a863cf752a1e994476dd9314743eed9f7b9310 Mon Sep 17 00:00:00 2001 From: CGodiksen <36046286+CGodiksen@users.noreply.github.com> Date: Sun, 7 Sep 2025 18:57:20 +0200 Subject: [PATCH 2/8] Add function to validate that remote tables are equivalent with local normal tables --- crates/modelardb_server/src/manager.rs | 41 ++++++++++++++++++++++++++ 1 file changed, 41 insertions(+) diff --git a/crates/modelardb_server/src/manager.rs b/crates/modelardb_server/src/manager.rs index 922440430..c30469703 100644 --- a/crates/modelardb_server/src/manager.rs +++ b/crates/modelardb_server/src/manager.rs @@ -21,6 +21,7 @@ use std::{env, str}; use arrow_flight::flight_service_client::FlightServiceClient; use arrow_flight::{Action, Result as FlightResult}; +use datafusion::arrow::datatypes::Schema; use datafusion::catalog::TableProvider; use modelardb_types::flight::protocol; use modelardb_types::types::{Node, ServerMode}; @@ -32,6 +33,7 @@ use tonic::transport::Channel; use crate::PORT; use crate::context::Context; +use crate::data_folders::DataFolder; use crate::error::{ModelarDbServerError, Result}; /// Manages metadata related to the manager and provides functionality for interacting with the manager. @@ -189,6 +191,45 @@ async fn do_action_and_extract_result( }) } +/// Validate that all normal tables in the local data folder exist in the remote data folder and have +/// the same schema. If all normal tables are valid, return a vector of tuples containing the +/// table name and schema of each normal table that is in the remote data folder but not in the local +/// data folder. If any normal table is invalid, return [`ModelarDbServerError`]. +async fn validate_normal_tables( + local_data_folder: &DataFolder, + remote_data_folder: &DataFolder, +) -> Result)>> { + let mut missing_normal_tables = vec![]; + + let remote_normal_tables = remote_data_folder + .table_metadata_manager + .normal_table_names() + .await?; + + for table_name in remote_normal_tables { + let remote_schema = normal_table_schema(remote_data_folder, &table_name).await?; + + if let Ok(local_schema) = normal_table_schema(local_data_folder, &table_name).await { + if remote_schema != local_schema { + return Err(ModelarDbServerError::InvalidState(format!( + "The normal table '{table_name}' has a different schema in the local data folder than in the remote data folder.", + ))); + } + } else { + missing_normal_tables.push((table_name, remote_schema)); + } + } + + Ok(missing_normal_tables) +} + +/// Retrieve the schema of a normal table from the Delta Lake in the data folder. If the table does +/// not exist, or the schema could not be retrieved, return [`ModelarDbServerError`]. +async fn normal_table_schema(data_folder: &DataFolder, table_name: &str) -> Result> { + let delta_table = data_folder.delta_lake.delta_table(table_name).await?; + Ok(TableProvider::schema(&delta_table)) +} + #[cfg(test)] mod tests { use super::*; From d98f75c295baefd79951381a855bdb9b96f2f5a3 Mon Sep 17 00:00:00 2001 From: CGodiksen <36046286+CGodiksen@users.noreply.github.com> Date: Sun, 7 Sep 2025 19:06:27 +0200 Subject: [PATCH 3/8] Add function to get unique and shared tables --- crates/modelardb_server/src/manager.rs | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/crates/modelardb_server/src/manager.rs b/crates/modelardb_server/src/manager.rs index c30469703..2c20825ef 100644 --- a/crates/modelardb_server/src/manager.rs +++ b/crates/modelardb_server/src/manager.rs @@ -16,6 +16,7 @@ //! Interface to connect to and interact with the manager, used if the server is started with a //! manager and needs to interact with it to initialize the metadata Delta Lake. +use std::collections::HashSet; use std::sync::Arc; use std::{env, str}; @@ -230,6 +231,22 @@ async fn normal_table_schema(data_folder: &DataFolder, table_name: &str) -> Resu Ok(TableProvider::schema(&delta_table)) } +/// Given the names of the tables in the local and remote data folders, return the unique tables in +/// the local data folder, the unique tables in the remote data folder, and the shared tables. +async fn unique_and_shared_tables( + local_table_names: Vec, + remote_table_names: Vec, +) -> (HashSet, HashSet, HashSet) { + let local_set: HashSet = local_table_names.into_iter().collect(); + let remote_set: HashSet = remote_table_names.into_iter().collect(); + + let unique_local_tables = local_set.difference(&remote_set).cloned().collect(); + let unique_remote_tables = remote_set.difference(&local_set).cloned().collect(); + let shared_tables = local_set.intersection(&remote_set).cloned().collect(); + + (unique_local_tables, unique_remote_tables, shared_tables) +} + #[cfg(test)] mod tests { use super::*; From cd87367989e4f8960ebdc93aa25d0d11785beee2 Mon Sep 17 00:00:00 2001 From: CGodiksen <36046286+CGodiksen@users.noreply.github.com> Date: Sun, 7 Sep 2025 19:36:57 +0200 Subject: [PATCH 4/8] Add function to remote tables are equivalent with local time series tables --- crates/modelardb_server/src/manager.rs | 51 +++++++++++++++++++------- crates/modelardb_types/src/types.rs | 2 +- 2 files changed, 39 insertions(+), 14 deletions(-) diff --git a/crates/modelardb_server/src/manager.rs b/crates/modelardb_server/src/manager.rs index 2c20825ef..f830d2ab2 100644 --- a/crates/modelardb_server/src/manager.rs +++ b/crates/modelardb_server/src/manager.rs @@ -213,7 +213,8 @@ async fn validate_normal_tables( if let Ok(local_schema) = normal_table_schema(local_data_folder, &table_name).await { if remote_schema != local_schema { return Err(ModelarDbServerError::InvalidState(format!( - "The normal table '{table_name}' has a different schema in the local data folder than in the remote data folder.", + "The normal table '{table_name}' has a different schema in the local data \ + folder than in the remote data folder.", ))); } } else { @@ -231,20 +232,44 @@ async fn normal_table_schema(data_folder: &DataFolder, table_name: &str) -> Resu Ok(TableProvider::schema(&delta_table)) } -/// Given the names of the tables in the local and remote data folders, return the unique tables in -/// the local data folder, the unique tables in the remote data folder, and the shared tables. -async fn unique_and_shared_tables( - local_table_names: Vec, - remote_table_names: Vec, -) -> (HashSet, HashSet, HashSet) { - let local_set: HashSet = local_table_names.into_iter().collect(); - let remote_set: HashSet = remote_table_names.into_iter().collect(); +/// Validate that all time series tables in the local data folder exist in the remote data folder +/// and have the same metadata. If all time series tables are valid, return a vector containing +/// the metadata of each time series table that is in the remote data folder but not in the local +/// data folder. If any time series table is invalid, return [`ModelarDbServerError`]. +async fn validate_time_series_tables( + local_data_folder: &DataFolder, + remote_data_folder: &DataFolder, +) -> Result> { + let mut missing_time_series_tables = vec![]; + + let remote_time_series_tables = remote_data_folder + .table_metadata_manager + .time_series_table_names() + .await?; - let unique_local_tables = local_set.difference(&remote_set).cloned().collect(); - let unique_remote_tables = remote_set.difference(&local_set).cloned().collect(); - let shared_tables = local_set.intersection(&remote_set).cloned().collect(); + for table_name in remote_time_series_tables { + let remote_metadata = remote_data_folder + .table_metadata_manager + .time_series_table_metadata_for_time_series_table(&table_name) + .await?; + + if let Ok(local_metadata) = local_data_folder + .table_metadata_manager + .time_series_table_metadata_for_time_series_table(&table_name) + .await + { + if remote_metadata != local_metadata { + return Err(ModelarDbServerError::InvalidState(format!( + "The time series table '{table_name}' has different metadata in the local data \ + folder than in the remote data folder.", + ))); + } + } else { + missing_time_series_tables.push(remote_metadata); + } + } - (unique_local_tables, unique_remote_tables, shared_tables) + Ok(missing_time_series_tables) } #[cfg(test)] diff --git a/crates/modelardb_types/src/types.rs b/crates/modelardb_types/src/types.rs index 3a770bfeb..7d27a2c4f 100644 --- a/crates/modelardb_types/src/types.rs +++ b/crates/modelardb_types/src/types.rs @@ -65,7 +65,7 @@ pub enum Table { } /// Metadata required to ingest data into a time series table and query a time series table. -#[derive(Debug, Clone)] +#[derive(Debug, Clone, PartialEq)] pub struct TimeSeriesTableMetadata { /// Name of the time series table. pub name: String, From 1191c017fc9b99d3e1d1df24b9e7a7e560f814f1 Mon Sep 17 00:00:00 2001 From: CGodiksen <36046286+CGodiksen@users.noreply.github.com> Date: Sun, 7 Sep 2025 19:43:25 +0200 Subject: [PATCH 5/8] Use utility functions to validate tables and create missing tables --- crates/modelardb_server/src/manager.rs | 46 ++++++++++---------------- 1 file changed, 18 insertions(+), 28 deletions(-) diff --git a/crates/modelardb_server/src/manager.rs b/crates/modelardb_server/src/manager.rs index f830d2ab2..e2ae4c05f 100644 --- a/crates/modelardb_server/src/manager.rs +++ b/crates/modelardb_server/src/manager.rs @@ -16,7 +16,6 @@ //! Interface to connect to and interact with the manager, used if the server is started with a //! manager and needs to interact with it to initialize the metadata Delta Lake. -use std::collections::HashSet; use std::sync::Arc; use std::{env, str}; @@ -25,7 +24,7 @@ use arrow_flight::{Action, Result as FlightResult}; use datafusion::arrow::datatypes::Schema; use datafusion::catalog::TableProvider; use modelardb_types::flight::protocol; -use modelardb_types::types::{Node, ServerMode}; +use modelardb_types::types::{Node, ServerMode, TimeSeriesTableMetadata}; use prost::Message; use tokio::sync::RwLock; use tonic::Request; @@ -90,10 +89,8 @@ impl Manager { /// retrieved from the remote data folder, or the tables could not be created, /// return [`ModelarDbServerError`]. pub(crate) async fn retrieve_and_create_tables(&self, context: &Arc) -> Result<()> { - let local_metadata_manager = &context - .data_folders - .local_data_folder - .table_metadata_manager; + let local_data_folder = &context.data_folders.local_data_folder; + let local_metadata_manager = &local_data_folder.table_metadata_manager; let remote_data_folder = &context .data_folders @@ -120,29 +117,22 @@ impl Manager { ))); } + // Validate that all tables that are in both the local and remote data folder are identical. + let missing_normal_tables = + validate_normal_tables(local_data_folder, remote_data_folder).await?; + + let missing_time_series_tables = + validate_time_series_tables(local_data_folder, remote_data_folder).await?; + // For each table that does not already exist locally, create the table. - let missing_cluster_tables = remote_table_names - .iter() - .filter(|table| !local_table_names.contains(table)); - - for table_name in missing_cluster_tables { - if remote_metadata_manager.is_normal_table(table_name).await? { - let delta_table = remote_data_folder - .delta_lake - .delta_table(table_name) - .await?; - - let schema = TableProvider::schema(&delta_table); - context.create_normal_table(table_name, &schema).await?; - } else { - let time_series_table_metadata = remote_metadata_manager - .time_series_table_metadata_for_time_series_table(table_name) - .await?; - - context - .create_time_series_table(&time_series_table_metadata) - .await?; - } + for (table_name, schema) in missing_normal_tables { + context.create_normal_table(&table_name, &schema).await?; + } + + for time_series_table_metadata in missing_time_series_tables { + context + .create_time_series_table(&time_series_table_metadata) + .await?; } Ok(()) From 24ea3de8136015bf5abd9ce54770ae7bd66ea0e1 Mon Sep 17 00:00:00 2001 From: CGodiksen <36046286+CGodiksen@users.noreply.github.com> Date: Sun, 7 Sep 2025 19:49:56 +0200 Subject: [PATCH 6/8] Add function to validate local tables exist remotely --- crates/modelardb_server/src/manager.rs | 50 ++++++++++++++++---------- 1 file changed, 32 insertions(+), 18 deletions(-) diff --git a/crates/modelardb_server/src/manager.rs b/crates/modelardb_server/src/manager.rs index e2ae4c05f..fb1ca8c37 100644 --- a/crates/modelardb_server/src/manager.rs +++ b/crates/modelardb_server/src/manager.rs @@ -90,7 +90,6 @@ impl Manager { /// return [`ModelarDbServerError`]. pub(crate) async fn retrieve_and_create_tables(&self, context: &Arc) -> Result<()> { let local_data_folder = &context.data_folders.local_data_folder; - let local_metadata_manager = &local_data_folder.table_metadata_manager; let remote_data_folder = &context .data_folders @@ -99,23 +98,8 @@ impl Manager { .ok_or(ModelarDbServerError::InvalidState( "Remote data folder is missing.".to_owned(), ))?; - let remote_metadata_manager = &remote_data_folder.table_metadata_manager; - - let local_table_names = local_metadata_manager.table_names().await?; - let remote_table_names = remote_metadata_manager.table_names().await?; - - // Check that all the local tables exist in the cluster's database schema already. - let invalid_node_tables: Vec = local_table_names - .iter() - .filter(|table| !remote_table_names.contains(table)) - .cloned() - .collect(); - - if !invalid_node_tables.is_empty() { - return Err(ModelarDbServerError::InvalidState(format!( - "The following tables do not exist in the cluster's database schema: {invalid_node_tables:?}.", - ))); - } + + validate_local_tables_exist_remotely(local_data_folder, remote_data_folder).await?; // Validate that all tables that are in both the local and remote data folder are identical. let missing_normal_tables = @@ -182,6 +166,36 @@ async fn do_action_and_extract_result( }) } +/// Validate that all tables in the local data folder exist in the remote data folder. If any table +/// does not exist in the remote data folder, return [`ModelarDbServerError`]. +async fn validate_local_tables_exist_remotely( + local_data_folder: &DataFolder, + remote_data_folder: &DataFolder, +) -> Result<()> { + let local_table_names = local_data_folder + .table_metadata_manager + .table_names() + .await?; + let remote_table_names = remote_data_folder + .table_metadata_manager + .table_names() + .await?; + + let invalid_tables: Vec = local_table_names + .iter() + .filter(|table| !remote_table_names.contains(table)) + .cloned() + .collect(); + + if !invalid_tables.is_empty() { + return Err(ModelarDbServerError::InvalidState(format!( + "The following tables do not exist in the remote data folder: {invalid_tables:?}.", + ))); + } + + Ok(()) +} + /// Validate that all normal tables in the local data folder exist in the remote data folder and have /// the same schema. If all normal tables are valid, return a vector of tuples containing the /// table name and schema of each normal table that is in the remote data folder but not in the local From 7f2c907cc49102e166627d89dc6d9e1f43561458 Mon Sep 17 00:00:00 2001 From: CGodiksen <36046286+CGodiksen@users.noreply.github.com> Date: Sun, 7 Sep 2025 20:14:39 +0200 Subject: [PATCH 7/8] Minor changes to error message --- crates/modelardb_server/src/manager.rs | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/crates/modelardb_server/src/manager.rs b/crates/modelardb_server/src/manager.rs index fb1ca8c37..4fc8f346c 100644 --- a/crates/modelardb_server/src/manager.rs +++ b/crates/modelardb_server/src/manager.rs @@ -189,7 +189,8 @@ async fn validate_local_tables_exist_remotely( if !invalid_tables.is_empty() { return Err(ModelarDbServerError::InvalidState(format!( - "The following tables do not exist in the remote data folder: {invalid_tables:?}.", + "The following tables do not exist in the remote data folder: {}.", + invalid_tables.join(", ") ))); } @@ -218,7 +219,7 @@ async fn validate_normal_tables( if remote_schema != local_schema { return Err(ModelarDbServerError::InvalidState(format!( "The normal table '{table_name}' has a different schema in the local data \ - folder than in the remote data folder.", + folder compared to the remote data folder.", ))); } } else { @@ -265,7 +266,7 @@ async fn validate_time_series_tables( if remote_metadata != local_metadata { return Err(ModelarDbServerError::InvalidState(format!( "The time series table '{table_name}' has different metadata in the local data \ - folder than in the remote data folder.", + folder compared to the remote data folder.", ))); } } else { From 421def76e1ed7bfa9140d03b888e72c8f35dd5cf Mon Sep 17 00:00:00 2001 From: CGodiksen <36046286+CGodiksen@users.noreply.github.com> Date: Sun, 7 Sep 2025 20:30:52 +0200 Subject: [PATCH 8/8] Update doc comment for clarity --- crates/modelardb_server/src/manager.rs | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/crates/modelardb_server/src/manager.rs b/crates/modelardb_server/src/manager.rs index 4fc8f346c..a0ceb8eb1 100644 --- a/crates/modelardb_server/src/manager.rs +++ b/crates/modelardb_server/src/manager.rs @@ -197,10 +197,10 @@ async fn validate_local_tables_exist_remotely( Ok(()) } -/// Validate that all normal tables in the local data folder exist in the remote data folder and have -/// the same schema. If all normal tables are valid, return a vector of tuples containing the -/// table name and schema of each normal table that is in the remote data folder but not in the local -/// data folder. If any normal table is invalid, return [`ModelarDbServerError`]. +/// For each normal table in the remote data folder, if the table also exists in the local data +/// folder, validate that the schemas are identical. If the schemas are not identical, return +/// [`ModelarDbServerError`]. Return a vector containing the name and schema of each normal table +/// that is in the remote data folder but not in the local data folder. async fn validate_normal_tables( local_data_folder: &DataFolder, remote_data_folder: &DataFolder, @@ -237,10 +237,10 @@ async fn normal_table_schema(data_folder: &DataFolder, table_name: &str) -> Resu Ok(TableProvider::schema(&delta_table)) } -/// Validate that all time series tables in the local data folder exist in the remote data folder -/// and have the same metadata. If all time series tables are valid, return a vector containing -/// the metadata of each time series table that is in the remote data folder but not in the local -/// data folder. If any time series table is invalid, return [`ModelarDbServerError`]. +/// For each time series table in the remote data folder, if the table also exists in the local +/// data folder, validate that the metadata is identical. If the metadata is not identical, return +/// [`ModelarDbServerError`]. Return a vector containing the metadata of each time series table +/// that is in the remote data folder but not in the local data folder. async fn validate_time_series_tables( local_data_folder: &DataFolder, remote_data_folder: &DataFolder,