From e33a9eb41ca5b11a7eaabe67eda0bbf32947961a Mon Sep 17 00:00:00 2001 From: ChilePiquin Date: Sun, 27 Sep 2026 12:43:45 -0700 Subject: [PATCH] feat(namespace): add in-memory directory manifest cache mode --- rust/lance-namespace-impls/src/dir.rs | 357 ++++++++++++++ .../lance-namespace-impls/src/dir/manifest.rs | 463 ++++++++++++++++-- 2 files changed, 784 insertions(+), 36 deletions(-) diff --git a/rust/lance-namespace-impls/src/dir.rs b/rust/lance-namespace-impls/src/dir.rs index e7d91d11e3d..30e6802c0d8 100644 --- a/rust/lance-namespace-impls/src/dir.rs +++ b/rust/lance-namespace-impls/src/dir.rs @@ -273,6 +273,7 @@ pub struct DirectoryNamespaceBuilder { storage_options: Option>, session: Option>, manifest_enabled: bool, + manifest_cache_mode: manifest::ManifestCacheMode, dir_listing_enabled: bool, inline_optimization_enabled: bool, table_version_tracking_enabled: bool, @@ -302,6 +303,7 @@ impl std::fmt::Debug for DirectoryNamespaceBuilder { .field("root", &self.root) .field("storage_options", &self.storage_options) .field("manifest_enabled", &self.manifest_enabled) + .field("manifest_cache_mode", &self.manifest_cache_mode) .field("dir_listing_enabled", &self.dir_listing_enabled) .field( "inline_optimization_enabled", @@ -344,6 +346,7 @@ impl DirectoryNamespaceBuilder { storage_options: None, session: None, manifest_enabled: true, + manifest_cache_mode: manifest::ManifestCacheMode::None, dir_listing_enabled: true, // Default to enabled for backwards compatibility inline_optimization_enabled: false, table_version_tracking_enabled: false, // Default to disabled @@ -366,6 +369,15 @@ impl DirectoryNamespaceBuilder { self } + /// Configure manifest read caching. + /// + /// Defaults to `none`. `in_memory` pins a snapshot of the manifest table + /// for the lifetime of this namespace instance. + pub fn manifest_cache_mode(mut self, mode: manifest::ManifestCacheMode) -> Self { + self.manifest_cache_mode = mode; + self + } + /// Enable or disable directory-based listing fallback. /// /// When enabled (default), falls back to directory scanning for tables not in the manifest. @@ -413,6 +425,7 @@ impl DirectoryNamespaceBuilder { /// It expects: /// - `root`: The root directory path (required) /// - `manifest_enabled`: Enable manifest-based table tracking (optional, default: true) + /// - `manifest_cache_mode`: Manifest read cache mode: `none` or `in_memory` (optional, default: none) /// - `dir_listing_enabled`: Enable directory listing for table discovery (optional, default: true) /// - `inline_optimization_enabled`: Enable replacement indices on __manifest rewrites (optional, default: false) /// - `storage.*`: Storage options (optional, prefix will be stripped) @@ -506,6 +519,12 @@ impl DirectoryNamespaceBuilder { .and_then(|v| str_to_bool(v)) .unwrap_or(true); + let manifest_cache_mode = properties + .get("manifest_cache_mode") + .map(|v| manifest::ManifestCacheMode::parse(v)) + .transpose()? + .unwrap_or_default(); + // Extract dir_listing_enabled (default: true) let dir_listing_enabled = properties .get("dir_listing_enabled") @@ -567,6 +586,7 @@ impl DirectoryNamespaceBuilder { storage_options, session, manifest_enabled, + manifest_cache_mode, dir_listing_enabled, inline_optimization_enabled, table_version_tracking_enabled, @@ -759,6 +779,7 @@ impl DirectoryNamespaceBuilder { self.dir_listing_enabled, self.inline_optimization_enabled, self.commit_retries, + self.manifest_cache_mode, ) .await { @@ -808,6 +829,7 @@ impl DirectoryNamespaceBuilder { manifest_ns: manifest_cell, write_manifest_ns: OnceCell::new(), manifest_enabled: self.manifest_enabled, + manifest_cache_mode: self.manifest_cache_mode, dir_listing_enabled: self.dir_listing_enabled, inline_optimization_enabled: self.inline_optimization_enabled, commit_retries: self.commit_retries, @@ -890,6 +912,7 @@ pub struct DirectoryNamespace { manifest_ns: OnceCell>, write_manifest_ns: OnceCell>, manifest_enabled: bool, + manifest_cache_mode: manifest::ManifestCacheMode, dir_listing_enabled: bool, inline_optimization_enabled: bool, commit_retries: Option, @@ -1020,10 +1043,32 @@ impl DirectoryNamespace { .or_else(|| self.manifest_ns.get()) } + fn is_in_memory_manifest_cache(&self) -> bool { + self.manifest_enabled + && matches!( + self.manifest_cache_mode, + manifest::ManifestCacheMode::InMemory + ) + } + + fn ensure_not_in_memory_cache_for_write(&self, operation: &str) -> Result<()> { + if self.is_in_memory_manifest_cache() { + return Err(NamespaceError::InvalidInput { + message: format!( + "{} is not supported when manifest_cache_mode is in_memory", + operation + ), + } + .into()); + } + Ok(()) + } + async fn manifest_ns_for_write(&self) -> Result>> { if !self.manifest_enabled { return Ok(None); } + self.ensure_not_in_memory_cache_for_write("write manifest")?; let manifest_ns = self .write_manifest_ns @@ -1037,6 +1082,7 @@ impl DirectoryNamespace { self.dir_listing_enabled, self.inline_optimization_enabled, self.commit_retries, + self.manifest_cache_mode, ) .await .map(Arc::new) @@ -1075,6 +1121,7 @@ impl DirectoryNamespace { self.dir_listing_enabled, self.inline_optimization_enabled, self.commit_retries, + self.manifest_cache_mode, ) .await .map(Arc::new) @@ -4017,6 +4064,7 @@ impl LanceNamespace for DirectoryNamespace { request: CreateTableVersionRequest, ) -> Result { self.record_op("create_table_version"); + self.ensure_not_in_memory_cache_for_write("create_table_version")?; let branch = Self::normalized_branch(request.branch.as_deref())?; let table_uri = self.resolve_table_location(&request.id).await?; let (table_uri, table_path, branch_parent_version) = match branch { @@ -4211,6 +4259,7 @@ impl LanceNamespace for DirectoryNamespace { request: BatchDeleteTableVersionsRequest, ) -> Result { self.record_op("batch_delete_table_versions"); + self.ensure_not_in_memory_cache_for_write("batch_delete_table_versions")?; let branch = Self::normalized_branch(request.branch.as_deref())?; // Single-table mode: use `id` (from path parameter) + `ranges` to delete // versions from one table. @@ -4265,6 +4314,7 @@ impl LanceNamespace for DirectoryNamespace { request: CreateTableIndexRequest, ) -> Result { self.record_op("create_table_index"); + self.ensure_not_in_memory_cache_for_write("create_table_index")?; let table_uri = self.resolve_table_location(&request.id).await?; let mut dataset = self .load_dataset(&table_uri, None, "create_table_index") @@ -4525,6 +4575,7 @@ impl LanceNamespace for DirectoryNamespace { request: AlterTransactionRequest, ) -> Result { self.record_op("alter_transaction"); + self.ensure_not_in_memory_cache_for_write("alter_transaction")?; // Parse the request ID: must include table id and transaction identifier let mut request_id = request.id.ok_or_else(|| { @@ -4748,6 +4799,7 @@ impl LanceNamespace for DirectoryNamespace { request: DropTableIndexRequest, ) -> Result { self.record_op("drop_table_index"); + self.ensure_not_in_memory_cache_for_write("drop_table_index")?; let table_uri = self.resolve_table_location(&request.id).await?; let index_name = request.index_name.as_deref().ok_or_else(|| { lance_core::Error::from(NamespaceError::InvalidInput { @@ -4818,6 +4870,7 @@ impl LanceNamespace for DirectoryNamespace { } async fn restore_table(&self, request: RestoreTableRequest) -> Result { + self.ensure_not_in_memory_cache_for_write("restore_table")?; let version = request.version; if version < 0 { return Err(Error::invalid_input_source( @@ -4883,6 +4936,7 @@ impl LanceNamespace for DirectoryNamespace { &self, request: UpdateTableSchemaMetadataRequest, ) -> Result { + self.ensure_not_in_memory_cache_for_write("update_table_schema_metadata")?; let table_uri = self.resolve_table_location(&request.id).await?; let mut dataset = self .load_dataset(&table_uri, None, "update_table_schema_metadata") @@ -5116,6 +5170,7 @@ impl LanceNamespace for DirectoryNamespace { request_data: Bytes, ) -> Result { self.record_op("insert_into_table"); + self.ensure_not_in_memory_cache_for_write("insert_into_table")?; let table_uri = self.resolve_table_location(&request.id).await?; let (reader, _num_rows) = Self::ipc_reader_from_request_data(&request_data, "insert_into_table")?; @@ -5155,6 +5210,7 @@ impl LanceNamespace for DirectoryNamespace { request_data: Bytes, ) -> Result { self.record_op("merge_insert_into_table"); + self.ensure_not_in_memory_cache_for_write("merge_insert_into_table")?; let table_uri = self.resolve_table_location(&request.id).await?; let on = merge_insert_on_columns(request.on.as_deref(), "merge_insert_into_table")?; @@ -5250,6 +5306,7 @@ impl LanceNamespace for DirectoryNamespace { async fn update_table(&self, request: UpdateTableRequest) -> Result { self.record_op("update_table"); + self.ensure_not_in_memory_cache_for_write("update_table")?; if request.updates.is_empty() { return Err(NamespaceError::InvalidInput { @@ -5338,6 +5395,7 @@ impl LanceNamespace for DirectoryNamespace { request: DeleteFromTableRequest, ) -> Result { self.record_op("delete_from_table"); + self.ensure_not_in_memory_cache_for_write("delete_from_table")?; if request.predicate.trim().is_empty() { return Err(NamespaceError::InvalidInput { @@ -5672,6 +5730,7 @@ impl LanceNamespace for DirectoryNamespace { request: CreateTableTagRequest, ) -> Result { self.record_op("create_table_tag"); + self.ensure_not_in_memory_cache_for_write("create_table_tag")?; if request.tag.is_empty() { return Err(NamespaceError::InvalidInput { message: "tag name must not be empty for create_table_tag".to_string(), @@ -5710,6 +5769,7 @@ impl LanceNamespace for DirectoryNamespace { request: DeleteTableTagRequest, ) -> Result { self.record_op("delete_table_tag"); + self.ensure_not_in_memory_cache_for_write("delete_table_tag")?; if request.tag.is_empty() { return Err(NamespaceError::InvalidInput { message: "tag name must not be empty for delete_table_tag".to_string(), @@ -5739,6 +5799,7 @@ impl LanceNamespace for DirectoryNamespace { request: UpdateTableTagRequest, ) -> Result { self.record_op("update_table_tag"); + self.ensure_not_in_memory_cache_for_write("update_table_tag")?; if request.tag.is_empty() { return Err(NamespaceError::InvalidInput { message: "tag name must not be empty for update_table_tag".to_string(), @@ -5777,6 +5838,7 @@ impl LanceNamespace for DirectoryNamespace { request: CreateTableBranchRequest, ) -> Result { self.record_op("create_table_branch"); + self.ensure_not_in_memory_cache_for_write("create_table_branch")?; if request.name.is_empty() { return Err(NamespaceError::InvalidInput { message: "branch name must not be empty for create_table_branch".to_string(), @@ -5893,6 +5955,7 @@ impl LanceNamespace for DirectoryNamespace { request: DeleteTableBranchRequest, ) -> Result { self.record_op("delete_table_branch"); + self.ensure_not_in_memory_cache_for_write("delete_table_branch")?; if request.name.is_empty() { return Err(NamespaceError::InvalidInput { message: "branch name must not be empty for delete_table_branch".to_string(), @@ -9783,9 +9846,40 @@ mod tests { let builder = DirectoryNamespaceBuilder::from_properties(properties, None).unwrap(); assert!(builder.manifest_enabled); assert!(builder.dir_listing_enabled); + assert_eq!( + builder.manifest_cache_mode, + manifest::ManifestCacheMode::None + ); assert!(!builder.inline_optimization_enabled); } + #[tokio::test] + async fn test_from_properties_manifest_cache_mode() { + let temp_dir = TempStdDir::default(); + + let mut properties = HashMap::new(); + properties.insert("root".to_string(), temp_dir.to_str().unwrap().to_string()); + properties.insert("manifest_cache_mode".to_string(), "in_memory".to_string()); + + let builder = DirectoryNamespaceBuilder::from_properties(properties, None).unwrap(); + assert_eq!( + builder.manifest_cache_mode, + manifest::ManifestCacheMode::InMemory + ); + } + + #[tokio::test] + async fn test_from_properties_rejects_invalid_manifest_cache_mode() { + let temp_dir = TempStdDir::default(); + + let mut properties = HashMap::new(); + properties.insert("root".to_string(), temp_dir.to_str().unwrap().to_string()); + properties.insert("manifest_cache_mode".to_string(), "pinned".to_string()); + + let err = DirectoryNamespaceBuilder::from_properties(properties, None).unwrap_err(); + assert!(err.to_string().contains("manifest_cache_mode")); + } + #[test] fn test_builder_disables_inline_optimization_by_default() { let builder = DirectoryNamespaceBuilder::new("memory://"); @@ -10107,6 +10201,269 @@ mod tests { assert_eq!(tables, vec!["table1".to_string()]); } + #[tokio::test] + async fn test_in_memory_manifest_cache_pins_namespace_snapshot() { + let temp_dir = TempStdDir::default(); + let root = temp_dir.to_str().unwrap(); + + let writer = DirectoryNamespaceBuilder::new(root) + .dir_listing_enabled(false) + .build() + .await + .unwrap(); + let ipc_data = create_test_ipc_data(&create_test_schema()); + let mut create_table_req = CreateTableRequest::new(); + create_table_req.id = Some(vec!["table1".to_string()]); + writer + .create_table(create_table_req, bytes::Bytes::from(ipc_data)) + .await + .unwrap(); + + let reader = DirectoryNamespaceBuilder::new(root) + .dir_listing_enabled(false) + .manifest_cache_mode(manifest::ManifestCacheMode::InMemory) + .build() + .await + .unwrap(); + let list_req = ListTablesRequest { + id: Some(vec![]), + ..Default::default() + }; + assert_eq!( + reader.list_tables(list_req.clone()).await.unwrap().tables, + vec!["table1".to_string()] + ); + + let ipc_data = create_test_ipc_data(&create_test_schema()); + let mut create_table_req = CreateTableRequest::new(); + create_table_req.id = Some(vec!["table2".to_string()]); + writer + .create_table(create_table_req, bytes::Bytes::from(ipc_data)) + .await + .unwrap(); + + assert_eq!( + reader.list_tables(list_req.clone()).await.unwrap().tables, + vec!["table1".to_string()] + ); + + let new_reader = DirectoryNamespaceBuilder::new(root) + .dir_listing_enabled(false) + .manifest_cache_mode(manifest::ManifestCacheMode::InMemory) + .build() + .await + .unwrap(); + assert_eq!( + new_reader.list_tables(list_req).await.unwrap().tables, + vec!["table1".to_string(), "table2".to_string()] + ); + } + + #[tokio::test] + async fn test_in_memory_manifest_cache_rejects_writes() { + let temp_dir = TempStdDir::default(); + let root = temp_dir.to_str().unwrap(); + + let writer = DirectoryNamespaceBuilder::new(root) + .dir_listing_enabled(false) + .build() + .await + .unwrap(); + let mut parent = CreateNamespaceRequest::new(); + parent.id = Some(vec!["parent".to_string()]); + writer.create_namespace(parent).await.unwrap(); + let mut empty = CreateNamespaceRequest::new(); + empty.id = Some(vec!["empty".to_string()]); + writer.create_namespace(empty).await.unwrap(); + + let cached = DirectoryNamespaceBuilder::new(root) + .dir_listing_enabled(false) + .manifest_cache_mode(manifest::ManifestCacheMode::InMemory) + .build() + .await + .unwrap(); + + let mut parent_id = NamespaceExistsRequest::new(); + parent_id.id = Some(vec!["parent".to_string()]); + cached + .namespace_exists(parent_id) + .await + .expect("cached reader should see parent in its pinned snapshot"); + + let mut child = CreateNamespaceRequest::new(); + child.id = Some(vec!["parent".to_string(), "child".to_string()]); + writer.create_namespace(child).await.unwrap(); + + let mut drop_parent = DropNamespaceRequest::new(); + drop_parent.id = Some(vec!["parent".to_string()]); + let err = cached + .drop_namespace(drop_parent) + .await + .expect_err("in_memory namespace instances must not orphan children"); + let err = err.to_string(); + assert!(err.contains("in_memory"), "{err}"); + + let verifier = DirectoryNamespaceBuilder::new(root) + .dir_listing_enabled(false) + .build() + .await + .unwrap(); + let mut parent_id = NamespaceExistsRequest::new(); + parent_id.id = Some(vec!["parent".to_string()]); + verifier.namespace_exists(parent_id).await.unwrap(); + let mut child_id = NamespaceExistsRequest::new(); + child_id.id = Some(vec!["parent".to_string(), "child".to_string()]); + verifier.namespace_exists(child_id).await.unwrap(); + + let mut drop_empty = DropNamespaceRequest::new(); + drop_empty.id = Some(vec!["empty".to_string()]); + let err = cached + .drop_namespace(drop_empty) + .await + .expect_err("in_memory namespace instances must reject namespace drops"); + let err = err.to_string(); + assert!(err.contains("in_memory"), "{err}"); + + let mut create_ns_req = CreateNamespaceRequest::new(); + create_ns_req.id = Some(vec!["sibling".to_string()]); + let err = cached + .create_namespace(create_ns_req) + .await + .expect_err("in_memory namespace instances must reject namespace creation"); + let err = err.to_string(); + assert!(err.contains("in_memory"), "{err}"); + + let mut create_table_req = CreateTableRequest::new(); + create_table_req.id = Some(vec!["table1".to_string()]); + let err = cached + .create_table( + create_table_req, + bytes::Bytes::from(create_test_ipc_data(&create_test_schema())), + ) + .await + .expect_err("in_memory namespace instances must reject table creation"); + let err = err.to_string(); + assert!(err.contains("in_memory"), "{err}"); + } + + #[tokio::test] + async fn test_in_memory_manifest_cache_rejects_data_write_after_deregister() { + let temp_dir = TempStdDir::default(); + let root = temp_dir.to_str().unwrap(); + + let writer = DirectoryNamespaceBuilder::new(root) + .dir_listing_enabled(false) + .build() + .await + .unwrap(); + let mut create = CreateTableRequest::new(); + create.id = Some(vec!["table1".to_string()]); + writer + .create_table(create, bytes::Bytes::from(create_non_empty_test_ipc_data())) + .await + .unwrap(); + + let cached = DirectoryNamespaceBuilder::new(root) + .dir_listing_enabled(false) + .manifest_cache_mode(manifest::ManifestCacheMode::InMemory) + .build() + .await + .unwrap(); + let mut list = ListTablesRequest::new(); + list.id = Some(vec![]); + assert_eq!( + cached.list_tables(list).await.unwrap().tables, + vec!["table1".to_string()] + ); + + let mut deregister = lance_namespace::models::DeregisterTableRequest::new(); + deregister.id = Some(vec!["table1".to_string()]); + let location = writer + .deregister_table(deregister) + .await + .unwrap() + .location + .unwrap(); + + let delete = DeleteFromTableRequest { + id: Some(vec!["table1".to_string()]), + predicate: "id = 1".to_string(), + ..Default::default() + }; + let err = cached + .delete_from_table(delete) + .await + .expect_err("in_memory namespace instances must reject data writes"); + let err = err.to_string(); + assert!(err.contains("in_memory"), "{err}"); + + let row_count = Dataset::open(&location) + .await + .unwrap() + .count_rows(None) + .await + .unwrap(); + assert_eq!(row_count, 2); + } + + #[tokio::test] + async fn test_in_memory_manifest_cache_failed_write_keeps_pinned_snapshot() { + let temp_dir = TempStdDir::default(); + let root = temp_dir.to_str().unwrap(); + + let writer = DirectoryNamespaceBuilder::new(root) + .dir_listing_enabled(false) + .build() + .await + .unwrap(); + let mut first = CreateTableRequest::new(); + first.id = Some(vec!["first".to_string()]); + writer + .create_table( + first, + bytes::Bytes::from(create_test_ipc_data(&create_test_schema())), + ) + .await + .unwrap(); + + let cached = DirectoryNamespaceBuilder::new(root) + .dir_listing_enabled(false) + .manifest_cache_mode(manifest::ManifestCacheMode::InMemory) + .build() + .await + .unwrap(); + let mut list = ListTablesRequest::new(); + list.id = Some(vec![]); + assert_eq!( + cached.list_tables(list.clone()).await.unwrap().tables, + vec!["first".to_string()] + ); + + let mut second = CreateTableRequest::new(); + second.id = Some(vec!["second".to_string()]); + writer + .create_table( + second, + bytes::Bytes::from(create_test_ipc_data(&create_test_schema())), + ) + .await + .unwrap(); + + let mut create_namespace = CreateNamespaceRequest::new(); + create_namespace.id = Some(vec!["other".to_string()]); + let err = cached + .create_namespace(create_namespace) + .await + .expect_err("in_memory namespace instances must reject catalog writes"); + let err = err.to_string(); + assert!(err.contains("in_memory"), "{err}"); + + assert_eq!( + cached.list_tables(list).await.unwrap().tables, + vec!["first".to_string()] + ); + } + /// Migration mode promises manifest-first lookup even at the root, so a /// reader built before `__manifest` existed must still resolve a /// manifest-only alias (`registered_table` -> `external_table.lance`) that diff --git a/rust/lance-namespace-impls/src/dir/manifest.rs b/rust/lance-namespace-impls/src/dir/manifest.rs index 123f7e314fe..b4ffe90e99e 100644 --- a/rust/lance-namespace-impls/src/dir/manifest.rs +++ b/rust/lance-namespace-impls/src/dir/manifest.rs @@ -227,6 +227,7 @@ struct ManifestTrainedIndex { created_index: CreatedIndex, } +#[derive(Debug, Clone)] struct ManifestRowValue { object_id: String, object_type: ObjectType, @@ -235,6 +236,143 @@ struct ManifestRowValue { base_objects: Option>, } +/// Read cache behavior for the directory manifest table. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub enum ManifestCacheMode { + /// Keep the existing behavior: check whether the manifest table has advanced + /// before each read and scan the current dataset. + #[default] + None, + /// Load the manifest rows into memory and reuse that snapshot for this + /// namespace instance. New namespace instances observe newer commits. + InMemory, +} + +impl ManifestCacheMode { + pub fn parse(value: &str) -> Result { + match value { + "none" => Ok(Self::None), + "in_memory" => Ok(Self::InMemory), + other => Err(NamespaceError::InvalidInput { + message: format!( + "Invalid manifest_cache_mode '{}'; expected 'none' or 'in_memory'", + other + ), + } + .into()), + } + } + + fn is_in_memory(self) -> bool { + matches!(self, Self::InMemory) + } +} + +#[derive(Debug)] +struct ManifestSnapshot { + rows: Arc<[ManifestRowValue]>, +} + +impl ManifestSnapshot { + async fn load(dataset: &Dataset) -> Result { + let mut scanner = dataset.scan(); + scanner + .project(&[ + "object_id", + "object_type", + "location", + "metadata", + "base_objects", + ]) + .map_err(|e| { + lance_core::Error::from(NamespaceError::Internal { + message: format!("Failed to project manifest columns: {:?}", e), + }) + })?; + + let batches = ManifestNamespace::execute_scanner(scanner).await?; + let mut rows = Vec::new(); + for batch in batches { + if batch.num_rows() == 0 { + continue; + } + let object_ids = ManifestNamespace::get_string_column(&batch, "object_id")?; + let object_types = ManifestNamespace::get_string_column(&batch, "object_type")?; + let locations = ManifestNamespace::get_string_column(&batch, "location")?; + let metadatas = ManifestNamespace::get_string_column(&batch, "metadata")?; + let base_objects = ManifestNamespace::base_objects_column_values(&batch)?; + for (row, base_objects) in base_objects.into_iter().enumerate().take(batch.num_rows()) { + rows.push(ManifestRowValue { + object_id: ManifestNamespace::required_string_value( + object_ids, + row, + "object_id", + )? + .to_string(), + object_type: ObjectType::parse(ManifestNamespace::required_string_value( + object_types, + row, + "object_type", + )?)?, + location: ManifestNamespace::optional_string_value(locations, row), + metadata: ManifestNamespace::optional_string_value(metadatas, row), + base_objects, + }); + } + } + + rows.sort_by(|left, right| left.object_id.cmp(&right.object_id)); + if let Some(duplicate) = rows + .windows(2) + .find(|pair| pair[0].object_id == pair[1].object_id) + .map(|pair| pair[0].object_id.clone()) + { + return Err(NamespaceError::Internal { + message: format!("Manifest contains duplicate object_id '{}'", duplicate), + } + .into()); + } + + Ok(Self { rows: rows.into() }) + } + + fn get(&self, object_id: &str) -> Option<&ManifestRowValue> { + self.rows + .binary_search_by(|row| row.object_id.as_str().cmp(object_id)) + .ok() + .map(|index| &self.rows[index]) + } + + fn direct_children<'a>( + &'a self, + parent: &[String], + object_type: ObjectType, + ) -> impl Iterator + 'a { + let parent_prefix = + (!parent.is_empty()).then(|| format!("{}{}", parent.join(DELIMITER), DELIMITER)); + self.rows.iter().filter(move |row| { + if row.object_type != object_type { + return false; + } + match &parent_prefix { + Some(prefix) => row + .object_id + .strip_prefix(prefix) + .is_some_and(|suffix| !suffix.contains(DELIMITER)), + None => !row.object_id.contains(DELIMITER), + } + }) + } + + fn descendant_count(&self, object_id: &str) -> usize { + let prefix = format!("{}{}", object_id, DELIMITER); + self.rows + .iter() + .filter(|row| row.object_id.starts_with(&prefix)) + .count() + } +} + struct ManifestOutputRow<'a> { object_id: &'a str, object_type: ObjectType, @@ -600,12 +738,19 @@ pub struct NamespaceInfo { /// The manifest dataset uses contiguous attached versions and this module never /// runs old-version cleanup on it, allowing reads to check only the immediate /// successor manifest before deciding whether a reload is needed. +#[derive(Debug)] +struct ManifestDatasetState { + dataset: Dataset, + snapshot: Option>, + cache_mode: ManifestCacheMode, +} + #[derive(Debug, Clone)] -pub struct DatasetConsistencyWrapper(Arc>); +pub struct DatasetConsistencyWrapper(Arc>); impl DatasetConsistencyWrapper { /// Create a new wrapper with the given dataset. - pub fn new(dataset: Dataset) -> Self { + pub fn new(dataset: Dataset, cache_mode: ManifestCacheMode) -> Self { debug_assert!( !dataset .manifest() @@ -614,13 +759,21 @@ impl DatasetConsistencyWrapper { .any(|key| key.starts_with("lance.auto_cleanup.")), "the directory manifest dataset must not enable old-version cleanup" ); - Self(Arc::new(RwLock::new(dataset))) + Self(Arc::new(RwLock::new(ManifestDatasetState { + dataset, + snapshot: None, + cache_mode, + }))) } /// Get an immutable reference to the dataset. /// Always reloads to ensure strong consistency. pub async fn get(&self) -> Result> { - self.reload().await?; + if self.cache_mode().await.is_in_memory() { + self.ensure_snapshot().await?; + } else { + self.reload().await?; + } let guard = DatasetReadGuard { guard: self.0.read().await, }; @@ -658,17 +811,51 @@ impl DatasetConsistencyWrapper { /// have the latest version. pub async fn set_latest(&self, dataset: Dataset) { let mut write_guard = self.0.write().await; - if dataset.manifest().version > write_guard.manifest().version { - *write_guard = dataset; + if dataset.manifest().version > write_guard.dataset.manifest().version { + write_guard.dataset = dataset; + write_guard.snapshot = None; } } + async fn snapshot(&self) -> Result>> { + if !self.cache_mode().await.is_in_memory() { + return Ok(None); + } + self.ensure_snapshot().await?; + Ok(self.0.read().await.snapshot.clone()) + } + + async fn cache_mode(&self) -> ManifestCacheMode { + self.0.read().await.cache_mode + } + + async fn is_in_memory_cache(&self) -> bool { + self.cache_mode().await.is_in_memory() + } + + async fn ensure_snapshot(&self) -> Result<()> { + { + let read_guard = self.0.read().await; + if !read_guard.cache_mode.is_in_memory() || read_guard.snapshot.is_some() { + return Ok(()); + } + } + + let mut write_guard = self.0.write().await; + if write_guard.cache_mode.is_in_memory() && write_guard.snapshot.is_none() { + write_guard.snapshot = Some(Arc::new( + ManifestSnapshot::load(&write_guard.dataset).await?, + )); + } + Ok(()) + } + /// Reload the dataset to the latest version. async fn reload(&self) -> Result<()> { // First check if we need to reload (with read lock) let read_guard = self.0.read().await; - let dataset_uri = read_guard.uri().to_string(); - let current_version = read_guard.version().version; + let dataset_uri = read_guard.dataset.uri().to_string(); + let current_version = read_guard.dataset.version().version; log::debug!( "Reload starting for uri={}, current_version={}", dataset_uri, @@ -678,11 +865,16 @@ impl DatasetConsistencyWrapper { // does not run old-version cleanup, so the immediate successor probe is // enough to detect changes without resolving or loading the latest // manifest on every namespace read. - let has_successor_version = read_guard.has_successor_version().await.map_err(|e| { - lance_core::Error::from(NamespaceError::Internal { - message: format!("Failed to check dataset staleness: {:?}", e), - }) - })?; + let has_successor_version = + read_guard + .dataset + .has_successor_version() + .await + .map_err(|e| { + lance_core::Error::from(NamespaceError::Internal { + message: format!("Failed to check dataset staleness: {:?}", e), + }) + })?; log::debug!( "Reload checked successor_version_exists={} for uri={}, current_version={}", has_successor_version, @@ -701,18 +893,30 @@ impl DatasetConsistencyWrapper { let mut write_guard = self.0.write().await; // Double-check after acquiring write lock (someone else might have reloaded) - let has_successor_version = write_guard.has_successor_version().await.map_err(|e| { - lance_core::Error::from(NamespaceError::Internal { - message: format!("Failed to check dataset staleness: {:?}", e), - }) - })?; + let has_successor_version = + write_guard + .dataset + .has_successor_version() + .await + .map_err(|e| { + lance_core::Error::from(NamespaceError::Internal { + message: format!("Failed to check dataset staleness: {:?}", e), + }) + })?; if has_successor_version { - write_guard.checkout_latest().await.map_err(|e| { + write_guard.dataset.checkout_latest().await.map_err(|e| { lance_core::Error::from(NamespaceError::Internal { message: format!("Failed to checkout latest: {:?}", e), }) })?; + write_guard.snapshot = None; + } + + if write_guard.cache_mode.is_in_memory() { + write_guard.snapshot = Some(Arc::new( + ManifestSnapshot::load(&write_guard.dataset).await?, + )); } Ok(()) @@ -720,32 +924,32 @@ impl DatasetConsistencyWrapper { } pub struct DatasetReadGuard<'a> { - guard: RwLockReadGuard<'a, Dataset>, + guard: RwLockReadGuard<'a, ManifestDatasetState>, } impl Deref for DatasetReadGuard<'_> { type Target = Dataset; fn deref(&self) -> &Self::Target { - &self.guard + &self.guard.dataset } } pub struct DatasetWriteGuard<'a> { - guard: RwLockWriteGuard<'a, Dataset>, + guard: RwLockWriteGuard<'a, ManifestDatasetState>, } impl Deref for DatasetWriteGuard<'_> { type Target = Dataset; fn deref(&self) -> &Self::Target { - &self.guard + &self.guard.dataset } } impl DerefMut for DatasetWriteGuard<'_> { fn deref_mut(&mut self) -> &mut Self::Target { - &mut self.guard + &mut self.guard.dataset } } @@ -858,10 +1062,15 @@ impl ManifestNamespace { dir_listing_enabled: bool, inline_optimization_enabled: bool, commit_retries: Option, + manifest_cache_mode: ManifestCacheMode, ) -> Result { - let manifest_dataset = - Self::ensure_manifest_table_up_to_date(&root, &storage_options, session.clone()) - .await?; + let manifest_dataset = Self::ensure_manifest_table_up_to_date( + &root, + &storage_options, + session.clone(), + manifest_cache_mode, + ) + .await?; Ok(Self::new( root, @@ -887,9 +1096,15 @@ impl ManifestNamespace { dir_listing_enabled: bool, inline_optimization_enabled: bool, commit_retries: Option, + manifest_cache_mode: ManifestCacheMode, ) -> Result { - let manifest_dataset = - Self::open_manifest_table(&root, &storage_options, session.clone()).await?; + let manifest_dataset = Self::open_manifest_table( + &root, + &storage_options, + session.clone(), + manifest_cache_mode, + ) + .await?; Ok(Self::new( root, @@ -1932,11 +2147,26 @@ impl ManifestNamespace { /// `rewrite_manifest` commit re-checks `ensure_writable` on each retry, so a /// concurrent upgrade in between is still caught. async fn ensure_manifest_writable(&self) -> Result<()> { + self.ensure_not_in_memory_cache_for_write("write manifest") + .await?; let dataset_guard = self.manifest_dataset.get().await?; ensure_can_write_manifest(dataset_guard.manifest())?; ensure_writable(dataset_guard.metadata()) } + async fn ensure_not_in_memory_cache_for_write(&self, operation: &str) -> Result<()> { + if self.manifest_dataset.is_in_memory_cache().await { + return Err(NamespaceError::InvalidInput { + message: format!( + "{} is not supported when manifest_cache_mode is in_memory", + operation + ), + } + .into()); + } + Ok(()) + } + async fn rewrite_manifest( &self, operation: &str, @@ -1946,6 +2176,8 @@ impl ManifestNamespace { M: ManifestStreamMutation + 'static, F: FnMut() -> M, { + self.ensure_not_in_memory_cache_for_write("rewrite manifest") + .await?; let _mutation_guard = self.manifest_mutation_lock.lock().await; let max_retries = self.manifest_rewrite_commit_retries(); let mut retries = 0; @@ -2114,6 +2346,10 @@ impl ManifestNamespace { /// Check if the manifest contains an object with the given ID async fn manifest_contains_object(&self, object_id: &str) -> Result { + if let Some(snapshot) = self.manifest_dataset.snapshot().await? { + return Ok(snapshot.get(object_id).is_some()); + } + let escaped_id = object_id.replace('\'', "''"); let filter = format!("object_id = '{}'", escaped_id); @@ -2146,6 +2382,14 @@ impl ManifestNamespace { /// Query the manifest for a table with the given object ID async fn query_manifest_for_table(&self, object_id: &str) -> Result> { + if let Some(snapshot) = self.manifest_dataset.snapshot().await? { + return snapshot + .get(object_id) + .filter(|row| row.object_type == ObjectType::Table) + .map(Self::table_info_from_manifest_row) + .transpose(); + } + let escaped_id = object_id.replace('\'', "''"); let filter = format!("object_id = '{}' AND object_type = 'table'", escaped_id); let mut scanner = self.manifest_scanner().await?; @@ -2215,6 +2459,58 @@ impl ManifestNamespace { Ok(found_result) } + fn deserialize_manifest_metadata( + object_type: &str, + object_id: &str, + metadata: Option<&str>, + ) -> Result>> { + let Some(metadata_str) = metadata else { + return Ok(None); + }; + serde_json::from_str::>(metadata_str) + .map(Some) + .map_err(|e| { + NamespaceError::Internal { + message: format!( + "Failed to deserialize metadata for {} '{}': {}", + object_type, object_id, e + ), + } + .into() + }) + } + + fn table_info_from_manifest_row(row: &ManifestRowValue) -> Result { + let location = row.location.clone().ok_or_else(|| { + lance_core::Error::from(NamespaceError::Internal { + message: format!("Manifest table '{}' has no location", row.object_id), + }) + })?; + let metadata = + Self::deserialize_manifest_metadata("table", &row.object_id, row.metadata.as_deref())?; + let (namespace, name) = Self::parse_object_id(&row.object_id); + Ok(TableInfo { + namespace, + name, + location, + metadata, + }) + } + + fn namespace_info_from_manifest_row(row: &ManifestRowValue) -> Result { + let metadata = Self::deserialize_manifest_metadata( + "namespace", + &row.object_id, + row.metadata.as_deref(), + )?; + let (namespace, name) = Self::parse_object_id(&row.object_id); + Ok(NamespaceInfo { + namespace, + name, + metadata, + }) + } + fn serialize_metadata( properties: Option<&HashMap>, object_type: &str, @@ -2277,6 +2573,13 @@ impl ManifestNamespace { /// List all table locations in the manifest (for root namespace only) /// Returns a set of table locations (e.g., "table_name.lance") pub async fn list_manifest_table_locations(&self) -> Result> { + if let Some(snapshot) = self.manifest_dataset.snapshot().await? { + return Ok(snapshot + .direct_children(&[], ObjectType::Table) + .filter_map(|row| row.location.clone()) + .collect()); + } + let filter = "object_type = 'table' AND NOT contains(object_id, '$')"; let mut scanner = self.manifest_scanner().await?; scanner.filter(filter).map_err(|e| { @@ -2404,6 +2707,14 @@ impl ManifestNamespace { /// Query the manifest for a namespace with the given object ID async fn query_manifest_for_namespace(&self, object_id: &str) -> Result> { + if let Some(snapshot) = self.manifest_dataset.snapshot().await? { + return snapshot + .get(object_id) + .filter(|row| row.object_type == ObjectType::Namespace) + .map(Self::namespace_info_from_manifest_row) + .transpose(); + } + let escaped_id = object_id.replace('\'', "''"); let filter = format!("object_id = '{}' AND object_type = 'namespace'", escaped_id); let mut scanner = self.manifest_scanner().await?; @@ -2476,6 +2787,7 @@ impl ManifestNamespace { root: &str, storage_options: &Option>, session: Option>, + manifest_cache_mode: ManifestCacheMode, ) -> Result { let manifest_path = format!("{}/{}", root, MANIFEST_TABLE_NAME); log::debug!("Attempting to load manifest from {}", manifest_path); @@ -2499,7 +2811,7 @@ impl ManifestNamespace { .load() .await?; ensure_readable(dataset.metadata())?; - Ok(DatasetConsistencyWrapper::new(dataset)) + Ok(DatasetConsistencyWrapper::new(dataset, manifest_cache_mode)) } /// Create or load the manifest dataset, ensuring it has the latest schema setup. @@ -2512,6 +2824,7 @@ impl ManifestNamespace { root: &str, storage_options: &Option>, session: Option>, + manifest_cache_mode: ManifestCacheMode, ) -> Result { let manifest_path = format!("{}/{}", root, MANIFEST_TABLE_NAME); log::debug!("Attempting to load manifest from {}", manifest_path); @@ -2576,7 +2889,7 @@ impl ManifestNamespace { })?; } - Ok(DatasetConsistencyWrapper::new(dataset)) + Ok(DatasetConsistencyWrapper::new(dataset, manifest_cache_mode)) } Err(err) if Self::is_not_found_load_error(&err) => { log::info!("Creating new manifest table at {}", manifest_path); @@ -2612,7 +2925,7 @@ impl ManifestNamespace { dataset.version().version, dataset.uri() ); - Ok(DatasetConsistencyWrapper::new(dataset)) + Ok(DatasetConsistencyWrapper::new(dataset, manifest_cache_mode)) } Err(ref e) if matches!( @@ -2655,7 +2968,7 @@ impl ManifestNamespace { ), }) })?; - Ok(DatasetConsistencyWrapper::new(dataset)) + Ok(DatasetConsistencyWrapper::new(dataset, manifest_cache_mode)) } Err(e) => Err(lance_core::Error::from(NamespaceError::Internal { message: format!("Failed to create manifest dataset: {:?}", e), @@ -2720,6 +3033,50 @@ impl LanceNamespace for ManifestNamespace { }) })?; + if let Some(snapshot) = self.manifest_dataset.snapshot().await? { + let table_entries: Vec<(String, String)> = snapshot + .direct_children(namespace_id, ObjectType::Table) + .map(|row| { + let (_namespace, name) = Self::parse_object_id(&row.object_id); + let location = row.location.clone().ok_or_else(|| { + lance_core::Error::from(NamespaceError::Internal { + message: format!("Manifest table '{}' has no location", row.object_id), + }) + })?; + Ok::<_, Error>((name, location)) + }) + .collect::>()?; + + let mut tables: Vec = if request.include_declared.unwrap_or(true) { + table_entries.into_iter().map(|(name, _)| name).collect() + } else { + let mut stream = futures::stream::iter(table_entries.into_iter().map( + |(name, location)| async move { + if self.location_has_actual_manifests(&location).await? { + Ok::, Error>(Some(name)) + } else { + Ok::, Error>(None) + } + }, + )) + .buffered(DECLARED_FILTER_CONCURRENCY); + + let mut filtered = Vec::new(); + while let Some(result) = stream.next().await { + if let Some(name) = result? { + filtered.push(name); + } + } + filtered + }; + + let next_page_token = + Self::apply_pagination(&mut tables, request.page_token, request.limit); + let mut response = ListTablesResponse::new(tables); + response.page_token = next_page_token; + return Ok(response); + } + // Build filter to find tables in this namespace let filter = if namespace_id.is_empty() { // Root namespace: find tables without a namespace prefix @@ -3222,6 +3579,21 @@ impl LanceNamespace for ManifestNamespace { }) })?; + if let Some(snapshot) = self.manifest_dataset.snapshot().await? { + let mut namespaces: Vec = snapshot + .direct_children(parent_namespace, ObjectType::Namespace) + .map(|row| { + let (_namespace, name) = Self::parse_object_id(&row.object_id); + name + }) + .collect(); + let next_page_token = + Self::apply_pagination(&mut namespaces, request.page_token, request.limit); + let mut response = ListNamespacesResponse::new(namespaces); + response.page_token = next_page_token; + return Ok(response); + } + // Build filter to find direct child namespaces let filter = if parent_namespace.is_empty() { // Root namespace: find all namespaces without a parent @@ -3385,6 +3757,18 @@ impl LanceNamespace for ManifestNamespace { .into()); } + if let Some(snapshot) = self.manifest_dataset.snapshot().await? { + let count = snapshot.descendant_count(&object_id); + if count > 0 { + return Err(NamespaceError::NamespaceNotEmpty { + message: format!("'{}' (contains {} child objects)", object_id, count), + } + .into()); + } + self.delete_from_manifest(&object_id).boxed().await?; + return Ok(DropNamespaceResponse::default()); + } + // Check for child namespaces let escaped_id = object_id.replace('\'', "''"); let prefix = format!("{}{}", escaped_id, DELIMITER); @@ -3681,6 +4065,8 @@ impl LanceNamespace for ManifestNamespace { &self, request: AlterTableAddColumnsRequest, ) -> Result { + self.ensure_not_in_memory_cache_for_write("alter_table_add_columns") + .await?; let table_id = request .id .as_ref() @@ -3746,6 +4132,8 @@ impl LanceNamespace for ManifestNamespace { &self, request: AlterTableAlterColumnsRequest, ) -> Result { + self.ensure_not_in_memory_cache_for_write("alter_table_alter_columns") + .await?; let table_id = request .id .as_ref() @@ -3798,6 +4186,8 @@ impl LanceNamespace for ManifestNamespace { &self, request: AlterTableDropColumnsRequest, ) -> Result { + self.ensure_not_in_memory_cache_for_write("alter_table_drop_columns") + .await?; let table_id = request .id .as_ref() @@ -3847,9 +4237,9 @@ mod tests { use super::{ BASE_OBJECTS_INDEX_NAME, ConflictResolution, CopyOnWriteMutation, DeleteObjectMutation, LANCE_DATA_DIR, LANCE_INDICES_DIR, MANIFEST_TABLE_NAME, ManifestBatchBuilder, - ManifestEntry, ManifestIndexAccumulator, ManifestNamespace, ManifestOutputRow, - ManifestRowValue, ManifestStreamMutation, OBJECT_ID_INDEX_NAME, OBJECT_TYPE_INDEX_NAME, - ObjectType, + ManifestCacheMode, ManifestEntry, ManifestIndexAccumulator, ManifestNamespace, + ManifestOutputRow, ManifestRowValue, ManifestStreamMutation, OBJECT_ID_INDEX_NAME, + OBJECT_TYPE_INDEX_NAME, ObjectType, }; use crate::DirectoryNamespaceBuilder; use arrow::datatypes::DataType; @@ -3897,6 +4287,7 @@ mod tests { true, inline_optimization_enabled, commit_retries, + ManifestCacheMode::None, ) .await .unwrap()