Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions java/lance-jni/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions python/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions rust/lance-io/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ serde = { workspace = true, features = ["derive"] }
tokio.workspace = true
tracing.workspace = true
url.workspace = true
uuid.workspace = true
Comment thread
dentiny marked this conversation as resolved.
path_abs.workspace = true
rand.workspace = true
tempfile.workspace = true
Expand Down
94 changes: 93 additions & 1 deletion rust/lance-io/src/object_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,10 @@ use object_store::ObjectStoreExt as OSObjectStoreExt;
use object_store::aws::AwsCredentialProvider;
#[cfg(any(feature = "aws", feature = "azure", feature = "gcp"))]
use object_store::{ClientOptions, HeaderMap, HeaderValue};
use object_store::{ListResult, ObjectMeta, ObjectStore as OSObjectStore, path::Path};
use object_store::{
ListResult, ObjectMeta, ObjectStore as OSObjectStore, PutMode, PutOptions, PutPayload,
path::Path,
};
use providers::local::FileStoreProvider;
use providers::memory::MemoryStoreProvider;
use tokio::io::AsyncWriteExt;
Expand Down Expand Up @@ -823,6 +826,57 @@ impl ObjectStore {
Writer::shutdown(writer.as_mut()).await
}

/// Atomically creates an object without replacing an existing object.
///
/// Local stores publish a uniquely named staging object with a conditional
/// rename. Other stores use their conditional create operation. Tencent COS
/// is rejected because it can silently ignore conditional create requests.
///
/// Returns [`object_store::Error::NotSupported`] without writing when the
/// backend cannot reliably provide put-if-absent semantics.
pub async fn put_if_absent(
&self,
path: &Path,
content: PutPayload,
) -> object_store::Result<()> {
if self.scheme == "cos" {
return Err(object_store::Error::NotSupported {
source: "Tencent COS does not reliably enforce put-if-absent after bucket \
versioning has ever been enabled"
.into(),
});
}

if self.is_local() {
let staging_path =
Path::from(format!("{}.tmp.{}", path, uuid::Uuid::new_v4().simple()));
self.inner.put(&staging_path, content).await?;
let result = self.inner.rename_if_not_exists(&staging_path, path).await;
if result.is_err()
&& let Err(error) = self.inner.delete(&staging_path).await
{
log::warn!(
"Failed to remove staging object {} after atomic create failed: {}",
staging_path,
error
);
}
result
} else {
self.inner
.put_opts(
path,
content,
PutOptions {
mode: PutMode::Create,
..Default::default()
},
)
.await
.map(|_| ())
}
}

