Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 5 additions & 20 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

12 changes: 12 additions & 0 deletions docs/ADMIN_API.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
```
Expand Down
292 changes: 292 additions & 0 deletions docs/ATOMIC_NAMESPACE_CREATE_DESIGN.md

Large diffs are not rendered by default.

8 changes: 4 additions & 4 deletions docs/STREAMING_DUMP_IMPORT_DESIGN.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).

---

Expand Down
2 changes: 1 addition & 1 deletion libsql-server/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
187 changes: 186 additions & 1 deletion libsql-server/src/namespace/configurator/helpers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,14 +20,101 @@ 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;
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<Arc<Path>>,
}

impl FreshDir {
/// Returns `None` when a namespace directory already exists at `db_path`.
pub(super) async fn begin(db_path: &Arc<Path>) -> crate::Result<Option<Self>> {
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/<namespace>` 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,
Expand Down Expand Up @@ -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<Path> = 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<Path> = 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());
}
}
12 changes: 12 additions & 0 deletions libsql-server/src/namespace/configurator/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,18 @@ pub trait ConfigureNamespace {
bottomless_db_id_init: NamespaceBottomlessDbIdInit,
) -> Pin<Box<dyn Future<Output = crate::Result<()>> + 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<Box<dyn Future<Output = crate::Result<()>> + Send + 'a>> {
let _ = namespace;
Box::pin(async { Ok(()) })
}

fn fork<'a>(
&'a self,
from_ns: &'a Namespace,
Expand Down
34 changes: 16 additions & 18 deletions libsql-server/src/namespace/configurator/primary.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -129,35 +129,33 @@ impl ConfigureNamespace for PrimaryConfigurator {
) -> Pin<Box<dyn Future<Output = crate::Result<Namespace>> + Send + 'a>> {
Box::pin(async move {
let db_path: Arc<Path> = 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<Box<dyn Future<Output = crate::Result<()>> + Send + 'a>> {
Box::pin(discard_incomplete(&self.base, namespace))
}

fn cleanup<'a>(
&'a self,
namespace: &'a NamespaceName,
Expand Down
Loading
Loading