diff --git a/Cargo.lock b/Cargo.lock index cc0d5ba7..5e90b596 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6455,6 +6455,7 @@ dependencies = [ "clap", "etcd-client", "futures", + "lance 9.0.0", "lance-context-api", "lance-context-core", "lance-context-merge", @@ -6466,6 +6467,7 @@ dependencies = [ "serde_json", "tempfile", "tokio", + "tower", "tower-http", "tracing", "tracing-subscriber", diff --git a/crates/lance-context-core/src/lib.rs b/crates/lance-context-core/src/lib.rs index 7e505616..6fce40f2 100644 --- a/crates/lance-context-core/src/lib.rs +++ b/crates/lance-context-core/src/lib.rs @@ -20,6 +20,7 @@ mod record; mod registry; mod registry_etcd; mod rollout; +pub mod rollout_append; mod rollout_store; pub mod serde; mod storage; diff --git a/crates/lance-context-core/src/merge_write_scope.rs b/crates/lance-context-core/src/merge_write_scope.rs index 5d5d19ae..e1f1d550 100644 --- a/crates/lance-context-core/src/merge_write_scope.rs +++ b/crates/lance-context-core/src/merge_write_scope.rs @@ -21,6 +21,12 @@ use std::{ use tokio::sync::{oneshot, Notify}; tokio::task_local! { static CURRENT: Arc; } +tokio::task_local! { static EXPECTED_BASE_VERSION: u64; } + +/// Pin the next publication without replacing a dataset's captured owner guard. +pub(crate) async fn at_base_version(version: u64, work: F) -> F::Output { + EXPECTED_BASE_VERSION.scope(version, work).await +} #[derive(Debug, Default)] struct Progress { @@ -114,10 +120,26 @@ impl MergeWriteScope { /// Count completed work, never timer ticks or admission attempts. This is a /// stall diagnostic, not evidence that WAL has been durably reclaimed. -pub(crate) fn checkpoint() { +pub fn checkpoint() { let _ = CURRENT.try_with(|scope| scope.completed_steps.fetch_add(1, Ordering::Relaxed)); } +/// Capture progress explicitly for Lance callbacks which can run on another task. +pub(crate) fn write_progress() -> lance::dataset::write::WriteProgressFn { + let scope = CURRENT.try_with(Arc::clone).ok(); + let previous = std::sync::Mutex::new((0u64, 0u64, 0u32)); + lance::dataset::write::WriteProgressFn::new(move |stats| { + let current = (stats.bytes_written, stats.rows_written, stats.files_written); + let mut old = previous.lock().unwrap(); + if current.0 > old.0 || current.1 > old.1 || current.2 > old.2 { + *old = current; + if let Some(scope) = &scope { + scope.completed_steps.fetch_add(1, Ordering::Relaxed); + } + } + }) +} + 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 { @@ -392,6 +414,12 @@ impl CommitHandler for GuardedCommit { transaction: Option, ) -> std::result::Result { let work = async { + if EXPECTED_BASE_VERSION + .try_with(|version| manifest.version != *version) + .unwrap_or(false) + { + return Err(CommitError::CommitConflict); + } authorize("base", manifest.version).await?; let inner = self.inner.clone(); let mut owned_manifest = manifest.clone(); diff --git a/crates/lance-context-core/src/rollout_append.rs b/crates/lance-context-core/src/rollout_append.rs new file mode 100644 index 00000000..b6b3ed13 --- /dev/null +++ b/crates/lance-context-core/src/rollout_append.rs @@ -0,0 +1,1079 @@ +//! Parallel, immutable rollout file preparation with one fenced metadata committer. +//! +//! The caller must hold the table's maintenance claim across planning and commit. +//! Workers need no write ownership: staging never publishes a table/WAL manifest. +use std::collections::{HashMap, HashSet}; +use std::sync::Arc; + +use arrow_array::{Array, BooleanArray, RecordBatch, StringArray}; +use arrow_schema::Schema; +use arrow_select::filter::filter_record_batch; +use datafusion::prelude::{col, lit}; +use futures::{StreamExt, TryStreamExt}; +use lance::dataset::{ + mem_wal::{DatasetMemWalExt, ShardManifestStore}, + transaction::{Operation, Transaction}, + CommitBuilder, Dataset, InsertBuilder, WriteMode, WriteParams, +}; +use lance::index::DatasetIndexExt; +use lance_index::mem_wal::{MemWalIndexDetails, MergedGeneration, MEM_WAL_INDEX_NAME}; +use lance_table::format::{pb, Fragment}; +use serde::{Deserialize, Serialize}; +use uuid::Uuid; + +use crate::store_base::{align_batch_to_schema, derive_shard_id, StorageBase}; +use crate::{LanceError as Error, MergeMemoryBudget, Session}; +use lance::Result; + +const CUTOVER_PREFIX: &str = "lance-context.rollout-append.cutover."; +pub const MAX_PLAN_GENERATIONS: usize = 256; +pub const MAX_STAGE_BYTES: usize = 256 * 1024 * 1024; + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct Generation { + pub number: u64, + pub path: String, +} + +/// A bounded, immutable prefix of one shard. URIs/credentials are never supplied +/// by an RPC caller: each process resolves the rollout name in its own registry. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct AppendPlan { + pub base_version: u64, + pub shard: Uuid, + pub generations: Vec, + pub max_bytes: usize, + pub legacy_through: u64, +} + +impl AppendPlan { + pub fn validate(&self) -> Result<()> { + if self.base_version == 0 + || self.generations.is_empty() + || self.generations.len() > MAX_PLAN_GENERATIONS + || self.max_bytes == 0 + || self.max_bytes > MAX_STAGE_BYTES + || self + .generations + .windows(2) + .any(|g| g[0].number >= g[1].number) + || self.generations.iter().any(|g| { + g.path.is_empty() || g.path.contains('/') || g.path.contains('\\') || g.path == ".." + }) + { + return Err(Error::invalid_input("invalid rollout append plan")); + } + Ok(()) + } +} + +/// Only file metadata crosses the worker/master boundary. No Arrow payloads. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct StagedAppend { + pub dataset_uri: String, + pub plan: AppendPlan, + pub completed: usize, + pub fragments: Vec, + pub rows: usize, +} + +/// Streaming RPC keeps completed work distinct from a live connection. +#[derive(Debug, Serialize, Deserialize)] +pub enum StageEvent { + Progress(u64), + Complete(StagedAppend), + Failed(String), +} + +pub(crate) async fn watermarks(dataset: &Dataset) -> Result> { + let indices = dataset.load_indices().await?; + let Some(index) = indices.iter().find(|i| i.name == MEM_WAL_INDEX_NAME) else { + return Ok(HashMap::new()); + }; + let details = index + .index_details + .as_ref() + .ok_or_else(|| Error::io("MemWAL index has no details"))?; + let details = MemWalIndexDetails::try_from(details.to_msg::()?)?; + Ok(details + .merged_generations + .into_iter() + .map(|g| (g.shard_id, g.generation)) + .collect()) +} + +async fn shard_store(dataset: &Dataset, shard: Uuid) -> Result { + Ok(ShardManifestStore::new( + dataset.object_store(None).await?, + &dataset.branch_location().path, + shard, + 16, + )) +} + +/// A previous append can be visible while its separate shard drain is not. +/// Reconcile solely from the atomic base watermark, without reading any WAL data. +pub(crate) async fn reconcile_shard(dataset: &Dataset, shard: Uuid) -> Result { + let Some(high) = watermarks(dataset).await?.get(&shard).copied() else { + return Ok(0); + }; + let store = shard_store(dataset, shard).await?; + let Some(manifest) = store.read_latest().await? else { + return Ok(0); + }; + let generations: HashSet<_> = manifest + .flushed_generations + .iter() + .filter(|g| g.generation <= high) + .map(|g| g.generation) + .collect(); + let count = generations.len(); + if count > 0 { + let paths: Vec<_> = manifest + .flushed_generations + .iter() + .filter(|g| generations.contains(&g.generation)) + .map(|g| g.path.clone()) + .collect(); + crate::merge_write_scope::drain_generations(store, manifest.writer_epoch, generations) + .await?; + let object_store = dataset.object_store(None).await?; + let root = dataset + .branch_location() + .path + .join("_mem_wal") + .join(shard.to_string().as_str()); + tokio::spawn(async move { + futures::stream::iter(paths) + .for_each_concurrent(16, |path| { + let object_store = object_store.clone(); + let path = root.clone().join(path.as_str()); + async move { + if let Err(error) = object_store.remove_dir_all(path).await { + tracing::warn!(%error, "merged rollout WAL directory cleanup failed"); + } + } + }) + .await; + }); + } + Ok(count) +} + +/// The only publisher. Open this inside the existing MergeWriteScope so both +/// table commits and shard drains retain ownership checks and version fencing. +pub struct AppendCoordinator { + dataset: Dataset, +} + +impl AppendCoordinator { + pub async fn open(uri: &str, session: Option>) -> Result { + // Do not open a resident shard writer or evolve schema on a staging worker. + let dataset = StorageBase::load_with_options(uri, None, session).await?; + let actual: Schema = dataset.schema().into(); + let expected = crate::rollout_schema(); + if dataset.manifest().should_use_legacy_format() { + return Err(Error::invalid_input( + "rollout staging requires Lance V2 files", + )); + } + if actual != expected { + return Err(Error::invalid_input( + "rollout append requires the current rollout schema", + )); + } + Ok(Self { dataset }) + } + + pub fn version(&self) -> u64 { + self.dataset.version().version + } + + pub async fn plan( + &mut self, + shards: &[String], + max_generations: usize, + max_bytes: usize, + ) -> Result<(Vec, usize)> { + if max_generations == 0 + || max_generations > MAX_PLAN_GENERATIONS + || max_bytes == 0 + || max_bytes > MAX_STAGE_BYTES + || shards.len() > 256 + { + return Err(Error::invalid_input("invalid rollout append limits")); + } + self.dataset.checkout_latest().await?; + let mut unique = HashSet::new(); + for name in shards { + if !unique.insert(derive_shard_id(Some(name))) { + return Err(Error::invalid_input("duplicate rollout append shard")); + } + } + // Metadata-only shard-directory enumeration also finds retired workers. + // An unavailable ingestion worker must not strand its immutable WAL. + unique.extend(self.dataset.list_mem_wal_latest_shard_ids().await?); + if unique.len() > 256 { + return Err(Error::invalid_input( + "rollout append supports at most 256 shards per table", + )); + } + let mut shard_ids: Vec<_> = unique.into_iter().collect(); + shard_ids.sort(); + let mut pending = Vec::new(); + let mut cutovers = HashMap::new(); + let mut reclaimed = 0; + for shard in shard_ids { + reclaimed += reconcile_shard(&self.dataset, shard).await?; + let Some(manifest) = shard_store(&self.dataset, shard) + .await? + .read_latest() + .await? + else { + continue; + }; + crate::merge_write_scope::checkpoint(); + let key = format!("{CUTOVER_PREFIX}{shard}"); + let cutoff = match self.dataset.metadata().get(&key) { + Some(value) => value + .parse::() + .map_err(|_| Error::io("invalid rollout cutover"))?, + None => { + // The historical prefix can contain a legacy append whose + // drain failed. Only that prefix needs a base ID lookup. + let high = manifest + .flushed_generations + .iter() + .map(|g| g.generation) + .max() + .unwrap_or(0); + cutovers.insert(key, high.to_string()); + high + } + }; + let mut generations: Vec<_> = manifest + .flushed_generations + .into_iter() + .map(|g| Generation { + number: g.generation, + path: g.path, + }) + .collect(); + generations.sort_by_key(|g| g.number); + generations.truncate(max_generations); + if !generations.is_empty() { + pending.push((shard, generations, cutoff)); + } + } + if !cutovers.is_empty() { + self.dataset + .update_metadata(cutovers.iter().map(|(k, v)| (k.as_str(), v.as_str()))) + .await?; + } + let version = self.version(); + let plans = pending + .into_iter() + .map(|(shard, generations, legacy_through)| AppendPlan { + base_version: version, + shard, + generations, + max_bytes, + legacy_through, + }) + .collect(); + Ok((plans, reclaimed)) + } + + /// Publication is a single transaction: immutable fragments + per-shard + /// merged prefix. No ID index maintenance, delete, or target payload scan. + /// Do not retry this Transaction blindly after an uncertain response: call + /// this method again so the freshly read watermarks decide what remains. + pub async fn commit(&mut self, staged: Vec) -> Result { + self.dataset.checkout_latest().await?; + let marks = watermarks(&self.dataset).await?; + let mut shards = HashSet::new(); + let mut fragments = Vec::new(); + let mut merged = Vec::new(); + for part in staged { + part.plan.validate()?; + if part.dataset_uri.trim_end_matches('/') != self.dataset.uri().trim_end_matches('/') { + return Err(Error::invalid_input( + "staging worker used a different dataset URI", + )); + } + if part.completed == 0 + || part.completed > part.plan.generations.len() + || !shards.insert(part.plan.shard) + { + return Err(Error::invalid_input("invalid or duplicate staged shard")); + } + let selected = &part.plan.generations[..part.completed]; + let high = selected.last().unwrap().number; + if let Some(old) = marks.get(&part.plan.shard) { + if high <= *old { + continue; + } + if selected[0].number <= *old { + return Err(Error::invalid_input( + "partially stale staged prefix; replan", + )); + } + } + // Require the version used to encode fields to have the same schema. + let snapshot = self + .dataset + .checkout_version(part.plan.base_version) + .await?; + // Dictionary values/offsets are loaded lazily and are file-local + // in V2. Raw Lance Schema equality compares that cache state too. + let options = lance::datatypes::SchemaCompareOptions { + compare_metadata: true, + compare_field_ids: true, + ..Default::default() + }; + if snapshot.manifest().data_storage_format.version + != self.dataset.manifest().data_storage_format.version + || Schema::from(snapshot.schema()) != Schema::from(self.dataset.schema()) + || snapshot + .schema() + .check_compatible(self.dataset.schema(), &options) + .is_err() + { + return Err(Error::invalid_input( + "schema changed during rollout staging; replan", + )); + } + let manifest = shard_store(&self.dataset, part.plan.shard) + .await? + .read_latest() + .await? + .ok_or_else(|| Error::io("staged shard disappeared"))?; + let mut pending: Vec<_> = manifest.flushed_generations.iter().collect(); + pending.sort_by_key(|g| g.generation); + if pending.len() < selected.len() + || pending + .iter() + .zip(selected) + .any(|(a, b)| a.generation != b.number || a.path != b.path) + { + return Err(Error::invalid_input("staged WAL prefix changed; replan")); + } + if part + .fragments + .iter() + .map(|f| f.physical_rows.unwrap_or(0)) + .sum::() + != part.rows + { + return Err(Error::invalid_input("staged row count mismatch")); + } + fragments.extend(part.fragments); + merged.push(MergedGeneration::new(part.plan.shard, high)); + } + if !merged.is_empty() { + commit_files(&mut self.dataset, fragments, merged).await?; + } + let mut reclaimed = 0; + for shard in shards { + reclaimed += reconcile_shard(&self.dataset, shard).await?; + } + Ok(reclaimed) + } +} + +/// Keep legacy fallback publications recoverable after an append-mode cutover. +pub(crate) fn has_cutover(dataset: &Dataset, shard: Uuid) -> bool { + dataset + .metadata() + .contains_key(&format!("{CUTOVER_PREFIX}{shard}")) +} + +pub(crate) async fn commit_files( + dataset: &mut Dataset, + fragments: Vec, + merged: Vec, +) -> Result<()> { + let operation = Operation::Update { + removed_fragment_ids: Vec::new(), + updated_fragments: Vec::new(), + new_fragments: fragments, + fields_modified: Vec::new(), + merged_generations: merged, + fields_for_preserving_frag_bitmap: Vec::new(), + update_mode: None, + inserted_rows_filter: None, + updated_fragment_offsets: None, + }; + // Even zero-retry Lance transactions can rebase before their first write. + let version = dataset.version().version; + crate::merge_write_scope::at_base_version( + version + 1, + CommitBuilder::new(Arc::new(dataset.clone())) + .with_max_retries(0) + .execute(Transaction::new(version, operation, None)), + ) + .await?; + dataset.checkout_latest().await?; + Ok(()) +} + +/// Encode files only. Holds the shared worker memory reservation until encoding +/// finishes, then returns small metadata. No shard epoch or manifest is changed. +pub async fn stage( + uri: &str, + plan: AppendPlan, + budget: Arc, + session: Option>, +) -> Result { + plan.validate()?; + let dataset = StorageBase::load_with_options(uri, None, session.clone()) + .await? + .checkout_version(plan.base_version) + .await?; + let dataset_uri = dataset.uri().to_owned(); + let schema: Arc = Arc::new(dataset.schema().into()); + if dataset.manifest().should_use_legacy_format() { + return Err(Error::invalid_input( + "rollout staging requires Lance V2 files", + )); + } + if *schema != crate::rollout_schema() { + return Err(Error::invalid_input("staging requires a rollout schema")); + } + let key = format!("{CUTOVER_PREFIX}{}", plan.shard); + if dataset.metadata().get(&key) != Some(&plan.legacy_through.to_string()) { + return Err(Error::invalid_input( + "staging plan does not match durable cutover", + )); + } + let mut reservation = budget.reserve(plan.max_bytes.min(budget.limit())).await; + let mut batches = Vec::new(); + let mut bytes = 0usize; + let mut completed = 0; + 'generations: for generation in &plan.generations { + let path = format!( + "{}/_mem_wal/{}/{}", + uri.trim_end_matches('/'), + plan.shard, + generation.path + ); + let source = StorageBase::load_with_options(&path, None, session.clone()).await?; + let mut stream = source.scan().try_into_stream().await?; + let mut current = Vec::new(); + while let Some(batch) = stream.try_next().await? { + let batch = align_batch_to_schema(batch, schema.clone())?; + bytes = bytes.saturating_add(batch.get_array_memory_size()); + if !reservation.try_grow_to(bytes.saturating_mul(2)) { + if completed == 0 { + return Err(Error::io( + "merge memory budget busy while reading first generation; retry", + )); + } + break 'generations; + } + current.push(batch); + crate::merge_write_scope::checkpoint(); + } + batches.extend(current); + completed += 1; + if bytes >= plan.max_bytes { + break; + } + } + // Rollout IDs are immutable. Dedup identical retries within this batch; + // during migration also exclude IDs already published by a legacy merge. + let mut seen = HashSet::new(); + if plan.generations[0].number <= plan.legacy_through { + let mut ids = HashSet::new(); + for batch in &batches { + let column = batch + .column_by_name("id") + .unwrap() + .as_any() + .downcast_ref::() + .ok_or_else(|| Error::invalid_input("rollout id must be Utf8"))?; + ids.extend(column.iter().flatten().map(str::to_owned)); + } + let ids: Vec<_> = ids.into_iter().collect(); + for chunk in ids.chunks(1024) { + let mut scanner = dataset.scan(); + scanner.project(&["id"])?; + scanner.filter_expr( + col("id").in_list(chunk.iter().map(|id| lit(id.clone())).collect(), false), + ); + let mut rows = scanner.try_into_stream().await?; + while let Some(batch) = rows.try_next().await? { + let column = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + seen.extend(column.iter().flatten().map(str::to_owned)); + crate::merge_write_scope::checkpoint(); + } + } + } + let mut output = Vec::new(); + for batch in batches { + let ids = batch + .column_by_name("id") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap(); + if ids.null_count() > 0 { + return Err(Error::invalid_input("rollout id cannot be null")); + } + let keep = BooleanArray::from( + ids.iter() + .map(|id| seen.insert(id.unwrap().to_owned())) + .collect::>(), + ); + if keep.true_count() == batch.num_rows() { + output.push(batch); + } else if keep.true_count() > 0 { + output.push(filter_record_batch(&batch, &keep)?); + } + } + let rows = output.iter().map(RecordBatch::num_rows).sum(); + let fragments = if output.is_empty() { + Vec::new() + } else { + let params = WriteParams { + mode: WriteMode::Append, + max_bytes_per_file: plan.max_bytes, + write_progress: Some(crate::merge_write_scope::write_progress()), + ..Default::default() + }; + let tx = InsertBuilder::new(Arc::new(dataset)) + .with_params(¶ms) + .execute_uncommitted(output) + .await?; + let Operation::Append { fragments } = tx.operation else { + return Err(Error::io("staging produced a non-append transaction")); + }; + fragments + }; + drop(reservation); + crate::merge_write_scope::checkpoint(); + Ok(StagedAppend { + dataset_uri, + plan, + completed, + fragments, + rows, + }) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::merge_write_scope::{CommitAuthorizer, MergeWriteScope}; + use crate::{RolloutRecord, RolloutStore, RolloutStoreOptions, ROLE_ARTIFACT}; + use chrono::{TimeZone, Utc}; + use serde_json::json; + use std::{ + future::Future, + pin::Pin, + sync::atomic::{AtomicBool, Ordering}, + }; + fn artifact_record(id: &str, bytes: &[u8]) -> RolloutRecord { + RolloutRecord { + id: id.to_string(), + rollout_id: "rollout-1".to_string(), + problem_id: "problem-1".to_string(), + dataset: None, + sequence_order: 1, + role: ROLE_ARTIFACT.to_string(), + created_at: Utc.timestamp_micros(1_700_000_000_500_000).unwrap(), + content: None, + content_type: "application/octet-stream".to_string(), + model_input_string: None, + model_output_string: None, + rationale: None, + problem_text: None, + user_metadata: None, + input_tokens: None, + output_tokens: None, + num_input_tokens: None, + num_output_tokens: None, + output_logprobs: None, + input_logprobs: None, + ref_logprobs: None, + loss_mask: None, + advantage: None, + reward: None, + raw_reward: None, + grader_id: None, + score: None, + include_in_training: None, + exclude_reason: None, + policy_version: None, + relationships: Vec::new(), + binary_payload: Some(bytes.to_vec()), + payload_size: Some(bytes.len() as i64), + payload_checksum: Some("sha256:cafef00d".to_string()), + artifact_type: Some("excel_grade_screenshot".to_string()), + metadata: Some(json!({"filename": "trace.bin"})), + } + } + + async fn writer(uri: &str, shard: &str) -> RolloutStore { + RolloutStore::open_with_options( + uri, + RolloutStoreOptions { + shard_id: Some(shard.into()), + merge_after_generations: Some(0), + ..Default::default() + }, + ) + .await + .unwrap() + } + async fn put(store: &RolloutStore, id: &str, bytes: usize) { + store + .add(&[artifact_record(id, &vec![42; bytes])]) + .await + .unwrap(); + store.flush().await.unwrap(); + } + async fn rows(uri: &str) -> usize { + Dataset::open(uri) + .await + .unwrap() + .count_rows(None) + .await + .unwrap() + } + fn budget() -> Arc { + MergeMemoryBudget::new(8 * 1024 * 1024) + } + + #[tokio::test] + async fn parallel_files_publish_once_and_preserve_new_writes() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_str().unwrap(); + let a = writer(uri, "a").await; + let b = writer(uri, "b").await; + put(&a, "a1", 4096).await; + put(&b, "b1", 4096).await; + let mut commit = AppendCoordinator::open(uri, None).await.unwrap(); + let (plans, _) = commit + .plan(&["a".into(), "b".into()], 64, 1024 * 1024) + .await + .unwrap(); + let before = commit.version(); + let memory = budget(); + let (a_part, b_part) = tokio::try_join!( + stage(uri, plans[0].clone(), memory.clone(), None), + stage(uri, plans[1].clone(), memory.clone(), None) + ) + .unwrap(); + assert_eq!(rows(uri).await, 0, "staging must not publish rows"); + assert_eq!(memory.reserved(), 0); + let mut wrong_store = a_part.clone(); + wrong_store.dataset_uri = format!("{uri}/other"); + assert!(commit + .commit(vec![wrong_store]) + .await + .unwrap_err() + .to_string() + .contains("different dataset URI")); + assert_eq!(rows(uri).await, 0); + put(&a, "a2", 4096).await; + assert_eq!( + commit + .commit(vec![a_part.clone(), b_part.clone()]) + .await + .unwrap(), + 2 + ); + assert_eq!(commit.version(), before + 1, "one version for both workers"); + assert_eq!(rows(uri).await, 2); + // Duplicate responses and retries do not create versions or rows. + assert_eq!(commit.commit(vec![a_part, b_part]).await.unwrap(), 0); + assert_eq!(commit.version(), before + 1); + let fresh = RolloutStore::open_existing_with_options(uri, RolloutStoreOptions::default()) + .await + .unwrap(); + assert_eq!(fresh.get_blob("a1").await.unwrap(), Some(vec![42; 4096])); + assert_eq!(fresh.get_blob("b1").await.unwrap(), Some(vec![42; 4096])); + assert_eq!(fresh.get_blob("a2").await.unwrap(), Some(vec![42; 4096])); + let (plans, _) = commit + .plan(&["a".into(), "b".into()], 64, 1024 * 1024) + .await + .unwrap(); + assert_eq!(plans.len(), 1); + let part = stage(uri, plans[0].clone(), memory, None).await.unwrap(); + commit.commit(vec![part]).await.unwrap(); + assert_eq!(rows(uri).await, 3); + assert_eq!( + commit + .dataset + .load_indices() + .await + .unwrap() + .iter() + .filter(|i| i.name != MEM_WAL_INDEX_NAME) + .count(), + 0, + "append merge must not build an ID index" + ); + } + + #[derive(Debug)] + struct FailDrain(AtomicBool); + impl CommitAuthorizer for FailDrain { + fn authorize<'a>( + &'a self, + resource: &'a str, + _version: u64, + ) -> Pin> + Send + 'a>> { + Box::pin(async move { + if resource.starts_with("shard:") && self.0.load(Ordering::SeqCst) { + Err(Error::io("injected failure after base commit")) + } else { + Ok(()) + } + }) + } + } + + #[tokio::test] + async fn restart_after_base_commit_repairs_drain_without_duplicate_append() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_str().unwrap(); + let a = writer(uri, "a").await; + put(&a, "a1", 4096).await; + let guard = Arc::new(FailDrain(AtomicBool::new(true))); + let scope = MergeWriteScope::with_pinned_authorizer(guard.clone()); + let (part, version) = scope + .run(async { + let mut commit = AppendCoordinator::open(uri, None).await.unwrap(); + let (plans, _) = commit.plan(&["a".into()], 64, 1024 * 1024).await.unwrap(); + let part = stage(uri, plans[0].clone(), budget(), None).await.unwrap(); + assert!(commit + .commit(vec![part.clone()]) + .await + .unwrap_err() + .to_string() + .contains("injected")); + (part, Dataset::open(uri).await.unwrap().version().version) + }) + .await; + scope.drain().await; + assert_eq!(rows(uri).await, 1); + put(&a, "a2", 4096).await; + // A fresh coordinator has no in-memory record of the first execution. + let mut restarted = AppendCoordinator::open(uri, None).await.unwrap(); + assert_eq!(restarted.commit(vec![part]).await.unwrap(), 1); + assert_eq!(restarted.version(), version); + assert_eq!(rows(uri).await, 1); + let (plans, _) = restarted + .plan(&["a".into()], 64, 1024 * 1024) + .await + .unwrap(); + assert_eq!(plans[0].generations.len(), 1); + let staged = stage(uri, plans[0].clone(), budget(), None).await.unwrap(); + restarted.commit(vec![staged]).await.unwrap(); + assert_eq!(rows(uri).await, 2); + } + + #[tokio::test] + async fn legacy_append_without_drain_is_not_duplicated_at_cutover() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_str().unwrap(); + let mut a = writer(uri, "a").await; + put(&a, "a1", 4096).await; + let guard = Arc::new(FailDrain(AtomicBool::new(true))); + let scope = MergeWriteScope::with_pinned_authorizer(guard); + // Open inside the scope so the legacy publisher uses its guard. + scope + .run(async { + a = writer(uri, "a").await; + assert!(a.cleanup_own_shard().await.is_err()); + }) + .await; + scope.drain().await; + assert_eq!(rows(uri).await, 1); + put(&a, "a2", 4096).await; + let mut commit = AppendCoordinator::open(uri, None).await.unwrap(); + let (plans, _) = commit.plan(&["a".into()], 64, 1024 * 1024).await.unwrap(); + let part = stage(uri, plans[0].clone(), budget(), None).await.unwrap(); + assert_eq!(part.rows, 1, "base ID lookup removes legacy retry"); + assert_eq!(commit.commit(vec![part]).await.unwrap(), 2); + assert_eq!(rows(uri).await, 2); + } + + #[tokio::test] + async fn legacy_merge_honors_committed_append_watermark() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_str().unwrap(); + let a = writer(uri, "a").await; + put(&a, "a1", 4096).await; + let scope = + MergeWriteScope::with_pinned_authorizer(Arc::new(FailDrain(AtomicBool::new(true)))); + scope + .run(async { + let mut commit = AppendCoordinator::open(uri, None).await.unwrap(); + let (plans, _) = commit.plan(&["a".into()], 64, 1024 * 1024).await.unwrap(); + let part = stage(uri, plans[0].clone(), budget(), None).await.unwrap(); + assert!(commit.commit(vec![part]).await.is_err()); + }) + .await; + scope.drain().await; + put(&a, "a2", 4096).await; + let mut legacy = writer(uri, "a").await; + legacy.cleanup_own_shard().await.unwrap(); + assert_eq!(rows(uri).await, 2); + let mut commit = AppendCoordinator::open(uri, None).await.unwrap(); + assert!(commit + .plan(&["a".into()], 64, 1024 * 1024) + .await + .unwrap() + .0 + .is_empty()); + } + + #[tokio::test] + async fn partial_stale_prefix_rejected_and_byte_cap_preserves_tail() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_str().unwrap(); + let a = writer(uri, "a").await; + for i in 0..3 { + put(&a, &format!("a{i}"), 128 * 1024).await; + } + let mut commit = AppendCoordinator::open(uri, None).await.unwrap(); + let (plans, _) = commit.plan(&["a".into()], 64, 1024 * 1024).await.unwrap(); + let large = stage(uri, plans[0].clone(), budget(), None).await.unwrap(); + let mut small_plan = plans[0].clone(); + small_plan.max_bytes = 64 * 1024; + let small = stage(uri, small_plan, budget(), None).await.unwrap(); + assert_eq!(small.completed, 1); + commit.commit(vec![small]).await.unwrap(); + assert!(commit + .commit(vec![large]) + .await + .unwrap_err() + .to_string() + .contains("partially stale")); + assert_eq!(rows(uri).await, 1); + assert_eq!( + commit.plan(&["a".into()], 64, 1024 * 1024).await.unwrap().0[0] + .generations + .len(), + 2 + ); + } + #[tokio::test] + async fn changed_schema_rejects_staged_files() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_str().unwrap(); + let a = writer(uri, "a").await; + put(&a, "a1", 4096).await; + let mut commit = AppendCoordinator::open(uri, None).await.unwrap(); + let (plans, _) = commit.plan(&["a".into()], 64, 1024 * 1024).await.unwrap(); + let part = stage(uri, plans[0].clone(), budget(), None).await.unwrap(); + let mut dataset = Dataset::open(uri).await.unwrap(); + dataset + .update_schema_metadata([("changed", "true")]) + .await + .unwrap(); + assert!(commit + .commit(vec![part]) + .await + .unwrap_err() + .to_string() + .contains("schema changed")); + assert_eq!(rows(uri).await, 0); + } + + #[tokio::test] + async fn exact_version_guard_rejects_lances_initial_rebase() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_str().unwrap(); + let _a = writer(uri, "a").await; + let old = Dataset::open(uri).await.unwrap(); + let mut current = old.clone(); + current + .update_metadata([("concurrent", "commit")]) + .await + .unwrap(); + let version = current.version().version; + let handler = lance_table::io::commit::commit_handler_from_url(uri, &None) + .await + .unwrap(); + let handler = Arc::new(crate::merge_write_scope::GuardedCommit::new(handler)); + let expected = old.version().version + 1; + let transaction = Transaction::new( + old.version().version, + Operation::Append { + fragments: Vec::new(), + }, + None, + ); + assert!(crate::merge_write_scope::at_base_version( + expected, + CommitBuilder::new(Arc::new(old)) + .with_commit_handler(handler) + .with_max_retries(0) + .execute(transaction) + ) + .await + .is_err()); + assert_eq!(Dataset::open(uri).await.unwrap().version().version, version); + } + + #[tokio::test] + async fn revoked_owner_cannot_publish_prepared_files() { + #[derive(Debug)] + struct Reject; + impl CommitAuthorizer for Reject { + fn authorize<'a>( + &'a self, + _resource: &'a str, + _version: u64, + ) -> Pin> + Send + 'a>> { + Box::pin(async { Err(Error::io("owner revoked")) }) + } + } + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_str().unwrap(); + let a = writer(uri, "a").await; + put(&a, "a1", 4096).await; + let mut commit = AppendCoordinator::open(uri, None).await.unwrap(); + let (plans, _) = commit.plan(&["a".into()], 64, 1024 * 1024).await.unwrap(); + let part = stage(uri, plans[0].clone(), budget(), None).await.unwrap(); + let version = commit.version(); + let scope = MergeWriteScope::with_pinned_authorizer(Arc::new(Reject)); + let mut publisher = scope.run(AppendCoordinator::open(uri, None)).await.unwrap(); + // A captured owner guard must remain effective after leaving its task-local scope. + let error = publisher.commit(vec![part.clone()]).await.unwrap_err(); + assert!(error.to_string().contains("owner revoked")); + assert_eq!(Dataset::open(uri).await.unwrap().version().version, version); + assert_eq!(rows(uri).await, 0); + // A replacement owner may reuse the immutable files. + commit.commit(vec![part]).await.unwrap(); + assert_eq!(rows(uri).await, 1); + } + /// Opt-in synthetic benchmark. Never scans an existing table: every fixture + /// gets a UUID path beneath the supplied scratch root. + #[tokio::test] + #[ignore = "opt-in rollout append benchmark"] + async fn benchmark_rollout_parallel_append() { + if std::env::var("ROLLOUT_APPEND_BENCH").as_deref() != Ok("1") { + return; + } + for repeat in 0..2 { + for parallel in if repeat == 0 { + [false, true] + } else { + [true, false] + } { + let dir = tempfile::tempdir().unwrap(); + let root = std::env::var("ROLLOUT_APPEND_BENCH_ROOT") + .unwrap_or_else(|_| dir.path().to_string_lossy().into()); + let uri = format!( + "{}/rollout-append-bench-{}", + root.trim_end_matches('/'), + Uuid::new_v4() + ); + let mut writers = Vec::new(); + let mut shards = Vec::new(); + for i in 0..4 { + let shard = format!("worker-{i}"); + let store = writer(&uri, &shard).await; + for generation in 0..8 { + let records: Vec<_> = (0..16) + .map(|row| { + artifact_record(&format!("{i}-{generation}-{row}"), &vec![42; 8192]) + }) + .collect(); + store.add(&records).await.unwrap(); + store.flush().await.unwrap(); + } + writers.push(store); + shards.push(shard); + } + let initial = Dataset::open(&uri).await.unwrap().version().version; + let start = std::time::Instant::now(); + if parallel { + let mut commit = AppendCoordinator::open(&uri, None).await.unwrap(); + let (plans, _) = commit.plan(&shards, 64, 2 * 1024 * 1024).await.unwrap(); + let memory = MergeMemoryBudget::new(64 * 1024 * 1024); + let results = futures::future::try_join_all( + plans + .into_iter() + .map(|plan| stage(&uri, plan, memory.clone(), None)), + ) + .await + .unwrap(); + assert_eq!(commit.commit(results).await.unwrap(), 32); + assert_eq!(memory.reserved(), 0); + } else { + for store in &mut writers { + store.cleanup_own_shard().await.unwrap(); + } + } + let seconds = start.elapsed().as_secs_f64(); + let dataset = Dataset::open(&uri).await.unwrap(); + assert_eq!(dataset.count_rows(None).await.unwrap(), 512); + let mut scan = dataset.scan(); + scan.project(&["id", "binary_payload"]).unwrap(); + let mut rows = scan.try_into_stream().await.unwrap(); + let mut ids = HashSet::new(); + while let Some(batch) = rows.try_next().await.unwrap() { + let keys = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let blobs = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + for row in 0..batch.num_rows() { + assert!(ids.insert(keys.value(row).to_owned())); + assert_eq!(blobs.value(row), vec![42; 8192]); + } + } + println!( + "APPEND_BENCH {}", + json!({"parallel":parallel,"repeat":repeat,"seconds":seconds, + "generations_per_second":32.0/seconds,"versions":dataset.version().version-initial,"rows":ids.len(),"uri":uri}) + ); + } + } + } + #[tokio::test] + async fn disable_then_reenable_preserves_fallback_commit_watermarks() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir.path().to_str().unwrap(); + let a = writer(uri, "a").await; + put(&a, "a1", 4096).await; + let mut commit = AppendCoordinator::open(uri, None).await.unwrap(); + let (plans, _) = commit.plan(&["a".into()], 64, 1024 * 1024).await.unwrap(); + let part = stage(uri, plans[0].clone(), budget(), None).await.unwrap(); + commit.commit(vec![part]).await.unwrap(); + put(&a, "a2", 4096).await; + let scope = + MergeWriteScope::with_pinned_authorizer(Arc::new(FailDrain(AtomicBool::new(true)))); + scope + .run(async { + let mut fallback = writer(uri, "a").await; + assert!(fallback.cleanup_own_shard().await.is_err()); + }) + .await; + scope.drain().await; + assert_eq!(rows(uri).await, 2); + let mut restarted = AppendCoordinator::open(uri, None).await.unwrap(); + let (plans, reclaimed) = restarted + .plan(&["a".into()], 64, 1024 * 1024) + .await + .unwrap(); + assert_eq!(reclaimed, 1); + assert!(plans.is_empty()); + assert_eq!(rows(uri).await, 2); + } +} diff --git a/crates/lance-context-core/src/store_base.rs b/crates/lance-context-core/src/store_base.rs index 95427404..3f6fbd58 100644 --- a/crates/lance-context-core/src/store_base.rs +++ b/crates/lance-context-core/src/store_base.rs @@ -961,6 +961,24 @@ impl StorageBase { /// **older** than the stored one. Reusing the shard's current epoch commits /// the drain and leaves the live writer untouched. async fn prepare_merge(&self, manifest: &ShardManifest) -> LanceResult> { + // A staged append publishes its recovery watermark atomically with its + // files. Honor it even when returning to the legacy merge implementation. + let latest = Self::load_with_options( + self.uri(), + self.storage_options.clone(), + self.session.clone(), + ) + .await?; + crate::rollout_append::reconcile_shard(&latest, self.write_shard).await?; + let marks = crate::rollout_append::watermarks(&latest).await?; + let mut manifest = manifest.clone(); + if let Some(high) = marks.get(&self.write_shard) { + manifest + .flushed_generations + .retain(|g| g.generation > *high); + } + manifest.flushed_generations.sort_by_key(|g| g.generation); + let manifest = &manifest; if manifest.flushed_generations.is_empty() { return Ok(None); } @@ -1044,11 +1062,28 @@ impl StorageBase { // before retrying, including generic stores without schema evolution. self.refresh_latest().await?; self.ensure_latest_schema().await?; + if crate::rollout_append::watermarks(&self.dataset) + .await? + .get(&self.write_shard) + .is_some_and(|high| merged_generations.iter().any(|g| g <= high)) + { + return Err(LanceError::io( + "prepared merge overlaps committed rollout watermark; reprepare", + )); + } + let watermark = if crate::rollout_append::has_cutover(&self.dataset, self.write_shard) { + merged_generations + .iter() + .max() + .map(|high| lance_index::mem_wal::MergedGeneration::new(self.write_shard, *high)) + } else { + None + }; if !batches.is_empty() { observe_phase!( "append", - Box::pin(self.merge_prepared_batches(batches, merge_schema)).await + Box::pin(self.merge_prepared_batches(batches, merge_schema, watermark)).await )?; self.pinned_version = None; } @@ -1281,6 +1316,7 @@ impl StorageBase { &mut self, batches: Vec, merge_schema: Arc, + watermark: Option, ) -> LanceResult<()> { observe_phase!("index", self.ensure_merge_key_index().await)?; let key_index = merge_schema.index_of(&self.key_column)?; @@ -1331,11 +1367,30 @@ impl StorageBase { self.dataset = Arc::unwrap_or_clone(dataset); } - let reader = RecordBatchIterator::new( - batches.into_iter().map(Ok::), - merge_schema, - ); - Box::pin(self.dataset.append(reader, None)).await?; + if let Some(watermark) = watermark { + let params = WriteParams { + mode: WriteMode::Append, + ..Default::default() + }; + let staged = lance::dataset::InsertBuilder::new(Arc::new(self.dataset.clone())) + .with_params(¶ms) + .execute_uncommitted(batches) + .await?; + let lance::dataset::transaction::Operation::Append { fragments } = staged.operation + else { + return Err(LanceError::io( + "legacy fallback did not stage append fragments", + )); + }; + crate::rollout_append::commit_files(&mut self.dataset, fragments, vec![watermark]) + .await?; + } else { + let reader = RecordBatchIterator::new( + batches.into_iter().map(Ok::), + merge_schema, + ); + Box::pin(self.dataset.append(reader, None)).await?; + } Ok(()) } diff --git a/crates/lance-context-master/src/catchup/executor.rs b/crates/lance-context-master/src/catchup/executor.rs index b026aeda..a8373041 100644 --- a/crates/lance-context-master/src/catchup/executor.rs +++ b/crates/lance-context-master/src/catchup/executor.rs @@ -79,6 +79,9 @@ pub async fn execute(mut config: MasterConfig, target: &str) -> Result<()> { } async fn merge_passes(state: &Arc, target: &str) -> Result { + if state.config.append.enabled(target) { + return crate::rollout_append::run(state, target).await; + } let config = &state.config.catchup; let session = RolloutStore::build_session(96 * 1024 * 1024, 32 * 1024 * 1024); let budget = MergeMemoryBudget::new(config.merge_memory_bytes); diff --git a/crates/lance-context-master/src/catchup/kubernetes.rs b/crates/lance-context-master/src/catchup/kubernetes.rs index b068144c..6708bb10 100644 --- a/crates/lance-context-master/src/catchup/kubernetes.rs +++ b/crates/lance-context-master/src/catchup/kubernetes.rs @@ -235,6 +235,23 @@ pub(super) fn render_job(config: &MasterConfig, record: &Record, mut spec: Value } let overrides = [ ("DATA_DIR", config.data_dir.clone()), + ("WORKER_ENDPOINTS", config.worker_endpoints.join(",")), + ( + "ROLLOUT_APPEND_TARGETS", + config.append.rollout_append_targets.join(","), + ), + ( + "ROLLOUT_APPEND_CONCURRENCY", + config.append.rollout_append_concurrency.to_string(), + ), + ( + "ROLLOUT_APPEND_MAX_GENERATIONS", + config.append.rollout_append_max_generations.to_string(), + ), + ( + "ROLLOUT_APPEND_MAX_BYTES", + config.append.rollout_append_max_bytes.to_string(), + ), ( "ROLLOUT_KEY_INDEX_TYPE", config.key_index_type.as_str().into(), diff --git a/crates/lance-context-master/src/config.rs b/crates/lance-context-master/src/config.rs index bdb51f23..ad925bb0 100644 --- a/crates/lance-context-master/src/config.rs +++ b/crates/lance-context-master/src/config.rs @@ -12,6 +12,8 @@ pub struct MasterConfig { pub maintenance: MaintenanceConfig, #[command(flatten)] pub catchup: crate::catchup::CatchupConfig, + #[command(flatten)] + pub append: crate::rollout_append::AppendConfig, /// Data directory / object-store prefix shared with the data-plane server. #[arg(long, env = "DATA_DIR", default_value = "./data")] pub data_dir: String, diff --git a/crates/lance-context-master/src/lib.rs b/crates/lance-context-master/src/lib.rs index 3fda31c2..34603d89 100644 --- a/crates/lance-context-master/src/lib.rs +++ b/crates/lance-context-master/src/lib.rs @@ -6,6 +6,7 @@ pub mod discovery; pub mod error; mod maintenance_execution; mod merge_execution; +pub mod rollout_append; pub mod routes; pub mod scanner; pub mod scheduler; diff --git a/crates/lance-context-master/src/maintenance_execution.rs b/crates/lance-context-master/src/maintenance_execution.rs index 59a31b6d..06d16fde 100644 --- a/crates/lance-context-master/src/maintenance_execution.rs +++ b/crates/lance-context-master/src/maintenance_execution.rs @@ -197,7 +197,7 @@ where .catchup .job_name .as_deref() - .ok_or("missing catch-up Job identity")? + .unwrap_or(&claim.task.id) } else { &claim.task.id }, diff --git a/crates/lance-context-master/src/merge_execution.rs b/crates/lance-context-master/src/merge_execution.rs index dffe3e62..7c374f01 100644 --- a/crates/lance-context-master/src/merge_execution.rs +++ b/crates/lance-context-master/src/merge_execution.rs @@ -59,6 +59,19 @@ pub(crate) async fn run_merge_wal( } return run_legacy(state, target).await; } + if state.config.append.enabled(target) { + if let Some(old) = coordinator.get(target).await? { + ensure_recovery_due(&coordinator, &old).await?; + recover_execution(state, &coordinator, &proof, old).await?; + } + return crate::maintenance_execution::run_as( + state, + claim, + lance_context_merge::MaintenanceKind::Catchup, + crate::rollout_append::run(state, target), + ) + .await; + } // Reconcile/fence before any other table mutation. One retry of the fan-out // after a completed barrier lets healthy shards progress in this task. for recovery_round in 0..2 { diff --git a/crates/lance-context-master/src/rollout_append.rs b/crates/lance-context-master/src/rollout_append.rs new file mode 100644 index 00000000..ae765a09 --- /dev/null +++ b/crates/lance-context-master/src/rollout_append.rs @@ -0,0 +1,595 @@ +//! Workers stage files concurrently; the master publishes bounded groups under +//! the existing table claim, idle watchdog, and manifest-version fence. +use crate::state::MasterState; +use futures::{stream, FutureExt, StreamExt}; +use lance_context_core::{ + merge_write_scope::checkpoint, + rollout_append::{AppendCoordinator, AppendPlan, StageEvent, StagedAppend}, +}; +use serde::Deserialize; +use std::{collections::HashSet, sync::Arc, time::Duration}; + +type Result = std::result::Result; + +#[derive(Clone, Debug, clap::Args)] +pub struct AppendConfig { + /// Opt-in immutable rollout targets. '*' selects all owned rollouts. + #[arg(long, env = "ROLLOUT_APPEND_TARGETS", value_delimiter = ',')] + pub rollout_append_targets: Vec, + /// Concurrent staging RPCs per table. Uses the workers' shared merge slots/budget. + #[arg(long, env = "ROLLOUT_APPEND_CONCURRENCY", default_value_t = 4)] + pub rollout_append_concurrency: usize, + #[arg(long, env = "ROLLOUT_APPEND_MAX_GENERATIONS", default_value_t = 64)] + pub rollout_append_max_generations: usize, + #[arg(long, env = "ROLLOUT_APPEND_MAX_BYTES", default_value_t = 67_108_864)] + pub rollout_append_max_bytes: usize, +} +impl Default for AppendConfig { + fn default() -> Self { + Self { + rollout_append_targets: Vec::new(), + rollout_append_concurrency: 4, + rollout_append_max_generations: 64, + rollout_append_max_bytes: 67_108_864, + } + } +} +impl AppendConfig { + pub fn enabled(&self, target: &str) -> bool { + !target.starts_with("generic:") + && self + .rollout_append_targets + .iter() + .any(|t| t == "*" || t == target) + } + fn validate(&self) -> Result<()> { + if !(1..=32).contains(&self.rollout_append_concurrency) + || !(1..=256).contains(&self.rollout_append_max_generations) + || self.rollout_append_max_bytes == 0 + || self.rollout_append_max_bytes > lance_context_core::rollout_append::MAX_STAGE_BYTES + { + return Err("invalid rollout append limits".into()); + } + Ok(()) + } +} + +#[derive(Deserialize)] +struct Capabilities { + #[serde(default)] + rollout_append_protocol: u32, + shard_name: Option, + #[serde(default)] + owned_targets: Vec, + #[serde(default)] + drain_targets: Vec, +} + +async fn stage_remote( + client: &reqwest::Client, + endpoint: &str, + target: &str, + plan: AppendPlan, +) -> Result { + let expected = plan.clone(); + let request = client + .post(format!( + "{}/api/v1/internal/rollout-append/{target}", + endpoint.trim_end_matches('/') + )) + .json(&plan); + let mut response = tokio::time::timeout(Duration::from_secs(10), request.send()) + .await + .map_err(|_| "staging admission timed out")? + .map_err(|e| e.to_string())? + .error_for_status() + .map_err(|e| e.to_string())?; + let mut pending = Vec::new(); + let mut sequence = 0; + loop { + // A silent socket is not evidence of worker progress. Completed read/ + // encoding steps are relayed to the existing maintenance idle watchdog. + let chunk = tokio::time::timeout(Duration::from_secs(10), response.chunk()) + .await + .map_err(|_| "staging progress stream disconnected")? + .map_err(|e| e.to_string())? + .ok_or("staging stream ended without a result")?; + pending.extend_from_slice(&chunk); + if pending.len() > 4 * 1024 * 1024 { + return Err("staged metadata exceeds 4 MiB".into()); + } + while let Some(end) = pending.iter().position(|b| *b == b'\n') { + let event: StageEvent = + serde_json::from_slice(&pending[..end]).map_err(|e| e.to_string())?; + pending.drain(..=end); + match event { + StageEvent::Progress(current) if current > sequence => { + sequence = current; + checkpoint(); + } + StageEvent::Progress(_) => {} + StageEvent::Failed(error) => return Err(error), + StageEvent::Complete(part) => { + if part.plan != expected { + return Err("worker returned a different staging plan".into()); + } + checkpoint(); + return Ok(part); + } + } + } + } +} + +pub(crate) async fn run(state: &Arc, target: &str) -> Result { + let config = &state.config.append; + config.validate()?; + let mut endpoints = Vec::new(); + let mut shards = Vec::new(); + let mut seen = HashSet::new(); + // Capability probing is cheap and bounded; never send a new protocol RPC to + // an older worker and silently fall back to its manifest-publishing merge. + // Dedicated Jobs contribute their own CPU/memory. The ordinary master + // remains metadata-only; it never falls back to local payload processing. + let dedicated = state.config.catchup.target.is_some(); + let probe_futures: Vec<_> = state + .config + .worker_endpoints + .iter() + .filter(|_| !dedicated) + .cloned() + .map(|endpoint| { + let client = state.http.clone(); + async move { + let response = client + .get(format!( + "{}/api/v1/internal/merge-executor", + endpoint.trim_end_matches('/') + )) + .timeout(Duration::from_secs(5)) + .send() + .await + .map_err(|e| e.to_string())? + .error_for_status() + .map_err(|e| e.to_string())?; + let caps: Capabilities = response.json().await.map_err(|e| e.to_string())?; + Ok::<_, String>((endpoint, caps)) + } + .boxed() + }) + .collect(); + let mut probes = stream::iter(probe_futures).buffer_unordered(8); + while let Some(result) = probes.next().await { + let (endpoint, caps) = match result { + Ok(value) => value, + Err(error) => { + tracing::warn!(%error, "staging worker unavailable; healthy workers will read its WAL"); + continue; + } + }; + if caps.rollout_append_protocol != 1 { + return Err(format!("worker {endpoint} lacks rollout append protocol 1")); + } + if !caps.owned_targets.iter().any(|t| t == "*" || t == target) + || caps.drain_targets.iter().any(|t| t == "*" || t == target) + { + return Err(format!( + "worker {endpoint} has not enabled owned merges for {target}" + )); + } + if let Some(shard) = caps.shard_name { + if seen.insert(shard.clone()) { + shards.push(shard); + } + } + endpoints.push(endpoint); + } + // Also visit historical shards whose worker no longer exists. Any staging + // worker can read their immutable WAL; no writer epoch is acquired. + for shard in &state.config.catchup.shards { + if seen.insert(shard.clone()) { + shards.push(shard.clone()); + } + } + if endpoints.is_empty() && !dedicated { + return Err("no staging workers available".into()); + } + let mut coordinator = AppendCoordinator::open( + &state.rollout_uri(target), + Some(lance_context_core::RolloutStore::build_session( + 32 * 1024 * 1024, + 32 * 1024 * 1024, + )), + ) + .await + .map_err(|e| e.to_string())?; + let started = tokio::time::Instant::now(); + let mut total = 0; + for _ in 0..16 { + if started.elapsed() >= Duration::from_secs(state.config.catchup.slice_secs) { + break; + } + let reclaimed = run_pass(state, target, &endpoints, &shards, &mut coordinator).await?; + total += reclaimed; + if reclaimed == 0 { + break; + } + } + Ok(format!("staged append reclaimed {total} generations")) +} + +async fn run_pass( + state: &Arc, + target: &str, + endpoints: &[String], + shards: &[String], + coordinator: &mut AppendCoordinator, +) -> Result { + let config = &state.config.append; + let (plans, mut reclaimed) = coordinator + .plan( + shards, + config.rollout_append_max_generations, + config.rollout_append_max_bytes, + ) + .await + .map_err(|e| e.to_string())?; + let local = state.config.catchup.target.as_ref().map(|_| { + ( + lance_context_core::MergeMemoryBudget::new(state.config.catchup.merge_memory_bytes), + lance_context_core::RolloutStore::build_session(96 * 1024 * 1024, 32 * 1024 * 1024), + ) + }); + let stage_futures: Vec<_> = plans.into_iter().enumerate().map(|(i, plan)| { + let client = state.http.clone(); + let routes = (!endpoints.is_empty()).then(|| ( + endpoints[i % endpoints.len()].clone(), endpoints[(i + 1) % endpoints.len()].clone())); + let target = target.to_owned(); + let uri = state.rollout_uri(&target); + let local = local.clone(); + async move { + let result = if let Some((budget, session)) = local { + lance_context_core::rollout_append::stage(&uri, plan.clone(), budget, Some(session)) + .await.map_err(|e| e.to_string()) + } else { + let (endpoint, fallback) = routes.expect("ordinary master requires staging endpoints"); + match stage_remote(&client, &endpoint, &target, plan.clone()).await { + Ok(part) => Ok(part), + Err(error) => { + tracing::warn!(%target, %endpoint, %error, "retrying immutable staging on another worker"); + stage_remote(&client, &fallback, &target, plan.clone()).await + } + } + }; + result.map_err(|error| (plan, error)) + }.boxed() + }).collect(); + let mut work = stream::iter(stage_futures).buffer_unordered(config.rollout_append_concurrency); + let mut group = Vec::new(); + let mut errors = Vec::new(); + let mut memory_retries = Vec::new(); + let mut flush = tokio::time::Instant::now() + Duration::from_secs(2); + loop { + tokio::select! { + next = work.next() => match next { + Some(Ok(part)) => { + if group.is_empty() { flush = tokio::time::Instant::now() + Duration::from_secs(2); } + group.push(part); + }, + Some(Err((plan, error))) => { + if local.is_some() && error.contains("merge memory budget busy while reading first generation; retry") { + memory_retries.push(plan); + } else { errors.push(error); } + }, + None => break, + }, + _ = tokio::time::sleep_until(flush), if !group.is_empty() => {}, + } + if group.len() >= config.rollout_append_concurrency + || (!group.is_empty() && tokio::time::Instant::now() >= flush) + { + reclaimed += coordinator + .commit(std::mem::take(&mut group)) + .await + .map_err(|e| e.to_string())?; + } + } + if !group.is_empty() { + reclaimed += coordinator.commit(group).await.map_err(|e| e.to_string())?; + } + // Every speculative reader has now dropped its reservation. Retry an + // indivisible oversized generation once alone, without Job-level backoff. + for plan in memory_retries { + let (budget, session) = local.as_ref().unwrap(); + match lance_context_core::rollout_append::stage( + &state.rollout_uri(target), + plan, + budget.clone(), + Some(session.clone()), + ) + .await + { + Ok(part) => { + reclaimed += coordinator + .commit(vec![part]) + .await + .map_err(|e| e.to_string())? + } + Err(error) => errors.push(error.to_string()), + } + } + metrics::counter!("master_rollout_append_generations_reclaimed_total") + .increment(reclaimed as u64); + tracing::info!(%target, reclaimed, failed = errors.len(), version = coordinator.version(), "parallel rollout append completed"); + if !errors.is_empty() { + return Err(format!( + "staged append reclaimed {reclaimed}; {} shard(s) failed: {}", + errors.len(), + errors[0] + )); + } + Ok(reclaimed) +} + +#[cfg(test)] +mod tests { + use super::*; + use axum::{ + extract::{Path, State}, + routing::{get, post}, + Json, Router, + }; + use clap::Parser; + use lance_context_core::{MergeMemoryBudget, RolloutStore, RolloutStoreOptions}; + use serde_json::json; + + #[derive(Clone)] + struct Workers { + uri: String, + barrier: Arc, + budget: Arc, + } + + async fn caps(Path(worker): Path) -> Json { + Json( + json!({"rollout_append_protocol": 1, "shard_name": worker, "owned_targets": ["hot"], "drain_targets": []}), + ) + } + async fn prepare(State(workers): State, Json(plan): Json) -> String { + // This test deadlocks if the master regresses to serial worker fan-out. + workers.barrier.wait().await; + let part = + lance_context_core::rollout_append::stage(&workers.uri, plan, workers.budget, None) + .await + .unwrap(); + format!( + "{}\n", + serde_json::to_string(&StageEvent::Complete(part)).unwrap() + ) + } + + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn master_stages_two_workers_in_parallel_and_commits_one_version() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir + .path() + .join("hot.rollout.lance") + .to_string_lossy() + .to_string(); + for shard in ["a", "b"] { + let store = RolloutStore::open_with_options( + &uri, + RolloutStoreOptions { + shard_id: Some(shard.into()), + merge_after_generations: Some(0), + ..Default::default() + }, + ) + .await + .unwrap(); + let dto = serde_json::from_value( + json!({"id": shard, "rollout_id": "r", "content": "payload"}), + ) + .unwrap(); + store + .add(&[lance_context_core::rollout_record_from_add_request(&dto)]) + .await + .unwrap(); + store.flush().await.unwrap(); + } + let before = lance::Dataset::open(&uri).await.unwrap().version().version; + let workers = Workers { + uri: uri.clone(), + barrier: Arc::new(tokio::sync::Barrier::new(2)), + budget: MergeMemoryBudget::new(8 * 1024 * 1024), + }; + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let router = Router::new() + .route("/{worker}/api/v1/internal/merge-executor", get(caps)) + .route( + "/{worker}/api/v1/internal/rollout-append/{target}", + post(prepare), + ) + .with_state(workers.clone()); + let server = tokio::spawn(async move { axum::serve(listener, router).await.unwrap() }); + let mut config = crate::config::MasterConfig::parse_from([ + "test", + "--data-dir", + dir.path().to_str().unwrap(), + ]); + config.worker_endpoints = + vec![format!("http://{address}/a"), format!("http://{address}/b")]; + config.append.rollout_append_targets = vec!["hot".into()]; + config.append.rollout_append_concurrency = 2; + config.etcd.etcd_endpoints = std::env::var("ETCD_TEST_ENDPOINTS") + .expect("ETCD_TEST_ENDPOINTS is required") + .split(',') + .map(str::to_owned) + .collect(); + config.etcd.etcd_prefix = + format!("/rollout-append-test/{}", lance_context_core::generate_id()); + let state = MasterState::new(config).await.unwrap(); + let result = tokio::time::timeout(Duration::from_secs(30), run(&state, "hot")) + .await + .unwrap() + .unwrap(); + assert!(result.contains("2 generations"), "{result}"); + let dataset = lance::Dataset::open(&uri).await.unwrap(); + assert_eq!(dataset.count_rows(None).await.unwrap(), 2); + assert_eq!( + dataset.version().version, + before + 2, + "one cutover metadata version and one combined append" + ); + assert_eq!(workers.budget.reserved(), 0); + server.abort(); + } + + #[tokio::test] + async fn unchanged_remote_progress_does_not_reset_watchdog() { + let dir = tempfile::tempdir().unwrap(); + let store = RolloutStore::open_with_options( + dir.path().to_str().unwrap(), + RolloutStoreOptions { + shard_id: Some("a".into()), + ..Default::default() + }, + ) + .await + .unwrap(); + let dto = serde_json::from_value(json!({"id":"a", "rollout_id":"r"})).unwrap(); + store + .add(&[lance_context_core::rollout_record_from_add_request(&dto)]) + .await + .unwrap(); + store.flush().await.unwrap(); + let mut coordinator = AppendCoordinator::open(dir.path().to_str().unwrap(), None) + .await + .unwrap(); + let plan = coordinator + .plan(&["a".into()], 1, 1024 * 1024) + .await + .unwrap() + .0 + .remove(0); + let part = StagedAppend { + dataset_uri: dir.path().to_string_lossy().into(), + plan: plan.clone(), + completed: 1, + fragments: Vec::new(), + rows: 0, + }; + let body = format!( + "{}\n{}\n{}\n{}\n", + serde_json::to_string(&StageEvent::Progress(0)).unwrap(), + serde_json::to_string(&StageEvent::Progress(1)).unwrap(), + serde_json::to_string(&StageEvent::Progress(1)).unwrap(), + serde_json::to_string(&StageEvent::Complete(part)).unwrap() + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let app = Router::new().route( + "/api/v1/internal/rollout-append/hot", + post(move || { + let body = body.clone(); + async move { body } + }), + ); + let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + let scope = Arc::new(lance_context_core::merge_write_scope::MergeWriteScope::default()); + scope + .run(stage_remote( + &reqwest::Client::new(), + &format!("http://{address}"), + "hot", + plan, + )) + .await + .unwrap(); + assert_eq!( + scope.completed_steps(), + 2, + "one actual progress update plus completion" + ); + server.abort(); + } + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn dedicated_job_uses_its_own_memory_for_oversized_generations() { + let dir = tempfile::tempdir().unwrap(); + let uri = dir + .path() + .join("hot.rollout.lance") + .to_string_lossy() + .to_string(); + for shard in ["a", "b"] { + let store = RolloutStore::open_with_options( + &uri, + RolloutStoreOptions { + shard_id: Some(shard.into()), + merge_after_generations: Some(0), + ..Default::default() + }, + ) + .await + .unwrap(); + let dto = serde_json::from_value( + json!({"id":shard,"rollout_id":"r","content":"x".repeat(2 * 1024 * 1024)}), + ) + .unwrap(); + store + .add(&[lance_context_core::rollout_record_from_add_request(&dto)]) + .await + .unwrap(); + store.flush().await.unwrap(); + } + let mut config = crate::config::MasterConfig::parse_from([ + "test", + "--data-dir", + dir.path().to_str().unwrap(), + ]); + config.catchup.target = Some("hot".into()); + config.catchup.job_name = Some("append-fixture".into()); + config.catchup.shards = vec!["a".into(), "b".into()]; + config.catchup.merge_max_bytes = 1024 * 1024; + config.catchup.merge_memory_bytes = 3 * 1024 * 1024; + config.append.rollout_append_concurrency = 2; + config.append.rollout_append_max_bytes = 1024 * 1024; + // No worker endpoints: this dedicated Pod must actually add compute. + config.etcd.etcd_endpoints = std::env::var("ETCD_TEST_ENDPOINTS") + .expect("ETCD_TEST_ENDPOINTS is required") + .split(',') + .map(str::to_owned) + .collect(); + config.etcd.etcd_prefix = + format!("/rollout-append-test/{}", lance_context_core::generate_id()); + let state = MasterState::new(config).await.unwrap(); + assert!( + tokio::time::timeout(Duration::from_secs(30), run(&state, "hot")) + .await + .unwrap() + .unwrap() + .contains("2 generations") + ); + assert_eq!( + lance::Dataset::open(&uri) + .await + .unwrap() + .count_rows(None) + .await + .unwrap(), + 2 + ); + let store = RolloutStore::open_existing_with_options(&uri, RolloutStoreOptions::default()) + .await + .unwrap(); + for id in ["a", "b"] { + assert_eq!( + store.get_by_id(id).await.unwrap().unwrap().content.unwrap(), + "x".repeat(2 * 1024 * 1024) + ); + } + } +} diff --git a/crates/lance-context-master/src/routes.rs b/crates/lance-context-master/src/routes.rs index 3177eab0..3a42c7af 100644 --- a/crates/lance-context-master/src/routes.rs +++ b/crates/lance-context-master/src/routes.rs @@ -778,6 +778,7 @@ mod tests { fn test_config(dir: &TempDir) -> MasterConfig { MasterConfig { + append: Default::default(), catchup: Default::default(), maintenance: Default::default(), merge_rollout: Default::default(), diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index 5c9a8128..f909f253 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -945,6 +945,7 @@ mod tests { fn config(dir: &TempDir) -> MasterConfig { MasterConfig { + append: Default::default(), catchup: Default::default(), maintenance: Default::default(), merge_rollout: lance_context_merge::rollout::MergeRollout { diff --git a/crates/lance-context-master/src/state.rs b/crates/lance-context-master/src/state.rs index 9d56f9b5..0e09c3d1 100644 --- a/crates/lance-context-master/src/state.rs +++ b/crates/lance-context-master/src/state.rs @@ -295,6 +295,7 @@ mod tests { fn test_config(dir: &TempDir) -> MasterConfig { MasterConfig { + append: Default::default(), catchup: Default::default(), maintenance: Default::default(), merge_rollout: Default::default(), diff --git a/crates/lance-context-master/src/task_store.rs b/crates/lance-context-master/src/task_store.rs index f42265f7..f7dbaf47 100644 --- a/crates/lance-context-master/src/task_store.rs +++ b/crates/lance-context-master/src/task_store.rs @@ -1505,6 +1505,7 @@ mod tests { fn config(dir: &TempDir) -> MasterConfig { MasterConfig { + append: Default::default(), catchup: Default::default(), maintenance: Default::default(), merge_rollout: Default::default(), diff --git a/crates/lance-context-server/Cargo.toml b/crates/lance-context-server/Cargo.toml index 6d838f41..55e92ed8 100644 --- a/crates/lance-context-server/Cargo.toml +++ b/crates/lance-context-server/Cargo.toml @@ -33,6 +33,8 @@ tracing-subscriber = { version = "0.3", features = ["env-filter"] } uuid = { version = "1", features = ["v4"] } [dev-dependencies] +lance = "9.0.0" +tower = { version = "0.5", features = ["util"] } etcd-client = { version = "0.19", features = ["tls"] } # Snapshotting recorder so tests can assert which metric series an operation # emitted, without installing a process-global Prometheus exporter. diff --git a/crates/lance-context-server/src/merge_execution.rs b/crates/lance-context-server/src/merge_execution.rs index 54237291..3b620e5c 100644 --- a/crates/lance-context-server/src/merge_execution.rs +++ b/crates/lance-context-server/src/merge_execution.rs @@ -142,6 +142,8 @@ pub async fn capabilities( "queue_timeout_secs": state.merge_executions.queue_timeout_secs, "idle_timeout_secs": state.merge_executions.idle_timeout_secs, "progress_protocol": 1, + "rollout_append_protocol": 1, + "shard_name": state.instance_id, "owned_targets": state.merge_executions.rollout.owned_targets, "drain_targets": state.merge_executions.rollout.drain_targets}), )) @@ -819,3 +821,173 @@ mod tests { client.delete(proof.key, None).await.unwrap(); } } + +/// This endpoint only writes unreachable immutable files. Dropping its HTTP +/// body cancels preparation; there is no detached manifest publisher to fence. +pub async fn stage_append( + State(state): State>, + axum::extract::Path(name): axum::extract::Path, + Json(plan): Json, +) -> Result { + use axum::response::IntoResponse; + use lance_context_core::{ + merge_write_scope::MergeWriteScope, + rollout_append::{stage, StageEvent}, + }; + lance_context_core::validate_store_name(&name).map_err(AppError::InvalidRequest)?; + plan.validate().map_err(AppError::from_lance)?; + if !state.merge_executions.owned(&name) || state.merge_executions.rollout.draining(&name) { + return Err(AppError::Overloaded( + "staged append requires an owned, non-draining rollout".into(), + )); + } + let budget = state.merge_budget.clone().ok_or_else(|| { + AppError::Overloaded("staged append requires a shared merge memory budget".into()) + })?; + if state.rollout_merge_max_bytes == 0 || plan.max_bytes > state.rollout_merge_max_bytes { + return Err(AppError::InvalidRequest( + "staged append exceeds worker byte limit".into(), + )); + } + let scope = Arc::new(MergeWriteScope::default()); + let worker_scope = scope.clone(); + let idle = Duration::from_secs(state.merge_executions.idle_timeout_secs); + let work = Box::pin(async move { + let _slot = state.acquire_merge_slot().await; + worker_scope + .run(stage( + &state.rollout_uri(&name), + plan, + budget, + state.rollout_session.clone(), + )) + .await + }); + let stream = futures::stream::unfold( + Some((work, scope, 0, tokio::time::Instant::now())), + move |next| async move { + let (mut work, scope, mut sequence, mut changed) = next?; + let event = tokio::select! { + result = &mut work => match result { + Ok(part) => StageEvent::Complete(part), + Err(error) => StageEvent::Failed(error.to_string()), + }, + _ = tokio::time::sleep(Duration::from_secs(1)) => { + let current = scope.completed_steps(); + if current > sequence { sequence = current; changed = tokio::time::Instant::now(); } + if changed.elapsed() >= idle { StageEvent::Failed("staged append made no progress".into()) } + else { StageEvent::Progress(sequence) } + } + }; + let terminal = !matches!(event, StageEvent::Progress(_)); + let mut bytes = serde_json::to_vec(&event).expect("serializable stage event"); + bytes.push(b'\n'); + let next = if terminal { + None + } else { + Some((work, scope, sequence, changed)) + }; + Some((Ok::<_, std::convert::Infallible>(bytes), next)) + }, + ); + Ok(( + [(axum::http::header::CONTENT_TYPE, "application/x-ndjson")], + axum::body::Body::from_stream(stream), + ) + .into_response()) +} + +#[cfg(test)] +mod append_tests { + use super::*; + use axum::{ + body::{to_bytes, Body}, + http::Request, + }; + use lance_context_core::{ + rollout_append::{AppendCoordinator, StageEvent}, + MergeMemoryBudget, RolloutStore, RolloutStoreOptions, + }; + use tower::ServiceExt; + + #[tokio::test] + async fn staging_http_writes_only_files_and_uses_shared_budget() { + let dir = tempfile::tempdir().unwrap(); + let mut state = + AppState::new_for_test_with_instance(dir.path().to_path_buf(), Some("a".into())).await; + state + .merge_executions + .rollout + .owned_targets + .push("hot".into()); + state.merge_budget = Some(MergeMemoryBudget::new(4 * 1024 * 1024)); + let uri = state.rollout_uri("hot"); + let store = RolloutStore::open_with_options( + &uri, + RolloutStoreOptions { + shard_id: Some("a".into()), + ..Default::default() + }, + ) + .await + .unwrap(); + let dto = serde_json::from_value( + serde_json::json!({"id":"a", "rollout_id":"r", "content":"hello"}), + ) + .unwrap(); + store + .add(&[lance_context_core::rollout_record_from_add_request(&dto)]) + .await + .unwrap(); + store.flush().await.unwrap(); + let mut coordinator = AppendCoordinator::open(&uri, None).await.unwrap(); + let plan = coordinator + .plan(&["a".into()], 1, 1024 * 1024) + .await + .unwrap() + .0 + .remove(0); + let version = coordinator.version(); + let state = Arc::new(state); + let request = Request::builder() + .method("POST") + .uri("/api/v1/internal/rollout-append/hot") + .header("content-type", "application/json") + .body(Body::from(serde_json::to_vec(&plan).unwrap())) + .unwrap(); + let response = crate::routes::router() + .with_state(state.clone()) + .oneshot(request) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let body = to_bytes(response.into_body(), 4 * 1024 * 1024) + .await + .unwrap(); + let event: StageEvent = serde_json::from_slice( + body.split(|b| *b == b'\n') + .rfind(|line| !line.is_empty()) + .unwrap(), + ) + .unwrap(); + let StageEvent::Complete(part) = event else { + panic!("staging failed: {event:?}") + }; + assert_eq!(part.rows, 1); + assert_eq!(state.merge_budget.as_ref().unwrap().reserved(), 0); + assert_eq!( + lance::Dataset::open(&uri).await.unwrap().version().version, + version + ); + assert_eq!( + lance::Dataset::open(&uri) + .await + .unwrap() + .count_rows(None) + .await + .unwrap(), + 0 + ); + assert_eq!(coordinator.commit(vec![part]).await.unwrap(), 1); + } +} diff --git a/crates/lance-context-server/src/routes/mod.rs b/crates/lance-context-server/src/routes/mod.rs index e23c2ab4..d5bd9774 100644 --- a/crates/lance-context-server/src/routes/mod.rs +++ b/crates/lance-context-server/src/routes/mod.rs @@ -30,6 +30,10 @@ pub fn router() -> Router> { "/api/v1/internal/merge-executor/cancel", post(crate::merge_execution::cancel), ) + .route( + "/api/v1/internal/rollout-append/{name}", + post(crate::merge_execution::stage_append), + ) .route("/api/v1/health", get(health::health_check)) .route("/api/v1/contexts", post(contexts::create_context)) .route("/api/v1/contexts", get(contexts::list_contexts)) diff --git a/docs/rollout-parallel-append.md b/docs/rollout-parallel-append.md new file mode 100644 index 00000000..131f76cc --- /dev/null +++ b/docs/rollout-parallel-append.md @@ -0,0 +1,141 @@ +# Parallel rollout WAL append + +Rollout rows are immutable and their IDs are never reused for different records. +This opt-in path moves WAL reading, encoding and immutable file uploads to workers, +then publishes several workers' files in one base-table transaction on the master. +Generic/context stores retain their existing keyed merge behavior. Staging requires +Lance V2 files, whose dictionary encodings are local to each file. + +## Configuration + +Deploy the compatible worker binary first. Select **owned, non-draining rollout +names** with `ROLLOUT_APPEND_TARGETS=table1,table2` on the master. `*` selects all +owned rollout targets; generic targets are always excluded. The empty default +keeps the existing path. Every staging worker must advertise append protocol 1 and owned-merge +configuration for the target. All ingestion writers must already honor the same +owned-target configuration; the table claim cannot exclude legacy, uncoordinated +writers. There is no fallback to an older +manifest-publishing worker RPC after a capability or staging failure. + +| Setting | Default | Bound | +| --- | --- | --- | +| `ROLLOUT_APPEND_CONCURRENCY` | 4 | 1–32 concurrent staging RPCs per table | +| `ROLLOUT_APPEND_MAX_GENERATIONS` | 64 | 1–256 generations per shard per pass | +| `ROLLOUT_APPEND_MAX_BYTES` | 64 MiB | 1–256 MiB before starting another generation | + +Worker merge slots and the existing process-wide `MergeMemoryBudget` are shared +with other merge work. A worker refuses staging if that budget is disabled or the +requested byte limit exceeds its configured limit. Generations remain indivisible; +an oversized first generation retains the existing exclusive-growth behavior. +Reservations account for twice buffered Arrow bytes to accommodate filtering. +Encoding and Lance caches also consume memory, so these settings are not an RSS cap. + +The master enumerates shard directories using metadata only, and also visits +advertised stable worker identities and additional `CATCHUP_SHARDS`. Unavailable +workers are skipped during staging admission; healthy workers can process their +WAL. An explicit incompatible capability response rejects the new protocol. Any compatible +worker can stage immutable generations from any shard; it does not claim its epoch. +Dedicated catch-up executors receive these settings through the generated Job +and perform parallel staging with their own CPU and shared catch-up memory budget; +they do not send their data work back to ordinary workers. A speculative read +which cannot grow its reservation retries once alone after other readers finish. +Both normal master merges and catch-up executions use the same publisher under +the existing maintenance ownership protocol. The ordinary master never reads WAL +payloads locally, including when no remote worker is available. + +## Commit and recovery + +1. Read bounded shard prefixes and capture the base version and schema. +2. Workers stage files with `InsertBuilder::execute_uncommitted`, returning only + fragment descriptors. Staging cannot publish base/shard manifests. +3. Publish a bounded group when it fills or two seconds after its first result. + A slow worker does not hold completed groups indefinitely. The two-second + interval governs batching, not execution expiry. +4. Use Lance's `Operation::Update` with **no removed/updated fragments**: its new + fragments and `merged_generations` watermarks are committed atomically. This + performs no row update, key-index maintenance, or target payload scan. +5. Drain only the WAL generations covered by the committed watermark, preserving + concurrent flushes and the live writer epoch. Delete drained directories + asynchronously, as with the previous merge implementation. + +An uncertain commit response is resolved by reopening the base and checking its +watermarks. Fully committed results become drain-only retries. Partially stale +prefixes must be replanned; appending their whole files would duplicate rows. +An unexpected worker dataset URI, schema/storage-format change, or changed WAL +prefix rejects the staged result before publication. + +The final commit handler pins the exact validated next version. Merely setting +Lance's retry count to zero is insufficient: it still attempts transaction rebase +before its first write. Ownership authorization and the existing shielded manifest +write scope remain active around the conditional storage write. Master failure or +revocation uses the existing recovery fence before another publisher takes over. + +A staging failure gets one retry on another worker. Old staging can leave orphan +files but cannot publish them. Persistent failures return to the existing durable +merge failure/backoff policy. HTTP progress is NDJSON: unchanged liveness messages +do not reset the master idle watchdog. Completed reads and increasing Lance +write statistics do; upload-buffer progress is not evidence of a committed table. + +## Existing tables and rollback + +At first use, persist a per-shard cutover boundary in table metadata. The historical +WAL prefix may contain rows from a legacy append whose subsequent drain failed. +While staging this prefix, probe the base for IDs in chunks of 1024, projecting +**only `id`**, and omit already-published immutable rows. This costs key lookups +until the historical prefix is consumed; it does not read old payloads or delete +base rows. Generations beyond the cutover do not require this migration lookup. + +Merge retries are idempotent through the atomic generation watermark. This is not +an arbitrary upsert or ingest-deduplication API: replaying the same ID into different +new generations/shards after cutover violates rollout's no-reappend contract. +Do not enable the path for workloads requiring that behavior or row updates. + +Disabling the master setting returns to legacy keyed merge. Updated legacy +preparation also reconciles committed watermarks before reading WAL, and its +append records a new atomic watermark after cutover. A staged commit interrupted +before drain can safely finish through that path, and disabling/re-enabling the +setting does not lose the fallback publisher's commit evidence. Rolling back to +an older binary without this compatibility code requires resolving pending legacy +writes before re-enabling staged append. Changes to +ownership/drain configuration still follow the existing ownership rollout rules. + +Staged files are not reader-visible until their manifest is committed. Files from +cancelled/failed stages remain subject to normal unverified-file retention; they +must not be eagerly deleted while another execution may still reference them. +No separate orphan-file vacuum or changes to compaction scheduling are introduced. +New fragments may be uncovered by the existing ID index until normal index +maintenance runs; point queries must retain their existing uncovered-fragment scan. + +## Verification + +Tests exercise parallel file preparation and a single combined commit, live writes +arriving between stage and publication, duplicate results, partial stale prefixes, +byte bounds, legacy append-without-drain migration, restart after base publication, +and fallback through the existing merge implementation. Master HTTP tests require +two workers to reach a barrier together, catching accidental serial fan-out. + +A local ARM64 debug-build benchmark (Lance 9.0.0, one process, local filesystem) +used four shards × eight generations × sixteen 8-KiB payload rows, about 4 MiB per +fixture. Two fresh fixtures per variant, with the second order reversed, produced: + +| Merge path | Mean time | Generations/s | Base version increments | +| --- | --- | --- | --- | +| Existing serial keyed merge | 1.141 s | 28.0 | 11 | +| Four concurrent stages + combined append | 0.724 s | 44.2 | 2 | + +Both variants consume the same 32 generations and verify all 512 final IDs and +payloads, with no duplicate base rows. The append path's two versions are the +one-time cutover metadata commit and one combined data/watermark commit. +This is approximately 1.58× throughput for this small synthetic workload, not an +Azure or multi-Pod production throughput measurement. Large-blob and highly +fragmented tables still need deployment-environment measurements. + +Run the opt-in benchmark with: + +```bash +ROLLOUT_APPEND_BENCH=1 cargo test -p lance-context-core --lib \ + benchmark_rollout_parallel_append -- --ignored --nocapture +``` + +`ROLLOUT_APPEND_BENCH_ROOT` optionally selects a scratch object-store prefix. Each +fixture creates a new UUID subdirectory; it never scans an existing table.