diff --git a/.github/workflows/rust-test.yml b/.github/workflows/rust-test.yml index d25bb6a..e30c5d6 100644 --- a/.github/workflows/rust-test.yml +++ b/.github/workflows/rust-test.yml @@ -81,7 +81,9 @@ jobs: - name: Run etcd-backed HA test suite env: ETCD_TEST_ENDPOINTS: http://127.0.0.1:2379 - run: cargo test -p lance-context-master -p lance-context-server -p lance-context-merge -- --ignored + run: | + cargo test -p lance-context-master -p lance-context-server -p lance-context-merge -- --ignored + cargo test -p lance-context-core registry --lib -- --ignored coverage: runs-on: ubuntu-24.04 @@ -132,6 +134,8 @@ jobs: run: | cargo llvm-cov --no-report --all-features \ -p lance-context-master -p lance-context-server -p lance-context-merge -- --ignored + cargo llvm-cov --no-report --all-features \ + -p lance-context-core --lib -- registry --ignored - name: Merge coverage reports run: cargo llvm-cov report --lcov --output-path lcov.info - name: Upload coverage to Codecov diff --git a/Cargo.lock b/Cargo.lock index 0ed20d8..cc0d5ba 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6352,7 +6352,9 @@ dependencies = [ "async-trait", "base64", "chrono", + "clap", "datafusion 54.1.0", + "etcd-client", "futures", "lance 9.0.0", "lance-context-api", diff --git a/crates/lance-context-core/Cargo.toml b/crates/lance-context-core/Cargo.toml index ef6be65..d1e3a42 100644 --- a/crates/lance-context-core/Cargo.toml +++ b/crates/lance-context-core/Cargo.toml @@ -43,6 +43,8 @@ serde = { version = "1", features = ["derive"] } serde_json = "1" futures = "0.3" tokio = { version = "1", features = ["sync", "time"] } +clap = { version = "4", features = ["derive", "env"] } +etcd-client = { version = "0.19", features = ["tls"] } tracing = "0.1" uuid = { version = "1.20.0", features = ["v4", "v5", "v7"] } diff --git a/crates/lance-context-core/src/etcd.rs b/crates/lance-context-core/src/etcd.rs new file mode 100644 index 0000000..2285efe --- /dev/null +++ b/crates/lance-context-core/src/etcd.rs @@ -0,0 +1,197 @@ +//! Shared etcd connection settings and client construction. +//! +//! The master has always needed etcd (task queue, locks); with the etcd store +//! registry the workers need it too. Both parse the same flags, so the flags +//! and the connect logic live here once. + +use std::time::Duration; + +use etcd_client::{Certificate, Client, ConnectOptions, Identity, TlsOptions}; +use lance::{Error as LanceError, Result as LanceResult}; + +/// etcd connection settings, one field per `ETCD_*` flag. +#[derive(Debug, Clone, Default, clap::Args)] +pub struct EtcdConfig { + /// Comma-separated etcd v3 endpoints. + #[arg(long, env = "ETCD_ENDPOINTS", value_delimiter = ',')] + pub etcd_endpoints: Vec, + + /// Namespace for every lance-context key in etcd. + #[arg(long, env = "ETCD_PREFIX", default_value = "/lance-context/master")] + pub etcd_prefix: String, + + /// Optional etcd username. `ETCD_PASSWORD` must also be set. + #[arg(long, env = "ETCD_USERNAME")] + pub etcd_username: Option, + + /// Optional etcd password. `ETCD_USERNAME` must also be set. + #[arg(long, env = "ETCD_PASSWORD")] + pub etcd_password: Option, + + /// Optional PEM CA certificate path for etcd TLS. + #[arg(long, env = "ETCD_CA_CERT")] + pub etcd_ca_cert: Option, + + /// Optional PEM client certificate path for etcd mutual TLS. + #[arg(long, env = "ETCD_CLIENT_CERT")] + pub etcd_client_cert: Option, + + /// Optional PEM client private-key path for etcd mutual TLS. + #[arg(long, env = "ETCD_CLIENT_KEY")] + pub etcd_client_key: Option, +} + +impl EtcdConfig { + /// Whether any endpoint is configured. + pub fn is_configured(&self) -> bool { + !self.etcd_endpoints.is_empty() + } + + /// `etcd_prefix` without a trailing slash, ready for `format!("{p}/…")`. + pub fn prefix(&self) -> &str { + self.etcd_prefix.trim_end_matches('/') + } + + /// Connect with the configured auth and TLS. + pub async fn connect(&self) -> LanceResult { + if self.etcd_endpoints.is_empty() { + return Err(LanceError::io("ETCD_ENDPOINTS is required")); + } + let mut options = ConnectOptions::new() + .with_connect_timeout(Duration::from_secs(5)) + .with_timeout(Duration::from_secs(10)) + .with_keep_alive(Duration::from_secs(10), Duration::from_secs(3)) + .with_require_leader(true); + match (&self.etcd_username, &self.etcd_password) { + (Some(username), Some(password)) => { + options = options.with_user(username, password); + } + (None, None) => {} + _ => { + return Err(LanceError::io( + "ETCD_USERNAME and ETCD_PASSWORD must be configured together", + )) + } + } + if let Some(path) = &self.etcd_ca_cert { + let pem = std::fs::read(path).map_err(|err| { + LanceError::io(format!("failed to read ETCD_CA_CERT '{path}': {err}")) + })?; + let mut tls = TlsOptions::new().ca_certificate(Certificate::from_pem(pem)); + match (&self.etcd_client_cert, &self.etcd_client_key) { + (Some(cert), Some(key)) => { + let cert_pem = std::fs::read(cert).map_err(|err| { + LanceError::io(format!("failed to read ETCD_CLIENT_CERT '{cert}': {err}")) + })?; + let key_pem = std::fs::read(key).map_err(|err| { + LanceError::io(format!("failed to read ETCD_CLIENT_KEY '{key}': {err}")) + })?; + tls = tls.identity(Identity::from_pem(cert_pem, key_pem)); + } + (None, None) => {} + _ => { + return Err(LanceError::io( + "ETCD_CLIENT_CERT and ETCD_CLIENT_KEY must be configured together", + )) + } + } + options = options.with_tls(tls); + } else if self.etcd_client_cert.is_some() || self.etcd_client_key.is_some() { + return Err(LanceError::io( + "ETCD_CA_CERT is required when configuring an etcd client certificate", + )); + } + Client::connect(self.etcd_endpoints.clone(), Some(options)) + .await + .map_err(|err| LanceError::io(format!("connect to etcd: {err}"))) + } +} + +/// Map an etcd client error into a Lance IO error with context. +pub fn etcd_error(context: &'static str) -> impl Fn(etcd_client::Error) -> LanceError { + move |err| LanceError::io(format!("{context}: {err}")) +} + +/// Which backend a store registry reads from, and which (if any) it mirrors +/// writes to. See `docs/design-registry-etcd.md` for the migration this drives. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, clap::ValueEnum)] +pub enum RegistryBackend { + /// The Lance table under `data_dir` (`_registry..lance`). + #[default] + Lance, + /// etcd, under `/registry//`. + Etcd, +} + +/// Registry backend selection flags, shared by the server and master. +#[derive(Debug, Clone, Default, clap::Args)] +pub struct RegistryConfig { + /// Backend the registries read from and write to. + #[arg(long, env = "REGISTRY_BACKEND", value_enum, default_value_t = RegistryBackend::Lance)] + pub registry_backend: RegistryBackend, + + /// Versioned etcd mirror for a Lance primary. Reverse mirroring is unsupported. + #[arg(long, env = "REGISTRY_MIRROR", value_enum)] + pub registry_mirror: Option, +} + +/// Open the registry for one store kind per `RegistryConfig`. +/// +/// `lance_uri` is the Lance table's location; `etcd` is a connected client when +/// `EtcdConfig` is configured. Returns the mirrored composite when a mirror is +/// requested. Fails when a requested backend is not available. +pub async fn open_registry( + kind: &'static str, + lance_uri: &str, + etcd: Option<(&Client, &str)>, + config: &RegistryConfig, +) -> LanceResult> { + use crate::{EtcdRegistry, LanceRegistry, MirroredRegistry, RolloutRegistry, StoreRegistry}; + use std::sync::Arc; + + async fn build( + backend: RegistryBackend, + kind: &'static str, + lance_uri: &str, + etcd: Option<(&Client, &str)>, + ) -> LanceResult> { + Ok(match backend { + RegistryBackend::Lance => Arc::new(LanceRegistry::new( + RolloutRegistry::open_or_create(lance_uri, None).await?, + )), + RegistryBackend::Etcd => { + let (client, prefix) = etcd.ok_or_else(|| { + LanceError::io(format!( + "REGISTRY_BACKEND/REGISTRY_MIRROR=etcd for the {kind} registry but ETCD_ENDPOINTS is required" + )) + })?; + EtcdRegistry::new(client.clone(), prefix, kind) + } + }) + } + + if config.registry_mirror.is_some() + && (config.registry_backend != RegistryBackend::Lance + || config.registry_mirror != Some(RegistryBackend::Etcd)) + { + return Err(LanceError::io("REGISTRY_MIRROR supports only Lance primary -> etcd mirror; reverse mirroring/rolling rollback is unsupported")); + } + let primary = build(config.registry_backend, kind, lance_uri, etcd).await?; + if let Some(etcd) = primary.as_any().downcast_ref::() { + etcd.activate_after_validation(lance_uri).await?; + } + let Some(mirror) = config.registry_mirror else { + return Ok(primary); + }; + if mirror == config.registry_backend { + return Err(LanceError::io( + "REGISTRY_MIRROR must differ from REGISTRY_BACKEND", + )); + } + let mirror = build(mirror, kind, lance_uri, etcd).await?; + Ok(Arc::new(MirroredRegistry { + primary, + mirror, + label: kind, + })) +} diff --git a/crates/lance-context-core/src/lib.rs b/crates/lance-context-core/src/lib.rs index c9666d4..17c3082 100644 --- a/crates/lance-context-core/src/lib.rs +++ b/crates/lance-context-core/src/lib.rs @@ -5,6 +5,7 @@ mod api_impl; mod context; mod datagen; mod datagen_store; +pub mod etcd; mod eval; mod export; pub mod generic_codec; @@ -16,6 +17,7 @@ pub mod metrics; mod namespace; mod record; mod registry; +mod registry_etcd; mod rollout; mod rollout_store; pub mod serde; @@ -64,7 +66,11 @@ pub use record::{ RetrieveResult, SearchResult, StateMetadata, UpdateResult, UpsertResult, LIFECYCLE_ACTIVE, LIFECYCLE_CONTRADICTED, }; -pub use registry::{RegistryEntry, RolloutRegistry}; +pub use registry::{ + backfill_registry, diff_registries, LanceRegistry, MirroredRegistry, RegistryDiff, + RegistryEntry, RegistryMismatch, RolloutRegistry, StoreRegistry, +}; +pub use registry_etcd::EtcdRegistry; pub use rollout::{RolloutRecord, ROLE_ARTIFACT, ROLE_ASSISTANT, ROLE_GRADE, ROLE_TOOL}; pub use rollout_store::{ rollout_schema, ListSource, PreparedMerge, RolloutFilters, RolloutObservation, RolloutPage, diff --git a/crates/lance-context-core/src/registry.rs b/crates/lance-context-core/src/registry.rs index 87c787f..ef2c849 100644 --- a/crates/lance-context-core/src/registry.rs +++ b/crates/lance-context-core/src/registry.rs @@ -32,7 +32,7 @@ use lance::io::{ObjectStoreParams, StorageOptionsAccessor}; use lance::{Error as LanceError, Result as LanceResult}; /// One entry in the rollout-store directory. -#[derive(Debug, Clone, PartialEq, Eq)] +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)] pub struct RegistryEntry { /// Logical store name (the experiment name); unique key. pub name: String, @@ -42,6 +42,231 @@ pub struct RegistryEntry { pub created_at: i64, } +/// A directory of stores: `name -> (uri, created_at)`. +/// +/// The name says which stores exist; the dataset on object storage is the +/// data. Implementations must be safe to share across tasks (`&self`). +#[async_trait::async_trait] +pub trait StoreRegistry: Send + Sync { + /// Whether a store named `name` exists. + async fn contains(&self, name: &str) -> LanceResult; + /// One directory entry, or `None` when absent. + async fn get(&self, name: &str) -> LanceResult>; + /// Every entry, in unspecified order. + async fn list(&self) -> LanceResult>; + /// Insert or replace the entry for `name`. Idempotent. + async fn upsert(&self, name: &str, uri: &str) -> LanceResult<()>; + /// Remove the entry for `name`. No-op when absent. + async fn remove(&self, name: &str) -> LanceResult<()>; + /// Insert every entry whose name is not already registered; returns how + /// many were inserted. Existing rows are left unchanged. + async fn insert_missing(&self, entries: &[(String, String)]) -> LanceResult; + /// The Lance table behind this registry, if any, so its owner can run + /// compaction and version pruning on it. `None` for backends that need no + /// maintenance (etcd). + fn lance_table(&self) -> Option<&tokio::sync::Mutex> { + None + } + /// For downcasting to a concrete backend (the master's migration routes + /// need the two halves of a [`MirroredRegistry`]). + fn as_any(&self) -> &dyn std::any::Any; +} + +/// [`StoreRegistry`] over a [`RolloutRegistry`]. The Lance handle must check +/// out the latest manifest before every read, hence `&mut self` inside and a +/// mutex here; every call is one manifest read on object storage. +pub struct LanceRegistry(pub tokio::sync::Mutex); + +impl LanceRegistry { + pub fn new(inner: RolloutRegistry) -> Self { + Self(tokio::sync::Mutex::new(inner)) + } +} + +#[async_trait::async_trait] +impl StoreRegistry for LanceRegistry { + async fn contains(&self, name: &str) -> LanceResult { + self.0.lock().await.contains(name).await + } + async fn get(&self, name: &str) -> LanceResult> { + self.0.lock().await.get(name).await + } + async fn list(&self) -> LanceResult> { + self.0.lock().await.list().await + } + async fn upsert(&self, name: &str, uri: &str) -> LanceResult<()> { + self.0.lock().await.upsert(name, uri).await + } + async fn remove(&self, name: &str) -> LanceResult<()> { + self.0.lock().await.remove(name).await + } + async fn insert_missing(&self, entries: &[(String, String)]) -> LanceResult { + self.0.lock().await.insert_missing(entries).await + } + fn lance_table(&self) -> Option<&tokio::sync::Mutex> { + Some(&self.0) + } + fn as_any(&self) -> &dyn std::any::Any { + self + } +} + +/// Lance primary with a versioned etcd mirror. Failed mirror writes do not +/// fail the primary mutation; a complete snapshot reconciliation repairs them. +/// Reverse mirroring is deliberately unsupported. +pub struct MirroredRegistry { + pub primary: Arc, + pub mirror: Arc, + pub label: &'static str, +} + +impl MirroredRegistry { + async fn sync_name(&self, name: &str) -> LanceResult<()> { + let (source, target) = migration_pair(&*self.primary, &*self.mirror)?; + let (uri, version, entry) = { + let mut source = source.0.lock().await; + let entry = source.get(name).await?; + (source.uri.clone(), source.dataset.version().version, entry) + }; + target + .apply_source(&uri, version, name, entry.as_ref()) + .await?; + Ok(()) + } +} + +#[async_trait::async_trait] +impl StoreRegistry for MirroredRegistry { + async fn contains(&self, name: &str) -> LanceResult { + self.primary.contains(name).await + } + async fn get(&self, name: &str) -> LanceResult> { + self.primary.get(name).await + } + async fn list(&self) -> LanceResult> { + self.primary.list().await + } + async fn upsert(&self, name: &str, uri: &str) -> LanceResult<()> { + self.primary.upsert(name, uri).await?; + if let Err(error) = self.sync_name(name).await { + tracing::warn!(registry = self.label, name, %error, "registry mirror upsert failed"); + } + Ok(()) + } + async fn remove(&self, name: &str) -> LanceResult<()> { + self.primary.remove(name).await?; + if let Err(error) = self.sync_name(name).await { + tracing::warn!(registry = self.label, name, %error, "registry mirror remove failed"); + } + Ok(()) + } + async fn insert_missing(&self, entries: &[(String, String)]) -> LanceResult { + let inserted = self.primary.insert_missing(entries).await?; + if let Err(error) = backfill_registry(&*self.primary, &*self.mirror).await { + tracing::warn!(registry = self.label, %error, "registry mirror reconciliation failed"); + } + Ok(inserted) + } + fn lance_table(&self) -> Option<&tokio::sync::Mutex> { + self.primary.lance_table() + } + fn as_any(&self) -> &dyn std::any::Any { + self + } +} + +fn migration_pair<'a>( + from: &'a dyn StoreRegistry, + to: &'a dyn StoreRegistry, +) -> LanceResult<(&'a LanceRegistry, &'a crate::EtcdRegistry)> { + match ( + from.as_any().downcast_ref::(), + to.as_any().downcast_ref::(), + ) { + (Some(source), Some(target)) => Ok((source, target)), + _ => Err(LanceError::io( + "registry reconciliation supports only Lance -> etcd", + )), + } +} + +/// Reconcile one consistent Lance snapshot, including changed values and +/// deletions. The source version fences delayed older snapshots and point +/// updates at the destination. The return value counts changed entries. +pub async fn backfill_registry( + from: &dyn StoreRegistry, + to: &dyn StoreRegistry, +) -> LanceResult { + let (source, target) = migration_pair(from, to)?; + let (uri, version, entries) = { + let mut source = source.0.lock().await; + let entries = source.list().await?; + ( + source.uri.clone(), + source.dataset.version().version, + entries, + ) + }; + target.reconcile_source(&uri, version, &entries).await +} + +#[derive(Debug, serde::Serialize)] +pub struct RegistryMismatch { + pub name: String, + pub primary: RegistryEntry, + pub mirror: RegistryEntry, +} + +/// A metadata comparison, not authorization to switch live writers. +#[derive(Debug, serde::Serialize)] +pub struct RegistryDiff { + pub only_in_primary: Vec, + pub only_in_mirror: Vec, + pub mismatched: Vec, +} + +impl RegistryDiff { + pub fn is_empty(&self) -> bool { + self.only_in_primary.is_empty() + && self.only_in_mirror.is_empty() + && self.mismatched.is_empty() + } +} + +pub async fn diff_registries( + a: &dyn StoreRegistry, + b: &dyn StoreRegistry, +) -> LanceResult { + let a: std::collections::BTreeMap<_, _> = a + .list() + .await? + .into_iter() + .map(|e| (e.name.clone(), e)) + .collect(); + let b: std::collections::BTreeMap<_, _> = b + .list() + .await? + .into_iter() + .map(|e| (e.name.clone(), e)) + .collect(); + Ok(RegistryDiff { + only_in_primary: a.keys().filter(|n| !b.contains_key(*n)).cloned().collect(), + only_in_mirror: b.keys().filter(|n| !a.contains_key(*n)).cloned().collect(), + mismatched: a + .iter() + .filter_map(|(name, primary)| { + b.get(name) + .filter(|mirror| *mirror != primary) + .map(|mirror| RegistryMismatch { + name: name.clone(), + primary: primary.clone(), + mirror: mirror.clone(), + }) + }) + .collect(), + }) +} + /// Durable directory of rollout stores, backed by a single Lance dataset. /// /// All operations take `&mut self` because Lance dataset handles are snapshots: diff --git a/crates/lance-context-core/src/registry_etcd.rs b/crates/lance-context-core/src/registry_etcd.rs new file mode 100644 index 0000000..a3398d5 --- /dev/null +++ b/crates/lance-context-core/src/registry_etcd.rs @@ -0,0 +1,708 @@ +//! [`StoreRegistry`] backed by etcd. +//! +//! One key per store: `/registry//` holding +//! `{"uri": ..., "created_at": ...}`. A point lookup is one round trip, an +//! upsert is one `put`, and nothing is retried against a shared commit point. +//! Compare `RolloutRegistry`, where every read is a manifest fetch from object +//! storage and every write is two Lance commits contended by every worker. +//! +//! Store names are validated to be portable path segments before they reach a +//! registry, so they are used as key segments unescaped. + +use std::sync::Arc; + +use etcd_client::{Client, Compare, CompareOp, DeleteOptions, GetOptions, Txn, TxnOp}; +use lance::{Error as LanceError, Result as LanceResult}; +use serde::{Deserialize, Serialize}; + +use crate::etcd::etcd_error; +use crate::registry::{RegistryEntry, StoreRegistry}; + +#[derive(Serialize, Deserialize)] +struct Value { + uri: String, + created_at: i64, + #[serde(default)] + deleted: bool, + #[serde(default)] + source_version: Option, +} + +#[derive(Serialize, Deserialize, Default)] +struct Migration { + source: String, + floor: u64, + active: bool, +} + +/// etcd-backed store directory for one store kind. +pub struct EtcdRegistry { + client: Client, + /// `/registry/` with no trailing slash. + prefix: String, + migration_key: String, +} + +impl EtcdRegistry { + /// `kind` is `rollout`, `generic` or `datagen`; `prefix` is the + /// deployment-wide etcd namespace (see `EtcdConfig::prefix`). + pub fn new(client: Client, prefix: &str, kind: &str) -> Arc { + Arc::new(Self { + client, + prefix: format!("{}/registry/{kind}", prefix.trim_end_matches('/')), + migration_key: format!( + "{}/registry-migrations/{kind}", + prefix.trim_end_matches('/') + ), + }) + } + + fn key(&self, name: &str) -> String { + format!("{}/{name}", self.prefix) + } + + fn entry(&self, key: &[u8], value: &[u8]) -> LanceResult> { + let key = String::from_utf8_lossy(key); + let name = key + .strip_prefix(&format!("{}/", self.prefix)) + .unwrap_or(&key) + .to_string(); + let value: Value = serde_json::from_slice(value) + .map_err(|e| LanceError::io(format!("decode registry entry '{name}': {e}")))?; + if value.deleted { + return Ok(None); + } + Ok(Some(RegistryEntry { + name, + uri: value.uri, + created_at: value.created_at, + })) + } + + fn encode(uri: &str, created_at: i64) -> LanceResult> { + serde_json::to_vec(&Value { + uri: uri.to_string(), + created_at, + deleted: false, + source_version: None, + }) + .map_err(|e| LanceError::io(format!("encode registry entry: {e}"))) + } + + async fn migration(&self) -> LanceResult<(Migration, i64)> { + let response = self + .client + .clone() + .get(self.migration_key.clone(), None) + .await + .map_err(etcd_error("registry migration read"))?; + match response.kvs().first() { + Some(kv) => Ok(( + serde_json::from_slice(kv.value()) + .map_err(|e| LanceError::io(format!("decode registry migration: {e}")))?, + kv.mod_revision(), + )), + None => Ok((Migration::default(), 0)), + } + } + + /// Seal this destination before opening it as the authoritative backend. + /// Operators must first quiesce all Lance registry writers and verify all + /// three kinds. Sealing rejects every delayed mirror request afterwards. + async fn activate(&self) -> LanceResult<()> { + for _ in 0..32 { + let (mut state, revision) = self.migration().await?; + if state.active { + return Ok(()); + } + state.active = true; + if self.put_migration(&state, revision).await? { + return Ok(()); + } + } + Err(LanceError::io("registry activation contention")) + } + + pub(crate) async fn activate_after_validation(&self, lance_uri: &str) -> LanceResult<()> { + let (state, _) = self.migration().await?; + if state.active { + return Ok(()); + } + if !state.source.is_empty() && state.source != lance_uri { + return Err(LanceError::io( + "registry migration source differs from configured data directory", + )); + } + let source = crate::LanceRegistry::new( + crate::RolloutRegistry::open_or_create(lance_uri, None).await?, + ); + let diff = crate::diff_registries(&source, self).await?; + if !diff.is_empty() { + return Err(LanceError::io(format!("registry cutover refused: reconcile names, URIs and timestamps from Lance first (missing={}, extra={}, mismatched={})", diff.only_in_primary.len(), diff.only_in_mirror.len(), diff.mismatched.len()))); + } + self.activate().await + } + + async fn ensure_writable(&self) -> LanceResult<()> { + for _ in 0..32 { + let (mut state, revision) = self.migration().await?; + if state.active { + return Ok(()); + } + if !state.source.is_empty() { + return Err(LanceError::io( + "registry is a migration mirror; activate through validated backend startup first", + )); + } + state.active = true; + // Do not seal a mirror that appeared after our initial read. + if self.put_migration(&state, revision).await? { + return Ok(()); + } + } + Err(LanceError::io("registry activation contention")) + } + + async fn put_migration(&self, state: &Migration, revision: i64) -> LanceResult { + let value = serde_json::to_vec(state).map_err(|e| LanceError::io(e.to_string()))?; + Ok(self + .client + .clone() + .txn( + Txn::new() + .when([Compare::mod_revision( + self.migration_key.clone(), + CompareOp::Equal, + revision, + )]) + .and_then([TxnOp::put(self.migration_key.clone(), value, None)]), + ) + .await + .map_err(etcd_error("registry migration CAS"))? + .succeeded()) + } + + async fn source_fence( + &self, + source: &str, + version: u64, + advance: bool, + ) -> LanceResult> { + for _ in 0..32 { + let (mut state, revision) = self.migration().await?; + if state.active { + return Err(LanceError::io( + "registry destination is active; mirror writes are sealed", + )); + } + if !state.source.is_empty() && state.source != source { + return Err(LanceError::io( + "registry migration source changed; use a fresh namespace", + )); + } + if version < state.floor { + return Ok(None); + } + if revision != 0 && (!advance || version == state.floor) { + return Ok(Some(revision)); + } + state.source = source.to_owned(); + if advance { + state.floor = version; + } + if self.put_migration(&state, revision).await? { + continue; + } + } + Err(LanceError::io("registry mirror fence contention")) + } + + /// Compare both the per-name watermark and the full-snapshot floor in + /// the same transaction as the write. Tombstones retain delete ordering. + pub(crate) async fn apply_source( + &self, + source: &str, + version: u64, + name: &str, + entry: Option<&RegistryEntry>, + ) -> LanceResult { + for _ in 0..32 { + let Some(fence) = self.source_fence(source, version, false).await? else { + return Ok(false); + }; + let key = self.key(name); + let response = self + .client + .clone() + .get(key.clone(), None) + .await + .map_err(etcd_error("registry mirror read"))?; + let revision = match response.kvs().first() { + Some(kv) => { + let old: Value = serde_json::from_slice(kv.value()) + .map_err(|e| LanceError::io(e.to_string()))?; + if old.source_version.is_some_and(|v| v >= version) { + return Ok(false); + } + kv.mod_revision() + } + None => 0, + }; + let value = Value { + uri: entry.map_or_else(String::new, |e| e.uri.clone()), + created_at: entry.map_or(0, |e| e.created_at), + deleted: entry.is_none(), + source_version: Some(version), + }; + let value = serde_json::to_vec(&value).map_err(|e| LanceError::io(e.to_string()))?; + let response = self + .client + .clone() + .txn( + Txn::new() + .when([ + Compare::mod_revision( + self.migration_key.clone(), + CompareOp::Equal, + fence, + ), + Compare::mod_revision(key.clone(), CompareOp::Equal, revision), + ]) + .and_then([TxnOp::put(key, value, None)]), + ) + .await + .map_err(etcd_error("registry mirror apply"))?; + if response.succeeded() { + return Ok(true); + } + } + Err(LanceError::io("registry mirror entry contention")) + } + + pub(crate) async fn reconcile_source( + &self, + source: &str, + version: u64, + entries: &[RegistryEntry], + ) -> LanceResult { + // Advance BEFORE listing destination names. This also fences an old + // snapshot inserting a name the newer snapshot never saw (e.g. a + // create/delete that happened before initial backfill). + if self.source_fence(source, version, true).await?.is_none() { + return Ok(0); + } + let names: std::collections::HashSet<_> = entries.iter().map(|e| e.name.as_str()).collect(); + let mut changed = 0; + for old in self.list().await? { + if !names.contains(old.name.as_str()) { + changed += usize::from(self.apply_source(source, version, &old.name, None).await?); + } + } + for entry in entries { + changed += usize::from( + self.apply_source(source, version, &entry.name, Some(entry)) + .await?, + ); + } + Ok(changed) + } +} + +#[async_trait::async_trait] +impl StoreRegistry for EtcdRegistry { + async fn contains(&self, name: &str) -> LanceResult { + Ok(self.get(name).await?.is_some()) + } + + async fn get(&self, name: &str) -> LanceResult> { + let response = self + .client + .clone() + .get(self.key(name), None) + .await + .map_err(etcd_error("registry get"))?; + response + .kvs() + .first() + .map(|kv| self.entry(kv.key(), kv.value())) + .transpose() + .map(Option::flatten) + } + + async fn list(&self) -> LanceResult> { + let response = self + .client + .clone() + .get( + format!("{}/", self.prefix), + Some(GetOptions::new().with_prefix()), + ) + .await + .map_err(etcd_error("registry list"))?; + response + .kvs() + .iter() + .map(|kv| self.entry(kv.key(), kv.value())) + .collect::>>() + .map(|entries| entries.into_iter().flatten().collect()) + } + + async fn upsert(&self, name: &str, uri: &str) -> LanceResult<()> { + self.ensure_writable().await?; + // Keep the original created_at when the entry already exists: a + // retried create must not look like a newer store. + let created_at = match self.get(name).await? { + Some(existing) => existing.created_at, + None => chrono::Utc::now().timestamp_millis(), + }; + self.client + .clone() + .put(self.key(name), Self::encode(uri, created_at)?, None) + .await + .map_err(etcd_error("registry put"))?; + Ok(()) + } + + async fn remove(&self, name: &str) -> LanceResult<()> { + self.ensure_writable().await?; + self.client + .clone() + .delete(self.key(name), Some(DeleteOptions::new())) + .await + .map_err(etcd_error("registry delete"))?; + Ok(()) + } + + async fn insert_missing(&self, entries: &[(String, String)]) -> LanceResult { + self.ensure_writable().await?; + // Each entry has its own compare-and-put transaction: live entries + // are preserved, and retained migration tombstones count as absent. + const BATCH: usize = 128; + let now = chrono::Utc::now().timestamp_millis(); + let mut inserted = 0usize; + let mut seen = std::collections::HashSet::new(); + for chunk in entries.chunks(BATCH) { + let mut client = self.client.clone(); + for (name, uri) in chunk { + if !seen.insert(name.clone()) { + continue; + } + let key = self.key(name); + let existing = client + .get(key.clone(), None) + .await + .map_err(etcd_error("registry insert_missing read"))?; + let revision = match existing.kvs().first() { + Some(kv) => { + if self.entry(kv.key(), kv.value())?.is_some() { + continue; + } + kv.mod_revision() + } + None => 0, + }; + let txn = Txn::new() + .when([Compare::mod_revision( + key.as_str(), + CompareOp::Equal, + revision, + )]) + .and_then([TxnOp::put(key.as_str(), Self::encode(uri, now)?, None)]); + if client + .txn(txn) + .await + .map_err(etcd_error("registry insert_missing"))? + .succeeded() + { + inserted += 1; + } + } + } + Ok(inserted) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } +} + +#[cfg(test)] +mod tests { + use super::*; + + async fn registry() -> Option> { + let endpoints = std::env::var("ETCD_TEST_ENDPOINTS").ok()?; + let client = Client::connect(endpoints.split(',').collect::>(), None) + .await + .expect("connect to test etcd"); + let prefix = format!("/lance-context-test/{}", uuid::Uuid::new_v4()); + Some(EtcdRegistry::new(client, &prefix, "rollout")) + } + + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn upsert_get_contains_remove() { + let r = registry().await.unwrap(); + assert!(!r.contains("a").await.unwrap()); + assert!(r.get("a").await.unwrap().is_none()); + + r.upsert("a", "/data/a.lance").await.unwrap(); + assert!(r.contains("a").await.unwrap()); + let first = r.get("a").await.unwrap().unwrap(); + assert_eq!(first.uri, "/data/a.lance"); + + // A retried create keeps the original created_at. + r.upsert("a", "/data/a2.lance").await.unwrap(); + let second = r.get("a").await.unwrap().unwrap(); + assert_eq!(second.uri, "/data/a2.lance"); + assert_eq!(second.created_at, first.created_at); + + r.remove("a").await.unwrap(); + assert!(!r.contains("a").await.unwrap()); + // Removing an absent name is a no-op. + r.remove("a").await.unwrap(); + } + + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn list_and_insert_missing() { + let r = registry().await.unwrap(); + r.upsert("keep", "/old").await.unwrap(); + let n = r + .insert_missing(&[ + ("keep".into(), "/new".into()), + ("x".into(), "/x".into()), + ("y".into(), "/y".into()), + ("y".into(), "/dup".into()), + ]) + .await + .unwrap(); + assert_eq!(n, 2, "only x and y were missing"); + let mut names: Vec = r + .list() + .await + .unwrap() + .into_iter() + .map(|e| e.name) + .collect(); + names.sort(); + assert_eq!(names, vec!["keep", "x", "y"]); + assert_eq!(r.get("keep").await.unwrap().unwrap().uri, "/old"); + } + + /// Migration step 1: Lance primary, etcd mirror. Writes reach both, reads + /// come from Lance, `diff` is empty after a backfill of pre-existing rows. + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn mirrored_registry_keeps_both_backends_in_step() { + use crate::registry::{ + backfill_registry, diff_registries, LanceRegistry, MirroredRegistry, + }; + let Some(etcd) = registry().await else { return }; + let dir = tempfile::TempDir::new().unwrap(); + let uri = dir.path().join("_registry.lance"); + let lance = crate::registry::RolloutRegistry::open_or_create(uri.to_str().unwrap(), None) + .await + .unwrap(); + let lance: Arc = Arc::new(LanceRegistry::new(lance)); + // Rows that existed before mirroring started. + lance.upsert("pre-1", "/p1").await.unwrap(); + lance.upsert("pre-2", "/p2").await.unwrap(); + + let etcd_dyn: Arc = etcd.clone(); + let diff = diff_registries(&*lance, &*etcd_dyn).await.unwrap(); + assert_eq!(diff.only_in_primary, vec!["pre-1", "pre-2"]); + assert!(diff.only_in_mirror.is_empty()); + assert_eq!(backfill_registry(&*lance, &*etcd_dyn).await.unwrap(), 2); + + let mirrored = MirroredRegistry { + primary: lance.clone(), + mirror: etcd_dyn.clone(), + label: "rollout", + }; + mirrored.upsert("new", "/n").await.unwrap(); + mirrored.remove("pre-1").await.unwrap(); + assert!(mirrored.contains("new").await.unwrap()); + + let diff = diff_registries(&*lance, &*etcd_dyn).await.unwrap(); + assert!(diff.is_empty(), "{diff:?}"); + let mut names: Vec = etcd + .list() + .await + .unwrap() + .into_iter() + .map(|e| e.name) + .collect(); + names.sort(); + assert_eq!(names, vec!["new", "pre-2"]); + } + + /// Names that are prefixes of other names must not collide: `list` is a + /// prefix scan on `/`, so `a` and `ab` are distinct keys and a + /// lookup of `a` must not see `ab`. + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn prefix_names_do_not_collide() { + let r = registry().await.unwrap(); + r.upsert("ab", "/ab").await.unwrap(); + assert!(!r.contains("a").await.unwrap()); + assert!(r.get("a").await.unwrap().is_none()); + r.upsert("a", "/a").await.unwrap(); + assert_eq!(r.list().await.unwrap().len(), 2); + r.remove("a").await.unwrap(); + assert!(r.contains("ab").await.unwrap()); + } + + fn source_entry(name: &str, uri: &str) -> RegistryEntry { + RegistryEntry { + name: name.into(), + uri: uri.into(), + created_at: 123456, + } + } + + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn reconciliation_repairs_missed_deletes_and_uri_changes_preserving_timestamps() { + use crate::{backfill_registry, diff_registries, LanceRegistry, RolloutRegistry}; + let mirror = registry().await.unwrap(); + let dir = tempfile::TempDir::new().unwrap(); + let source = LanceRegistry::new( + RolloutRegistry::open_or_create( + dir.path().join("registry.lance").to_str().unwrap(), + None, + ) + .await + .unwrap(), + ); + source.upsert("deleted", "/deleted").await.unwrap(); + source.upsert("changed", "/old").await.unwrap(); + backfill_registry(&source, &*mirror).await.unwrap(); + // Primary commits succeeded while mirror writes were unavailable. + source.remove("deleted").await.unwrap(); + source.upsert("changed", "/current").await.unwrap(); + let diff = diff_registries(&source, &*mirror).await.unwrap(); + assert_eq!(diff.only_in_mirror, ["deleted"]); + assert_eq!(diff.mismatched.len(), 1); + assert_eq!(diff.mismatched[0].primary.uri, "/current"); + assert_eq!(diff.mismatched[0].mirror.uri, "/old"); + assert_eq!(backfill_registry(&source, &*mirror).await.unwrap(), 2); + assert!(diff_registries(&source, &*mirror).await.unwrap().is_empty()); + assert_eq!( + source.get("changed").await.unwrap(), + mirror.get("changed").await.unwrap() + ); + assert!(!mirror.contains("deleted").await.unwrap()); + assert_eq!(backfill_registry(&source, &*mirror).await.unwrap(), 0); + } + + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn stale_snapshots_and_point_updates_cannot_resurrect_or_remove_newer_entries() { + let r = registry().await.unwrap(); + let old = source_entry("x", "/old"); + let new = source_entry("x", "/new"); + // New snapshot never saw x: a per-name watermark alone is insufficient. + r.reconcile_source("source", 20, &[]).await.unwrap(); + assert_eq!( + r.reconcile_source("source", 10, std::slice::from_ref(&old)) + .await + .unwrap(), + 0 + ); + assert!(!r.apply_source("source", 10, "x", Some(&old)).await.unwrap()); + assert!(r.get("x").await.unwrap().is_none()); + r.apply_source("source", 30, "x", Some(&new)).await.unwrap(); + r.reconcile_source("source", 25, &[]).await.unwrap(); + assert_eq!(r.get("x").await.unwrap(), Some(new.clone())); + r.apply_source("source", 40, "x", None).await.unwrap(); + assert!(!r.apply_source("source", 39, "x", Some(&old)).await.unwrap()); + assert!(r.get("x").await.unwrap().is_none()); + r.apply_source("source", 41, "x", Some(&new)).await.unwrap(); + assert_eq!(r.get("x").await.unwrap(), Some(new)); + } + + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn interrupted_reconciliation_retries_same_snapshot_and_seals_at_cutover() { + let r = registry().await.unwrap(); + let entries = [source_entry("a", "/a"), source_entry("b", "/b")]; + // Failure after the floor and one entry committed, before the rest. + r.source_fence("source", 10, true).await.unwrap(); + r.apply_source("source", 10, "a", Some(&entries[0])) + .await + .unwrap(); + assert_eq!(r.reconcile_source("source", 10, &entries).await.unwrap(), 1); + assert_eq!(r.list().await.unwrap().len(), 2); + // A mirror must not silently accept authoritative writes before cutover. + assert!(r.upsert("a", "/unsafe").await.is_err()); + r.activate().await.unwrap(); + r.upsert("a", "/authoritative").await.unwrap(); + assert!(r.reconcile_source("source", 100, &[]).await.is_err()); + assert!(r.apply_source("source", 100, "a", None).await.is_err()); + assert_eq!(r.get("a").await.unwrap().unwrap().uri, "/authoritative"); + } + + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn cutover_rejects_missing_or_mismatched_entries_and_timestamp_only_differences() { + use crate::{backfill_registry, diff_registries, LanceRegistry, RolloutRegistry}; + let r = registry().await.unwrap(); + let dir = tempfile::TempDir::new().unwrap(); + let uri = dir + .path() + .join("registry.lance") + .to_string_lossy() + .to_string(); + let source = LanceRegistry::new(RolloutRegistry::open_or_create(&uri, None).await.unwrap()); + source.upsert("a", "/a").await.unwrap(); + assert!(r.activate_after_validation(&uri).await.is_err()); + backfill_registry(&source, &*r).await.unwrap(); + let original = source.get("a").await.unwrap().unwrap(); + // Inject the unversioned value produced by the previous mirror implementation. + r.client + .clone() + .put( + r.key("a"), + EtcdRegistry::encode("/a", original.created_at + 1).unwrap(), + None, + ) + .await + .unwrap(); + let diff = diff_registries(&source, &*r).await.unwrap(); + assert_eq!(diff.mismatched.len(), 1); + assert!(r.activate_after_validation(&uri).await.is_err()); + backfill_registry(&source, &*r).await.unwrap(); + r.activate_after_validation(&uri).await.unwrap(); + assert_eq!(r.get("a").await.unwrap(), Some(original)); + } + + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn authoritative_discovery_can_replace_a_migration_tombstone() { + let r = registry().await.unwrap(); + r.apply_source("source", 1, "a", None).await.unwrap(); + r.activate().await.unwrap(); + assert_eq!( + r.insert_missing(&[("a".into(), "/recreated".into())]) + .await + .unwrap(), + 1 + ); + assert_eq!(r.get("a").await.unwrap().unwrap().uri, "/recreated"); + } + + #[tokio::test] + async fn reverse_mirroring_is_rejected_before_opening_backends() { + use crate::etcd::{open_registry, RegistryBackend, RegistryConfig}; + let config = RegistryConfig { + registry_backend: RegistryBackend::Etcd, + registry_mirror: Some(RegistryBackend::Lance), + }; + let err = open_registry("datagen", "/not-opened", None, &config) + .await + .err() + .unwrap(); + assert!(err.to_string().contains("reverse mirroring")); + } +} diff --git a/crates/lance-context-master/src/config.rs b/crates/lance-context-master/src/config.rs index fe7ed96..cc12362 100644 --- a/crates/lance-context-master/src/config.rs +++ b/crates/lance-context-master/src/config.rs @@ -181,35 +181,15 @@ pub struct MasterConfig { #[arg(long, env = "MERGE_WAL_CONCURRENCY", default_value_t = 4)] pub merge_wal_concurrency: usize, - /// Comma-separated etcd v3 endpoints. Scheduler state (task queue, - /// lease-based claims, per-experiment write locks) lives in etcd so several - /// stateless master replicas can share one queue. Required. - #[arg(long, env = "ETCD_ENDPOINTS", value_delimiter = ',')] - pub etcd_endpoints: Vec, - - /// Namespace for all lance-context master keys in etcd. - #[arg(long, env = "ETCD_PREFIX", default_value = "/lance-context/master")] - pub etcd_prefix: String, - - /// Optional etcd username. `ETCD_PASSWORD` must also be set. - #[arg(long, env = "ETCD_USERNAME")] - pub etcd_username: Option, - - /// Optional etcd password. `ETCD_USERNAME` must also be set. - #[arg(long, env = "ETCD_PASSWORD")] - pub etcd_password: Option, - - /// Optional PEM CA certificate path for etcd TLS. - #[arg(long, env = "ETCD_CA_CERT")] - pub etcd_ca_cert: Option, - - /// Optional PEM client certificate path for etcd mutual TLS. - #[arg(long, env = "ETCD_CLIENT_CERT")] - pub etcd_client_cert: Option, - - /// Optional PEM client private-key path for etcd mutual TLS. - #[arg(long, env = "ETCD_CLIENT_KEY")] - pub etcd_client_key: Option, + /// etcd connection. Scheduler state (task queue, lease-based claims, + /// per-experiment write locks) lives in etcd so several stateless master + /// replicas can share one queue. `ETCD_ENDPOINTS` is required. + #[command(flatten)] + pub etcd: lance_context_core::etcd::EtcdConfig, + + /// Which backend the store registries (rollout, generic) live in. + #[command(flatten)] + pub registry: lance_context_core::etcd::RegistryConfig, /// TTL for etcd task claims and distributed locks. The master renews leases /// while work is running; orphaned tasks are requeued after expiry. diff --git a/crates/lance-context-master/src/discovery.rs b/crates/lance-context-master/src/discovery.rs index c4713fe..bbc1c93 100644 --- a/crates/lance-context-master/src/discovery.rs +++ b/crates/lance-context-master/src/discovery.rs @@ -1,7 +1,7 @@ //! Startup discovery for rollout datasets created before the registry existed. use lance::io::ObjectStore; -use lance_context_core::{join_uri, RolloutRegistry}; +use lance_context_core::{join_uri, StoreRegistry}; const ROLLOUT_SUFFIX: &str = ".rollout.lance"; @@ -9,7 +9,7 @@ const ROLLOUT_SUFFIX: &str = ".rollout.lance"; /// registry rows in one batch. Returns the number of rows inserted. pub async fn backfill_registry( data_dir: &str, - registry: &mut RolloutRegistry, + registry: &dyn StoreRegistry, ) -> lance::Result { let (store, base_path) = ObjectStore::from_uri(data_dir).await?; let mut children = store.read_dir(base_path).await?; @@ -32,6 +32,7 @@ pub async fn backfill_registry( #[cfg(test)] mod tests { use super::*; + use lance_context_core::RolloutRegistry; use lance_context_core::RolloutStore; use tempfile::TempDir; @@ -54,22 +55,24 @@ mod tests { .unwrap(); let registry_uri = dir.path().join("_registry.rollout.lance"); - let mut registry = RolloutRegistry::open_or_create(registry_uri.to_str().unwrap(), None) - .await - .unwrap(); + let registry = lance_context_core::LanceRegistry::new( + RolloutRegistry::open_or_create(registry_uri.to_str().unwrap(), None) + .await + .unwrap(), + ); registry .upsert("existing", existing_uri.to_str().unwrap()) .await .unwrap(); assert_eq!( - backfill_registry(dir.path().to_str().unwrap(), &mut registry) + backfill_registry(dir.path().to_str().unwrap(), ®istry) .await .unwrap(), 1 ); assert_eq!( - backfill_registry(dir.path().to_str().unwrap(), &mut registry) + backfill_registry(dir.path().to_str().unwrap(), ®istry) .await .unwrap(), 0 diff --git a/crates/lance-context-master/src/routes.rs b/crates/lance-context-master/src/routes.rs index c9d7cd8..d2ce1a3 100644 --- a/crates/lance-context-master/src/routes.rs +++ b/crates/lance-context-master/src/routes.rs @@ -190,8 +190,6 @@ pub async fn list_experiments( let known: HashSet = experiments.iter().map(|e| e.name.clone()).collect(); let matches: Vec<_> = state .registry - .write() - .await .list() .await .map_err(MasterError::from_lance)? @@ -239,8 +237,6 @@ pub async fn get_experiment( // side effect, then read it back. let entry = state .registry - .write() - .await .get(&name) .await .map_err(MasterError::from_lance)? @@ -265,8 +261,6 @@ pub async fn get_experiment( // registry and observe on demand rather than reporting it as missing. let entry = state .registry - .write() - .await .get(&name) .await .map_err(MasterError::from_lance)? @@ -423,8 +417,6 @@ pub async fn rescan_experiment( ) -> Result, MasterError> { let entry = state .registry - .write() - .await .get(&name) .await .map_err(MasterError::from_lance)? @@ -509,8 +501,6 @@ pub async fn compact_experiment( // Only enqueue known experiments. let exists = state .registry - .write() - .await .contains(&name) .await .map_err(MasterError::from_lance)?; @@ -551,8 +541,6 @@ pub async fn enqueue_task( ) -> Result<(StatusCode, Json), MasterError> { let exists = state .registry - .write() - .await .contains(&req.target) .await .map_err(MasterError::from_lance)?; @@ -598,6 +586,77 @@ pub async fn list_repairs( .map_err(MasterError::from_lance) } +#[derive(Debug, serde::Deserialize)] +pub struct RegistryParams { + /// `rollout` (default), `generic`, or `datagen`. + #[serde(default)] + pub kind: Option, +} + +fn registry_pair<'a>( + state: &'a MasterState, + kind: Option<&str>, +) -> Result<(&'a Arc, &'static str), MasterError> { + match kind.unwrap_or("rollout") { + "rollout" => Ok((&state.registry, "rollout")), + "generic" => Ok((&state.generic_registry, "generic")), + "datagen" => Ok((&state.datagen_registry, "datagen")), + other => Err(MasterError::InvalidRequest(format!( + "unknown registry kind '{other}' (rollout|generic|datagen)" + ))), + } +} + +/// Compare names, URIs and creation timestamps. A live comparison is only an +/// observation; freeze all registry writers and reconcile all three kinds +/// before cutover. An empty diff is not permission for a rolling backend flip. +pub async fn registry_diff( + State(state): State>, + Query(params): Query, +) -> Result, MasterError> { + let (registry, kind) = registry_pair(&state, params.kind.as_deref())?; + let Some(mirrored) = registry + .as_any() + .downcast_ref::() + else { + return Err(MasterError::InvalidRequest( + "REGISTRY_MIRROR is not configured; nothing to diff".to_string(), + )); + }; + let diff = lance_context_core::diff_registries(&*mirrored.primary, &*mirrored.mirror) + .await + .map_err(MasterError::from_lance)?; + Ok(Json(serde_json::json!({ + "kind": kind, + "only_in_primary": diff.only_in_primary, + "only_in_mirror": diff.only_in_mirror, + "mismatched": diff.mismatched, + "cutover_requires_quiescence": true, + }))) +} + +/// `POST /api/v1/registry/backfill?kind=` — reconcile the complete primary +/// snapshot, including changed values and deletions. The master also does this on +/// startup and every maintenance round; this is for forcing it. +pub async fn registry_backfill( + State(state): State>, + Query(params): Query, +) -> Result, MasterError> { + let (registry, kind) = registry_pair(&state, params.kind.as_deref())?; + let Some(mirrored) = registry + .as_any() + .downcast_ref::() + else { + return Err(MasterError::InvalidRequest( + "REGISTRY_MIRROR is not configured; nothing to backfill".to_string(), + )); + }; + let copied = lance_context_core::backfill_registry(&*mirrored.primary, &*mirrored.mirror) + .await + .map_err(MasterError::from_lance)?; + Ok(Json(serde_json::json!({ "kind": kind, "copied": copied }))) +} + /// `GET /api/v1/tasks` — paginated tasks (queue + recent history), newest first. pub async fn list_tasks( State(state): State>, @@ -653,6 +712,8 @@ pub fn api_router() -> Router> { .route("/tasks", post(enqueue_task).get(list_tasks)) .route("/scheduler/cooldowns", get(list_cooldowns)) .route("/scheduler/repairs", get(list_repairs)) + .route("/registry/diff", get(registry_diff)) + .route("/registry/backfill", post(registry_backfill)) .route("/tasks/{id}", get(get_task)) .route("/rescan", post(rescan)) } @@ -663,8 +724,6 @@ async fn open_registered_store( ) -> Result>, MasterError> { let entry = state .registry - .write() - .await .get(name) .await .map_err(MasterError::from_lance)? @@ -740,13 +799,12 @@ mod tests { worker_endpoints: vec![], task_concurrency: 4, merge_wal_concurrency: 4, - etcd_endpoints: test_etcd_endpoints(), - etcd_prefix: format!("/lance-context/test/{}", generate_id()), - etcd_username: None, - etcd_password: None, - etcd_ca_cert: None, - etcd_client_cert: None, - etcd_client_key: None, + etcd: lance_context_core::etcd::EtcdConfig { + etcd_endpoints: test_etcd_endpoints(), + etcd_prefix: format!("/lance-context/test/{}", generate_id()), + ..Default::default() + }, + registry: lance_context_core::etcd::RegistryConfig::default(), etcd_lease_ttl_secs: 5, task_history_limit: 1_000, task_history_ttl_secs: 86_400, @@ -818,13 +876,7 @@ mod tests { let uri = state.rollout_uri(&name); // Creating the store materializes an (empty) base table on disk. RolloutStore::open(&uri).await.unwrap(); - state - .registry - .write() - .await - .upsert(&name, &uri) - .await - .unwrap(); + state.registry.upsert(&name, &uri).await.unwrap(); } let scanned = scanner::scan_once(&state).await.unwrap(); @@ -877,13 +929,7 @@ mod tests { let name = format!("exp-{i}"); let uri = state.rollout_uri(&name); RolloutStore::open(&uri).await.unwrap(); - state - .registry - .write() - .await - .upsert(&name, &uri) - .await - .unwrap(); + state.registry.upsert(&name, &uri).await.unwrap(); } // Populates both the stats table and the in-memory cache. scanner::scan_once(&state).await.unwrap(); @@ -975,18 +1021,12 @@ mod tests { let name = "gone"; let uri = state.rollout_uri(name); RolloutStore::open(&uri).await.unwrap(); - state - .registry - .write() - .await - .upsert(name, &uri) - .await - .unwrap(); + state.registry.upsert(name, &uri).await.unwrap(); scanner::scan_once(&state).await.unwrap(); assert!(state.stats.lock().await.get(name).await.unwrap().is_some()); // Remove from registry -> next scan drops the stats row. - state.registry.write().await.remove(name).await.unwrap(); + state.registry.remove(name).await.unwrap(); scanner::scan_once(&state).await.unwrap(); assert!(state.stats.lock().await.get(name).await.unwrap().is_none()); } @@ -1015,13 +1055,7 @@ mod tests { let name = "behind"; let uri = state.rollout_uri(name); let store = RolloutStore::open(&uri).await.unwrap(); - state - .registry - .write() - .await - .upsert(name, &uri) - .await - .unwrap(); + state.registry.upsert(name, &uri).await.unwrap(); scanner::scan_once(&state).await.unwrap(); let before = state.stats.lock().await.get(name).await.unwrap().unwrap(); assert_eq!(before.pending_wal_generations, 0); @@ -1051,13 +1085,7 @@ mod tests { for name in ["target", "other"] { let uri = state.rollout_uri(name); RolloutStore::open(&uri).await.unwrap(); - state - .registry - .write() - .await - .upsert(name, &uri) - .await - .unwrap(); + state.registry.upsert(name, &uri).await.unwrap(); } let Json(detail) = get_experiment( @@ -1105,13 +1133,7 @@ mod tests { // the memtable is sealed into a committed WAL generation. Flush rather // than relaxing the assertions below. store.flush().await.unwrap(); - state - .registry - .write() - .await - .upsert("records", &uri) - .await - .unwrap(); + state.registry.upsert("records", &uri).await.unwrap(); let Json(page) = list_experiment_records( State(state.clone()), @@ -1180,13 +1202,7 @@ mod tests { // the memtable is sealed into a committed WAL generation. Flush rather // than relaxing the assertions below. store.flush().await.unwrap(); - state - .registry - .write() - .await - .upsert("records", &uri) - .await - .unwrap(); + state.registry.upsert("records", &uri).await.unwrap(); // A valid SELECT returns rows over the merged view. let Json(result) = query_experiment_sql( @@ -1258,13 +1274,7 @@ mod tests { // See the note in `records_endpoint_...`: rows are only visible to the // handle the endpoint opens after the memtable is sealed. store.flush().await.unwrap(); - state - .registry - .write() - .await - .upsert("blobs", &uri) - .await - .unwrap(); + state.registry.upsert("blobs", &uri).await.unwrap(); let response = download_experiment_blob( State(state.clone()), diff --git a/crates/lance-context-master/src/scanner.rs b/crates/lance-context-master/src/scanner.rs index 6336fb4..0f7d5f9 100644 --- a/crates/lance-context-master/src/scanner.rs +++ b/crates/lance-context-master/src/scanner.rs @@ -14,8 +14,8 @@ use std::time::Duration; use chrono::Utc; use futures::stream::{self, StreamExt}; use lance_context_core::{ - CompactionConfig, GenericStore, GenericStoreOptions, RolloutRegistry, RolloutStore, - RolloutStoreOptions, + CompactionConfig, GenericStore, GenericStoreOptions, RolloutStore, RolloutStoreOptions, + StoreRegistry, }; use tokio::sync::Semaphore; use tokio::task::JoinHandle; @@ -167,23 +167,27 @@ pub async fn maintain_stats(state: &Arc) -> lance::Result<()> { async fn maintain_registry( state: &Arc, label: &'static str, - registry: &tokio::sync::RwLock, + registry: &Arc, ) -> lance::Result<()> { + // etcd-backed registries need no maintenance. + let Some(table) = registry.lance_table() else { + return Ok(()); + }; let ttl = Duration::from_secs(state.config.stats_history_ttl_secs); let start = std::time::Instant::now(); let (compaction, cleaner) = { - let mut registry = registry.write().await; - tokio::time::timeout(MAINTENANCE_TIMEOUT, registry.compact()) + let mut table = table.lock().await; + tokio::time::timeout(MAINTENANCE_TIMEOUT, table.compact()) .await .unwrap_or_else(|_| Err(lance::Error::io("registry compaction timed out")))? }; let removal = cleaner.cleanup(ttl).await; - let reload = registry.write().await.reload().await; + let reload = table.lock().await.reload().await; let removal = match (removal, reload) { (Ok(removal), Ok(())) => removal, (Err(e), _) | (Ok(_), Err(e)) => return Err(e), }; - let version = registry.read().await.version(); + let version = table.lock().await.version(); metrics::counter!("master_registry_versions_removed_total", "registry" => label) .increment(removal.old_versions); metrics::gauge!("master_registry_version", "registry" => label).set(version as f64); @@ -232,10 +236,25 @@ async fn try_scan_once(state: &Arc, maintain: bool) -> lance::Resul for (label, registry) in [ ("rollout", &state.registry), ("generic", &state.generic_registry), + ("datagen", &state.datagen_registry), ] { if let Err(e) = maintain_registry(state, label, registry).await { tracing::warn!(registry = label, error = %e, "registry maintenance failed"); } + if let Some(m) = registry + .as_any() + .downcast_ref::() + { + match lance_context_core::backfill_registry(&*m.primary, &*m.mirror).await { + Ok(copied) if copied > 0 => { + tracing::info!(registry = label, copied, "healed registry mirror") + } + Ok(_) => {} + Err(e) => { + tracing::warn!(registry = label, error = %e, "registry mirror heal failed") + } + } + } } } let release = state.task_store.release_coordination_lock(guard).await; @@ -254,8 +273,6 @@ async fn scan_once_inner(state: &Arc) -> lance::Result { // depend on this: the sweep reads nothing but this table. let mut entries: Vec = state .registry - .write() - .await .list() .await? .into_iter() @@ -268,8 +285,6 @@ async fn scan_once_inner(state: &Arc) -> lance::Result { entries.extend( state .generic_registry - .write() - .await .list() .await? .into_iter() diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index 0fed321..b953536 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -907,15 +907,14 @@ mod tests { task_cooldown_after_failures: 3, task_cooldown_base_secs: 600, task_cooldown_max_secs: 21_600, - etcd_endpoints: std::env::var("ETCD_TEST_ENDPOINTS") - .map(|value| value.split(',').map(str::to_string).collect()) - .unwrap_or_default(), - etcd_prefix: format!("/lance-context/test/{}", generate_id()), - etcd_username: None, - etcd_password: None, - etcd_ca_cert: None, - etcd_client_cert: None, - etcd_client_key: None, + etcd: lance_context_core::etcd::EtcdConfig { + etcd_endpoints: std::env::var("ETCD_TEST_ENDPOINTS") + .map(|value| value.split(',').map(str::to_string).collect()) + .unwrap_or_default(), + etcd_prefix: format!("/lance-context/test/{}", generate_id()), + ..Default::default() + }, + registry: lance_context_core::etcd::RegistryConfig::default(), etcd_lease_ttl_secs: 5, task_history_limit: 1_000, task_history_ttl_secs: 86_400, @@ -934,13 +933,7 @@ mod tests { let mut store = RolloutStore::open(&uri).await.unwrap(); store.add(&[rollout_record("a")]).await.unwrap(); store.cleanup_own_shard().await.unwrap(); - state - .registry - .write() - .await - .upsert("exp", &uri) - .await - .unwrap(); + state.registry.upsert("exp", &uri).await.unwrap(); crate::scanner::scan_once(&state).await.unwrap(); let metrics = compact_inner(&state, "exp").await.unwrap(); assert_eq!(metrics.fragments_removed, 0); @@ -1114,13 +1107,7 @@ mod tests { store.cleanup_own_shard().await.unwrap(); } } - state - .registry - .write() - .await - .upsert(name, &uri) - .await - .unwrap(); + state.registry.upsert(name, &uri).await.unwrap(); // Seed a stats row so post-compaction upsert has a prior counter. crate::scanner::scan_once(&state).await.unwrap(); @@ -1168,13 +1155,7 @@ mod tests { store.cleanup_own_shard().await.unwrap(); } } - state - .registry - .write() - .await - .upsert(name, &uri) - .await - .unwrap(); + state.registry.upsert(name, &uri).await.unwrap(); crate::scanner::scan_once(&state).await.unwrap(); let compact = enqueue(&state, TaskKind::Compact, name).await.unwrap(); @@ -1288,13 +1269,7 @@ mod tests { store.cleanup_own_shard().await.unwrap(); assert!(!store.has_id_btree_index().await.unwrap()); } - state - .registry - .write() - .await - .upsert(name, &uri) - .await - .unwrap(); + state.registry.upsert(name, &uri).await.unwrap(); let rec = enqueue(&state, TaskKind::MergeWal, name).await.unwrap(); assert_eq!(await_terminal(&state, &rec.id).await.state, TaskState::Done); @@ -1347,13 +1322,7 @@ mod tests { .join(format!("{name}.rollout.lance/data/{victim_path}")), ) .unwrap(); - state - .registry - .write() - .await - .upsert(name, &uri) - .await - .unwrap(); + state.registry.upsert(name, &uri).await.unwrap(); crate::scanner::scan_once(&state).await.unwrap(); let compact = enqueue(&state, TaskKind::Compact, name).await.unwrap(); @@ -1423,13 +1392,7 @@ mod tests { store.cleanup_own_shard().await.unwrap(); } } - state - .registry - .write() - .await - .upsert(name, &uri) - .await - .unwrap(); + state.registry.upsert(name, &uri).await.unwrap(); let rec = enqueue(&state, TaskKind::IndexId, name).await.unwrap(); let status = await_terminal(&state, &rec.id).await; @@ -1475,10 +1438,10 @@ mod tests { }; use lance_context_merge::{Coordinator, Execution}; use tower::ServiceExt; - let client = etcd_client::Client::connect(cfg.etcd_endpoints.clone(), None) + let client = etcd_client::Client::connect(cfg.etcd.etcd_endpoints.clone(), None) .await .unwrap(); - let coordinator = Coordinator::new(client, cfg.etcd_prefix.clone()); + let coordinator = Coordinator::new(client, cfg.etcd.etcd_prefix.clone()); let owned_targets = cfg.merge_rollout.owned_targets.clone(); axum::Router::new() .route( @@ -1836,13 +1799,7 @@ mod tests { store.cleanup_own_shard().await.unwrap(); } } - state - .registry - .write() - .await - .upsert(name, &uri) - .await - .unwrap(); + state.registry.upsert(name, &uri).await.unwrap(); crate::scanner::scan_once(&state).await.unwrap(); // Saturate and over-fill the merge queue *before* the Compact. @@ -1903,13 +1860,7 @@ mod tests { store.cleanup_own_shard().await.unwrap(); } } - state - .registry - .write() - .await - .upsert(name, &uri) - .await - .unwrap(); + state.registry.upsert(name, &uri).await.unwrap(); crate::scanner::scan_once(&state).await.unwrap(); let compact = enqueue(&state, TaskKind::Compact, name).await.unwrap(); diff --git a/crates/lance-context-master/src/state.rs b/crates/lance-context-master/src/state.rs index 703a40e..28b5142 100644 --- a/crates/lance-context-master/src/state.rs +++ b/crates/lance-context-master/src/state.rs @@ -3,9 +3,10 @@ use std::num::NonZeroUsize; use std::sync::Arc; +use lance_context_core::etcd::open_registry; use lance_context_core::{ - join_uri, CompactionConfig, GenericStoreOptions, RolloutRegistry, RolloutStore, - RolloutStoreOptions, Session, + backfill_registry, join_uri, CompactionConfig, GenericStoreOptions, MirroredRegistry, + RolloutStore, RolloutStoreOptions, Session, StoreRegistry, }; use lru::LruCache; use tokio::sync::{Mutex, RwLock, Semaphore}; @@ -60,7 +61,7 @@ fn build_compaction_config(config: &MasterConfig) -> CompactionConfig { /// mutating methods take `&mut self`. pub struct MasterState { /// Durable directory of which rollout stores exist. - pub registry: RwLock, + pub registry: Arc, /// Durable directory of which generic stores exist. Same registry format /// as rollout (the data plane writes `_registry.generic.lance` with the /// same `RolloutRegistry` type); read here so the stats scan and the @@ -69,7 +70,9 @@ pub struct MasterState { /// to commit into one base table and lost to `Too many concurrent writers` /// while the master -- whose etcd-locked MergeWal task exists to serialize /// exactly that -- never heard of them. - pub generic_registry: RwLock, + pub generic_registry: Arc, + /// Datagen participates in the same registry migration as the worker. + pub datagen_registry: Arc, /// Periodically-refreshed per-experiment metrics (master-owned). pub stats: Mutex, /// Last snapshot written to the stats table, kept in memory so @@ -136,20 +139,50 @@ impl MasterState { let registry_uri = join_uri(&base_uri, "_registry.rollout.lance"); let generic_registry_uri = join_uri(&base_uri, "_registry.generic.lance"); let stats_uri = join_uri(&base_uri, "_stats.rollout.lance"); - let mut registry = RolloutRegistry::open_or_create(®istry_uri, None).await?; - let generic_registry = RolloutRegistry::open_or_create(&generic_registry_uri, None).await?; - let backfilled = discovery::backfill_registry(&config.data_dir, &mut registry).await?; + let etcd = Some((task_store.etcd_client(), config.etcd.prefix())); + let registry = open_registry("rollout", ®istry_uri, etcd, &config.registry).await?; + let generic_registry = + open_registry("generic", &generic_registry_uri, etcd, &config.registry).await?; + let datagen_registry = open_registry( + "datagen", + &join_uri(&base_uri, "_registry.datagen.lance"), + etcd, + &config.registry, + ) + .await?; + let backfilled = discovery::backfill_registry(&config.data_dir, &*registry).await?; if backfilled > 0 { tracing::info!( experiments = backfilled, "backfilled rollout registry from data directory" ); } + // With a mirror configured, seed it from the primary so a backend + // switch finds every store already there. Repeated on every + // maintenance round; this covers the very first start. + for (label, reg) in [ + ("rollout", ®istry), + ("generic", &generic_registry), + ("datagen", &datagen_registry), + ] { + if let Some(m) = reg.as_any().downcast_ref::() { + match backfill_registry(&*m.primary, &*m.mirror).await { + Ok(copied) if copied > 0 => { + tracing::info!(registry = label, copied, "seeded registry mirror") + } + Ok(_) => {} + Err(error) => { + tracing::warn!(registry = label, %error, "registry mirror seed failed") + } + } + } + } let stats = StatsStore::open_or_create(&stats_uri, None).await?; task_store.release_coordination_lock(init_guard).await?; let state = Arc::new(Self { - registry: RwLock::new(registry), - generic_registry: RwLock::new(generic_registry), + registry, + generic_registry, + datagen_registry, stats: Mutex::new(stats), stats_cache: RwLock::new(Arc::new(Vec::new())), record_stores: Mutex::new(LruCache::new( @@ -265,15 +298,14 @@ mod tests { worker_endpoints: vec![], task_concurrency: 4, merge_wal_concurrency: 4, - etcd_endpoints: std::env::var("ETCD_TEST_ENDPOINTS") - .map(|value| value.split(',').map(str::to_string).collect()) - .unwrap_or_default(), - etcd_prefix: format!("/lance-context/test/{}", generate_id()), - etcd_username: None, - etcd_password: None, - etcd_ca_cert: None, - etcd_client_cert: None, - etcd_client_key: None, + etcd: lance_context_core::etcd::EtcdConfig { + etcd_endpoints: std::env::var("ETCD_TEST_ENDPOINTS") + .map(|value| value.split(',').map(str::to_string).collect()) + .unwrap_or_default(), + etcd_prefix: format!("/lance-context/test/{}", generate_id()), + ..Default::default() + }, + registry: lance_context_core::etcd::RegistryConfig::default(), etcd_lease_ttl_secs: 5, task_history_limit: 1_000, task_history_ttl_secs: 86_400, @@ -315,13 +347,62 @@ mod tests { RolloutStore::open(uri.to_str().unwrap()).await.unwrap(); let state = MasterState::new(test_config(&dir)).await.unwrap(); - let entries = state.registry.write().await.list().await.unwrap(); + let entries = state.registry.list().await.unwrap(); assert_eq!(entries.len(), 1); assert_eq!(entries[0].name, "legacy"); assert_eq!(entries[0].uri, uri.to_string_lossy().to_string()); } + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn startup_and_migration_routes_include_datagen() { + use lance_context_core::{etcd::RegistryBackend, RolloutRegistry}; + let dir = TempDir::new().unwrap(); + for kind in ["rollout", "generic", "datagen"] { + let uri = dir.path().join(format!("_registry.{kind}.lance")); + let mut registry = RolloutRegistry::open_or_create(uri.to_str().unwrap(), None) + .await + .unwrap(); + registry + .upsert("existing", &format!("/data/{kind}")) + .await + .unwrap(); + } + let mut config = test_config(&dir); + config.registry.registry_mirror = Some(RegistryBackend::Etcd); + let state = MasterState::new(config).await.unwrap(); + for kind in ["rollout", "generic", "datagen"] { + let axum::Json(diff) = crate::routes::registry_diff( + axum::extract::State(state.clone()), + axum::extract::Query(crate::routes::RegistryParams { + kind: Some(kind.into()), + }), + ) + .await + .unwrap(); + assert_eq!(diff["only_in_primary"], serde_json::json!([])); + assert_eq!(diff["only_in_mirror"], serde_json::json!([])); + assert_eq!(diff["mismatched"], serde_json::json!([])); + } + let m = state + .datagen_registry + .as_any() + .downcast_ref::() + .unwrap(); + m.primary.remove("existing").await.unwrap(); + let axum::Json(result) = crate::routes::registry_backfill( + axum::extract::State(state.clone()), + axum::extract::Query(crate::routes::RegistryParams { + kind: Some("datagen".into()), + }), + ) + .await + .unwrap(); + assert_eq!(result["copied"], 1); + assert!(m.mirror.get("existing").await.unwrap().is_none()); + } + /// A master that dies mid-task leaves the record in `Running`. Once its /// etcd lease expires the claim key vanishes, and the next master's startup /// `recover_orphaned` pass must return the task to `Queued` — otherwise an diff --git a/crates/lance-context-master/src/task_store.rs b/crates/lance-context-master/src/task_store.rs index aecea41..afd768b 100644 --- a/crates/lance-context-master/src/task_store.rs +++ b/crates/lance-context-master/src/task_store.rs @@ -12,10 +12,7 @@ use std::sync::Arc; use std::time::Duration; use chrono::Utc; -use etcd_client::{ - Certificate, Client, Compare, CompareOp, ConnectOptions, GetOptions, Identity, PutOptions, - TlsOptions, Txn, TxnOp, -}; +use etcd_client::{Client, Compare, CompareOp, GetOptions, PutOptions, Txn, TxnOp}; use lance_context_api::{RepairRecord, TaskCooldown, TaskKind, TaskRecord, TaskState}; use lance_context_core::generate_id; use tokio::sync::oneshot; @@ -226,6 +223,12 @@ impl TaskStore { ) } + /// The underlying etcd client, for other etcd-backed components (the + /// store registries) that should share one connection. + pub fn etcd_client(&self) -> &Client { + &self.inner.client + } + pub async fn open(config: &MasterConfig) -> lance::Result { let store = Self { inner: Arc::new(EtcdTaskStore::connect(config).await?), @@ -478,7 +481,7 @@ fn prunable_terminal_ids( impl EtcdTaskStore { async fn connect(config: &MasterConfig) -> lance::Result { - if config.etcd_endpoints.is_empty() { + if !config.etcd.is_configured() { return Err(lance::Error::io( "ETCD_ENDPOINTS is required to run the master", )); @@ -486,56 +489,10 @@ impl EtcdTaskStore { if config.etcd_lease_ttl_secs < 5 { return Err(lance::Error::io("ETCD_LEASE_TTL_SECS must be at least 5")); } - let mut options = ConnectOptions::new() - .with_connect_timeout(Duration::from_secs(5)) - .with_timeout(Duration::from_secs(10)) - .with_keep_alive(Duration::from_secs(10), Duration::from_secs(3)) - .with_require_leader(true); - match (&config.etcd_username, &config.etcd_password) { - (Some(username), Some(password)) => { - options = options.with_user(username, password); - } - (None, None) => {} - _ => { - return Err(lance::Error::io( - "ETCD_USERNAME and ETCD_PASSWORD must be configured together", - )) - } - } - if let Some(path) = &config.etcd_ca_cert { - let pem = std::fs::read(path).map_err(|err| { - lance::Error::io(format!("failed to read ETCD_CA_CERT '{path}': {err}")) - })?; - let mut tls = TlsOptions::new().ca_certificate(Certificate::from_pem(pem)); - match (&config.etcd_client_cert, &config.etcd_client_key) { - (Some(cert), Some(key)) => { - let cert_pem = std::fs::read(cert).map_err(|err| { - lance::Error::io(format!("failed to read ETCD_CLIENT_CERT '{cert}': {err}")) - })?; - let key_pem = std::fs::read(key).map_err(|err| { - lance::Error::io(format!("failed to read ETCD_CLIENT_KEY '{key}': {err}")) - })?; - tls = tls.identity(Identity::from_pem(cert_pem, key_pem)); - } - (None, None) => {} - _ => { - return Err(lance::Error::io( - "ETCD_CLIENT_CERT and ETCD_CLIENT_KEY must be configured together", - )) - } - } - options = options.with_tls(tls); - } else if config.etcd_client_cert.is_some() || config.etcd_client_key.is_some() { - return Err(lance::Error::io( - "ETCD_CA_CERT is required when configuring an etcd client certificate", - )); - } - let client = Client::connect(config.etcd_endpoints.clone(), Some(options)) - .await - .map_err(etcd_error("connect to etcd"))?; + let client = config.etcd.connect().await?; Ok(Self { client, - prefix: config.etcd_prefix.trim_end_matches('/').to_string(), + prefix: config.etcd.prefix().to_string(), lease_ttl: config.etcd_lease_ttl_secs, }) } @@ -1477,12 +1434,12 @@ mod tests { async fn compact_noop_lease_expiry_allows_bounded_recheck() { let dir = TempDir::new().unwrap(); let mut cfg = config(&dir); - cfg.etcd_endpoints = std::env::var("ETCD_TEST_ENDPOINTS") + cfg.etcd.etcd_endpoints = std::env::var("ETCD_TEST_ENDPOINTS") .unwrap() .split(',') .map(str::to_owned) .collect(); - cfg.etcd_prefix = format!("/test/{}", generate_id()); + cfg.etcd.etcd_prefix = format!("/test/{}", generate_id()); let store = TaskStore::open(&cfg).await.unwrap(); store .record_compact_noop("table", "/table", 10, "options") @@ -1500,7 +1457,7 @@ mod tests { .compact_is_unchanged("table", "/table", -1, "options") .await .unwrap()); - let key = lance_context_merge::execution_key(&cfg.etcd_prefix, "table") + let key = lance_context_merge::execution_key(&cfg.etcd.etcd_prefix, "table") .replace("/merge-executions/", "/compact-noops/"); let mut client = store.inner.client.clone(); let kv = client.get(key, None).await.unwrap(); @@ -1541,13 +1498,12 @@ mod tests { worker_endpoints: vec![], task_concurrency: 4, merge_wal_concurrency: 4, - etcd_endpoints: vec![], - etcd_prefix: "/test".to_string(), - etcd_username: None, - etcd_password: None, - etcd_ca_cert: None, - etcd_client_cert: None, - etcd_client_key: None, + etcd: lance_context_core::etcd::EtcdConfig { + etcd_endpoints: vec![], + etcd_prefix: "/test".to_string(), + ..Default::default() + }, + registry: lance_context_core::etcd::RegistryConfig::default(), etcd_lease_ttl_secs: 30, task_history_limit: 1_000, task_history_ttl_secs: 86_400, @@ -1563,12 +1519,12 @@ mod tests { async fn merge_serializes_other_writers_and_survives_claim_loss() { let dir = TempDir::new().unwrap(); let mut cfg = config(&dir); - cfg.etcd_endpoints = std::env::var("ETCD_TEST_ENDPOINTS") + cfg.etcd.etcd_endpoints = std::env::var("ETCD_TEST_ENDPOINTS") .unwrap() .split(',') .map(str::to_string) .collect(); - cfg.etcd_prefix = format!("/merge-lock-test/{}", generate_id()); + cfg.etcd.etcd_prefix = format!("/merge-lock-test/{}", generate_id()); let store = TaskStore::open(&cfg).await.unwrap(); store .enqueue(TaskKind::MergeWal, "shared", Vec::new()) @@ -1697,7 +1653,7 @@ mod tests { .client .clone() .delete( - cfg.etcd_prefix, + cfg.etcd.etcd_prefix, Some(etcd_client::DeleteOptions::new().with_prefix()), ) .await @@ -1792,8 +1748,8 @@ mod tests { .expect("ETCD_TEST_ENDPOINTS must point to a test etcd"); let dir = TempDir::new().unwrap(); let mut cfg = config(&dir); - cfg.etcd_endpoints = endpoint.split(',').map(str::to_string).collect(); - cfg.etcd_prefix = format!("/lance-context/test/{}", generate_id()); + cfg.etcd.etcd_endpoints = endpoint.split(',').map(str::to_string).collect(); + cfg.etcd.etcd_prefix = format!("/lance-context/test/{}", generate_id()); cfg.task_cooldown_after_failures = 100; // stay below threshold; count only let store = TaskStore::open(&cfg).await.unwrap(); @@ -1829,8 +1785,8 @@ mod tests { .expect("ETCD_TEST_ENDPOINTS must point to a test etcd"); let dir = TempDir::new().unwrap(); let mut cfg = config(&dir); - cfg.etcd_endpoints = endpoint.split(',').map(str::to_string).collect(); - cfg.etcd_prefix = format!("/lance-context/test/{}", generate_id()); + cfg.etcd.etcd_endpoints = endpoint.split(',').map(str::to_string).collect(); + cfg.etcd.etcd_prefix = format!("/lance-context/test/{}", generate_id()); cfg.etcd_lease_ttl_secs = 5; let first = TaskStore::open(&cfg).await.unwrap(); @@ -1915,10 +1871,12 @@ mod tests { .await .unwrap(); - let mut client = Client::connect(cfg.etcd_endpoints, None).await.unwrap(); + let mut client = Client::connect(cfg.etcd.etcd_endpoints, None) + .await + .unwrap(); client .delete( - cfg.etcd_prefix, + cfg.etcd.etcd_prefix, Some(etcd_client::DeleteOptions::new().with_prefix()), ) .await @@ -1939,8 +1897,8 @@ mod tests { .expect("ETCD_TEST_ENDPOINTS must point to a test etcd"); let dir = TempDir::new().unwrap(); let mut cfg = config(&dir); - cfg.etcd_endpoints = endpoint.split(',').map(str::to_string).collect(); - cfg.etcd_prefix = format!("/lance-context/test/{}", generate_id()); + cfg.etcd.etcd_endpoints = endpoint.split(',').map(str::to_string).collect(); + cfg.etcd.etcd_prefix = format!("/lance-context/test/{}", generate_id()); cfg.etcd_lease_ttl_secs = 5; let store = TaskStore::open(&cfg).await.unwrap(); @@ -1997,10 +1955,12 @@ mod tests { assert_eq!(merge.task.kind, TaskKind::MergeWal); store.finish(merge, Ok("done".to_string())).await.unwrap(); - let mut client = Client::connect(cfg.etcd_endpoints, None).await.unwrap(); + let mut client = Client::connect(cfg.etcd.etcd_endpoints, None) + .await + .unwrap(); client .delete( - cfg.etcd_prefix, + cfg.etcd.etcd_prefix, Some(etcd_client::DeleteOptions::new().with_prefix()), ) .await diff --git a/crates/lance-context-server/src/config.rs b/crates/lance-context-server/src/config.rs index aa387ea..641cf3e 100644 --- a/crates/lance-context-server/src/config.rs +++ b/crates/lance-context-server/src/config.rs @@ -4,10 +4,6 @@ use clap::Parser; #[command(name = "lance-context-server")] #[command(about = "REST API server for lance-context")] pub struct ServerConfig { - /// Connection used lazily by explicitly enabled owned merge targets. - #[command(flatten)] - pub merge_etcd: lance_context_merge::EtcdConfig, - #[command(flatten)] pub merge_rollout: lance_context_merge::rollout::MergeRollout, @@ -211,9 +207,30 @@ pub struct ServerConfig { /// default session (the pre-fix, leak-prone behavior). #[arg(long, env = "ROLLOUT_CACHE_BYTES", default_value = "2147483648")] pub rollout_cache_bytes: usize, + + /// Shared etcd settings for registries and lazily connected owned merges. + #[command(flatten)] + pub etcd: lance_context_core::etcd::EtcdConfig, + + /// Which backend the store registries (rollout, generic, datagen) live in. + #[command(flatten)] + pub registry: lance_context_core::etcd::RegistryConfig, } impl ServerConfig { + /// Use the same namespace and credentials for lazily connected merge ownership. + pub(crate) fn merge_etcd_config(&self) -> lance_context_merge::EtcdConfig { + lance_context_merge::EtcdConfig { + etcd_endpoints: self.etcd.etcd_endpoints.clone(), + etcd_prefix: self.etcd.etcd_prefix.clone(), + etcd_username: self.etcd.etcd_username.clone(), + etcd_password: self.etcd.etcd_password.clone(), + etcd_ca_cert: self.etcd.etcd_ca_cert.clone(), + etcd_client_cert: self.etcd.etcd_client_cert.clone(), + etcd_client_key: self.etcd.etcd_client_key.clone(), + } + } + /// Resolve the instance id used for server-managed MemWAL sharding: the /// explicit `--instance-id`/`INSTANCE_ID` if provided, otherwise the /// `HOSTNAME` environment variable (stable per-pod under a StatefulSet). diff --git a/crates/lance-context-server/src/routes/datagen.rs b/crates/lance-context-server/src/routes/datagen.rs index 8974613..c714ebf 100644 --- a/crates/lance-context-server/src/routes/datagen.rs +++ b/crates/lance-context-server/src/routes/datagen.rs @@ -29,8 +29,6 @@ pub async fn create_datagen_store( let mut handle = state.datagen_handles.lock(&req.name).await; if state .datagen_registry - .write() - .await .contains(&req.name) .await .map_err(AppError::from_lance)? @@ -78,8 +76,6 @@ pub async fn list_datagen_stores( ) -> Result, AppError> { let entries = state .datagen_registry - .write() - .await .list() .await .map_err(AppError::from_lance)?; diff --git a/crates/lance-context-server/src/routes/generic.rs b/crates/lance-context-server/src/routes/generic.rs index 087f0fd..ce4f6ec 100644 --- a/crates/lance-context-server/src/routes/generic.rs +++ b/crates/lance-context-server/src/routes/generic.rs @@ -40,8 +40,6 @@ pub async fn create_generic_store( if state .generic_registry - .write() - .await .contains(&req.name) .await .map_err(AppError::from_lance)? @@ -95,8 +93,6 @@ pub async fn list_generic_stores( ) -> Result, AppError> { let entries = state .generic_registry - .write() - .await .list() .await .map_err(AppError::from_lance)?; diff --git a/crates/lance-context-server/src/routes/rollouts.rs b/crates/lance-context-server/src/routes/rollouts.rs index cae4e05..d8777eb 100644 --- a/crates/lance-context-server/src/routes/rollouts.rs +++ b/crates/lance-context-server/src/routes/rollouts.rs @@ -179,8 +179,6 @@ pub async fn create_rollout_store( // Existence is tracked durably in the registry, not by cache membership. if state .rollout_registry - .write() - .await .contains(&req.name) .await .map_err(AppError::from_lance)? @@ -236,8 +234,6 @@ pub async fn list_rollout_stores( // opening every dataset. let entries = state .rollout_registry - .write() - .await .list() .await .map_err(AppError::from_lance)?; @@ -1067,13 +1063,7 @@ mod tests { .await .is_err() ); - assert!(state - .rollout_registry - .write() - .await - .contains("rl") - .await - .unwrap()); + assert!(state.rollout_registry.contains("rl").await.unwrap()); std::fs::remove_file(&path).unwrap(); std::fs::rename(&saved, &path).unwrap(); assert_eq!( @@ -1083,13 +1073,7 @@ mod tests { StatusCode::NO_CONTENT ); assert!(!path.exists()); - assert!(!state - .rollout_registry - .write() - .await - .contains("rl") - .await - .unwrap()); + assert!(!state.rollout_registry.contains("rl").await.unwrap()); let row = rollout_record_from_add_request(&record_with_size("late-write", None)); assert!(cached.read().await.add(&[row]).await.is_err()); } @@ -1359,14 +1343,7 @@ mod tests { assert!(matches!(err, AppError::InvalidRequest(_)), "{name}"); } assert!(state.rollout_stores.lock().await.is_empty()); - assert!(state - .rollout_registry - .write() - .await - .list() - .await - .unwrap() - .is_empty()); + assert!(state.rollout_registry.list().await.unwrap().is_empty()); } #[tokio::test] diff --git a/crates/lance-context-server/src/state.rs b/crates/lance-context-server/src/state.rs index ba8e326..b17f502 100644 --- a/crates/lance-context-server/src/state.rs +++ b/crates/lance-context-server/src/state.rs @@ -5,10 +5,11 @@ use std::path::PathBuf; use std::sync::{Arc, Weak}; use std::time::Duration; +use lance_context_core::etcd::open_registry; use lance_context_core::{ join_uri, validate_store_name, ContextStore, ContextStoreOptions, DatagenStore, - DatagenStoreOptions, GenericStore, GenericStoreOptions, MergeMemoryBudget, RolloutRegistry, - RolloutStore, RolloutStoreOptions, Session, + DatagenStoreOptions, GenericStore, GenericStoreOptions, MergeMemoryBudget, RolloutStore, + RolloutStoreOptions, Session, StoreRegistry, }; use lru::LruCache; use tokio::sync::{Mutex, OwnedMutexGuard, RwLock, Semaphore}; @@ -106,7 +107,7 @@ pub struct AppState { /// miss (existence check) and to back the list endpoint. Guarded by a lock /// because every operation refreshes the snapshot and therefore takes /// `&mut`. - pub rollout_registry: RwLock, + pub rollout_registry: Arc, pub base_uri: String, /// Stable identity of this server instance, used as the MemWAL shard key for /// every server-managed store so each instance owns exactly one shard. @@ -158,7 +159,7 @@ pub struct AppState { /// Durable directory of which datagen stores exist. A separate registry /// dataset from [`Self::rollout_registry`] so the two store kinds never /// collide on a shared name. - pub datagen_registry: RwLock, + pub datagen_registry: Arc, /// Bounded LRU of resident generic-store handles, mirroring /// [`Self::rollout_stores`]. pub generic_stores: Mutex>>>, @@ -171,7 +172,7 @@ pub struct AppState { /// refactor exists to remove. The cost is that listing stores cannot report /// their schemas without opening each dataset, which is why /// `GenericStoreInfo::schema` is `None` in list responses. - pub generic_registry: RwLock, + pub generic_registry: Arc, } /// Process-wide admission control for the total artifact-blob payload held in @@ -305,30 +306,51 @@ impl AppState { pub async fn new(config: ServerConfig) -> Result { let instance_id = config.resolved_instance_id(); let base_uri = config.data_dir.clone(); - let registry_uri = join_uri(&base_uri, "_registry.rollout.lance"); - let registry = RolloutRegistry::open_or_create(®istry_uri, None) - .await - .map_err(AppError::from_lance)?; + let etcd_client = if config.etcd.is_configured() + && (matches!( + config.registry.registry_backend, + lance_context_core::etcd::RegistryBackend::Etcd + ) || config.registry.registry_mirror.is_some()) + { + Some(config.etcd.connect().await.map_err(AppError::from_lance)?) + } else { + None + }; + let etcd = etcd_client.as_ref().map(|c| (c, config.etcd.prefix())); + let registry = open_registry( + "rollout", + &join_uri(&base_uri, "_registry.rollout.lance"), + etcd, + &config.registry, + ) + .await + .map_err(AppError::from_lance)?; let capacity = NonZeroUsize::new(config.rollout_cache_capacity) .unwrap_or_else(|| NonZeroUsize::new(DEFAULT_ROLLOUT_CACHE_CAPACITY).unwrap()); let blob_budget = (config.rollout_max_inflight_blob_bytes > 0) .then(|| BlobBudget::new(config.rollout_max_inflight_blob_bytes)); let rollout_session = build_rollout_session(config.rollout_cache_bytes); - let generic_registry_uri = join_uri(&base_uri, "_registry.generic.lance"); - let generic_registry = RolloutRegistry::open_or_create(&generic_registry_uri, None) - .await - .map_err(AppError::from_lance)?; - let datagen_registry_uri = join_uri(&base_uri, "_registry.datagen.lance"); - let datagen_registry = RolloutRegistry::open_or_create(&datagen_registry_uri, None) - .await - .map_err(AppError::from_lance)?; + let generic_registry = open_registry( + "generic", + &join_uri(&base_uri, "_registry.generic.lance"), + etcd, + &config.registry, + ) + .await + .map_err(AppError::from_lance)?; + let datagen_registry = open_registry( + "datagen", + &join_uri(&base_uri, "_registry.datagen.lance"), + etcd, + &config.registry, + ) + .await + .map_err(AppError::from_lance)?; config .merge_rollout .validate() .map_err(AppError::InvalidRequest)?; - if !config.merge_rollout.owned_targets.is_empty() - && config.merge_etcd.etcd_endpoints.is_empty() - { + if !config.merge_rollout.owned_targets.is_empty() && config.etcd.etcd_endpoints.is_empty() { return Err(AppError::InvalidRequest( "owned targets require ETCD_ENDPOINTS".into(), )); @@ -343,7 +365,7 @@ impl AppState { } Ok(Self { merge_executions: crate::merge_execution::Executions::configured( - config.merge_etcd.clone(), + config.merge_etcd_config(), config.merge_rollout.clone(), config.merge_execution_timeout_secs, config.merge_queue_timeout_secs, @@ -352,7 +374,7 @@ impl AppState { stores: RwLock::new(std::collections::HashMap::new()), rollout_stores: Mutex::new(LruCache::new(capacity)), rollout_handles: StoreHandles::default(), - rollout_registry: RwLock::new(registry), + rollout_registry: registry, base_uri, instance_id, rollout_merge_after_generations: config.rollout_merge_after_generations, @@ -371,10 +393,10 @@ impl AppState { rollout_session, datagen_stores: Mutex::new(LruCache::new(capacity)), datagen_handles: StoreHandles::default(), - datagen_registry: RwLock::new(datagen_registry), + datagen_registry, generic_stores: Mutex::new(LruCache::new(capacity)), generic_handles: StoreHandles::default(), - generic_registry: RwLock::new(generic_registry), + generic_registry, }) } @@ -405,18 +427,31 @@ impl AppState { instance_id: Option, ) -> Self { let base_uri = base_path.to_string_lossy().to_string(); - let registry_uri = join_uri(&base_uri, "_registry.rollout.lance"); - let registry = RolloutRegistry::open_or_create(®istry_uri, None) - .await - .expect("open test registry"); - let generic_registry_uri = join_uri(&base_uri, "_registry.generic.lance"); - let generic_registry = RolloutRegistry::open_or_create(&generic_registry_uri, None) - .await - .expect("open test generic registry"); - let datagen_registry_uri = join_uri(&base_uri, "_registry.datagen.lance"); - let datagen_registry = RolloutRegistry::open_or_create(&datagen_registry_uri, None) - .await - .expect("open test datagen registry"); + let lance_only = lance_context_core::etcd::RegistryConfig::default(); + let registry = open_registry( + "rollout", + &join_uri(&base_uri, "_registry.rollout.lance"), + None, + &lance_only, + ) + .await + .expect("open test registry"); + let generic_registry = open_registry( + "generic", + &join_uri(&base_uri, "_registry.generic.lance"), + None, + &lance_only, + ) + .await + .expect("open test generic registry"); + let datagen_registry = open_registry( + "datagen", + &join_uri(&base_uri, "_registry.datagen.lance"), + None, + &lance_only, + ) + .await + .expect("open test datagen registry"); Self { merge_executions: crate::merge_execution::Executions::new(None, 600), stores: RwLock::new(std::collections::HashMap::new()), @@ -424,7 +459,7 @@ impl AppState { NonZeroUsize::new(DEFAULT_ROLLOUT_CACHE_CAPACITY).unwrap(), )), rollout_handles: StoreHandles::default(), - rollout_registry: RwLock::new(registry), + rollout_registry: registry, base_uri, instance_id, rollout_merge_after_generations: 0, @@ -443,12 +478,12 @@ impl AppState { std::num::NonZeroUsize::new(DEFAULT_ROLLOUT_CACHE_CAPACITY).unwrap(), )), generic_handles: StoreHandles::default(), - generic_registry: RwLock::new(generic_registry), + generic_registry, datagen_stores: Mutex::new(LruCache::new( NonZeroUsize::new(DEFAULT_ROLLOUT_CACHE_CAPACITY).unwrap(), )), datagen_handles: StoreHandles::default(), - datagen_registry: RwLock::new(datagen_registry), + datagen_registry, } } @@ -493,8 +528,6 @@ impl AppState { ))); } self.rollout_registry - .write() - .await .upsert(name, uri) .await .map_err(AppError::from_lance)?; @@ -516,7 +549,7 @@ impl AppState { let mut handle = self.rollout_handles.lock(name).await; // Keep the registry entry and cached handle until physical deletion // succeeds, so failures can be retried without losing the store's URI. - let mut registry = self.rollout_registry.write().await; + let registry = &self.rollout_registry; if !registry .contains(name) .await @@ -576,8 +609,6 @@ impl AppState { // Existence is the registry's job, not the cache's. let exists = self .rollout_registry - .write() - .await .contains(name) .await .map_err(AppError::from_lance)?; @@ -596,7 +627,7 @@ impl AppState { .map_err(AppError::from_lance)?; let opened = Arc::new(RwLock::new(opened)); // Also observe registry removal by another server during the open. - let mut registry = self.rollout_registry.write().await; + let registry = &self.rollout_registry; if !registry .contains(name) .await @@ -684,8 +715,6 @@ impl AppState { ))); } self.datagen_registry - .write() - .await .upsert(name, uri) .await .map_err(AppError::from_lance)?; @@ -704,7 +733,7 @@ impl AppState { let mut handle = self.datagen_handles.lock(name).await; // Keep the registry entry and cached handle until physical deletion // succeeds, so failures can be retried without losing the store's URI. - let mut registry = self.datagen_registry.write().await; + let registry = &self.datagen_registry; if !registry .contains(name) .await @@ -755,8 +784,6 @@ impl AppState { let exists = self .datagen_registry - .write() - .await .contains(name) .await .map_err(AppError::from_lance)?; @@ -773,7 +800,7 @@ impl AppState { .map_err(AppError::from_lance)?; let opened = Arc::new(RwLock::new(opened)); // Also observe registry removal by another server during the open. - let mut registry = self.datagen_registry.write().await; + let registry = &self.datagen_registry; if !registry .contains(name) .await @@ -829,8 +856,6 @@ impl AppState { ))); } self.generic_registry - .write() - .await .upsert(name, uri) .await .map_err(AppError::from_lance)?; @@ -849,7 +874,7 @@ impl AppState { let mut handle = self.generic_handles.lock(name).await; // Keep the registry entry and cached handle until physical deletion // succeeds, so failures can be retried without losing the store's URI. - let mut registry = self.generic_registry.write().await; + let registry = &self.generic_registry; if !registry .contains(name) .await @@ -901,8 +926,6 @@ impl AppState { let exists = self .generic_registry - .write() - .await .contains(name) .await .map_err(AppError::from_lance)?; @@ -921,7 +944,7 @@ impl AppState { .map_err(AppError::from_lance)?; let opened = Arc::new(RwLock::new(opened)); // Also observe registry removal by another server during the open. - let mut registry = self.generic_registry.write().await; + let registry = &self.generic_registry; if !registry .contains(name) .await diff --git a/docs/src/design/registry-etcd.md b/docs/src/design/registry-etcd.md new file mode 100644 index 0000000..34403d5 --- /dev/null +++ b/docs/src/design/registry-etcd.md @@ -0,0 +1,102 @@ +# Store registry migration to etcd + +The rollout, generic and datagen registries map a name to a dataset URI and its +original creation timestamp. Payload data stays in object storage. The default +backend remains Lance. etcd provides point lookups without reading a Lance +manifest; existing worker store caches remain unchanged. + +## Supported configurations + +| REGISTRY_BACKEND | REGISTRY_MIRROR | Behavior | +|---|---|---| +| lance | unset | Existing Lance registry | +| lance | etcd | Lance authority, versioned etcd migration mirror | +| etcd | unset | Authoritative etcd registry; mirror writes sealed | + +Other mirror combinations fail startup. In particular, etcd-to-Lance mirroring +and live rollback by flipping environment variables are **not supported**. +Both processes share the existing `ETCD_*` connection/TLS settings. No registry +backend is changed by merely deploying the binary. + +## Versioned reconciliation + +Each etcd entry lives at `/registry//`. Its JSON carries +`uri`, `created_at`, and, during migration, a Lance `source_version` and a +`deleted` tombstone flag. Tombstones are invisible to `get`, `contains` and +`list`. They retain ordering so a delayed upsert cannot resurrect a deletion. + +A primary mutation commits to Lance first. The mirror then receives the entry +(or its absence) and version from a fresh source snapshot. Mirror failures are +logged without failing a successful primary mutation. Master startup and periodic +registry maintenance reconcile **all three kinds**, including datagen. +`POST /api/v1/registry/backfill?kind=rollout|generic|datagen` runs the same repair. +The `copied` response counts changed entries, including deletions and updates. + +Reconciliation reads one consistent Lance snapshot, preserving URIs and creation +timestamps. Before applying it, it advances a persistent snapshot floor at +`/registry-migrations/`. Each destination write atomically +compares both that floor's etcd revision and the destination key's revision, +and rejects an older per-name source version. The floor also covers names never +seen by the destination: an old snapshot cannot insert a name that a newer +snapshot already observed as absent. A newer point update is protected against +an older snapshot's deletion pass. A crashed/partially applied reconciliation +can resume at the same source version; a newer snapshot can supersede it. + +The source URI is bound into migration state. Do not recreate the Lance registry +at that URI or reuse an etcd namespace for a different source dataset. Use a +fresh namespace for a new source lineage. Ordinary version pruning does not +reset Lance versions. Retain migration state and tombstones; deleting them +removes the delayed-write protection. + +`ETCD_PREFIX` also scopes master task records and locks. Do not rotate that +shared prefix merely to reset a registry mirror: preserving those records +requires a separate coordinated migration. + +Once etcd is opened as authority, an atomic seal rejects all subsequent mirror +writes, including requests delayed in transit. Authoritative creation/discovery +can replace a migration tombstone. There is no read fallback to stale Lance. + +## Cutover procedure + +Mirroring repairs concurrent traffic; it does **not** make a rolling switch +between two authoritative backends safe. An empty live diff is an observation, +not cutover authorization. + +1. Deploy compatible binaries with `REGISTRY_BACKEND=lance`, + `REGISTRY_MIRROR=etcd`. Include every worker and auxiliary master. Keep one + authoritative backend throughout this phase. +2. Stop admitting registry mutations across **every** writer and drain in-flight + operations. This includes store create/delete, discovery/backfill and cold + retirement in auxiliary masters. Stop all old writer processes before changing + authority. External routing alone does not stop background registry writers. + The application does not implement an automatic fleet-wide quiescence barrier. +3. While writers are quiescent, reconcile each of rollout, generic and datagen. + Verify `GET /api/v1/registry/diff?kind=...` for each: `only_in_primary`, + `only_in_mirror`, and `mismatched` must all be empty. `mismatched` includes + both URI differences and creation-timestamp differences. +4. Start the new fleet with `REGISTRY_BACKEND=etcd` and **omit** + `REGISTRY_MIRROR`. Before first activation, startup compares the destination + against the Lance source and refuses mismatches. It then seals mirroring. + Already active registries do not reconnect to Lance on each restart. + Validate all three directories before reopening registry mutations. + +Startup comparison cannot prove the absence of old in-flight Lance writes; +step 2 is required. First activation is per registry kind, so a startup failure +can leave some kinds sealed. Keep writers quiescent, inspect the failure and +repair the remaining unsealed kinds before retrying. Never clear a seal while +etcd writers might exist. + +## Recovery and limitations + +Losing an etcd mirror before cutover does not lose authoritative registry state: +rebuild it from Lance in a fresh etcd namespace. Snapshot repair handles missed +inserts, changed mappings, missed deletes and interrupted repair passes. + +After cutover the Lance directories are historical snapshots, **not** rollback +copies. Returning to Lance requires a separate, verified export/import while +all registry writers are quiescent. Recreating entries from dataset directory +names alone cannot reproduce arbitrary URIs, original timestamps or deletions. +Back up etcd and do not describe directory discovery as a lossless rollback. + +This changes registry metadata only. It does not fix WAL merge liveness, +compaction scheduling, or payload read failures.