From 2e2e460e133f2ca06ed9c42a602531842b0314a2 Mon Sep 17 00:00:00 2001 From: Beinan Date: Tue, 29 Sep 2026 20:11:00 +0000 Subject: [PATCH 1/4] feat: etcd-backed store registries behind a StoreRegistry trait, with a mirrored migration path The store registries (rollout, generic, datagen) are a name -> uri directory every worker consults on a cache miss and writes on every create/delete. As a Lance table each read is a manifest fetch from object storage and each write is two contended commits; even after #274 keeps it compact, a create costs seconds and throttling turns it into failures. Introduce a StoreRegistry trait (&self, so the RwLock every request went through goes away) with three impls: LanceRegistry (today's table), EtcdRegistry (one key per store under /registry//, put-if-absent backfill), and MirroredRegistry (read one, write both, mirror best-effort). REGISTRY_BACKEND / REGISTRY_MIRROR select them; defaults keep Lance with no mirror, so this change is behaviour-neutral until the flags are set. The master seeds the mirror from the primary at startup and heals it every maintenance round, and exposes GET /registry/diff and POST /registry/backfill for verifying a migration step. ETCD_* flags move to core (EtcdConfig) so the server can share them; the master's etcd prefix is unchanged. Co-Authored-By: Claude Fable 5 --- Cargo.lock | 2 + crates/lance-context-core/Cargo.toml | 2 + crates/lance-context-core/src/etcd.rs | 190 +++++++++++ crates/lance-context-core/src/lib.rs | 8 +- crates/lance-context-core/src/registry.rs | 152 +++++++++ .../lance-context-core/src/registry_etcd.rs | 303 ++++++++++++++++++ crates/lance-context-master/src/config.rs | 38 +-- crates/lance-context-master/src/discovery.rs | 17 +- crates/lance-context-master/src/routes.rs | 165 +++++----- crates/lance-context-master/src/scanner.rs | 36 ++- crates/lance-context-master/src/scheduler.rs | 77 +---- crates/lance-context-master/src/state.rs | 56 ++-- crates/lance-context-master/src/task_store.rs | 106 ++---- crates/lance-context-server/src/config.rs | 25 +- .../src/routes/datagen.rs | 4 - .../src/routes/generic.rs | 4 - .../src/routes/rollouts.rs | 29 +- crates/lance-context-server/src/state.rs | 137 ++++---- docs/src/design/registry-etcd.md | 209 ++++++++++++ 19 files changed, 1187 insertions(+), 373 deletions(-) create mode 100644 crates/lance-context-core/src/etcd.rs create mode 100644 crates/lance-context-core/src/registry_etcd.rs create mode 100644 docs/src/design/registry-etcd.md 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..89d1af1 --- /dev/null +++ b/crates/lance-context-core/src/etcd.rs @@ -0,0 +1,190 @@ +//! 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, + + /// Optional second backend that also receives every write, best-effort. + /// Used during a migration so the backend being moved to (or kept as a + /// fallback) stays current. Must differ from `registry_backend`. + #[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) + } + }) + } + + let primary = build(config.registry_backend, kind, lance_uri, etcd).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..0969206 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, RegistryEntry, + 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..e40e286 100644 --- a/crates/lance-context-core/src/registry.rs +++ b/crates/lance-context-core/src/registry.rs @@ -42,6 +42,158 @@ 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 + } +} + +/// A [`StoreRegistry`] that reads from `primary` and writes to both `primary` +/// and `mirror`. The mirror write is best-effort: it is logged on failure and +/// never fails the call, because the mirror exists to be caught up by a +/// periodic backfill during a backend migration, not to be authoritative. +pub struct MirroredRegistry { + pub primary: Arc, + pub mirror: Arc, + pub label: &'static str, +} + +#[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.mirror.upsert(name, uri).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.mirror.remove(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) = self.mirror.insert_missing(entries).await { + tracing::warn!(registry = self.label, %error, "registry mirror backfill failed"); + } + Ok(inserted) + } + fn lance_table(&self) -> Option<&tokio::sync::Mutex> { + // Whichever side is the Lance table still needs maintaining while it + // is written to. + self.primary + .lance_table() + .or_else(|| self.mirror.lance_table()) + } + fn as_any(&self) -> &dyn std::any::Any { + self + } +} + +/// Copy every entry of `from` that `to` lacks. Returns how many were copied. +/// Used to seed a new backend from the old one and to heal mirror misses. +pub async fn backfill_registry( + from: &dyn StoreRegistry, + to: &dyn StoreRegistry, +) -> LanceResult { + let entries: Vec<(String, String)> = from + .list() + .await? + .into_iter() + .map(|e| (e.name, e.uri)) + .collect(); + to.insert_missing(&entries).await +} + +/// Names present in exactly one of two registries: `(only_in_a, only_in_b)`. +pub async fn diff_registries( + a: &dyn StoreRegistry, + b: &dyn StoreRegistry, +) -> LanceResult<(Vec, Vec)> { + let a_names: HashSet = a.list().await?.into_iter().map(|e| e.name).collect(); + let b_names: HashSet = b.list().await?.into_iter().map(|e| e.name).collect(); + let mut only_a: Vec = a_names.difference(&b_names).cloned().collect(); + let mut only_b: Vec = b_names.difference(&a_names).cloned().collect(); + only_a.sort(); + only_b.sort(); + Ok((only_a, only_b)) +} + /// 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..d181e1a --- /dev/null +++ b/crates/lance-context-core/src/registry_etcd.rs @@ -0,0 +1,303 @@ +//! [`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, +} + +/// etcd-backed store directory for one store kind. +pub struct EtcdRegistry { + client: Client, + /// `/registry/` with no trailing slash. + prefix: 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('/')), + }) + } + + 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}")))?; + Ok(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, + }) + .map_err(|e| LanceError::io(format!("encode registry entry: {e}"))) + } +} + +#[async_trait::async_trait] +impl StoreRegistry for EtcdRegistry { + async fn contains(&self, name: &str) -> LanceResult { + let response = self + .client + .clone() + .get(self.key(name), Some(GetOptions::new().with_count_only())) + .await + .map_err(etcd_error("registry get"))?; + Ok(response.count() > 0) + } + + 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() + } + + 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() + } + + async fn upsert(&self, name: &str, uri: &str) -> LanceResult<()> { + // 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.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 { + // One put-if-absent transaction per entry, batched so a 10k-row + // backfill is a few hundred round trips rather than 10k. A txn with + // several compares is all-or-nothing, which is not what we want here, + // so each entry is its own txn inside the batch loop. + 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 txn = Txn::new() + .when([Compare::version(key.as_str(), CompareOp::Equal, 0)]) + .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 (only_lance, only_etcd) = diff_registries(&*lance, &*etcd_dyn).await.unwrap(); + assert_eq!(only_lance, vec!["pre-1", "pre-2"]); + assert!(only_etcd.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 (only_lance, only_etcd) = diff_registries(&*lance, &*etcd_dyn).await.unwrap(); + assert!(only_lance.is_empty(), "{only_lance:?}"); + assert!(only_etcd.is_empty(), "{only_etcd:?}"); + 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()); + } +} 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..2ccf5eb 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,76 @@ pub async fn list_repairs( .map_err(MasterError::from_lance) } +#[derive(Debug, serde::Deserialize)] +pub struct RegistryParams { + /// `rollout` (default) or `generic`. + #[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")), + other => Err(MasterError::InvalidRequest(format!( + "unknown registry kind '{other}' (rollout|generic)" + ))), + } +} + +/// `GET /api/v1/registry/diff?kind=` — names present in the primary backend +/// but not the mirror, and vice versa. Empty on both sides means the two +/// backends agree and a migration step can proceed. 400 when no mirror is +/// configured, because then there is nothing to compare. +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 (only_primary, only_mirror) = + 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": only_primary, + "only_in_mirror": only_mirror, + }))) +} + +/// `POST /api/v1/registry/backfill?kind=` — copy every primary entry the +/// mirror lacks into the mirror. Idempotent. 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 +711,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 +723,6 @@ async fn open_registered_store( ) -> Result>, MasterError> { let entry = state .registry - .write() - .await .get(name) .await .map_err(MasterError::from_lance)? @@ -740,13 +798,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 +875,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 +928,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 +1020,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 +1054,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 +1084,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 +1132,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 +1201,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 +1273,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..bbcfd16 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); @@ -236,6 +240,20 @@ async fn try_scan_once(state: &Arc, maintain: bool) -> lance::Resul 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 +272,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 +284,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..c5e81bd 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, @@ -1114,13 +1113,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 +1161,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 +1275,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 +1328,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 +1398,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 +1444,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 +1805,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 +1866,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..e65fa13 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,7 @@ 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, /// Periodically-refreshed per-experiment metrics (master-owned). pub stats: Mutex, /// Last snapshot written to the stats table, kept in memory so @@ -136,20 +137,38 @@ 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 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)] { + 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, stats: Mutex::new(stats), stats_cache: RwLock::new(Arc::new(Vec::new())), record_stores: Mutex::new(LruCache::new( @@ -265,15 +284,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,7 +333,7 @@ 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"); diff --git a/crates/lance-context-master/src/task_store.rs b/crates/lance-context-master/src/task_store.rs index aecea41..bbb93d3 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, }) } @@ -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..d7ed185 --- /dev/null +++ b/docs/src/design/registry-etcd.md @@ -0,0 +1,209 @@ +# Design: move the store registries from Lance tables to etcd + +Status: proposal. Follows #274 (which makes the Lance registries sustainable +enough that this is not urgent). + +## What the registry is and who touches it + +Three tables, one per store kind, each a directory `{name -> uri, created_at}` +shared by every worker and the master: + +| table | rows (prod) | writers | readers | +|---|---|---|---| +| `_registry.rollout.lance` | ~10.7k | worker create/delete, master retire | every worker miss on `get_or_open_*`, master scan (every 300 s, `list()`), discovery backfill | +| `_registry.generic.lance` | ~600 | same | same | +| `_registry.datagen.lance` | small | same | same | + +Every call site today goes through `RolloutRegistry` (`crates/lance-context-core/src/registry.rs`): +`contains`, `get`, `list`, `upsert`, `remove`, `insert_missing`. Each read does +`checkout_latest` first, i.e. one manifest read from ADLS per call. Each `upsert` +is delete + append, two Lance commits with optimistic concurrency across 20 +workers. + +### Why the Lance table is the wrong tool here + +- A directory is a key-value set with point lookups and rare listing. Lance is + a columnar table optimised for scans; every point lookup pays a manifest read + plus a scan, and every write pays a full commit protocol. +- Without maintenance it grows quadratically (#274 explains the 29 GB). With + maintenance it still costs a manifest read per `contains`, ~0.7 s on ADLS + today, 1.5 s+ under throttling, and 20 workers do it on every cache miss. +- Writes from 20 workers contend on one commit point. Lance retries, but a + create under ADLS throttling took 47 s on average and failed 79 times in the + last 15 h. +- Existence checks and store opens are the hot path for *every* API call that + misses the worker's LRU. Registry latency is user-visible latency. + +### What etcd gives + +- Point `contains`/`get`: one round trip, single-digit ms, from any pod. +- `upsert`/`remove`: one transaction, no manifest, no commit retries. +- `list` with prefix: fine at 10k keys (etcd pages by default; the master + scans every 300 s and already tolerates seconds). +- Watch: workers can invalidate their caches on create/delete instead of + polling; not needed for v1 but free. +- The master already runs etcd (task store, coordination locks, cooldowns) and + the cluster exposes `rocketkeep-etcd:2379` to the namespace. Workers do not + link `etcd-client` yet. + +### What etcd does not give, and how we cope + +- **Not the source of truth for data.** The dataset on ADLS is. The registry + only says "this name exists and lives here". If etcd is lost we must be able + to rebuild it. The master already has `discovery.rs` (`insert_missing` from a + directory listing); keep it and point it at etcd. +- **Size.** Value is `{uri, created_at}` (~200 B). 11k keys is ~2 MB. etcd is + comfortable to a few hundred MB; irrelevant here. +- **Availability coupling.** A worker that cannot reach etcd cannot create or + open a store it does not have cached. Today a worker that cannot reach ADLS + is equally dead, and etcd is a 3-node cluster inside the same k8s cluster, so + this is not a new failure domain in practice. Still: reads fall back to the + in-process cache, and we keep the Lance table as a read-only fallback during + the migration window (below). + +## Data model + +``` +/registry// -> {"uri": "...", "created_at": 1790000000000} +``` + +- `` is the same `ETCD_PREFIX` the master uses (`/lance-context` by + default), so one etcd holds one deployment's tasks and registries together. +- `` is `rollout`, `generic`, `datagen`. +- `` is the store name, already validated to be a portable path segment + (`validate_name`), so it is a safe key segment with no escaping. +- Value is JSON. No lease: registry entries are durable until explicitly + removed. + +## Code shape + +A trait in `lance-context-core` so the server and master do not care which +backend they have: + +```rust +#[async_trait] +pub trait StoreRegistry: Send + Sync { + async fn contains(&self, name: &str) -> LanceResult; + async fn get(&self, name: &str) -> LanceResult>; + async fn list(&self) -> LanceResult>; + async fn upsert(&self, name: &str, uri: &str) -> LanceResult<()>; + async fn remove(&self, name: &str) -> LanceResult; + async fn insert_missing(&self, entries: &[(String, String)]) -> LanceResult; +} +``` + +Note `&self`, not `&mut self`: etcd needs no handle mutation, and the +`RwLock` everywhere today exists only because Lance handles +must `checkout_latest`. Dropping the lock removes a serialisation point every +worker request goes through. + +Two impls: + +- `LanceRegistry` — today's `RolloutRegistry`, wrapped. Keeps `maintain()`. +- `EtcdRegistry` — `etcd_client::Client` + prefix. `upsert` is a plain `put` + (idempotent by construction; the delete-then-append dance goes away). + `insert_missing` is a batch of `txn(version(key)==0 → put)`. + +Selection by config: `REGISTRY_BACKEND=lance|etcd` (default `lance` until the +migration below is done), plus the existing `ETCD_ENDPOINTS`/`ETCD_PREFIX`/TLS +flags lifted from the master's config into core so the server can reuse them. +`etcd-client` becomes a dependency of `lance-context-core` (feature-gated +`etcd`, on by default). + +Worker-side cache: `get_or_open_*` already has an LRU of open stores and only +asks the registry on a miss. Keep that. Add a small negative cache +(name → not-found, 5 s TTL) so a client hammering a non-existent name does not +turn into an etcd storm; cheap and safe because creates go through the same +worker fleet and invalidate locally, and a 5 s stale "not found" on another +worker is the same window Lance gives today. + +## Migration + +Zero-downtime, three deploys, each independently revertable. + +1. **Dual-write** (`REGISTRY_BACKEND=lance`, `REGISTRY_MIRROR=etcd`): reads hit + Lance as today; every `upsert`/`remove` also goes to etcd, best-effort with + a warning on failure. Master runs a one-shot backfill on startup + (`list()` from Lance → `insert_missing` into etcd) under the existing + `state-init` coordination lock, and re-runs it every maintenance round so + any mirror miss heals. Ship this, let it run a day, then compare: a + `GET /api/v1/registry/diff` on the master lists keys present in one backend + and not the other. Expect empty. +2. **Read from etcd** (`REGISTRY_BACKEND=etcd`, `REGISTRY_MIRROR=lance`): + reads hit etcd; writes still mirror to Lance so a revert to step 1 loses + nothing. Watch worker `rollout_store_cache_misses_total` latency and + `POST /rollouts` latency drop. Run a few days. +3. **etcd only** (`REGISTRY_BACKEND=etcd`, no mirror): stop writing the Lance + tables. Leave them on disk for a while as a cold backup; `discovery.rs` can + rebuild etcd from the ADLS directory listing regardless. + +Each step is env-only in mango; no image change between steps once the code is +in. Revert at any step is flipping the env back. + +## Master changes + +- `MasterState.registry` / `generic_registry` become `Arc`. + The scanner's `list()` every 300 s is unchanged in shape. +- Registry maintenance from #274 stays for the Lance backend and becomes a + no-op for etcd. +- New `GET /api/v1/registry/diff` (step 1 verification) and + `POST /api/v1/registry/backfill` (manual re-sync). +- Retirement (`retire_cold_experiments`) calls `remove()` on the trait; no + change. + +## Worker changes + +- `AppState.{rollout,generic,datagen}_registry` become `Arc`; + the `RwLock` goes away. +- `create_*_store`: `contains` → open/create dataset → `upsert`. Same order; + the dataset is still created before the registry row so a crash between them + leaves an orphan dataset (discovery picks it up) rather than a dangling + registry entry. Unchanged semantics, ~45 s less latency. +- `get_or_open_*`: `contains` on miss, as today. The double `contains` (before + and after the open) can stay; it is now two 2 ms calls. +- `unregister_*`: `contains` → delete data → `remove`. Unchanged. + +## Consistency + +- **Create/create race** on the same name from two workers: today, both pass + `contains`, both create the dataset (Lance `Create` is put-if-absent so one + fails), the winner upserts. With etcd, same: dataset creation is still the + arbiter; the loser's `contains` re-check after the failed create returns + true and it answers 409. No change in outcome. +- **Create/delete race**: today serialised per name only within one worker + (`rollout_handles.lock(name)`); across workers it is last-writer-wins on the + Lance table. etcd is the same but with a smaller window. If we want it + strictly ordered we can use `txn(mod_revision == expected)`; not proposed + for v1 because the current behaviour has not caused a problem. +- **Read-your-writes**: etcd reads are linearisable by default. A worker that + just created a store sees it from any other worker immediately, which is + *better* than today (Lance `checkout_latest` is also strongly consistent, so + no regression, but etcd removes the 0.7 s manifest read). + +## Rollback + +- Step 1 → nothing to roll back; the mirror is additive. +- Step 2 → flip `REGISTRY_BACKEND=lance`. Lance was kept current by the mirror. +- Step 3 → re-enable the mirror, run backfill etcd → Lance (the same + `insert_missing` path in reverse, exposed as + `POST /api/v1/registry/backfill?to=lance`), flip. + +## Effort + +- core: trait + `EtcdRegistry` + config plumbing, ~400 lines incl. tests + (tests run against the etcd the master tests already require). +- server: swap type, drop the `RwLock`, mirror logic, negative cache, ~150 lines. +- master: swap type, backfill, diff/backfill routes, ~150 lines. +- mango: three env flips. + +Roughly two days of work plus the staged rollout. Not blocking anything +today: #274 takes the registry from "quadratic and failing" to "linear and +slow", and the etcd move takes it to "fast". + +## Not in scope + +- Moving `_stats` to etcd. It is a 10k-row table rewritten as one snapshot per + scan round and read for range queries (`list_above_fragment_count`, + `list_above_pending_wal`), which etcd cannot do; it stays in Lance. +- Moving MemWAL shard manifests. They are per-shard, written by one writer, + and read by the LSM scanner; different problem. From 21c3481d754d6f05c231baf2f4e1d429f0c4ceb0 Mon Sep 17 00:00:00 2001 From: Beinan Wang <> Date: Fri, 2 Oct 2026 06:42:46 +0000 Subject: [PATCH 2/4] fix: reconcile all registry kinds safely during etcd migration --- .github/workflows/rust-test.yml | 6 +- crates/lance-context-core/src/etcd.rs | 13 +- crates/lance-context-core/src/lib.rs | 4 +- crates/lance-context-core/src/registry.rs | 137 ++++-- .../lance-context-core/src/registry_etcd.rs | 449 +++++++++++++++++- crates/lance-context-master/src/routes.rs | 29 +- crates/lance-context-master/src/scanner.rs | 1 + crates/lance-context-master/src/state.rs | 65 ++- docs/src/design/registry-etcd.md | 311 ++++-------- 9 files changed, 731 insertions(+), 284 deletions(-) diff --git a/.github/workflows/rust-test.yml b/.github/workflows/rust-test.yml index d25bb6a..350c46c 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/crates/lance-context-core/src/etcd.rs b/crates/lance-context-core/src/etcd.rs index 89d1af1..2285efe 100644 --- a/crates/lance-context-core/src/etcd.rs +++ b/crates/lance-context-core/src/etcd.rs @@ -130,9 +130,7 @@ pub struct RegistryConfig { #[arg(long, env = "REGISTRY_BACKEND", value_enum, default_value_t = RegistryBackend::Lance)] pub registry_backend: RegistryBackend, - /// Optional second backend that also receives every write, best-effort. - /// Used during a migration so the backend being moved to (or kept as a - /// fallback) stays current. Must differ from `registry_backend`. + /// Versioned etcd mirror for a Lance primary. Reverse mirroring is unsupported. #[arg(long, env = "REGISTRY_MIRROR", value_enum)] pub registry_mirror: Option, } @@ -172,7 +170,16 @@ pub async fn open_registry( }) } + 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); }; diff --git a/crates/lance-context-core/src/lib.rs b/crates/lance-context-core/src/lib.rs index 0969206..17c3082 100644 --- a/crates/lance-context-core/src/lib.rs +++ b/crates/lance-context-core/src/lib.rs @@ -67,8 +67,8 @@ pub use record::{ LIFECYCLE_CONTRADICTED, }; pub use registry::{ - backfill_registry, diff_registries, LanceRegistry, MirroredRegistry, RegistryEntry, - RolloutRegistry, StoreRegistry, + 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}; diff --git a/crates/lance-context-core/src/registry.rs b/crates/lance-context-core/src/registry.rs index e40e286..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, @@ -111,16 +111,30 @@ impl StoreRegistry for LanceRegistry { } } -/// A [`StoreRegistry`] that reads from `primary` and writes to both `primary` -/// and `mirror`. The mirror write is best-effort: it is logged on failure and -/// never fails the call, because the mirror exists to be caught up by a -/// periodic backfill during a backend migration, not to be authoritative. +/// 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 { @@ -134,64 +148,123 @@ impl StoreRegistry for MirroredRegistry { } async fn upsert(&self, name: &str, uri: &str) -> LanceResult<()> { self.primary.upsert(name, uri).await?; - if let Err(error) = self.mirror.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.mirror.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) = self.mirror.insert_missing(entries).await { - tracing::warn!(registry = self.label, %error, "registry mirror backfill failed"); + 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> { - // Whichever side is the Lance table still needs maintaining while it - // is written to. - self.primary - .lance_table() - .or_else(|| self.mirror.lance_table()) + self.primary.lance_table() } fn as_any(&self) -> &dyn std::any::Any { self } } -/// Copy every entry of `from` that `to` lacks. Returns how many were copied. -/// Used to seed a new backend from the old one and to heal mirror misses. +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 entries: Vec<(String, String)> = from - .list() - .await? - .into_iter() - .map(|e| (e.name, e.uri)) - .collect(); - to.insert_missing(&entries).await + 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() + } } -/// Names present in exactly one of two registries: `(only_in_a, only_in_b)`. pub async fn diff_registries( a: &dyn StoreRegistry, b: &dyn StoreRegistry, -) -> LanceResult<(Vec, Vec)> { - let a_names: HashSet = a.list().await?.into_iter().map(|e| e.name).collect(); - let b_names: HashSet = b.list().await?.into_iter().map(|e| e.name).collect(); - let mut only_a: Vec = a_names.difference(&b_names).cloned().collect(); - let mut only_b: Vec = b_names.difference(&a_names).cloned().collect(); - only_a.sort(); - only_b.sort(); - Ok((only_a, only_b)) +) -> 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. diff --git a/crates/lance-context-core/src/registry_etcd.rs b/crates/lance-context-core/src/registry_etcd.rs index d181e1a..a3398d5 100644 --- a/crates/lance-context-core/src/registry_etcd.rs +++ b/crates/lance-context-core/src/registry_etcd.rs @@ -22,6 +22,17 @@ use crate::registry::{RegistryEntry, StoreRegistry}; 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. @@ -29,6 +40,7 @@ pub struct EtcdRegistry { client: Client, /// `/registry/` with no trailing slash. prefix: String, + migration_key: String, } impl EtcdRegistry { @@ -38,6 +50,10 @@ impl EtcdRegistry { Arc::new(Self { client, prefix: format!("{}/registry/{kind}", prefix.trim_end_matches('/')), + migration_key: format!( + "{}/registry-migrations/{kind}", + prefix.trim_end_matches('/') + ), }) } @@ -45,7 +61,7 @@ impl EtcdRegistry { format!("{}/{name}", self.prefix) } - fn entry(&self, key: &[u8], value: &[u8]) -> LanceResult { + fn entry(&self, key: &[u8], value: &[u8]) -> LanceResult> { let key = String::from_utf8_lossy(key); let name = key .strip_prefix(&format!("{}/", self.prefix)) @@ -53,32 +69,249 @@ impl EtcdRegistry { .to_string(); let value: Value = serde_json::from_slice(value) .map_err(|e| LanceError::io(format!("decode registry entry '{name}': {e}")))?; - Ok(RegistryEntry { + 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_trait::async_trait] -impl StoreRegistry for EtcdRegistry { - async fn contains(&self, name: &str) -> LanceResult { + async fn migration(&self) -> LanceResult<(Migration, i64)> { let response = self .client .clone() - .get(self.key(name), Some(GetOptions::new().with_count_only())) + .get(self.migration_key.clone(), None) .await - .map_err(etcd_error("registry get"))?; - Ok(response.count() > 0) + .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> { @@ -93,6 +326,7 @@ impl StoreRegistry for EtcdRegistry { .first() .map(|kv| self.entry(kv.key(), kv.value())) .transpose() + .map(Option::flatten) } async fn list(&self) -> LanceResult> { @@ -109,10 +343,12 @@ impl StoreRegistry for EtcdRegistry { .kvs() .iter() .map(|kv| self.entry(kv.key(), kv.value())) - .collect() + .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? { @@ -128,6 +364,7 @@ impl StoreRegistry for EtcdRegistry { } async fn remove(&self, name: &str) -> LanceResult<()> { + self.ensure_writable().await?; self.client .clone() .delete(self.key(name), Some(DeleteOptions::new())) @@ -137,10 +374,9 @@ impl StoreRegistry for EtcdRegistry { } async fn insert_missing(&self, entries: &[(String, String)]) -> LanceResult { - // One put-if-absent transaction per entry, batched so a 10k-row - // backfill is a few hundred round trips rather than 10k. A txn with - // several compares is all-or-nothing, which is not what we want here, - // so each entry is its own txn inside the batch loop. + 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; @@ -152,8 +388,25 @@ impl StoreRegistry for EtcdRegistry { 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::version(key.as_str(), CompareOp::Equal, 0)]) + .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) @@ -257,9 +510,9 @@ mod tests { lance.upsert("pre-2", "/p2").await.unwrap(); let etcd_dyn: Arc = etcd.clone(); - let (only_lance, only_etcd) = diff_registries(&*lance, &*etcd_dyn).await.unwrap(); - assert_eq!(only_lance, vec!["pre-1", "pre-2"]); - assert!(only_etcd.is_empty()); + 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 { @@ -271,9 +524,8 @@ mod tests { mirrored.remove("pre-1").await.unwrap(); assert!(mirrored.contains("new").await.unwrap()); - let (only_lance, only_etcd) = diff_registries(&*lance, &*etcd_dyn).await.unwrap(); - assert!(only_lance.is_empty(), "{only_lance:?}"); - assert!(only_etcd.is_empty(), "{only_etcd:?}"); + let diff = diff_registries(&*lance, &*etcd_dyn).await.unwrap(); + assert!(diff.is_empty(), "{diff:?}"); let mut names: Vec = etcd .list() .await @@ -300,4 +552,157 @@ mod tests { 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/routes.rs b/crates/lance-context-master/src/routes.rs index 2ccf5eb..d2ce1a3 100644 --- a/crates/lance-context-master/src/routes.rs +++ b/crates/lance-context-master/src/routes.rs @@ -588,7 +588,7 @@ pub async fn list_repairs( #[derive(Debug, serde::Deserialize)] pub struct RegistryParams { - /// `rollout` (default) or `generic`. + /// `rollout` (default), `generic`, or `datagen`. #[serde(default)] pub kind: Option, } @@ -600,16 +600,16 @@ fn registry_pair<'a>( 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)" + "unknown registry kind '{other}' (rollout|generic|datagen)" ))), } } -/// `GET /api/v1/registry/diff?kind=` — names present in the primary backend -/// but not the mirror, and vice versa. Empty on both sides means the two -/// backends agree and a migration step can proceed. 400 when no mirror is -/// configured, because then there is nothing to compare. +/// 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, @@ -623,19 +623,20 @@ pub async fn registry_diff( "REGISTRY_MIRROR is not configured; nothing to diff".to_string(), )); }; - let (only_primary, only_mirror) = - lance_context_core::diff_registries(&*mirrored.primary, &*mirrored.mirror) - .await - .map_err(MasterError::from_lance)?; + 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": only_primary, - "only_in_mirror": only_mirror, + "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=` — copy every primary entry the -/// mirror lacks into the mirror. Idempotent. The master also does this on +/// `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>, diff --git a/crates/lance-context-master/src/scanner.rs b/crates/lance-context-master/src/scanner.rs index bbcfd16..0f7d5f9 100644 --- a/crates/lance-context-master/src/scanner.rs +++ b/crates/lance-context-master/src/scanner.rs @@ -236,6 +236,7 @@ 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"); diff --git a/crates/lance-context-master/src/state.rs b/crates/lance-context-master/src/state.rs index e65fa13..28b5142 100644 --- a/crates/lance-context-master/src/state.rs +++ b/crates/lance-context-master/src/state.rs @@ -71,6 +71,8 @@ pub struct MasterState { /// while the master -- whose etcd-locked MergeWal task exists to serialize /// exactly that -- never heard of them. 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 @@ -141,6 +143,13 @@ impl MasterState { 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!( @@ -151,7 +160,11 @@ impl MasterState { // 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)] { + 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 => { @@ -169,6 +182,7 @@ impl MasterState { let state = Arc::new(Self { registry, generic_registry, + datagen_registry, stats: Mutex::new(stats), stats_cache: RwLock::new(Arc::new(Vec::new())), record_stores: Mutex::new(LruCache::new( @@ -340,6 +354,55 @@ mod tests { 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/docs/src/design/registry-etcd.md b/docs/src/design/registry-etcd.md index d7ed185..34403d5 100644 --- a/docs/src/design/registry-etcd.md +++ b/docs/src/design/registry-etcd.md @@ -1,209 +1,102 @@ -# Design: move the store registries from Lance tables to etcd - -Status: proposal. Follows #274 (which makes the Lance registries sustainable -enough that this is not urgent). - -## What the registry is and who touches it - -Three tables, one per store kind, each a directory `{name -> uri, created_at}` -shared by every worker and the master: - -| table | rows (prod) | writers | readers | -|---|---|---|---| -| `_registry.rollout.lance` | ~10.7k | worker create/delete, master retire | every worker miss on `get_or_open_*`, master scan (every 300 s, `list()`), discovery backfill | -| `_registry.generic.lance` | ~600 | same | same | -| `_registry.datagen.lance` | small | same | same | - -Every call site today goes through `RolloutRegistry` (`crates/lance-context-core/src/registry.rs`): -`contains`, `get`, `list`, `upsert`, `remove`, `insert_missing`. Each read does -`checkout_latest` first, i.e. one manifest read from ADLS per call. Each `upsert` -is delete + append, two Lance commits with optimistic concurrency across 20 -workers. - -### Why the Lance table is the wrong tool here - -- A directory is a key-value set with point lookups and rare listing. Lance is - a columnar table optimised for scans; every point lookup pays a manifest read - plus a scan, and every write pays a full commit protocol. -- Without maintenance it grows quadratically (#274 explains the 29 GB). With - maintenance it still costs a manifest read per `contains`, ~0.7 s on ADLS - today, 1.5 s+ under throttling, and 20 workers do it on every cache miss. -- Writes from 20 workers contend on one commit point. Lance retries, but a - create under ADLS throttling took 47 s on average and failed 79 times in the - last 15 h. -- Existence checks and store opens are the hot path for *every* API call that - misses the worker's LRU. Registry latency is user-visible latency. - -### What etcd gives - -- Point `contains`/`get`: one round trip, single-digit ms, from any pod. -- `upsert`/`remove`: one transaction, no manifest, no commit retries. -- `list` with prefix: fine at 10k keys (etcd pages by default; the master - scans every 300 s and already tolerates seconds). -- Watch: workers can invalidate their caches on create/delete instead of - polling; not needed for v1 but free. -- The master already runs etcd (task store, coordination locks, cooldowns) and - the cluster exposes `rocketkeep-etcd:2379` to the namespace. Workers do not - link `etcd-client` yet. - -### What etcd does not give, and how we cope - -- **Not the source of truth for data.** The dataset on ADLS is. The registry - only says "this name exists and lives here". If etcd is lost we must be able - to rebuild it. The master already has `discovery.rs` (`insert_missing` from a - directory listing); keep it and point it at etcd. -- **Size.** Value is `{uri, created_at}` (~200 B). 11k keys is ~2 MB. etcd is - comfortable to a few hundred MB; irrelevant here. -- **Availability coupling.** A worker that cannot reach etcd cannot create or - open a store it does not have cached. Today a worker that cannot reach ADLS - is equally dead, and etcd is a 3-node cluster inside the same k8s cluster, so - this is not a new failure domain in practice. Still: reads fall back to the - in-process cache, and we keep the Lance table as a read-only fallback during - the migration window (below). - -## Data model - -``` -/registry// -> {"uri": "...", "created_at": 1790000000000} -``` - -- `` is the same `ETCD_PREFIX` the master uses (`/lance-context` by - default), so one etcd holds one deployment's tasks and registries together. -- `` is `rollout`, `generic`, `datagen`. -- `` is the store name, already validated to be a portable path segment - (`validate_name`), so it is a safe key segment with no escaping. -- Value is JSON. No lease: registry entries are durable until explicitly - removed. - -## Code shape - -A trait in `lance-context-core` so the server and master do not care which -backend they have: - -```rust -#[async_trait] -pub trait StoreRegistry: Send + Sync { - async fn contains(&self, name: &str) -> LanceResult; - async fn get(&self, name: &str) -> LanceResult>; - async fn list(&self) -> LanceResult>; - async fn upsert(&self, name: &str, uri: &str) -> LanceResult<()>; - async fn remove(&self, name: &str) -> LanceResult; - async fn insert_missing(&self, entries: &[(String, String)]) -> LanceResult; -} -``` - -Note `&self`, not `&mut self`: etcd needs no handle mutation, and the -`RwLock` everywhere today exists only because Lance handles -must `checkout_latest`. Dropping the lock removes a serialisation point every -worker request goes through. - -Two impls: - -- `LanceRegistry` — today's `RolloutRegistry`, wrapped. Keeps `maintain()`. -- `EtcdRegistry` — `etcd_client::Client` + prefix. `upsert` is a plain `put` - (idempotent by construction; the delete-then-append dance goes away). - `insert_missing` is a batch of `txn(version(key)==0 → put)`. - -Selection by config: `REGISTRY_BACKEND=lance|etcd` (default `lance` until the -migration below is done), plus the existing `ETCD_ENDPOINTS`/`ETCD_PREFIX`/TLS -flags lifted from the master's config into core so the server can reuse them. -`etcd-client` becomes a dependency of `lance-context-core` (feature-gated -`etcd`, on by default). - -Worker-side cache: `get_or_open_*` already has an LRU of open stores and only -asks the registry on a miss. Keep that. Add a small negative cache -(name → not-found, 5 s TTL) so a client hammering a non-existent name does not -turn into an etcd storm; cheap and safe because creates go through the same -worker fleet and invalidate locally, and a 5 s stale "not found" on another -worker is the same window Lance gives today. - -## Migration - -Zero-downtime, three deploys, each independently revertable. - -1. **Dual-write** (`REGISTRY_BACKEND=lance`, `REGISTRY_MIRROR=etcd`): reads hit - Lance as today; every `upsert`/`remove` also goes to etcd, best-effort with - a warning on failure. Master runs a one-shot backfill on startup - (`list()` from Lance → `insert_missing` into etcd) under the existing - `state-init` coordination lock, and re-runs it every maintenance round so - any mirror miss heals. Ship this, let it run a day, then compare: a - `GET /api/v1/registry/diff` on the master lists keys present in one backend - and not the other. Expect empty. -2. **Read from etcd** (`REGISTRY_BACKEND=etcd`, `REGISTRY_MIRROR=lance`): - reads hit etcd; writes still mirror to Lance so a revert to step 1 loses - nothing. Watch worker `rollout_store_cache_misses_total` latency and - `POST /rollouts` latency drop. Run a few days. -3. **etcd only** (`REGISTRY_BACKEND=etcd`, no mirror): stop writing the Lance - tables. Leave them on disk for a while as a cold backup; `discovery.rs` can - rebuild etcd from the ADLS directory listing regardless. - -Each step is env-only in mango; no image change between steps once the code is -in. Revert at any step is flipping the env back. - -## Master changes - -- `MasterState.registry` / `generic_registry` become `Arc`. - The scanner's `list()` every 300 s is unchanged in shape. -- Registry maintenance from #274 stays for the Lance backend and becomes a - no-op for etcd. -- New `GET /api/v1/registry/diff` (step 1 verification) and - `POST /api/v1/registry/backfill` (manual re-sync). -- Retirement (`retire_cold_experiments`) calls `remove()` on the trait; no - change. - -## Worker changes - -- `AppState.{rollout,generic,datagen}_registry` become `Arc`; - the `RwLock` goes away. -- `create_*_store`: `contains` → open/create dataset → `upsert`. Same order; - the dataset is still created before the registry row so a crash between them - leaves an orphan dataset (discovery picks it up) rather than a dangling - registry entry. Unchanged semantics, ~45 s less latency. -- `get_or_open_*`: `contains` on miss, as today. The double `contains` (before - and after the open) can stay; it is now two 2 ms calls. -- `unregister_*`: `contains` → delete data → `remove`. Unchanged. - -## Consistency - -- **Create/create race** on the same name from two workers: today, both pass - `contains`, both create the dataset (Lance `Create` is put-if-absent so one - fails), the winner upserts. With etcd, same: dataset creation is still the - arbiter; the loser's `contains` re-check after the failed create returns - true and it answers 409. No change in outcome. -- **Create/delete race**: today serialised per name only within one worker - (`rollout_handles.lock(name)`); across workers it is last-writer-wins on the - Lance table. etcd is the same but with a smaller window. If we want it - strictly ordered we can use `txn(mod_revision == expected)`; not proposed - for v1 because the current behaviour has not caused a problem. -- **Read-your-writes**: etcd reads are linearisable by default. A worker that - just created a store sees it from any other worker immediately, which is - *better* than today (Lance `checkout_latest` is also strongly consistent, so - no regression, but etcd removes the 0.7 s manifest read). - -## Rollback - -- Step 1 → nothing to roll back; the mirror is additive. -- Step 2 → flip `REGISTRY_BACKEND=lance`. Lance was kept current by the mirror. -- Step 3 → re-enable the mirror, run backfill etcd → Lance (the same - `insert_missing` path in reverse, exposed as - `POST /api/v1/registry/backfill?to=lance`), flip. - -## Effort - -- core: trait + `EtcdRegistry` + config plumbing, ~400 lines incl. tests - (tests run against the etcd the master tests already require). -- server: swap type, drop the `RwLock`, mirror logic, negative cache, ~150 lines. -- master: swap type, backfill, diff/backfill routes, ~150 lines. -- mango: three env flips. - -Roughly two days of work plus the staged rollout. Not blocking anything -today: #274 takes the registry from "quadratic and failing" to "linear and -slow", and the etcd move takes it to "fast". - -## Not in scope - -- Moving `_stats` to etcd. It is a 10k-row table rewritten as one snapshot per - scan round and read for range queries (`list_above_fragment_count`, - `list_above_pending_wal`), which etcd cannot do; it stays in Lance. -- Moving MemWAL shard manifests. They are per-shard, written by one writer, - and read by the LSM scanner; different problem. +# 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. From 29fa598d1810974cf093df3aec643125110618cf Mon Sep 17 00:00:00 2001 From: Beinan Wang <> Date: Fri, 2 Oct 2026 06:58:11 +0000 Subject: [PATCH 3/4] ci: pass registry coverage filter to the test harness --- .github/workflows/rust-test.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/rust-test.yml b/.github/workflows/rust-test.yml index 350c46c..e30c5d6 100644 --- a/.github/workflows/rust-test.yml +++ b/.github/workflows/rust-test.yml @@ -135,7 +135,7 @@ jobs: 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 + -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 From 6a332f9e87d9085f8fa0fbf921c9e89843757292 Mon Sep 17 00:00:00 2001 From: Beinan Wang <> Date: Fri, 2 Oct 2026 07:27:47 +0000 Subject: [PATCH 4/4] Adapt no-op regression tests to shared registry configuration --- crates/lance-context-master/src/scheduler.rs | 8 +------- crates/lance-context-master/src/task_store.rs | 6 +++--- 2 files changed, 4 insertions(+), 10 deletions(-) diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index c5e81bd..b953536 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -933,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); diff --git a/crates/lance-context-master/src/task_store.rs b/crates/lance-context-master/src/task_store.rs index bbb93d3..afd768b 100644 --- a/crates/lance-context-master/src/task_store.rs +++ b/crates/lance-context-master/src/task_store.rs @@ -1434,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") @@ -1457,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();