Skip to content
Merged
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
110 changes: 61 additions & 49 deletions Cargo.lock

Large diffs are not rendered by default.

3 changes: 2 additions & 1 deletion crates/lance-context-core/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,12 @@ default = ["metrics"]
metrics = ["dep:metrics"]

[dependencies]
async-trait = "0.1"
async-trait = "0.1.92"
lance-table = "9.0.0"
lance-io = "9.0.0"
object_store = "0.13.2"
base64 = "0.22"
bytes = "1"
arrow-array = "58"
arrow-ipc = "58"
arrow-json = "58"
Expand Down
3 changes: 2 additions & 1 deletion crates/lance-context-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ pub mod merge_budget;
pub mod merge_write_scope;
pub mod metrics;
mod namespace;
mod preparation_io;
mod record;
mod registry;
mod registry_etcd;
Expand Down Expand Up @@ -103,4 +104,4 @@ pub use lance::Error as LanceError;
// capacity-bounded cache session across all resident rollout stores.
pub use lance::session::Session;

pub use store_base::PreparedCompaction;
pub use store_base::{PreparedCompaction, PreparedKeyIndex};
69 changes: 67 additions & 2 deletions crates/lance-context-core/src/merge_write_scope.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,9 @@ pub trait CommitAuthorizer: std::fmt::Debug + Send + Sync {

#[derive(Debug, Default)]
pub struct MergeWriteScope {
completed_steps: AtomicU64,
completed_steps: Arc<AtomicU64>,
observe_file_io: bool,
local_io: Mutex<Vec<(lance_io::utils::tracking_store::IOTracker, u64)>>,
progress: Mutex<Progress>,
changed: Notify,
authorizer: Option<Arc<dyn CommitAuthorizer>>,
Expand All @@ -55,7 +57,18 @@ pub struct MergeWriteScope {

impl MergeWriteScope {
pub fn completed_steps(&self) -> u64 {
self.completed_steps.load(Ordering::Relaxed)
self.local_io.lock().unwrap().iter().fold(
self.completed_steps.load(Ordering::Relaxed),
|steps, (tracker, baseline)| {
let stats = tracker.stats();
steps.saturating_add(
stats
.read_bytes
.saturating_add(stats.written_bytes)
.saturating_sub(*baseline),
)
},
)
}
pub fn new() -> Arc<Self> {
Arc::new(Self::default())
Expand All @@ -79,6 +92,17 @@ impl MergeWriteScope {
})
}

/// Preparation reads/writes immutable files. Observe successful file IO so
/// long scans do not appear idle when Lance omits training callbacks.
pub fn with_preparation_authorizer(authorizer: Arc<dyn CommitAuthorizer>) -> Arc<Self> {
Arc::new(Self {
authorizer: Some(authorizer),
pin_opened_handles: true,
observe_file_io: true,
..Self::default()
})
}

/// Only after dropping the merge future and after a durable storage
/// barrier fenced every admitted version. Aborting before that evidence
/// would abandon an uncertain storage write.
Expand Down Expand Up @@ -140,6 +164,47 @@ pub(crate) fn write_progress() -> lance::dataset::write::WriteProgressFn {
})
}

/// Capture only the preparation scope. Cached handles from ordinary writers
/// must never acquire another task's diagnostic counter.
pub(crate) fn preparation_store_wrapper(
) -> Option<Arc<dyn lance_io::object_store::WrappingObjectStore>> {
CURRENT
.try_with(|scope| {
scope.observe_file_io.then(|| {
Arc::new(crate::preparation_io::ProgressWrapper(
scope.completed_steps.clone(),
)) as Arc<dyn lance_io::object_store::WrappingObjectStore>
})
})
.ok()
.flatten()
}

pub(crate) fn observes_preparation_io() -> bool {
CURRENT
.try_with(|scope| scope.observe_file_io)
.unwrap_or(false)
}

/// Lance's direct local/uring readers and writers bypass ObjectStore wrappers.
/// Their tracker records completed byte IO. The unique preparation wrapper in
/// the store params isolates this handle's tracker from other cached stores.
pub(crate) fn register_local_preparation_io(store: &ObjectStore) {
if !store.is_local() {
return;
}
let _ = CURRENT.try_with(|scope| {
if scope.observe_file_io {
let tracker = store.io_tracker().clone();
let stats = tracker.stats();
scope.local_io.lock().unwrap().push((
tracker,
stats.read_bytes.saturating_add(stats.written_bytes),
));
}
});
}

pub(crate) async fn authorize(resource: &str, version: u64) -> Result<()> {
if let Ok(scope) = CURRENT.try_with(Arc::clone) {
if let Some(authorizer) = &scope.authorizer {
Expand Down
212 changes: 212 additions & 0 deletions crates/lance-context-core/src/preparation_io.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,212 @@
//! Diagnostics for immutable preparation IO, never write authorization.
use async_trait::async_trait;
use bytes::Bytes;
use futures::{stream::BoxStream, StreamExt};
use lance_io::object_store::WrappingObjectStore;
use object_store::{path::Path, *};
use std::sync::{
atomic::{AtomicU64, Ordering},
Arc,
};

#[derive(Debug)]
pub(crate) struct ProgressWrapper(pub Arc<AtomicU64>);

impl WrappingObjectStore for ProgressWrapper {
fn wrap(&self, _: &str, original: Arc<dyn ObjectStore>) -> Arc<dyn ObjectStore> {
Arc::new(ProgressStore {
inner: original,
steps: self.0.clone(),
})
}
}

#[derive(Debug)]
struct ProgressStore {
inner: Arc<dyn ObjectStore>,
steps: Arc<AtomicU64>,
}

impl std::fmt::Display for ProgressStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "PreparationProgress({})", self.inner)
}
}

// Metadata/lease polling must not mask stalled data work. Match path segments
// so both bucket-relative and dataset-relative stores work.
fn data_file(path: &Path) -> bool {
path.as_ref()
.split('/')
.any(|s| matches!(s, "data" | "_indices"))
}

#[async_trait]
impl ObjectStore for ProgressStore {
async fn put_opts(
&self,
path: &Path,
payload: PutPayload,
opts: PutOptions,
) -> Result<PutResult> {
let count = data_file(path) && payload.content_length() > 0;
let result = self.inner.put_opts(path, payload, opts).await?;
if count {
self.steps.fetch_add(1, Ordering::Relaxed);
}
Ok(result)
}

async fn put_multipart_opts(
&self,
path: &Path,
opts: PutMultipartOptions,
) -> Result<Box<dyn MultipartUpload>> {
let inner = self.inner.put_multipart_opts(path, opts).await?;
if !data_file(path) {
return Ok(inner);
}
Ok(Box::new(ProgressUpload {
inner,
steps: self.steps.clone(),
}))
}

async fn get_opts(&self, path: &Path, opts: GetOptions) -> Result<GetResult> {
let count = data_file(path) && !opts.head;
let result = self.inner.get_opts(path, opts).await?;
if !count {
return Ok(result);
}
let meta = result.meta.clone();
let range = result.range.clone();
let attributes = result.attributes.clone();
let steps = self.steps.clone();
// Count body chunks after they arrive, not successful response headers.
// into_stream also observes actual local file reads without buffering.
let payload = GetResultPayload::Stream(
result
.into_stream()
.inspect(move |chunk| {
if chunk.as_ref().is_ok_and(|b| !b.is_empty()) {
steps.fetch_add(1, Ordering::Relaxed);
}
})
.boxed(),
);
Ok(GetResult {
payload,
meta,
range,
attributes,
})
}

async fn get_ranges(&self, path: &Path, ranges: &[std::ops::Range<u64>]) -> Result<Vec<Bytes>> {
let result = self.inner.get_ranges(path, ranges).await?;
if data_file(path) && result.iter().any(|b| !b.is_empty()) {
self.steps.fetch_add(1, Ordering::Relaxed);
}
Ok(result)
}

fn delete_stream(
&self,
paths: BoxStream<'static, Result<Path>>,
) -> BoxStream<'static, Result<Path>> {
self.inner.delete_stream(paths)
}
fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
self.inner.list(prefix)
}
fn list_with_offset(
&self,
prefix: Option<&Path>,
offset: &Path,
) -> BoxStream<'static, Result<ObjectMeta>> {
self.inner.list_with_offset(prefix, offset)
}
async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult> {
self.inner.list_with_delimiter(prefix).await
}
async fn copy_opts(&self, from: &Path, to: &Path, opts: CopyOptions) -> Result<()> {
self.inner.copy_opts(from, to, opts).await
}
async fn rename_opts(&self, from: &Path, to: &Path, opts: RenameOptions) -> Result<()> {
self.inner.rename_opts(from, to, opts).await
}
}

