From 81e7e71ddea20ab016c616dbb202991c8af6c593 Mon Sep 17 00:00:00 2001 From: Yang Cen Date: Tue, 11 Aug 2026 02:18:01 +0800 Subject: [PATCH] fix(storage): recover mirrored manifest creates --- rust/lancedb/src/io/object_store.rs | 72 ++++++++++++++++++++++++++++- 1 file changed, 70 insertions(+), 2 deletions(-) diff --git a/rust/lancedb/src/io/object_store.rs b/rust/lancedb/src/io/object_store.rs index d594bd857..f889de7f1 100644 --- a/rust/lancedb/src/io/object_store.rs +++ b/rust/lancedb/src/io/object_store.rs @@ -8,7 +8,7 @@ use std::{fmt::Formatter, sync::Arc}; use futures::{StreamExt, TryFutureExt, stream::BoxStream}; use lance::io::WrappingObjectStore; use object_store::{ - CopyOptions, Error, GetOptions, GetResult, ListResult, MultipartUpload, ObjectMeta, + CopyMode, CopyOptions, Error, GetOptions, GetResult, ListResult, MultipartUpload, ObjectMeta, ObjectStore, ObjectStoreExt, PutMultipartOptions, PutOptions, PutPayload, PutResult, Result, UploadPart, path::Path, }; @@ -52,7 +52,7 @@ impl PrimaryOnly for Path { /// store. We have primary store that is durable but slow, and a secondary /// store that is fast but not asdurable /// -/// Note: this object store does not mirror writes to *.manifest files +/// Note: this object store does not mirror writes to `_latest.manifest`. #[async_trait] impl ObjectStore for MirroringObjectStore { async fn put_opts( @@ -137,6 +137,16 @@ impl ObjectStore for MirroringObjectStore { // or may be evicted before the copy begins. match self.secondary.copy_opts(from, to, options.clone()).await { Ok(()) | Err(Error::NotFound { .. }) => {} + // The secondary is a non-authoritative cache. A process can + // leave an orphaned manifest there if it exits before creating + // the durable primary object, so let the primary decide the + // outcome of create-only manifest copies. + Err(Error::AlreadyExists { .. } | Error::Precondition { .. }) + if options.mode == CopyMode::Create + && to + .filename() + .map(|name| name.ends_with(".manifest")) + .unwrap_or(false) => {} Err(err) => return Err(err), } self.primary.copy_opts(from, to, options).await @@ -340,6 +350,64 @@ mod test { )); } + #[tokio::test] + async fn test_create_manifest_recovers_from_orphaned_secondary() { + let primary: Arc = Arc::new(InMemory::new()); + let secondary: Arc = Arc::new(InMemory::new()); + let store = MirroringObjectStore { + primary: primary.clone(), + secondary: secondary.clone(), + }; + let staging = Path::from("_versions/1.manifest-staging"); + let finalized = Path::from("_versions/1.manifest"); + + primary + .put(&staging, "manifest contents".into()) + .await + .unwrap(); + secondary + .put(&staging, "manifest contents".into()) + .await + .unwrap(); + secondary + .copy_if_not_exists(&staging, &finalized) + .await + .expect("simulate a crash after secondary create and before primary create"); + + store + .copy_if_not_exists(&staging, &finalized) + .await + .expect("an orphaned secondary manifest must not block primary creation"); + + let copied = primary + .get(&finalized) + .await + .unwrap() + .bytes() + .await + .unwrap(); + assert_eq!(copied, "manifest contents"); + + assert!(matches!( + store.copy_if_not_exists(&staging, &finalized).await, + Err(Error::AlreadyExists { .. } | Error::Precondition { .. }) + )); + + let non_manifest = Path::from("data/existing.lance"); + secondary + .copy_if_not_exists(&staging, &non_manifest) + .await + .unwrap(); + assert!(matches!( + store.copy_if_not_exists(&staging, &non_manifest).await, + Err(Error::AlreadyExists { .. } | Error::Precondition { .. }) + )); + assert!(matches!( + primary.head(&non_manifest).await, + Err(Error::NotFound { .. }) + )); + } + // This test is ignored because lance 3.0 introduced LocalWriter optimization // that bypasses the object store wrapper for local writes. The mirroring feature // still works for remote/cloud storage, but can't be tested with local storage.