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/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). --- 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() };