#[derive(Debug)]
struct ProgressUpload {
inner: Box<dyn MultipartUpload>,
steps: Arc<AtomicU64>,
}
#[async_trait]
impl MultipartUpload for ProgressUpload {
fn put_part(&mut self, payload: PutPayload) -> UploadPart {
let nonempty = payload.content_length() > 0;
let part = self.inner.put_part(payload);
let steps = self.steps.clone();
Box::pin(async move {
part.await?;
if nonempty {
steps.fetch_add(1, Ordering::Relaxed);
}
Ok(())
})
}
async fn complete(&mut self) -> Result<PutResult> {
self.inner.complete().await
}
async fn abort(&mut self) -> Result<()> {
self.inner.abort().await
}
}

#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn counts_data_bodies_and_upload_parts_but_not_metadata_or_failures() {
let steps = Arc::new(AtomicU64::new(0));
let store = ProgressStore {
inner: Arc::new(memory::InMemory::new()),
steps: steps.clone(),
};
let data = Path::from("table/data/part.lance");
let metadata = Path::from("table/_versions/1.manifest");
store.put(&metadata, "metadata".into()).await.unwrap();
store.get(&metadata).await.unwrap().bytes().await.unwrap();
assert_eq!(steps.load(Ordering::Relaxed), 0);
store.put(&data, "payload".into()).await.unwrap();
let result = store.get(&data).await.unwrap();
store.head(&data).await.unwrap();
assert_eq!(
steps.load(Ordering::Relaxed),
1,
"headers and HEAD are not data progress"
);
assert_eq!(result.bytes().await.unwrap(), "payload");
assert_eq!(steps.load(Ordering::Relaxed), 2);
assert!(store.get(&Path::from("table/data/missing")).await.is_err());
assert!(store
.put_opts(&data, "conflict".into(), PutOptions::from(PutMode::Create))
.await
.is_err());
assert_eq!(steps.load(Ordering::Relaxed), 2);
let mut upload = store
.put_multipart(&Path::from("table/_indices/new/index.idx"))
.await
.unwrap();
upload.put_part("index".into()).await.unwrap();
upload.complete().await.unwrap();
assert_eq!(steps.load(Ordering::Relaxed), 3);
let ranges = store.get_ranges(&data, &[0..3, 3..7]).await.unwrap();
assert_eq!(
ranges,
vec![Bytes::from_static(b"pay"), Bytes::from_static(b"load")]
);
assert_eq!(steps.load(Ordering::Relaxed), 4);
}
}
Loading
Loading