From 02c36b3282655b210c2ff80b76b2930c85b30569 Mon Sep 17 00:00:00 2001 From: Tomasz Szymczyszyn Date: Fri, 9 Oct 2026 21:28:36 +0200 Subject: [PATCH 1/2] libsql-server: make namespace creation all-or-nothing `NamespaceStore::create` persisted the namespace config before setting the namespace up, and nothing undid it: a dump that failed to import left a "ghost" namespace (config without data) whose name could not be reused and which the next user request lazily turned into an empty database. A request cancelled mid-import ran no cleanup at all, since the error branch of the future is never reached when it is dropped, and the schema configurator had no cleanup even for errors. The metastore write is now the commit point and everything before it is reversible: - `Reservation` (store) write-locks the namespace's cache entry for the duration of the creation and keeps the config in memory only. `publish` installs the finished namespace; dropping the reservation unpublished (error or cancellation) removes the in-memory config and only then releases the lock, so waiters observe a namespace that does not exist, never a partial one. `fork` uses it too, replacing its ad-hoc guard, which blocked in place (panics on a current-thread runtime) and published before flushing. - `FreshDir` (configurators) removes a brand-new namespace directory when setup fails or its future is dropped; used by the primary and schema configurators. - An `.incomplete` marker is written into a new directory and removed once the config is persisted. A creation that finds a marked directory for a name the metastore doesn't know discards it, so a process crash during an import never blocks a retry; the metastore's filesystem recovery skips marked directories. Directories without the marker are left alone. Requests for a namespace that is being created wait for the outcome, as during `fork` and `reset`. In single-namespace mode `create` of the default namespace remains a config upsert, but creating it *from a dump* while it exists is now rejected instead of silently skipping the import. Contract for callers: 2xx means complete; anything else, including a lost connection, means retry; `already exists` on the retry means the earlier attempt did complete. Requires async-lock 3 for `RwLock::write_arc` (already in the lockfile; the bump removes the duplicate 2.x version). Tests: 9 lifecycle tests (failed/cancelled creation leaves no trace for both importers and for shared-schema namespaces, concurrent requests wait, crash remnants are discarded, unmarked directories and published namespaces are untouched, single-namespace dump create is rejected) plus unit tests for the guards. Nine of the eleven integration tests fail on the previous code. --- Cargo.lock | 25 +- libsql-server/Cargo.toml | 2 +- .../src/namespace/configurator/helpers.rs | 187 ++++++++- .../src/namespace/configurator/mod.rs | 12 + .../src/namespace/configurator/primary.rs | 34 +- .../src/namespace/configurator/schema.rs | 24 +- libsql-server/src/namespace/meta_store.rs | 7 + libsql-server/src/namespace/mod.rs | 15 + libsql-server/src/namespace/store.rs | 172 +++++--- libsql-server/tests/namespaces/dumps.rs | 121 +++--- libsql-server/tests/namespaces/lifecycle.rs | 366 ++++++++++++++++++ libsql-server/tests/namespaces/mod.rs | 15 +- 12 files changed, 834 insertions(+), 146 deletions(-) create mode 100644 libsql-server/tests/namespaces/lifecycle.rs diff --git a/Cargo.lock b/Cargo.lock index 436299524e..1b951752e8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -224,22 +224,13 @@ dependencies = [ "zstd-safe 7.2.0", ] -[[package]] -name = "async-lock" -version = "2.8.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "287272293e9d8c41773cec55e365490fe034813a2f172f502d6ddcf75b2f582b" -dependencies = [ - "event-listener 2.5.3", -] - [[package]] name = "async-lock" version = "3.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff6e472cdea888a4bd64f342f09b3f50e1886d32afe8df3d663c01140b811b18" dependencies = [ - "event-listener 5.3.1", + "event-listener", "event-listener-strategy", "pin-project-lite", ] @@ -1920,12 +1911,6 @@ dependencies = [ "windows-sys 0.52.0", ] -[[package]] -name = "event-listener" -version = "2.5.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0206175f82b8d6bf6652ff7d71a1e27fd2e4efde587fd368662814d6ec1d9ce0" - [[package]] name = "event-listener" version = "5.3.1" @@ -1943,7 +1928,7 @@ version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0f214dc438f977e6d4e3500aaa277f5ad94ca83fbbd9b1a15713ce2344ccc5a1" dependencies = [ - "event-listener 5.3.1", + "event-listener", "pin-project-lite", ] @@ -3010,7 +2995,7 @@ dependencies = [ "aes", "anyhow", "arbitrary", - "async-lock 2.8.0", + "async-lock", "async-recursion", "async-stream", "async-tempfile", @@ -3423,12 +3408,12 @@ version = "0.12.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32cf62eb4dd975d2dde76432fb1075c49e3ee2331cf36f1f8fd4b66550d32b6f" dependencies = [ - "async-lock 3.4.0", + "async-lock", "async-trait", "crossbeam-channel", "crossbeam-epoch", "crossbeam-utils", - "event-listener 5.3.1", + "event-listener", "futures-util", "once_cell", "parking_lot", diff --git a/libsql-server/Cargo.toml b/libsql-server/Cargo.toml index e9721a32ba..6c11402828 100644 --- a/libsql-server/Cargo.toml +++ b/libsql-server/Cargo.toml @@ -11,7 +11,7 @@ path = "src/main.rs" [dependencies] anyhow = "1.0.66" -async-lock = "2.6.0" +async-lock = "3" async-stream = "0.3.5" async-tempfile = "0.4.0" async-trait = "0.1.58" diff --git a/libsql-server/src/namespace/configurator/helpers.rs b/libsql-server/src/namespace/configurator/helpers.rs index 023333be30..9ff8560a6f 100644 --- a/libsql-server/src/namespace/configurator/helpers.rs +++ b/libsql-server/src/namespace/configurator/helpers.rs @@ -20,7 +20,7 @@ use crate::namespace::meta_store::MetaStoreHandle; use crate::namespace::replication_wal::{make_replication_wal_wrapper, ReplicationWalWrapper}; use crate::namespace::{ NamespaceBottomlessDbId, NamespaceBottomlessDbIdInit, NamespaceName, ResolveNamespacePathFn, - RestoreOption, + RestoreOption, INCOMPLETE_MARKER, }; use crate::replication::ReplicationLogger; use crate::stats::Stats; @@ -28,6 +28,93 @@ use crate::{StatsSender, BLOCKING_RT, DB_CREATE_TIMEOUT, DEFAULT_AUTO_CHECKPOINT use super::{BaseNamespaceConfig, PrimaryConfig}; +/// The directory of a namespace that does not exist yet. +/// +/// Created (with the [`INCOMPLETE_MARKER`]) by [`begin`](Self::begin) and removed again when the +/// guard is dropped before [`keep`](Self::keep) is called: on error, or when the request that was +/// creating the namespace is cancelled and its future dropped. The marker itself is removed by +/// the namespace store once the namespace's config is persisted. +pub(super) struct FreshDir { + path: Option>, +} + +impl FreshDir { + /// Returns `None` when a namespace directory already exists at `db_path`. + pub(super) async fn begin(db_path: &Arc) -> crate::Result> { + if db_path.try_exists()? { + return Ok(None); + } + tokio::fs::create_dir_all(db_path).await?; + tokio::fs::File::create(db_path.join(INCOMPLETE_MARKER)).await?; + Ok(Some(Self { + path: Some(db_path.clone()), + })) + } + + /// The namespace was set up; the directory stays. + pub(super) fn keep(mut self) { + self.path.take(); + } +} + +impl Drop for FreshDir { + fn drop(&mut self) { + let Some(path) = self.path.take() else { + return; + }; + tracing::warn!( + path = %path.display(), + "namespace creation did not complete; removing its directory" + ); + // Dropping may happen while the creating future is being torn down (cancellation), so + // the removal is handed to the runtime rather than awaited here. Should it fail, the + // directory still carries the marker and the next creation of this namespace discards it. + match tokio::runtime::Handle::try_current() { + Ok(rt) => { + rt.spawn(async move { + if let Err(e) = tokio::fs::remove_dir_all(&*path).await { + log_remove_error(&path, e); + } + }); + } + Err(_) => { + if let Err(e) = std::fs::remove_dir_all(&*path) { + log_remove_error(&path, e); + } + } + } + } +} + +fn log_remove_error(path: &Path, e: std::io::Error) { + if e.kind() != std::io::ErrorKind::NotFound { + tracing::error!( + path = %path.display(), + "failed to remove unfinished namespace directory: {e}" + ); + } +} + +/// Remove `dbs/` if it was left behind by a creation that never completed, i.e. if it +/// still carries the [`INCOMPLETE_MARKER`]. Directories without the marker are left alone. +pub(super) async fn discard_incomplete( + base: &BaseNamespaceConfig, + namespace: &NamespaceName, +) -> crate::Result<()> { + let db_path = base.base_path.join("dbs").join(namespace.as_str()); + if !db_path.join(INCOMPLETE_MARKER).try_exists()? { + return Ok(()); + } + tracing::warn!( + namespace = %namespace, + path = %db_path.display(), + "discarding namespace directory left by an interrupted creation" + ); + metrics::increment_counter!("libsql_server_namespace_create_discarded_incomplete"); + tokio::fs::remove_dir_all(&db_path).await?; + Ok(()) +} + #[tracing::instrument(skip_all)] pub(super) async fn make_primary_connection_maker( primary_config: &PrimaryConfig, @@ -405,3 +492,101 @@ pub(super) async fn cleanup_primary( Ok(()) } + +#[cfg(test)] +mod test { + use super::*; + + async fn settle() { + // the removal is spawned; give it a chance to run + for _ in 0..50 { + tokio::task::yield_now().await; + tokio::time::sleep(Duration::from_millis(2)).await; + } + } + + #[tokio::test] + async fn fresh_dir_is_removed_unless_kept() { + let tmp = tempfile::tempdir().unwrap(); + let path: Arc = tmp.path().join("dbs").join("ns").into(); + + let fresh = FreshDir::begin(&path).await.unwrap().expect("fresh"); + assert!(path.join(INCOMPLETE_MARKER).exists()); + drop(fresh); + settle().await; + assert!(!path.exists(), "dropped before keep: removed"); + + let fresh = FreshDir::begin(&path).await.unwrap().expect("fresh again"); + fresh.keep(); + settle().await; + assert!( + path.join(INCOMPLETE_MARKER).exists(), + "kept: stays, marker included" + ); + + assert!( + FreshDir::begin(&path).await.unwrap().is_none(), + "an existing directory is not fresh" + ); + } + + #[tokio::test] + async fn fresh_dir_is_removed_when_its_future_is_dropped() { + let tmp = tempfile::tempdir().unwrap(); + let path: Arc = tmp.path().join("dbs").join("ns").into(); + + let setup = { + let path = path.clone(); + async move { + let fresh = FreshDir::begin(&path).await.unwrap(); + std::future::pending::<()>().await; // "the import" + fresh.map(FreshDir::keep); + } + }; + let aborted = tokio::time::timeout(Duration::from_millis(50), setup).await; + assert!(aborted.is_err()); + settle().await; + assert!(!path.exists()); + } + + #[tokio::test] + async fn discard_incomplete_only_removes_marked_directories() { + let tmp = tempfile::tempdir().unwrap(); + let base = BaseNamespaceConfig { + base_path: tmp.path().into(), + extensions: Arc::new([]), + stats_sender: tokio::sync::mpsc::channel(1).0, + max_response_size: 0, + max_total_response_size: 0, + max_concurrent_connections: Arc::new(tokio::sync::Semaphore::new(1)), + max_concurrent_requests: 0, + encryption_config: None, + connection_creation_timeout: None, + disable_intelligent_throttling: false, + dump_import: Default::default(), + }; + let marked = tmp.path().join("dbs").join("marked"); + let unmarked = tmp.path().join("dbs").join("unmarked"); + for dir in [&marked, &unmarked] { + std::fs::create_dir_all(dir).unwrap(); + std::fs::write(dir.join("data"), b"junk").unwrap(); + } + std::fs::write(marked.join(INCOMPLETE_MARKER), b"").unwrap(); + + discard_incomplete(&base, &NamespaceName::from_string("marked".into()).unwrap()) + .await + .unwrap(); + discard_incomplete( + &base, + &NamespaceName::from_string("unmarked".into()).unwrap(), + ) + .await + .unwrap(); + discard_incomplete(&base, &NamespaceName::from_string("absent".into()).unwrap()) + .await + .unwrap(); + + assert!(!marked.exists()); + assert!(unmarked.join("data").exists()); + } +} diff --git a/libsql-server/src/namespace/configurator/mod.rs b/libsql-server/src/namespace/configurator/mod.rs index 064f9a8f7c..ecb384fc66 100644 --- a/libsql-server/src/namespace/configurator/mod.rs +++ b/libsql-server/src/namespace/configurator/mod.rs @@ -131,6 +131,18 @@ pub trait ConfigureNamespace { bottomless_db_id_init: NamespaceBottomlessDbIdInit, ) -> Pin> + Send + 'a>>; + /// Remove whatever a previous, interrupted creation of `namespace` left on disk, so that a + /// new creation starts from nothing. Only called for a name that is neither loaded nor known + /// to the metastore; must leave anything that does not carry the + /// [`INCOMPLETE_MARKER`](super::INCOMPLETE_MARKER) alone. + fn discard_incomplete<'a>( + &'a self, + namespace: &'a NamespaceName, + ) -> Pin> + Send + 'a>> { + let _ = namespace; + Box::pin(async { Ok(()) }) + } + fn fork<'a>( &'a self, from_ns: &'a Namespace, diff --git a/libsql-server/src/namespace/configurator/primary.rs b/libsql-server/src/namespace/configurator/primary.rs index f68405fad6..f918ce7a7b 100644 --- a/libsql-server/src/namespace/configurator/primary.rs +++ b/libsql-server/src/namespace/configurator/primary.rs @@ -21,7 +21,7 @@ use crate::namespace::{ use crate::run_periodic_checkpoint; use crate::schema::{has_pending_migration_task, setup_migration_table}; -use super::helpers::cleanup_primary; +use super::helpers::{cleanup_primary, discard_incomplete, FreshDir}; use super::{BaseNamespaceConfig, ConfigureNamespace, PrimaryConfig}; pub struct PrimaryConfigurator { @@ -129,35 +129,33 @@ impl ConfigureNamespace for PrimaryConfigurator { ) -> Pin> + Send + 'a>> { Box::pin(async move { let db_path: Arc = self.base.base_path.join("dbs").join(name.as_str()).into(); - let fresh_namespace = !db_path.try_exists()?; - // FIXME: make that truly atomic. explore the idea of using temp directories, and it's implications - match self + // A brand-new directory is removed again if setup fails or this future is dropped. + let fresh = FreshDir::begin(&db_path).await?; + let ns = self .try_new_primary( name.clone(), meta_store_handle, restore_option, resolve_attach_path, - db_path.clone(), + db_path, broadcaster, self.base.encryption_config.clone(), ) - .await - { - Ok(this) => Ok(this), - Err(e) if fresh_namespace => { - tracing::error!( - "an error occured while deleting creating namespace, cleaning..." - ); - if let Err(e) = tokio::fs::remove_dir_all(&db_path).await { - tracing::error!("failed to remove dirty namespace directory: {e}") - } - Err(e) - } - Err(e) => Err(e), + .await?; + if let Some(fresh) = fresh { + fresh.keep(); } + Ok(ns) }) } + fn discard_incomplete<'a>( + &'a self, + namespace: &'a NamespaceName, + ) -> Pin> + Send + 'a>> { + Box::pin(discard_incomplete(&self.base, namespace)) + } + fn cleanup<'a>( &'a self, namespace: &'a NamespaceName, diff --git a/libsql-server/src/namespace/configurator/schema.rs b/libsql-server/src/namespace/configurator/schema.rs index 275fd71e93..40c7614849 100644 --- a/libsql-server/src/namespace/configurator/schema.rs +++ b/libsql-server/src/namespace/configurator/schema.rs @@ -13,7 +13,9 @@ use crate::namespace::{ }; use crate::schema::SchedulerHandle; -use super::helpers::{cleanup_primary, make_primary_connection_maker}; +use super::helpers::{ + cleanup_primary, discard_incomplete, make_primary_connection_maker, FreshDir, +}; use super::{BaseNamespaceConfig, ConfigureNamespace, PrimaryConfig}; pub struct SchemaConfigurator { @@ -52,9 +54,10 @@ impl ConfigureNamespace for SchemaConfigurator { ) -> std::pin::Pin> + Send + 'a>> { Box::pin(async move { let mut join_set = JoinSet::new(); - let db_path = self.base.base_path.join("dbs").join(name.as_str()); - - tokio::fs::create_dir_all(&db_path).await?; + let db_path: Arc = + self.base.base_path.join("dbs").join(name.as_str()).into(); + // A brand-new directory is removed again if setup fails or this future is dropped. + let fresh = FreshDir::begin(&db_path).await?; let (connection_maker, wal_manager, stats) = make_primary_connection_maker( &self.primary_config, @@ -72,6 +75,10 @@ impl ConfigureNamespace for SchemaConfigurator { ) .await?; + if let Some(fresh) = fresh { + fresh.keep(); + } + Ok(Namespace { db: Database::Schema(SchemaDatabase::new( self.migration_scheduler.clone(), @@ -89,11 +96,18 @@ impl ConfigureNamespace for SchemaConfigurator { tasks: join_set, stats, db_config_store: db_config.clone(), - path: db_path.into(), + path: db_path, }) }) } + fn discard_incomplete<'a>( + &'a self, + namespace: &'a NamespaceName, + ) -> std::pin::Pin> + Send + 'a>> { + Box::pin(discard_incomplete(&self.base, namespace)) + } + fn cleanup<'a>( &'a self, namespace: &'a NamespaceName, diff --git a/libsql-server/src/namespace/meta_store.rs b/libsql-server/src/namespace/meta_store.rs index 70b419ebe9..5840ae9dfc 100644 --- a/libsql-server/src/namespace/meta_store.rs +++ b/libsql-server/src/namespace/meta_store.rs @@ -223,6 +223,13 @@ impl MetaStoreInner { if !entry.path().is_dir() { continue; } + if entry.path().join(super::INCOMPLETE_MARKER).try_exists()? { + tracing::warn!( + "skipping `{}`: left by an interrupted namespace creation", + entry.path().display() + ); + continue; + } let config_path = entry.path().join("config.json"); let name = NamespaceName::from_string(entry.file_name().to_str().unwrap().to_string())?; diff --git a/libsql-server/src/namespace/mod.rs b/libsql-server/src/namespace/mod.rs index 2cea87ed75..f39b3c8555 100644 --- a/libsql-server/src/namespace/mod.rs +++ b/libsql-server/src/namespace/mod.rs @@ -57,6 +57,12 @@ pub enum NamespaceBottomlessDbIdInit { } /// A namespace isolates the resources pertaining to a database of type T +/// Marker file present in a namespace directory from its creation until the namespace's config +/// has been persisted to the metastore. A directory that still carries it belongs to a creation +/// that did not complete (the process crashed, or cleanup failed); the next attempt to create a +/// namespace with that name discards it. See `docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md`. +pub(crate) const INCOMPLETE_MARKER: &str = ".incomplete"; + #[derive(Debug)] pub struct Namespace { pub db: Database, @@ -73,6 +79,15 @@ impl Namespace { &self.name } + /// Remove the [`INCOMPLETE_MARKER`]: the namespace is now persisted in the metastore (or was + /// opened from it, which implies the same). + pub(crate) async fn mark_complete(&self) -> std::io::Result<()> { + match tokio::fs::remove_file(self.path.join(INCOMPLETE_MARKER)).await { + Err(e) if e.kind() != std::io::ErrorKind::NotFound => Err(e), + _ => Ok(()), + } + } + async fn destroy(mut self) -> anyhow::Result<()> { self.tasks.shutdown().await; self.db.destroy(); diff --git a/libsql-server/src/namespace/store.rs b/libsql-server/src/namespace/store.rs index 86e9438ccd..f76843caba 100644 --- a/libsql-server/src/namespace/store.rs +++ b/libsql-server/src/namespace/store.rs @@ -1,7 +1,7 @@ use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; -use async_lock::RwLock; +use async_lock::{RwLock, RwLockWriteGuardArc}; use chrono::NaiveDateTime; use futures::TryFutureExt; use moka::future::Cache; @@ -27,6 +27,86 @@ use super::{Namespace, ResetCb, ResetOp, ResolveNamespacePathFn, RestoreOption}; type NamespaceEntry = Arc>>; +/// A namespace name reserved for a creation (or fork) in progress. +/// +/// Holding a reservation holds the name's cache entry write-locked, so every other access to +/// the name — user requests, a second create, destroy, shutdown — waits until the outcome is +/// known. [`publish`](Self::publish) installs the finished namespace. Dropping the reservation +/// unpublished (error, or the creating request was cancelled) removes the namespace's in-memory +/// config and only then releases the lock, so waiters observe a namespace that does not exist, +/// never a half-built one. +/// +/// The metastore row is only written after setup succeeds (see `create`), so the row's presence +/// implies a complete namespace. See `docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md`. +struct Reservation { + guard: Option>>, + metadata: MetaStore, + name: NamespaceName, +} + +impl Reservation { + /// Fails with `NamespaceAlreadyExist` if the namespace is loaded or known to the metastore. + /// + /// Must be called before `MetaStore::handle(name)`: `handle` inserts the name into the + /// in-memory config map that `exists` reads, so the name becomes visible only while it is + /// already locked here. + async fn acquire(store: &NamespaceStore, name: NamespaceName) -> crate::Result { + let entry = store + .inner + .store + .get_with(name.clone(), async { Default::default() }) + .await; + let guard = entry.write_arc().await; + if guard.is_some() || store.inner.metadata.exists(&name).await { + return Err(Error::NamespaceAlreadyExist(name.to_string())); + } + Ok(Self { + guard: Some(guard), + metadata: store.inner.metadata.clone(), + name, + }) + } + + /// Install the finished namespace and release the lock. + fn publish(mut self, ns: Namespace) { + self.guard + .take() + .expect("a reservation is published at most once") + .replace(ns); + } +} + +impl Drop for Reservation { + fn drop(&mut self) { + let Some(guard) = self.guard.take() else { + return; + }; + tracing::warn!( + namespace = %self.name, + "namespace creation did not complete; discarding it" + ); + metrics::increment_counter!("libsql_server_namespace_create_aborted"); + let metadata = self.metadata.clone(); + let name = self.name.clone(); + let cleanup = move || { + // Removes the in-memory config and the metastore row, if any was written. + if let Err(e) = metadata.remove(name.clone()) { + tracing::error!(namespace = %name, "failed to discard namespace config: {e}"); + } + // Release the lock only once the name no longer exists. + drop(guard); + }; + // `MetaStore::remove` takes blocking locks, and we may be running inside a future that + // is being dropped: hand the work to the blocking pool instead of blocking in place. + match tokio::runtime::Handle::try_current() { + Ok(rt) => { + rt.spawn_blocking(cleanup); + } + Err(_) => cleanup(), + } + } +} + /// Stores and manage a set of namespaces. pub struct NamespaceStore { pub inner: Arc, @@ -225,15 +305,7 @@ impl NamespaceStore { return Err(crate::error::Error::NamespaceDoesntExist(from.to_string())); } - let to_entry = self - .inner - .store - .get_with(to.clone(), async { Default::default() }) - .await; - let mut to_lock = to_entry.write().await; - if to_lock.is_some() { - return Err(crate::error::Error::NamespaceAlreadyExist(to.to_string())); - } + let reservation = Reservation::acquire(self, to.clone()).await?; // FIXME: we could potentially delete the namespace while trying to fork it if !self.inner.metadata.exists(&from).await { @@ -249,33 +321,8 @@ impl NamespaceStore { return Err(crate::error::Error::NamespaceDoesntExist(from.to_string())); }; - struct Bomb { - store: MetaStore, - ns: NamespaceName, - should_delete: bool, - } - - impl Drop for Bomb { - fn drop(&mut self) { - if self.should_delete { - // we need to block in place because the inner connection may blocking, or - // unsing tokio's blocking methods (bottomless), which would cause a panic. - if let Err(e) = - tokio::task::block_in_place(|| self.store.remove(self.ns.clone())) - { - tracing::error!("failed to clean handle while forking: {e}"); - } - } - } - } - - let mut bomb = Bomb { - store: self.inner.metadata.clone(), - ns: to.clone(), - should_delete: true, - }; - let handle = self.inner.metadata.handle(to.clone()).await; + // In memory only: the fork reads it; nothing is persisted until it succeeded. handle .store_and_maybe_flush(Some(to_config.into()), false) .await?; @@ -291,10 +338,8 @@ impl NamespaceStore { ) .await?; - to_lock.replace(to_ns); handle.flush().await?; - // defuse - bomb.should_delete = false; + reservation.publish(to_ns); Ok(()) } @@ -394,6 +439,11 @@ impl NamespaceStore { let ns = self .make_namespace(namespace, db_config, restore_option) .await?; + // A namespace opened from the metastore is complete by definition; this only removes + // a marker left by a crash between persisting the config and removing it. + if let Err(e) = ns.mark_complete().await { + tracing::warn!(namespace = %namespace, "could not remove creation marker: {e}"); + } Ok(Some(ns)) }; @@ -430,24 +480,46 @@ impl NamespaceStore { .await; }; - // With namespaces disabled, the default namespace can be auto-created, - // otherwise it's an error. + // With namespaces disabled, the default namespace always exists and `create` is an + // idempotent config upsert of it. Creating *from a dump* is not: the dump must never be + // skipped because the namespace happens to be loaded already, so that case takes the + // reserving path below like any other creation. // FIXME: move the default namespace check out of this function. - if self.inner.allow_lazy_creation || namespace == NamespaceName::default() { + if (self.inner.allow_lazy_creation || namespace == NamespaceName::default()) + && matches!(restore_option, RestoreOption::Latest) + { tracing::trace!("auto-creating the namespace"); - } else if self.inner.metadata.exists(&namespace).await { - return Err(Error::NamespaceAlreadyExist(namespace.to_string())); + let handle = self.inner.metadata.handle(namespace.clone()).await; + handle.store(Arc::new(db_config)).await?; + self.load_namespace(&namespace, handle, restore_option) + .await?; + return Ok(()); } - let db_config = Arc::new(db_config); + // Nothing below is visible or durable until the namespace is complete: + // - the reservation locks the name; dropping it unpublished undoes the in-memory config; + // - the configurator removes a brand-new directory if setup fails or is cancelled; + // - the metastore row is written last, so its presence implies a complete namespace. + let reservation = Reservation::acquire(self, namespace.clone()).await?; let handle = self.inner.metadata.handle(namespace.clone()).await; - tracing::debug!("storing db config"); - handle.store(db_config).await?; - tracing::debug!("completed storing db config, loading namespace"); - self.load_namespace(&namespace, handle, restore_option) + handle + .store_and_maybe_flush(Some(Arc::new(db_config)), false) + .await?; + self.get_configurator(&handle.get()) + .discard_incomplete(&namespace) .await?; + tracing::debug!("setting up namespace"); + let ns = self + .make_namespace(&namespace, handle.clone(), restore_option) + .await?; + tracing::debug!("storing db config"); + handle.flush().await?; + if let Err(e) = ns.mark_complete().await { + tracing::warn!(namespace = %namespace, "could not remove creation marker: {e}"); + } + reservation.publish(ns); - tracing::debug!("completed loading namespace"); + tracing::debug!("completed creating namespace"); Ok(()) } diff --git a/libsql-server/tests/namespaces/dumps.rs b/libsql-server/tests/namespaces/dumps.rs index 958d24daab..1c3ca7869e 100644 --- a/libsql-server/tests/namespaces/dumps.rs +++ b/libsql-server/tests/namespaces/dumps.rs @@ -17,11 +17,11 @@ use crate::common::http::{Client, Response}; use crate::common::net::{TurmoilAcceptor, TurmoilConnector}; use crate::namespaces::{make_primary, make_primary_with_db_config}; -const BUFFERED: Option<&str> = Some("buffered"); -const STREAMING: Option<&str> = Some("streaming"); +pub(super) const BUFFERED: Option<&str> = Some("buffered"); +pub(super) const STREAMING: Option<&str> = Some("streaming"); /// `POST /v1/namespaces/:ns/create` with `dump_url` and, optionally, `dump_importer`. -async fn create_from_dump( +pub(super) async fn create_from_dump( client: &Client, ns: &str, dump_url: &str, @@ -39,17 +39,17 @@ async fn create_from_dump( .await } -fn file_url(path: &Path) -> String { +pub(super) fn file_url(path: &Path) -> String { format!("file:{}", path.display()) } -fn sim() -> Sim<'static> { +pub(super) fn sim() -> Sim<'static> { Builder::new() .simulation_duration(Duration::from_secs(1000)) .build() } -async fn count_rows(ns: &str, table: &str) -> anyhow::Result { +pub(super) async fn count_rows(ns: &str, table: &str) -> anyhow::Result { let db = Database::open_remote_with_connector( &format!("http://{ns}.primary:8080"), "", @@ -69,7 +69,7 @@ fn make_dump_store(sim: &mut Sim, chunks: Vec>) { } /// Like [`make_dump_store`], pausing `pause` (simulated time) before each chunk. -fn make_dump_store_paced( +pub(super) fn make_dump_store_paced( sim: &mut Sim, chunks: Vec>, pause: Duration, @@ -124,7 +124,39 @@ fn make_dump_store_chunked(sim: &mut Sim, dump: &'static str, chunk_size: usize) ); } -const SIMPLE_DUMP: &str = r#" +/// Serve `dump` from `dump-store:8080` in 64-byte chunks, 100ms (simulated) apart: with +/// turmoil's up-to-100ms message latency this keeps a creation in flight for a couple of +/// seconds, long enough for other requests to be issued while it runs. +pub(super) fn make_slow_dump_store(sim: &mut Sim, dump: &'static str) { + make_dump_store_paced( + sim, + dump.as_bytes() + .chunks(64) + .map(|c| Ok(Bytes::copy_from_slice(c))) + .collect(), + Duration::from_millis(100), + ); +} + +/// Ten rows, importable by both importers (no "attach" anywhere), ~1.1 KB so that +/// [`make_slow_dump_store`] delivers it in ~18 chunks. +pub(super) const SLOW_DUMP: &str = r#"PRAGMA foreign_keys=OFF; +BEGIN TRANSACTION; +CREATE TABLE test (id INTEGER PRIMARY KEY, name TEXT, note TEXT); +INSERT INTO test VALUES (1, 'one', 'the first row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (2, 'two', 'the second row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (3, 'three', 'the third row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (4, 'four', 'the fourth row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (5, 'five', 'the fifth row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (6, 'six', 'the sixth row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (7, 'seven', 'the seventh row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (8, 'eight', 'the eighth row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (9, 'nine', 'the ninth row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (10, 'ten', 'the tenth row of a dump that takes a while to arrive'); +COMMIT; +"#; + +pub(super) const SIMPLE_DUMP: &str = r#" PRAGMA foreign_keys=OFF; BEGIN TRANSACTION; CREATE TABLE test (x); @@ -745,7 +777,7 @@ fn load_dump_with_nested_case_streaming() { /// Exercises framing across arbitrary HTTP body chunk boundaries, including inside multibyte /// characters, string literals with semicolons, comments and a trigger body. -const CHUNKY_DUMP: &str = r#"PRAGMA foreign_keys=OFF; +pub(super) const CHUNKY_DUMP: &str = r#"PRAGMA foreign_keys=OFF; BEGIN TRANSACTION; -- a comment; with a semicolon CREATE TABLE test (id INTEGER PRIMARY KEY, name TEXT); @@ -1437,34 +1469,26 @@ fn streaming_failure_under_backpressure() { sim.run().unwrap(); } -/// The admin request is abandoned while the import is still streaming in. The executor thread -/// must notice the closed channel, roll back and release the connection, so the server stays -/// healthy and the namespace can be used afterwards instead of being wedged by a half-open -/// transaction. -#[test] -fn streaming_cancelled_request_rolls_back() { +/// The admin request is abandoned while the dump is still being transferred. The namespace +/// must not survive in any form: the importer rolls back and releases its connection, the +/// configurator removes the directory and the store never persists the config, so the name is +/// immediately reusable (see `docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md`). +fn cancelled_create_leaves_no_trace_with(importer: Option<&'static str>) { let mut sim = sim(); let tmp = tempdir().unwrap(); + let tmp_path = tmp.path().to_path_buf(); make_primary(&mut sim, tmp.path().to_path_buf()); - // Deliver the dump slowly (64-byte chunks, 50ms apart) so the request is still in flight - // when it is abandoned. Few, large chunks also keep the number of unread segments below - // turmoil's simulated socket buffer once nobody reads them anymore. - make_dump_store_paced( - &mut sim, - CHUNKY_DUMP - .as_bytes() - .chunks(64) - .map(|c| Ok(Bytes::copy_from_slice(c))) - .collect(), - Duration::from_millis(50), - ); + // Deliver the dump slowly so the request is still in flight when it is abandoned. Few, + // large chunks also keep the number of unread segments below turmoil's simulated socket + // buffer once nobody reads them anymore. + make_slow_dump_store(&mut sim, SLOW_DUMP); sim.client("client", async move { let client = Client::new(); let aborted = tokio::time::timeout( - Duration::from_millis(100), - create_from_dump(&client, "foo", "http://dump-store:8080/", STREAMING), + Duration::from_millis(500), + create_from_dump(&client, "foo", "http://dump-store:8080/", importer), ) .await; assert!(aborted.is_err(), "the request should still be in flight"); @@ -1473,32 +1497,31 @@ fn streaming_cancelled_request_rolls_back() { tokio::time::sleep(Duration::from_secs(2)).await; std::thread::sleep(Duration::from_millis(300)); - // Depending on where the request was cut, the namespace was registered (pre-existing - // behavior: the metastore row survives a failed create) or not. Either way it must be - // usable now: no lingering write transaction, no partial data. + // The namespace does not exist, for the admin API, for users, and on disk. let resp = client - .post_raw("http://primary:9090/v1/namespaces/foo/create", json!({})) + .get("http://primary:9090/v1/namespaces/foo/config") .await?; - assert!( - resp.status() == StatusCode::OK || resp.status() == StatusCode::BAD_REQUEST, - "unexpected status {}", - resp.status() - ); - assert!(count_rows("foo", "test").await.is_err()); - let db = - Database::open_remote_with_connector("http://foo.primary:8080", "", TurmoilConnector)?; - let conn = db.connect()?; - conn.execute("create table after_cancel (x)", ()).await?; - conn.execute("insert into after_cancel values (1)", ()) - .await?; - assert_eq!(count_rows("foo", "after_cancel").await?, 1); + assert_eq!(resp.status(), StatusCode::NOT_FOUND); + let err = count_rows("foo", "test").await.unwrap_err().to_string(); + assert!(err.contains("doesn't exist"), "unexpected error: {err}"); + super::lifecycle::wait_until_gone(&tmp_path.join("dbs").join("foo")); - // and other namespaces are unaffected - let resp = create_from_dump(&client, "bar", "http://dump-store:8080/", STREAMING).await?; + // ... so the same name can be created again, from the same dump. + let resp = create_from_dump(&client, "foo", "http://dump-store:8080/", importer).await?; assert_eq!(resp.status(), StatusCode::OK); - assert_eq!(count_rows("bar", "test").await?, 3); + assert_eq!(count_rows("foo", "test").await?, 10); Ok(()) }); sim.run().unwrap(); } + +#[test] +fn cancelled_create_leaves_no_trace() { + cancelled_create_leaves_no_trace_with(BUFFERED); +} + +#[test] +fn cancelled_create_leaves_no_trace_streaming() { + cancelled_create_leaves_no_trace_with(STREAMING); +} diff --git a/libsql-server/tests/namespaces/lifecycle.rs b/libsql-server/tests/namespaces/lifecycle.rs new file mode 100644 index 0000000000..7f8659dad7 --- /dev/null +++ b/libsql-server/tests/namespaces/lifecycle.rs @@ -0,0 +1,366 @@ +//! Namespace creation is all-or-nothing: after a failed, cancelled or crashed +//! `POST /v1/namespaces/:ns/create` the namespace either does not exist (no config, no +//! directory, name reusable) or is complete. See `docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md`. + +use std::path::Path; +use std::time::Duration; + +use hyper::StatusCode; +use serde_json::json; +use tempfile::tempdir; + +use crate::common::http::Client; +use crate::namespaces::dumps::{ + count_rows, create_from_dump, file_url, make_slow_dump_store, sim, BUFFERED, SIMPLE_DUMP, + SLOW_DUMP, STREAMING, +}; +use crate::namespaces::{make_primary, make_single_namespace_primary}; + +const INCOMPLETE_MARKER: &str = ".incomplete"; + +const BROKEN_DUMP: &str = r#" + BEGIN TRANSACTION; + CREATE TABLE test (x); + INSERT INTO test VALUES(1); + THIS IS NOT SQL; + COMMIT;"#; + +/// [`SLOW_DUMP`] with a syntax error right before the end. +const SLOW_BROKEN_DUMP: &str = r#"PRAGMA foreign_keys=OFF; +BEGIN TRANSACTION; +CREATE TABLE test (id INTEGER PRIMARY KEY, name TEXT, note TEXT); +INSERT INTO test VALUES (1, 'one', 'the first row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (2, 'two', 'the second row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (3, 'three', 'the third row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (4, 'four', 'the fourth row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (5, 'five', 'the fifth row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (6, 'six', 'the sixth row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (7, 'seven', 'the seventh row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (8, 'eight', 'the eighth row of a dump that takes a while to arrive'); +INSERT INTO test VALUES (9, 'nine', 'the ninth row of a dump that takes a while to arrive'); +THIS IS NOT SQL; +COMMIT; +"#; + +/// Directory removal after an aborted creation is handed to the runtime; poll (real time, the +/// blocking pool is not driven by the simulated clock) until it is gone. +pub(super) fn wait_until_gone(path: &Path) { + for _ in 0..100 { + if !path.exists() { + return; + } + std::thread::sleep(Duration::from_millis(20)); + } + panic!("{} still exists", path.display()); +} + +async fn assert_absent(client: &Client, tmp: &Path, ns: &str) -> anyhow::Result<()> { + let resp = client + .get(&format!("http://primary:9090/v1/namespaces/{ns}/config")) + .await?; + assert_eq!( + resp.status(), + StatusCode::NOT_FOUND, + "{ns} should not exist" + ); + let err = count_rows(ns, "test").await.unwrap_err().to_string(); + assert!(err.contains("doesn't exist"), "unexpected error: {err}"); + wait_until_gone(&tmp.join("dbs").join(ns)); + Ok(()) +} + +/// A dump that fails to import leaves nothing behind, and the name can be reused right away. +/// Before this change the config survived and the next access lazily created an empty database. +fn failed_create_leaves_no_trace_with(importer: Option<&'static str>, shared_schema: bool) { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + let tmp_path = tmp.path().to_path_buf(); + std::fs::write(tmp_path.join("broken.sql"), BROKEN_DUMP).unwrap(); + std::fs::write(tmp_path.join("good.sql"), SIMPLE_DUMP).unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf()); + + sim.client("client", async move { + let client = Client::new(); + let create = |dump: &str| { + let mut body = json!({ "dump_url": file_url(&tmp_path.join(dump)), "shared_schema": shared_schema }); + if let Some(importer) = importer { + body["dump_importer"] = json!(importer); + } + client.post_raw("http://primary:9090/v1/namespaces/foo/create", body) + }; + + let resp = create("broken.sql").await?; + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); + assert_absent(&client, &tmp_path, "foo").await?; + + let resp = create("good.sql").await?; + assert_eq!( + resp.status(), + StatusCode::OK, + "{}", + resp.body_string().await.unwrap_or_default() + ); + assert_eq!(count_rows("foo", "test").await?, 1); + assert!(!tmp_path.join("dbs/foo").join(INCOMPLETE_MARKER).exists()); + Ok(()) + }); + + sim.run().unwrap(); +} + +#[test] +fn failed_create_leaves_no_trace() { + failed_create_leaves_no_trace_with(BUFFERED, false); +} + +#[test] +fn failed_create_leaves_no_trace_streaming() { + failed_create_leaves_no_trace_with(STREAMING, false); +} + +#[test] +fn failed_create_leaves_no_trace_shared_schema() { + failed_create_leaves_no_trace_with(STREAMING, true); +} + +/// Requests for a namespace that is being created wait for the outcome: a second `create` +/// for the same name is rejected once the first one succeeded, and a user request issued +/// mid-import sees the complete data. +#[test] +fn concurrent_requests_wait_for_creation() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf()); + make_slow_dump_store(&mut sim, SLOW_DUMP); + + sim.client("client", async move { + let client = Client::new(); + let slow = create_from_dump(&client, "foo", "http://dump-store:8080/", STREAMING); + // issued while the import is in progress (it takes ~2s of simulated time) + let second = async { + tokio::time::sleep(Duration::from_millis(500)).await; + client + .post_raw("http://primary:9090/v1/namespaces/foo/create", json!({})) + .await + }; + let user = async { + tokio::time::sleep(Duration::from_millis(500)).await; + count_rows("foo", "test").await + }; + let (slow, second, user) = tokio::join!(slow, second, user); + assert_eq!(slow?.status(), StatusCode::OK); + let second = second?; + assert_eq!(second.status(), StatusCode::BAD_REQUEST); + assert!(second.body_string().await?.contains("already exists")); + assert_eq!(user?, 10); + Ok(()) + }); + + sim.run().unwrap(); +} + +/// If the creation that others are waiting on fails, the waiters observe a namespace that does +/// not exist, and a subsequent `create` of the same name succeeds. +#[test] +fn concurrent_create_succeeds_after_failure() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + let tmp_path = tmp.path().to_path_buf(); + std::fs::write(tmp_path.join("good.sql"), SIMPLE_DUMP).unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf()); + make_slow_dump_store(&mut sim, SLOW_BROKEN_DUMP); + + sim.client("client", async move { + let client = Client::new(); + let slow = create_from_dump(&client, "foo", "http://dump-store:8080/", STREAMING); + // issued while the import is in progress (it takes ~2s of simulated time) + let user = async { + tokio::time::sleep(Duration::from_millis(500)).await; + count_rows("foo", "test").await + }; + let (slow, user) = tokio::join!(slow, user); + assert_eq!(slow?.status(), StatusCode::BAD_REQUEST); + let err = user.unwrap_err().to_string(); + assert!(err.contains("doesn't exist"), "unexpected error: {err}"); + assert_absent(&client, &tmp_path, "foo").await?; + + let resp = create_from_dump( + &client, + "foo", + &file_url(&tmp_path.join("good.sql")), + STREAMING, + ) + .await?; + assert_eq!(resp.status(), StatusCode::OK); + assert_eq!(count_rows("foo", "test").await?, 1); + Ok(()) + }); + + sim.run().unwrap(); +} + +/// A directory left by a creation the process died in the middle of (it still carries the +/// marker, there is no config for it) is discarded by the next creation of that name, and is +/// not adopted as a namespace when the metastore is rebuilt from the filesystem. +#[test] +fn crash_remnant_is_discarded() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + let tmp_path = tmp.path().to_path_buf(); + std::fs::write(tmp_path.join("good.sql"), SIMPLE_DUMP).unwrap(); + let remnant = tmp_path.join("dbs/foo"); + std::fs::create_dir_all(&remnant).unwrap(); + for f in [INCOMPLETE_MARKER, ".sentinel", "data", "wallog"] { + std::fs::write(remnant.join(f), b"junk").unwrap(); + } + make_primary(&mut sim, tmp.path().to_path_buf()); + + sim.client("client", async move { + let client = Client::new(); + // not recovered as a namespace at startup (the metastore was empty) + let resp = client + .get("http://primary:9090/v1/namespaces/foo/config") + .await?; + assert_eq!(resp.status(), StatusCode::NOT_FOUND); + + let resp = create_from_dump( + &client, + "foo", + &file_url(&tmp_path.join("good.sql")), + STREAMING, + ) + .await?; + assert_eq!( + resp.status(), + StatusCode::OK, + "{}", + resp.body_string().await.unwrap_or_default() + ); + assert_eq!(count_rows("foo", "test").await?, 1); + assert!(!remnant.join(INCOMPLETE_MARKER).exists()); + Ok(()) + }); + + sim.run().unwrap(); +} + +/// Directories without the marker are never touched by a creation, whatever is in them. +#[test] +fn unmarked_directory_is_left_alone() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + let tmp_path = tmp.path().to_path_buf(); + std::fs::write(tmp_path.join("good.sql"), SIMPLE_DUMP).unwrap(); + let dir = tmp_path.join("dbs/foo"); + std::fs::create_dir_all(&dir).unwrap(); + std::fs::write(dir.join("keep.txt"), b"precious").unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf()); + + sim.client("client", async move { + let client = Client::new(); + let _ = create_from_dump( + &client, + "foo", + &file_url(&tmp_path.join("good.sql")), + STREAMING, + ) + .await?; + assert_eq!(std::fs::read(dir.join("keep.txt"))?, b"precious"); + Ok(()) + }); + + sim.run().unwrap(); +} + +/// A marker on a namespace that *is* known to the metastore (the process died between +/// persisting the config and removing the marker) is harmless: `create` rejects the name, +/// the data is served, and the marker is dropped when the namespace is next loaded. +#[test] +fn marker_on_published_namespace_is_ignored() { + let tmp = tempdir().unwrap(); + let tmp_path = tmp.path().to_path_buf(); + std::fs::write(tmp_path.join("good.sql"), SIMPLE_DUMP).unwrap(); + let marker = tmp_path.join("dbs/foo").join(INCOMPLETE_MARKER); + + let mut first = sim(); + make_primary(&mut first, tmp.path().to_path_buf()); + first.client("client", { + let tmp_path = tmp_path.clone(); + let marker = marker.clone(); + async move { + let client = Client::new(); + let resp = create_from_dump( + &client, + "foo", + &file_url(&tmp_path.join("good.sql")), + STREAMING, + ) + .await?; + assert_eq!(resp.status(), StatusCode::OK); + assert!(!marker.exists()); + + // simulate a crash right after the config was persisted + std::fs::write(&marker, b"").unwrap(); + let resp = create_from_dump( + &client, + "foo", + &file_url(&tmp_path.join("good.sql")), + STREAMING, + ) + .await?; + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); + assert!(resp.body_string().await?.contains("already exists")); + assert_eq!(count_rows("foo", "test").await?, 1); + assert!(marker.exists(), "the marker is only removed on load"); + Ok(()) + } + }); + first.run().unwrap(); + + // restart: the namespace is loaded from the metastore and the stale marker goes away + let mut restarted = sim(); + make_primary(&mut restarted, tmp.path().to_path_buf()); + restarted.client("client", async move { + assert_eq!(count_rows("foo", "test").await?, 1); + wait_until_gone(&marker); + Ok(()) + }); + restarted.run().unwrap(); +} + +/// Single-namespace mode: `create` of the default namespace is a config upsert, but creating +/// it *from a dump* while it already exists must be rejected rather than silently skipping the +/// import (which is what happened before). +#[test] +fn single_namespace_mode_rejects_dump_into_existing_namespace() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + let tmp_path = tmp.path().to_path_buf(); + std::fs::write(tmp_path.join("good.sql"), SIMPLE_DUMP).unwrap(); + make_single_namespace_primary(&mut sim, tmp.path().to_path_buf()); + + sim.client("client", async move { + let client = Client::new(); + // config upsert of the (already created) default namespace still works + let resp = client + .post_raw( + "http://primary:9090/v1/namespaces/default/create", + json!({ "max_db_size": "1mb" }), + ) + .await?; + assert_eq!(resp.status(), StatusCode::OK); + + let resp = create_from_dump( + &client, + "default", + &file_url(&tmp_path.join("good.sql")), + STREAMING, + ) + .await?; + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); + assert!(resp.body_string().await?.contains("already exists")); + Ok(()) + }); + + sim.run().unwrap(); +} diff --git a/libsql-server/tests/namespaces/mod.rs b/libsql-server/tests/namespaces/mod.rs index c33cc229c3..fa03c9b765 100644 --- a/libsql-server/tests/namespaces/mod.rs +++ b/libsql-server/tests/namespaces/mod.rs @@ -1,6 +1,7 @@ #![allow(deprecated)] mod dumps; +mod lifecycle; mod meta; mod shared_schema; @@ -20,6 +21,16 @@ fn make_primary(sim: &mut Sim, path: PathBuf) { } fn make_primary_with_db_config(sim: &mut Sim, path: PathBuf, db_config: DbConfig) { + make_primary_with(sim, path, db_config, false); +} + +/// A primary running in single-namespace mode (`--disable-namespaces`): only the default +/// namespace exists, and it is created at startup. +fn make_single_namespace_primary(sim: &mut Sim, path: PathBuf) { + make_primary_with(sim, path, DbConfig::default(), true); +} + +fn make_primary_with(sim: &mut Sim, path: PathBuf, db_config: DbConfig, single_namespace: bool) { init_tracing(); sim.host("primary", move || { let path = path.clone(); @@ -41,8 +52,8 @@ fn make_primary_with_db_config(sim: &mut Sim, path: PathBuf, db_config: DbConfig acceptor: TurmoilAcceptor::bind(([0, 0, 0, 0], 4567)).await?, tls_config: None, }), - disable_namespaces: false, - disable_default_namespace: true, + disable_namespaces: single_namespace, + disable_default_namespace: !single_namespace, ..Default::default() }; From aef8a74fdc0274236f4165fafd1d78d18774a841 Mon Sep 17 00:00:00 2001 From: Tomasz Szymczyszyn Date: Fri, 9 Oct 2026 21:28:38 +0200 Subject: [PATCH 2/2] docs: atomic namespace creation design and admin API contract Adds docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md, states the all-or-nothing creation contract in ADMIN_API.md, and marks the lifecycle items of the streaming dump import design (section 13) as fixed. --- docs/ADMIN_API.md | 12 + docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md | 292 +++++++++++++++++++++++++ docs/STREAMING_DUMP_IMPORT_DESIGN.md | 8 +- 3 files changed, 308 insertions(+), 4 deletions(-) create mode 100644 docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md diff --git a/docs/ADMIN_API.md b/docs/ADMIN_API.md index 651f43ebae..d9e985afea 100644 --- a/docs/ADMIN_API.md +++ b/docs/ADMIN_API.md @@ -53,6 +53,18 @@ Both importers produce the same data. Known differences: - invalid UTF-8 or NUL bytes yield `400` (buffered: `500`), and statements that return rows are executed with their rows discarded (buffered: `500`). +Creation is all-or-nothing. The namespace becomes visible, and its configuration is persisted, +only after it has been fully set up (including the dump import). If the request fails, is +cancelled by the client, or the server dies while it is in flight, the namespace either does not +exist (its name can be created again right away) or is complete — never partially imported. +Requests addressed to a namespace while it is being created wait for the outcome. For callers: + +- `2xx`: the namespace is complete; +- any other outcome, including a lost connection: retry the same request; a + `400 Namespace already exists` on the retry means the earlier attempt did complete. + +See `docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md`. + ```HTTP DELETE /v1/namespaces/:namespace ``` diff --git a/docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md b/docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md new file mode 100644 index 0000000000..73ac223d81 --- /dev/null +++ b/docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md @@ -0,0 +1,292 @@ +# Atomic namespace creation for `libsql-server` + +Stacked on `tszymczyszyn/streaming-dump-importer` (PR #51). Fixes the pre-existing lifecycle gap described in PR #51's design doc §13 and the review walkthrough: a failed, cancelled or crashed `POST /v1/namespaces/:ns/create` could leave a *ghost* namespace (metastore row without data), a directory without a row, or both. Sections marked *as built* note where the implementation differs from the sketches. + +--- + +## 1. Problem + +`NamespaceStore::create` persists the namespace config **first** and only then runs the configurator's `setup` (directory, logger, connections, dump import). Three resources are mutated without a common commit point: + +| Resource | Written at | Cleaned on error | Cleaned on cancellation | Cleaned after crash | +|---|---|---|---|---| +| metastore row (`namespace_configs`) | start of `create` | **no** | **no** | n/a (persisted) | +| `dbs//` + `.sentinel` + SQLite/WAL/log files | during `setup` | primary only, if the dir was fresh | **no** (future dropped, `Err` branch never runs) | **no** | +| cache entry (`NamespaceEntry`) | end of `load_namespace` | n/a | n/a | n/a | + +Consequences: a failed create reserves the name forever (`NamespaceAlreadyExist` on retry) while user traffic lazily materialises an **empty** database; a cancelled create can leave any prefix of the resources; a client cannot tell from a non-200 whether the namespace now exists. `SchemaConfigurator::setup` has no cleanup at all. `fork` has a drop guard for the row but uses `block_in_place`, which panics on a current-thread runtime. + +Both dump importers run inside this sequence; the problem is independent of PR #51. + +## 2. Goals + +- **Contract:** after `create` returns a non-2xx, or the connection is lost, the namespace is either *absent* (no row, no data, name reusable) or *complete* (row persisted, data complete, serving). Never partial. +- Cancellation-safe (dropping the request future cleans up) and crash-tolerant (a process crash mid-import never blocks a retry). +- No lock-step changes to clients, replication, the metastore schema or the `DatabaseConfig` protobuf. +- Small: reuse the pattern `fork` already uses; one mechanism for `create` and `fork`. + +Non-goals: a server-side import timeout; idempotency keys / async create jobs; bottomless remote cleanup for aborted imports (remote objects under a never-published db_id are inert junk, as today); the moka eviction hazard shared with `reset`/`fork` (§8). + +## 3. Design in one paragraph + +Make the **metastore write the commit point** and make everything before it reversible. `create` reserves the name by write-locking its cache entry, puts the config in the metastore's *in-memory* map only (so `setup` can read it), runs `setup`, and only on success **flushes the row** and installs the namespace. A `Reservation` drop guard undoes the in-memory config if the future ends any other way — error *or* cancellation — and releases the lock only after that, so waiters never observe a half-built namespace. The filesystem gets the same treatment one layer down: `setup` wraps a brand-new directory in a `FreshDir` guard that removes it on drop. For crashes, a brand-new directory carries an `.incomplete` marker until the row is flushed; the next `create` of that name discards a marked directory. + +``` +create(name, restore, config) +│ +├─ Reservation::acquire(name) lock cache entry; reject if loaded or row exists +├─ handle = metadata.handle(name) ─┐ in-memory only; nothing on disk +├─ handle.store_and_maybe_flush(cfg, false) ┘ +├─ configurator.discard_incomplete(name) remove `dbs/` if it carries `.incomplete` +├─ ns = make_namespace(...) setup: FreshDir::begin → mkdir + marker … import … keep() +├─ handle.flush() ◄── COMMIT POINT: row persisted +├─ ns.mark_complete() remove `.incomplete` (best effort) +└─ reservation.publish(ns) entry := Some(ns); lock released + any Err / drop above ⇒ Reservation::drop ⇒ metadata.remove(name), then release lock + any Err / drop inside setup ⇒ FreshDir::drop ⇒ remove_dir_all(dbs/) +``` + +## 4. Components + +### 4.1 `Reservation` (store.rs) — replaces `fork`'s ad-hoc `Bomb` + +```rust +/// A namespace name reserved for a creation or fork in progress. +/// +/// Holding it holds the name's cache entry write-locked: user requests, a second create, +/// destroy and shutdown all wait until the outcome is known. `publish` installs the finished +/// namespace. Dropping it unpublished (error or cancelled request) removes the in-memory +/// config and only then releases the lock, so waiters see "doesn't exist", never a partial one. +struct Reservation { + guard: Option>>, + metadata: MetaStore, + name: NamespaceName, +} + +impl Reservation { + async fn acquire(store: &NamespaceStore, name: NamespaceName) -> crate::Result { + let entry = store.inner.store.get_with(name.clone(), async { Default::default() }).await; + let guard = entry.write_arc().await; + if guard.is_some() || store.inner.metadata.exists(&name).await { + return Err(Error::NamespaceAlreadyExist(name.to_string())); + } + Ok(Self { guard: Some(guard), metadata: store.inner.metadata.clone(), name }) + } + + fn publish(mut self, ns: Namespace) { + self.guard.take().expect("published once").replace(ns); + } +} + +impl Drop for Reservation { + fn drop(&mut self) { + let Some(guard) = self.guard.take() else { return }; + let (metadata, name) = (self.metadata.clone(), self.name.clone()); + tracing::warn!(namespace = %name, "namespace creation did not complete; discarding it"); + let cleanup = move || { + if let Err(e) = metadata.remove(name.clone()) { + tracing::error!(namespace = %name, "failed to discard namespace config: {e}"); + } + drop(guard); // release only once the name no longer exists + }; + match tokio::runtime::Handle::try_current() { + Ok(rt) => { rt.spawn_blocking(cleanup); } + Err(_) => cleanup(), + } + } +} +``` + +Why `write_arc`: the guard must outlive the `create` future (it moves into the cleanup task). Requires `async-lock = "3"` in `libsql-server/Cargo.toml` (3.4.0 is already in the lockfile transitively; `proxy.rs`'s `RwLockUpgradableReadGuard::upgrade` is unchanged in 3.x). + +`MetaStore::remove` takes blocking locks, hence `spawn_blocking` (no `block_in_place`, so it also works on turmoil's current-thread runtime). It deletes the row if one exists and the in-memory entry; both are what we want in every abort path, including "flush succeeded but publish didn't" (then the row is deleted again and the complete directory becomes a marked orphan, see §6). + +### 4.2 `create` (store.rs) + +```rust +pub async fn create(&self, namespace, restore_option, db_config) -> crate::Result<()> { + shutdown check — unchanged + shared_schema_name ⇒ fork(…) — unchanged (fork now uses Reservation, §4.6) + + // Single-namespace mode / default namespace: `create` is an idempotent config upsert of a + // namespace that always exists. Keep today's path for it, but only for `Latest`: creating + // *from a dump* must never silently skip the import because the namespace is already loaded + // (today it does), so dump creates always take the reserving path below. + if (self.inner.allow_lazy_creation || namespace == NamespaceName::default()) + && matches!(restore_option, RestoreOption::Latest) + { + let handle = self.inner.metadata.handle(namespace.clone()).await; + handle.store(Arc::new(db_config)).await?; + self.load_namespace(&namespace, handle, restore_option).await?; + return Ok(()); + } + + let reservation = Reservation::acquire(self, namespace.clone()).await?; + let handle = self.inner.metadata.handle(namespace.clone()).await; + handle.store_and_maybe_flush(Some(Arc::new(db_config)), false).await?; + self.get_configurator(&handle.get()).discard_incomplete(&namespace).await?; + let ns = self.make_namespace(&namespace, handle.clone(), restore_option).await?; + handle.flush().await?; // commit point + if let Err(e) = ns.mark_complete().await { // cosmetic from here on + tracing::warn!(namespace = %namespace, "could not remove creation marker: {e}"); + } + reservation.publish(ns); + Ok(()) +} +``` + +Ordering that matters: `acquire` before `handle()` (`handle()` inserts into the in-memory map, which is what `exists()` reads — so the name becomes visible only while it is already locked); `flush` before `publish` (durable before visible); `publish` consumes the reservation (no double cleanup). + +### 4.3 `FreshDir` (configurator/helpers.rs) — cancellation-safe directory + +```rust +/// Directory of a namespace that does not exist yet. Removed again if dropped before `keep`: +/// on error, or when the request creating the namespace is cancelled. +pub(super) struct FreshDir { path: Option> } + +impl FreshDir { + /// `None` if a namespace already lives at `db_path`. + pub(super) async fn begin(db_path: &Arc) -> crate::Result> { + if db_path.try_exists()? { return Ok(None); } + tokio::fs::create_dir_all(db_path).await?; + tokio::fs::File::create(db_path.join(INCOMPLETE_MARKER)).await?; + Ok(Some(Self { path: Some(db_path.clone()) })) + } + pub(super) fn keep(mut self) { self.path.take(); } +} + +impl Drop for FreshDir { + fn drop(&mut self) { + let Some(path) = self.path.take() else { return }; + let rm = async move { + match tokio::fs::remove_dir_all(&*path).await { + Ok(()) | Err(e) if e.kind() == ErrorKind::NotFound => {} + Err(e) => tracing::error!(path = %path.display(), "failed to remove unfinished namespace directory: {e}"), + } + }; + match tokio::runtime::Handle::try_current() { + Ok(rt) => { rt.spawn(rm); } + Err(_) => { let _ = std::fs::remove_dir_all(&*path); } + } + } +} +``` + +Used by `PrimaryConfigurator::setup` and `SchemaConfigurator::setup` (which gains cleanup it never had): + +```rust +let db_path: Arc = self.base.base_path.join("dbs").join(name.as_str()).into(); +let fresh = FreshDir::begin(&db_path).await?; +let ns = self.try_new_primary(…).await?; // on Err or drop, `fresh` removes the directory +if let Some(fresh) = fresh { fresh.keep(); } +Ok(ns) +``` + +`try_new_primary`'s own `create_dir_all` becomes redundant for the fresh case and stays harmless for existing directories. The `remove_dir_all`-on-error branch in `setup` is replaced by the guard (drop order: the `try_new_primary` future and its `JoinSet` are dropped before `fresh`, so background tasks are aborted before the directory goes). + +### 4.4 `.incomplete` marker — crash tolerance + +- `pub(crate) const INCOMPLETE_MARKER: &str = ".incomplete";` in `namespace/mod.rs`. +- Created by `FreshDir::begin` for every brand-new directory. +- Removed by `Namespace::mark_complete(&self)` (`remove_file(self.path.join(MARKER))`, `NotFound` ok), called (a) by `create` after `flush`, (b) by `load_namespace`'s init after `make_namespace` — a namespace opened from the metastore is complete by definition, so lazily loaded namespaces (and single-namespace-mode auto-creation) never keep a marker. +- Consulted in exactly one place: `ConfigureNamespace::discard_incomplete(name)` (new trait method, default no-op; primary and schema share `helpers::discard_incomplete(base, name)` = `if dbs//.incomplete exists → remove_dir_all`). `create` calls it after `acquire`, i.e. only for a name with **no row and no loaded namespace** — the only situation in which a marked directory is unambiguously garbage. Lazy loads never see it, so a stale marker on a published namespace (crash between `flush` and `mark_complete`) is removed, never acted upon. +- `MetaStoreInner::maybe_recover_from_fs` (metastore empty ⇒ adopt `dbs/*`) skips directories carrying the marker, so a crash during the very first creation on a fresh server does not get adopted as a namespace at the next start. + +A directory **without** a marker and without a row keeps today's semantics (`Dump` ⇒ `LoadDumpExistingDb`; `Latest` ⇒ adopted). Such directories can only come from pre-marker servers or manual intervention; being destructive there is not worth it. + +### 4.5 `Namespace::mark_complete` (namespace/mod.rs) + +```rust +pub(crate) async fn mark_complete(&self) -> std::io::Result<()> { + match tokio::fs::remove_file(self.path.join(INCOMPLETE_MARKER)).await { + Err(e) if e.kind() != ErrorKind::NotFound => Err(e), + _ => Ok(()), + } +} +``` + +### 4.6 `fork` (store.rs) + +Replace the `Bomb` + manual `to_lock` with `Reservation::acquire(to)`; order becomes `configurator.fork(…)` → `handle.flush()` → `reservation.publish(ns)` (today it publishes *then* flushes, so a flush failure leaves a running namespace without a row). Net code removal; same semantics otherwise. `ForkTask`'s own temp-dir + rename stays as is. + +## 5. Flow for a dump import (streaming or buffered) + +1. Admin handler resolves `dump_url` (errors here touch nothing). +2. `create` → `Reservation::acquire` (409-class error if the name is loaded or has a row). +3. In-memory config; `discard_incomplete` (removes a crashed previous attempt, if any). +4. `setup`: `FreshDir::begin` (mkdir + marker) → logger, WAL, connection maker → `load_dump` → `keep()`. +5. `flush` — the namespace now exists durably. +6. `mark_complete`, `publish` — the namespace is now visible; waiters proceed. + +## 6. Failure matrix + +| Interrupted at | Error (returned) | Cancellation (future dropped) | Crash (process dies) | +|---|---|---|---| +| 1 | nothing to undo | nothing to undo | nothing | +| 2–3 (before `setup`) | `Reservation` removes in-memory config | same | nothing persisted; nothing on disk | +| 4, inside `setup` (incl. import) | `FreshDir` removes dir; `Reservation` removes config; executor rolls back SQLite | same (importer executor is uncancelled: closes channel → ROLLBACK → connection released; dir removal races with it benignly on POSIX, see §9) | marked orphan dir, no row ⇒ invisible; next `create` discards it | +| between `setup` and `flush` | config removed; **complete dir, marked** ⇒ next `create` discards and redoes | same | same | +| `flush` fails | row never written (single SQLite statement); as above | — | — | +| between `flush` and `publish` | `Reservation` deletes the row again; marked complete dir ⇒ redone next time | same | row + marker + complete data ⇒ namespace **exists and serves**; marker removed on first lazy load | +| after `publish` | — | response may be lost ⇒ client retries ⇒ `NamespaceAlreadyExist` ⇒ by the contract, complete | same | + +Client rule: **2xx ⇒ complete; anything else ⇒ retry; `NamespaceAlreadyExist` on retry ⇒ complete.** No inspection of tables needed. + +## 7. Concurrency + +- **Second `create` of the same name:** waits on the write lock; then `AlreadyExist` (first succeeded) or proceeds (first failed). Previously: `AlreadyExist` even after the first failed. +- **User request during creation (`with`)**: `exists()` is true (in-memory config), `try_get_with` finds the locked entry, `read()` waits; then sees the namespace or `None ⇒ NamespaceDoesntExist`. Same waiting behaviour as during `fork`/`reset` today. +- **`destroy` during creation:** `metadata.remove` succeeds (in-memory only), then waits on the lock; when creation finishes, `flush` is a no-op (config gone), `publish` installs, `destroy` takes and destroys it. Net: created then destroyed, in that order. Acceptable; no worse than today. +- **Store shutdown during creation:** `shutdown` waits on the entry lock, then shuts the finished namespace down normally. +- **`handle()` side effect:** `MetaStore::handle` inserts a default config into the in-memory map (pre-existing). Keeping `acquire` strictly before `handle()` is what makes the reservation airtight; documented in code. + +## 8. Known limitation inherited from `fork`/`reset` + +The reservation is the cache entry's lock. If moka evicts that entry during a long import (capacity pressure; TTI is 24 h), a concurrent `with()` re-inserts a fresh entry and lazily opens the directory being written. This exists today for `reset` and `fork` and is unchanged here. The structural fix is an in-flight-operations map outside the cache; separate change. + +## 9. Interaction with the streaming importer + +On cancellation the reader future is dropped, the executor thread sees the closed channel, rolls back and drops its connection — concurrently with `FreshDir`'s `remove_dir_all`. On POSIX, unlinking files that SQLite still has open is harmless (writes go to unlinked inodes; `-shm`/log unlinks return `ENOENT`, which SQLite tolerates). If removal fails because a background task recreated a file in the window, the directory stays **marked** and is discarded by the next `create` — the marker makes cleanup eventually consistent, so no ordering between the guard and the executor is required. + +## 10. Observability + +- `warn!` on every abort path (`Reservation`: "namespace creation did not complete; discarding it"; `FreshDir`: "… removing its directory"; `discard_incomplete`: "discarding namespace directory left by an interrupted creation"), with the namespace / path. +- Counters `libsql_server_namespace_create_aborted` and `libsql_server_namespace_create_discarded_incomplete` (*as built:* no `stage` label — the logs carry the detail). +- The contract of §6 is stated in `docs/ADMIN_API.md` under namespace creation. + +## 11. Tests (turmoil, `tests/namespaces/lifecycle.rs` and `dumps.rs`) + +*As built.* 9 of the 11 integration tests fail on the unfixed code; the other two pin behaviour that must hold both before and after. + +1. `failed_create_leaves_no_trace{,_streaming,_shared_schema}`: invalid dump ⇒ 400; `GET /v1/namespaces/foo/config` ⇒ 404 (was 200); hrana ⇒ "doesn't exist" (was "no such table"); `dbs/foo` gone; `create` again with a valid dump ⇒ 200, data present, no marker. +2. `cancelled_create_leaves_no_trace{,_streaming}` (replaces `streaming_cancelled_request_rolls_back`, which tolerated "OK or BAD_REQUEST" on retry): request abandoned mid-transfer ⇒ config 404, hrana "doesn't exist", directory gone, retry from the same dump ⇒ 200 with all rows. +3. `concurrent_requests_wait_for_creation`: a second `create` and a user query issued during a ~2 s import wait; afterwards the second gets `already exists` and the query sees all rows. `concurrent_create_succeeds_after_failure`: the import fails ⇒ the waiting query gets "doesn't exist" and a new `create` succeeds. +4. `crash_remnant_is_discarded`: `dbs/foo/{.incomplete,.sentinel,data,wallog}` exists before the server starts ⇒ not adopted by the metastore (config 404), `create foo` from a dump ⇒ 200. `unmarked_directory_is_left_alone`: a marker-less `dbs/foo/keep.txt` survives a `create`. `marker_on_published_namespace_is_ignored`: marker added to a published namespace ⇒ `create` ⇒ `already exists`, data served; after a restart the marker is gone. +5. `single_namespace_mode_rejects_dump_into_existing_namespace`: config upsert of `default` still works; `create default` from a dump ⇒ `already exists` (was: 200, dump silently skipped). +6. Unit (`configurator::helpers::test`): `FreshDir` removed on drop and when its enclosing future is dropped, kept on `keep()`; `discard_incomplete` removes only marked directories. +7. `fork_namespace`, `shared_schema::*` and all dump tests unchanged. + +Slow-transfer tests use `make_slow_dump_store` (64-byte chunks, 100 ms apart): with turmoil's random 0–100 ms message latency, concurrent requests are issued 500 ms in and the import lasts ~2 s. + +## 12. Scope + +| File | Change | +|---|---| +| `libsql-server/Cargo.toml` | `async-lock = "3"` | +| `namespace/store.rs` | `Reservation`; `create` reserving path; `fork` on `Reservation`; `load_namespace` init calls `mark_complete` | +| `namespace/mod.rs` | `INCOMPLETE_MARKER`, `Namespace::mark_complete` | +| `namespace/configurator/mod.rs` | `ConfigureNamespace::discard_incomplete` (default no-op) | +| `namespace/configurator/helpers.rs` | `FreshDir`, `discard_incomplete` helper | +| `namespace/configurator/{primary,schema}.rs` | use `FreshDir`; implement `discard_incomplete`; drop ad-hoc cleanup | +| `namespace/meta_store.rs` | skip marked dirs in `maybe_recover_from_fs` | +| `docs/ADMIN_API.md`, PR #51 design doc §13 | contract; mark the lifecycle item as fixed | +| tests | §11 | + +Branch `tszymczyszyn/atomic-namespace-create` off `tszymczyszyn/streaming-dump-importer`; the PR targets the importer branch until #51 merges, then retargets `v0.9.30-shopify-patches`. + +## 13. Alternatives considered + +- **Row first with a lifecycle state (`creating`/`ready`).** Needs a metastore schema or protobuf change that replicates to replicas via the config handshake, plus a startup sweep that knows configurator paths. More moving parts for the same guarantee. +- **Import into a temp directory and rename.** The logger, stats and connection maker capture the path at open; publishing would require close → rename → reopen, and bottomless would see a restore decision on reopen. `fork` can do it because it only writes a raw `data` file before opening. +- **Dedup via moka `try_get_with` instead of an explicit lock.** When the initialising future is dropped, moka lets another waiter run *its* init — a lazy `with()` would then create an empty namespace. The explicit reservation makes waiters observe the outcome instead. +- **Making `with()` fail fast ("namespace is being created") instead of waiting.** Cleaner UX for minute-long imports, but changes behaviour for `reset`/`fork` too; can be layered on later without touching this design. diff --git a/docs/STREAMING_DUMP_IMPORT_DESIGN.md b/docs/STREAMING_DUMP_IMPORT_DESIGN.md index 525f947071..f8358fffe8 100644 --- a/docs/STREAMING_DUMP_IMPORT_DESIGN.md +++ b/docs/STREAMING_DUMP_IMPORT_DESIGN.md @@ -585,12 +585,12 @@ Metrics (`metrics` 0.21 macros with labels): ## 13. Pre-existing issues this design does not fix (state them in the PR) -1. `NamespaceStore::create` stores the namespace config in the metastore **before** loading. On import failure the directory is removed (`PrimaryConfigurator::setup`) but the metastore row stays, so a later request to the namespace lazily creates an **empty** database. Existing tests rely on this (`select … from test` errors because the table is missing, not because the namespace is missing). -2. If the Admin HTTP request is cancelled mid-import, `setup`'s cleanup does not run (the future is dropped). The streaming executor still rolls back SQLite state, but the directory and `.sentinel` remain. -3. No server-side timeout exists for dump imports; the Admin HTTP request stays open for the whole import. Clients must not time out, or must tolerate (2). +1. ~~`NamespaceStore::create` stores the namespace config in the metastore **before** loading. On import failure the directory is removed (`PrimaryConfigurator::setup`) but the metastore row stays, so a later request to the namespace lazily creates an **empty** database.~~ **Fixed** by the stacked change described in `docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md`: the config is persisted only after setup succeeds, and a failed or cancelled creation leaves nothing behind. +2. ~~If the Admin HTTP request is cancelled mid-import, `setup`'s cleanup does not run (the future is dropped). The streaming executor still rolls back SQLite state, but the directory and `.sentinel` remain.~~ **Fixed**, same change (drop guards for the directory and the reservation of the name). +3. No server-side timeout exists for dump imports; the Admin HTTP request stays open for the whole import. Clients must not time out; if they do, the creation is discarded and can be retried (see the creation contract in `docs/ADMIN_API.md`). 4. Bottomless, if enabled, has its own frame buffering/backpressure outside this design. -These are tracked by the DB Mover quarantine/fence work (Retail #35846/#35848). +(3) and (4) are tracked by the DB Mover quarantine/fence work (Retail #35846/#35848). ---