pub async fn delete(&self, path: &Path) -> Result<()> {
self.inner.delete(path).await?;
Ok(())
Expand Down Expand Up @@ -1297,6 +1351,44 @@ mod tests {
Ok(contents)
}

#[tokio::test]
async fn test_put_if_absent() {
let temp_dir = TempStrDir::default();
let path = Path::from(format!("{}/atomic-create", temp_dir.as_str()));
let store = ObjectStore::local();
store
.put_if_absent(&path, Bytes::from_static(b"first").into())
.await
.unwrap();
let error = store
.put_if_absent(&path, Bytes::from_static(b"second").into())
.await
.unwrap_err();
assert!(matches!(
error,
object_store::Error::AlreadyExists { .. } | object_store::Error::Precondition { .. }
));
assert_eq!(
store.read_one_all(&path).await.unwrap(),
b"first".as_slice()
);
}

#[tokio::test]
async fn test_put_if_absent_rejects_cos() {
let mut store = ObjectStore::memory();
store.scheme = "cos".to_string();
let path = Path::from("atomic-create");

let error = store
.put_if_absent(&path, Bytes::from_static(b"value").into())
.await
.unwrap_err();

assert!(matches!(error, object_store::Error::NotSupported { .. }));
assert!(!store.exists(&path).await.unwrap());
}

#[test]
fn test_io_parallelism_clamped_to_nonzero() {
// `io_parallelism()` feeds `buffered`/`buffer_unordered` windows; a value of 0 makes those
Expand Down
76 changes: 16 additions & 60 deletions rust/lance/src/dataset/mem_wal/manifest.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,8 +39,6 @@ use lance_index::mem_wal::{ShardManifest, ShardStatus};
use lance_io::object_store::ObjectStore;
use lance_table::format::pb;
use log::{info, warn};
use object_store::PutMode;
use object_store::PutOptions;
use object_store::path::Path;
use prost::Message;
use serde::{Deserialize, Serialize};
Expand Down Expand Up @@ -184,68 +182,26 @@ impl ShardManifestStore {
let pb_manifest = pb::ShardManifest::from(manifest);
let bytes = pb_manifest.encode_to_vec();

if self.object_store.is_local() {
// Local storage: Use temp file + atomic rename for fencing
let temp_filename = format!("{}.tmp.{}", filename, uuid::Uuid::new_v4());
let temp_path = self.manifest_dir.clone().join(temp_filename.as_str());

// Write to temp file
self.object_store
.inner
.put(&temp_path, Bytes::from(bytes).into())
.await
.map_err(|e| Error::io(format!("Failed to write temp manifest: {}", e)))?;

// Atomically rename to final path
match self
.object_store
.inner
.rename_if_not_exists(&temp_path, &path)
.await
{
Ok(()) => {}
Err(object_store::Error::AlreadyExists { .. }) => {
// Clean up temp file
let _ = self.object_store.delete(&temp_path).await;
return Err(Error::io(format!(
self.object_store
.put_if_absent(&path, Bytes::from(bytes).into())
.await
.map_err(|error| {
if matches!(
error,
object_store::Error::AlreadyExists { .. }
| object_store::Error::Precondition { .. }
) {
Error::io(format!(
"Manifest version {} already exists for shard {}",
version, self.shard_id
)));
}
Err(e) => {
// Clean up temp file
let _ = self.object_store.delete(&temp_path).await;
return Err(Error::io(format!(
))
} else {
Error::io(format!(
"Failed to write manifest version {} for shard {}: {}",
version, self.shard_id, e
)));
version, self.shard_id, error
))
}
}
} else {
// Cloud storage: Use PUT-IF-NOT-EXISTS
let put_opts = PutOptions {
mode: PutMode::Create,
..Default::default()
};

self.object_store
.inner
.put_opts(&path, Bytes::from(bytes).into(), put_opts)
.await
.map_err(|e| {
if matches!(e, object_store::Error::AlreadyExists { .. }) {
Error::io(format!(
"Manifest version {} already exists for shard {}",
version, self.shard_id
))
} else {
Error::io(format!(
"Failed to write manifest version {} for shard {}: {}",
version, self.shard_id, e
))
}
})?;
}
})?;

// Best-effort update version hint (failures are logged as warnings)
self.write_version_hint(version).await;
Expand Down
59 changes: 11 additions & 48 deletions rust/lance/src/dataset/mem_wal/wal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,6 @@ use lance_core::{Error, FenceReason, Result};
use lance_io::object_store::ObjectStore;
use object_store::ObjectStoreExt;
use object_store::path::Path;
use object_store::{PutMode, PutOptions};
use tokio::sync::{Mutex, mpsc, watch};

use tracing::instrument;
Expand Down Expand Up @@ -1579,53 +1578,17 @@ async fn atomic_put(
bytes: Bytes,
) -> std::result::Result<(), AtomicPutError> {
let path = dir.clone().join(filename);
if object_store.is_local() {
let temp = dir
.clone()
.join(format!("{}.tmp.{}", filename, Uuid::new_v4()));
object_store
.inner
.put(&temp, bytes.into())
.await
.map_err(|e| {
AtomicPutError::Other(Error::io(format!("failed to write temp file: {}", e)))
})?;
match object_store.inner.rename_if_not_exists(&temp, &path).await {
Ok(()) => Ok(()),
Err(object_store::Error::AlreadyExists { .. }) => {
let _ = object_store.delete(&temp).await;
Err(AtomicPutError::AlreadyExists)
}
Err(e) => {
let _ = object_store.delete(&temp).await;
Err(AtomicPutError::Other(Error::io(format!(
"failed to create {} atomically: {}",
path, e
))))
}
}
} else {
object_store
.inner
.put_opts(
&path,
bytes.into(),
PutOptions {
mode: PutMode::Create,
..Default::default()
},
)
.await
.map_err(|e| match e {
object_store::Error::AlreadyExists { .. }
| object_store::Error::Precondition { .. } => AtomicPutError::AlreadyExists,
_ => AtomicPutError::Other(Error::io(format!(
"failed to create {} atomically: {}",
path, e
))),
})?;
Ok(())
}
object_store
.put_if_absent(&path, bytes.into())
.await
.map_err(|error| match error {
object_store::Error::AlreadyExists { .. }
| object_store::Error::Precondition { .. } => AtomicPutError::AlreadyExists,
_ => AtomicPutError::Other(Error::io(format!(
"failed to create {} atomically: {}",
path, error
))),
})
}

/// Probe forward from a hint position to find the next unwritten position.
Expand Down
Loading
Loading