From dd47a069ab6536bae8fc4b77c8eeb381eeb52c3a Mon Sep 17 00:00:00 2001 From: River Date: Tue, 29 Sep 2026 17:20:13 +0000 Subject: [PATCH 1/3] libsql-server: source read fence for SQL with read leases SetSourceReadFence now closes read admission in memory, persists SOURCE_READ_DRAINING, and waits for every read lease already held before it persists SOURCE_READ_FENCED. Leases are taken only after checking the gate under the controller's lease lock, so once admission is closed the set of leases can only shrink. Each SQL program holds a lease for as long as it runs (including a Hrana cursor still producing rows), as does describe and each admin shell query. ATTACH of a namespace is a read of that namespace: the attaching connection keeps a lease on it for every later program until it is detached. A connection idle inside a transaction holds no lease; its next program is refused with MIGRATION_READ_FENCED and rolled back. /beta/listen is refused where reads are denied and ends when reads are fenced. At the deadline (--namespace-fence-default-read-drain-ms, 30s) running programs are cancelled through the connection's progress-handler cancel flag and report the read fence; the drain still waits for the actual releases and answers DRAINING if they do not come, and a replay of the same command resumes it. ClearSourceReadFence reopens reads with writes still fenced. Dump and replication stream leases follow. Co-authored-by: Tomasz Szymczyszyn --- libsql-server/src/admin_shell.rs | 114 +++- libsql-server/src/config.rs | 3 + .../src/connection/connection_core.rs | 67 +- libsql-server/src/connection/legacy.rs | 2 +- libsql-server/src/connection/program.rs | 9 +- libsql-server/src/http/user/listen.rs | 40 +- libsql-server/src/main.rs | 9 + .../src/namespace/fence/controller.rs | 332 +++++++++- libsql-server/src/namespace/fence/drain.rs | 42 +- libsql-server/src/namespace/fence/hooks.rs | 4 + libsql-server/src/namespace/fence/mod.rs | 7 +- libsql-server/src/namespace/fence/read.rs | 588 ++++++++++++++++++ libsql-server/src/namespace/meta_store.rs | 15 + libsql-server/src/namespace/mod.rs | 13 +- libsql-server/src/namespace/store.rs | 26 +- 15 files changed, 1230 insertions(+), 41 deletions(-) create mode 100644 libsql-server/src/namespace/fence/read.rs diff --git a/libsql-server/src/admin_shell.rs b/libsql-server/src/admin_shell.rs index e97b272d72..84f11e7fe8 100644 --- a/libsql-server/src/admin_shell.rs +++ b/libsql-server/src/admin_shell.rs @@ -10,6 +10,8 @@ use tonic::metadata::{AsciiMetadataValue, BinaryMetadataValue}; use crate::connection::Connection as _; use crate::database::Connection; +use crate::namespace::fence::controller::{FenceController, LeaseKind}; +use crate::namespace::fence::state::OperationClass; use crate::namespace::{NamespaceName, NamespaceStore}; use self::rpc::admin_shell_service_server::{AdminShellService, AdminShellServiceServer}; @@ -41,17 +43,20 @@ impl AdminShell { queries: impl Stream>, ) -> anyhow::Result>> { let namespace = NamespaceName::from_bytes(ns).unwrap(); - let connection_maker = self + let (connection_maker, fence) = self .namespace_store - .with(namespace, |ns| ns.db.connection_maker()) + .with(namespace, |ns| { + (ns.db.connection_maker(), ns.fence().clone()) + }) .await?; let connection = connection_maker.create().await?; - Ok(run_shell(connection, queries)) + Ok(run_shell(connection, fence, queries)) } } fn run_shell( conn: Connection, + fence: std::sync::Arc, queries: impl Stream>, ) -> impl Stream> { async_stream::stream! { @@ -59,9 +64,7 @@ fn run_shell( while let Some(q) = queries.next().await { let Ok(q) = q else { break }; let res = tokio::task::block_in_place(|| { - conn.with_raw(move |conn| { - run_one(conn, q.query) - }) + conn.with_raw(|conn| run_admitted(&fence, conn, q.query)) }); yield res @@ -69,6 +72,39 @@ fn run_shell( } } +/// Run one shell query as a read of the namespace (`docs/NAMESPACE_FENCE.md` section 9): the +/// shell runs raw SQL, so the read gate is checked, and a read lease held, per query. Writes are +/// refused by the WAL gate like any other connection's. +fn run_admitted( + fence: &std::sync::Arc, + conn: &mut rusqlite::Connection, + q: String, +) -> Result { + let interrupt = conn.get_interrupt_handle(); + let lease = + match fence.acquire_read_lease(OperationClass::NormalRead, LeaseKind::Sql, move || { + interrupt.interrupt() + }) { + Ok(lease) => lease, + Err(e) => { + return Ok(rpc::Response { + resp: Some(Resp::Error(rpc::Error { + error: e.to_string(), + })), + }) + } + }; + let res = run_one(conn, q); + if lease.cancelled_by_fence() { + return Ok(rpc::Response { + resp: Some(Resp::Error(rpc::Error { + error: "the query was cancelled by the namespace read fence".into(), + })), + }); + } + res +} + fn run_one(conn: &mut rusqlite::Connection, q: String) -> Result { match try_run_one(conn, q) { Ok(resp) => Ok(resp), @@ -242,3 +278,69 @@ impl Display for RowsFormatter { Ok(()) } } + +#[cfg(test)] +mod fence_tests { + use uuid::Uuid; + + use super::*; + use crate::namespace::fence::command::{FenceCommand, FenceRequest}; + use crate::namespace::fence::drain::tests::{fence_outcome, raw, Source, LONG, OP}; + use crate::namespace::fence::outcome::FenceOutcome; + + fn error(resp: &rpc::Response) -> &str { + match &resp.resp { + Some(Resp::Error(e)) => &e.error, + other => panic!("expected an error response, got {other:?}"), + } + } + + /// The shell runs raw SQL, so it checks the read gate per query: reads are served while + /// the source is only write-fenced (writes are refused by the WAL gate), and refused once + /// its reads are fenced. + #[tokio::test(flavor = "multi_thread")] + async fn admin_shell_read_denied() { + let s = Source::new().await; + raw(&s.conn().await, "insert into t values (1)") + .await + .unwrap(); + let acquired = s.execute(s.acquire(OP, 1, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&acquired), FenceOutcome::Applied); + + let conn = s.conn().await; + let rows = conn + .with_raw(|c| run_admitted(&s.fence, c, "select count(*) from t".into())) + .unwrap(); + assert!(matches!(rows.resp, Some(Resp::Rows(ref r)) if r.rows.len() == 1)); + let write = conn + .with_raw(|c| run_admitted(&s.fence, c, "insert into t values (2)".into())) + .unwrap(); + assert!(error(&write).contains("authoriz"), "{}", error(&write)); + + let gate = s.fence.gate(); + let fenced = s + .execute(FenceRequest { + namespace: "ns".into(), + operation_id: OP, + command_id: Uuid::from_u128(2), + expected_state: gate.state(), + expected_revision: gate.revision(), + command: FenceCommand::SetSourceReadFence { + drain_policy: Some(LONG), + }, + }) + .await + .unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + + let refused = conn + .with_raw(|c| run_admitted(&s.fence, c, "select count(*) from t".into())) + .unwrap(); + assert!( + error(&refused).starts_with("MIGRATION_READ_FENCED"), + "{}", + error(&refused) + ); + assert_eq!(s.fence.read_lease_counts().total(), 0); + } +} diff --git a/libsql-server/src/config.rs b/libsql-server/src/config.rs index 4457232d03..48dee38dfa 100644 --- a/libsql-server/src/config.rs +++ b/libsql-server/src/config.rs @@ -196,6 +196,9 @@ pub struct MetaStoreConfig { /// How long `AcquireSourceWriteFence` waits for active writers when the request names no /// drain policy. `None` is the default of 30 seconds. pub namespace_fence_default_write_drain: Option, + /// How long `SetSourceReadFence` waits for running reads and streams before it cancels + /// them, when the request names no drain policy. `None` is the default of 30 seconds. + pub namespace_fence_default_read_drain: Option, } #[derive(Debug, Clone)] diff --git a/libsql-server/src/connection/connection_core.rs b/libsql-server/src/connection/connection_core.rs index 212e5c2b3d..821428ecc9 100644 --- a/libsql-server/src/connection/connection_core.rs +++ b/libsql-server/src/connection/connection_core.rs @@ -11,7 +11,8 @@ use crate::connection::legacy::open_conn_active_checkpoint; use crate::error::Error; use crate::metrics::{PROGRAM_EXEC_COUNT, QUERY_CANCELED, VACUUM_COUNT, WAL_CHECKPOINT_COUNT}; use crate::namespace::broadcasters::BroadcasterHandle; -use crate::namespace::fence::controller::FenceConnState; +use crate::namespace::fence::controller::{FenceConnState, LeaseKind}; +use crate::namespace::fence::outcome::{FenceError, FenceOutcome}; use crate::namespace::fence::state::OperationClass; use crate::namespace::meta_store::MetaStoreHandle; use crate::namespace::ResolveNamespacePathFn; @@ -102,6 +103,8 @@ impl CoreConnection { ); let canceled = Arc::new(AtomicBool::new(false)); + // The read drain cancels a running program at its deadline through the same flag. + fence.set_cancel_flag(canceled.clone()); conn.progress_handler(100, { let canceled = canceled.clone(); @@ -218,6 +221,22 @@ impl CoreConnection { // The program is admitted under the gate's current write generation; a write // transaction it opens must start under the same one (section 8.1). fence.begin_program(); + // ...and admitted for reading, holding a read lease for as long as it runs, including + // while a cursor is still producing rows (section 9). A connection left idle in a + // transaction held no lease; its next program is refused here and its transaction is + // rolled back. + let read_lease = { + let lock = this.lock(); + match fence.begin_read_program(|| attached_schemas(lock.raw())) { + Ok(lease) => lease, + Err(e) => { + if !lock.conn.is_autocommit() { + lock.rollback(); + } + return Err(Error::NamespaceFence(e)); + } + } + }; builder.init(&this.lock().builder_config)?; let mut vm = Vm::new( @@ -272,16 +291,40 @@ impl CoreConnection { vm.step(&conn.raw())?; } + if read_lease.cancelled_by_fence() { + // The read drain reached its deadline and interrupted the program: report the fence, + // not the interruption, and leave no transaction behind. + let lock = this.lock(); + if !lock.conn.is_autocommit() { + lock.rollback(); + } + return Err(Error::NamespaceFence(FenceError::new( + FenceOutcome::MigrationReadFenced, + "the program was cancelled by the namespace read fence", + ))); + } + { let lock = this.lock(); let is_autocommit = lock.conn.is_autocommit(); let current_fno = (lock.get_current_frame_no)(); vm.builder().finish(current_fno, is_autocommit)?; } + drop(read_lease); Ok(vm.into_builder()) } + pub(super) fn describe_admitted(&self, sql: &str) -> crate::Result { + // Describing prepares the statement, which reads the schema: it is a read. + let _lease = self.fence.controller().acquire_read_lease( + self.fence.read_class(), + LeaseKind::Sql, + || (), + )?; + self.describe(sql) + } + fn rollback(&self) { if let Err(e) = self.conn.execute("ROLLBACK", ()) { tracing::error!("failed to rollback: {e}"); @@ -425,6 +468,28 @@ impl CoreConnection { } } +/// The schema aliases attached on `conn` (other than `main` and `temp`). +/// `None` when they cannot be listed. +fn attached_schemas(conn: &libsql_sys::Connection) -> Option> { + let mut aliases = Vec::new(); + let result = conn.prepare("PRAGMA database_list").and_then(|mut stmt| { + let mut rows = stmt.query(())?; + while let Some(row) = rows.next()? { + let name: String = row.get(1)?; + if name != "main" && name != "temp" { + aliases.push(name); + } + } + Ok(()) + }); + if let Err(e) = result { + // Keep every recorded attachment: over-leasing is safe, missing one is not. + tracing::warn!("could not list attached schemas: {e}"); + return None; + } + Some(aliases) +} + #[cfg(test)] mod test { use itertools::Itertools; diff --git a/libsql-server/src/connection/legacy.rs b/libsql-server/src/connection/legacy.rs index 71763e199e..efc176cae8 100644 --- a/libsql-server/src/connection/legacy.rs +++ b/libsql-server/src/connection/legacy.rs @@ -438,7 +438,7 @@ where DESCRIBE_COUNT.increment(1); check_describe_auth(ctx)?; let conn = self.inner.clone(); - let res = tokio::task::spawn_blocking(move || conn.lock().describe(&sql)) + let res = tokio::task::spawn_blocking(move || conn.lock().describe_admitted(&sql)) .await .unwrap(); diff --git a/libsql-server/src/connection/program.rs b/libsql-server/src/connection/program.rs index 8cafe681e7..f285343137 100644 --- a/libsql-server/src/connection/program.rs +++ b/libsql-server/src/connection/program.rs @@ -205,10 +205,15 @@ where let attached = attached.strip_prefix('"').unwrap_or(attached); let attached = attached.strip_suffix('"').unwrap_or(attached); let attached = NamespaceName::from_string(attached.into())?; - let path = (self.resolve_attach_path)(&attached)?; + let target = (self.resolve_attach_path)(&attached)?; + // Attaching a namespace reads it: the attachment is admitted by that namespace's fence, + // and the connection holds a read lease on it while its programs run (section 9). + if let Some(fence) = &self.fence { + fence.attach(attached_alias.trim_matches('"'), target.fence)?; + } let query = format!( "ATTACH DATABASE 'file:{}?mode=ro' AS \"{attached_alias}\"", - path.join("data").display() + target.path.join("data").display() ); tracing::trace!("ATTACH rewritten to: {query}"); Ok(query) diff --git a/libsql-server/src/http/user/listen.rs b/libsql-server/src/http/user/listen.rs index 04f95fdfc2..3a63ad32a4 100644 --- a/libsql-server/src/http/user/listen.rs +++ b/libsql-server/src/http/user/listen.rs @@ -1,6 +1,9 @@ use crate::broadcaster::BroadcastMsg; use crate::error::Error; use crate::metrics::{LISTEN_EVENTS_DROPPED, LISTEN_EVENTS_SENT}; +use crate::namespace::fence::controller::{FenceController, GateSnapshot}; +use crate::namespace::fence::outcome::FenceError; +use crate::namespace::fence::state::OperationClass; use crate::{ auth::Authenticated, namespace::{NamespaceName, NamespaceStore}, @@ -18,7 +21,9 @@ use serde::{Deserialize, Serialize}; use std::boxed::Box; use std::convert::Infallible; use std::pin::Pin; +use std::sync::Arc; use std::time::Duration; +use tokio::sync::watch; use tokio_stream::wrappers::errors::BroadcastStreamRecvError; use super::db_factory::namespace_from_headers; @@ -93,11 +98,19 @@ pub(super) async fn handle_listen( return Ok(ListenResponse::Redirect(Redirect::temporary(&url))); } + // Change notifications are a read of the namespace: refused where reads are denied, and + // ended when the namespace's reads are fenced (`docs/NAMESPACE_FENCE.md` section 9). + let fence = state.namespaces.fence_gate(&namespace).await?; + if let Some(fence) = &fence { + fence.permits(OperationClass::NormalRead)?; + } + let stream = sse_stream( state.namespaces.clone(), namespace, query.table.clone(), query.action.clone(), + fence, ) .await; @@ -115,9 +128,10 @@ async fn sse_stream( namespace: NamespaceName, table: String, actions: Option>, + fence: Option>, ) -> SseStream { Box::pin( - listen_stream(store, namespace, table, actions) + listen_stream(store, namespace, table, actions, fence) .await .map(|result| { Ok(match result { @@ -136,12 +150,19 @@ async fn listen_stream( namespace: NamespaceName, table: String, actions: Option>, + fence: Option>, ) -> impl Stream> { async_stream::try_stream! { let _sub = Subscription::new(store.clone(), namespace.clone(), table.clone()); let mut stream = store.subscribe(namespace.clone(), table.clone()); + let mut gate = fence.as_ref().map(|fence| fence.subscribe()); - while let Some(item) = stream.next().await { + loop { + let item = tokio::select! { + item = stream.next() => Ok(item), + denied = read_denied(&mut gate) => Err(Error::NamespaceFence(denied)), + }; + let Some(item) = item? else { break }; match item { Ok(msg) => if filter_actions(&msg, &actions) { LISTEN_EVENTS_SENT.increment(1); @@ -156,6 +177,21 @@ async fn listen_stream( } } +/// Resolves once the gate denies normal reads, with the denial; never without a gate. +async fn read_denied(gate: &mut Option>) -> FenceError { + let Some(gate) = gate else { + return std::future::pending().await; + }; + loop { + if let Err(e) = gate.borrow_and_update().permits(OperationClass::NormalRead) { + return e; + } + if gate.changed().await.is_err() { + return std::future::pending().await; + } + } +} + fn filter_actions(msg: &BroadcastMsg, actions: &Option>) -> bool { actions.as_ref().map_or(true, |actions| { actions.iter().any(|action| { diff --git a/libsql-server/src/main.rs b/libsql-server/src/main.rs index a1ab51c4c4..f8f56e7494 100644 --- a/libsql-server/src/main.rs +++ b/libsql-server/src/main.rs @@ -274,6 +274,12 @@ struct Cli { #[clap(long, env = "SQLD_NAMESPACE_FENCE_DEFAULT_WRITE_DRAIN_MS")] namespace_fence_default_write_drain_ms: Option, + /// How long, in milliseconds, setting a namespace read fence waits for running reads and + /// streams when the request names no drain policy, before it cancels them. Defaults to 30 + /// seconds. + #[clap(long, env = "SQLD_NAMESPACE_FENCE_DEFAULT_READ_DRAIN_MS")] + namespace_fence_default_read_drain_ms: Option, + /// Shutdown timeout duration in seconds, defaults to 30 seconds. #[clap(long, env = "SQLD_SHUTDOWN_TIMEOUT")] shutdown_timeout: Option, @@ -673,6 +679,9 @@ fn make_meta_store_config(config: &Cli) -> anyhow::Result { namespace_fence_default_write_drain: config .namespace_fence_default_write_drain_ms .map(Duration::from_millis), + namespace_fence_default_read_drain: config + .namespace_fence_default_read_drain_ms + .map(Duration::from_millis), }) } diff --git a/libsql-server/src/namespace/fence/controller.rs b/libsql-server/src/namespace/fence/controller.rs index 7e718e5516..ff946554c6 100644 --- a/libsql-server/src/namespace/fence/controller.rs +++ b/libsql-server/src/namespace/fence/controller.rs @@ -10,11 +10,12 @@ //! namespace cache, so evicting and reloading a namespace hands the reloaded namespace the //! same controller. -use std::sync::atomic::{AtomicU64, Ordering}; +use std::collections::HashMap; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::Arc; use parking_lot::Mutex; -use tokio::sync::{watch, OwnedMutexGuard}; +use tokio::sync::{watch, Notify, OwnedMutexGuard}; use uuid::Uuid; use crate::connection::connection_manager::{ConnectionManager, WeakConnectionManager}; @@ -52,6 +53,11 @@ pub struct GateSnapshot { /// (section 8.3, step 2): write admission is closed on top of whatever `fence` allows. /// Never persisted. pub installing: Option, + /// The in-memory read-closing gate of a `SetSourceReadFence` that is being persisted + /// (section 9, step 2): normal reads and streams are refused on top of whatever `fence` + /// allows. Never persisted, and cleared by every publication of a commit. It does not move + /// the write generation: write admission is already closed wherever a read fence can be set. + pub closing_reads: Option, } impl GateSnapshot { @@ -61,6 +67,7 @@ impl GateSnapshot { write_generation: 0, indeterminate: None, installing: None, + closing_reads: None, } } @@ -115,6 +122,17 @@ impl GateSnapshot { )); } } + if let Some((operation_id, command_id)) = self.closing_reads { + if matches!(class, OperationClass::NormalRead | OperationClass::Stream) { + return Err(FenceError::new( + FenceOutcome::MigrationReadFenced, + format!( + "{class:?} is not permitted: fence command {command_id} of operation \ + {operation_id} is closing read admission" + ), + )); + } + } Ok(()) } @@ -177,6 +195,84 @@ pub(crate) struct LiveWriteDrain { pub(crate) current_frame_no: GetCurrentFrameNo, } +/// What a read lease covers (section 9). +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub enum LeaseKind { + /// A running SQL program (including a Hrana cursor producing rows), or an ATTACH of the + /// namespace by a program running on another namespace. + Sql, + /// A `/dump` stream. + Dump, + /// A replication `log_entries` or `snapshot` stream. + Replication, +} + +/// The number of read leases held on a namespace, by kind. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct ReadLeaseCounts { + pub sql: usize, + pub dump: usize, + pub replication: usize, +} + +impl ReadLeaseCounts { + pub fn total(&self) -> usize { + self.sql + self.dump + self.replication + } +} + +struct LeaseEntry { + kind: LeaseKind, + cancel: Box, + cancelled: Arc, +} + +#[derive(Default)] +struct ReadLeaseSet { + next_id: u64, + live: HashMap, +} + +/// Read work admitted by the gate and counted by the read drain until it is dropped +/// ([`FenceController::acquire_read_lease`]). +pub struct ReadLease { + controller: Arc, + id: u64, + kind: LeaseKind, + cancelled: Arc, +} + +impl ReadLease { + pub fn kind(&self) -> LeaseKind { + self.kind + } + + pub fn controller(&self) -> &Arc { + &self.controller + } + + /// Whether the read drain cancelled this lease's work at its deadline. + pub fn cancelled_by_fence(&self) -> bool { + self.cancelled.load(Ordering::Acquire) + } +} + +impl std::fmt::Debug for ReadLease { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ReadLease") + .field("namespace", &self.controller.namespace) + .field("id", &self.id) + .field("kind", &self.kind) + .finish() + } +} + +impl Drop for ReadLease { + fn drop(&mut self) { + self.controller.release_read_lease(self.id); + } +} + /// The fence controller of one namespace. pub struct FenceController { namespace: NamespaceName, @@ -187,6 +283,10 @@ pub struct FenceController { /// What the write drain needs from each of the namespace's primary connection makers /// (section 8.3). write_drains: Mutex>, + /// The read leases held on the namespace (section 9). + read_leases: Mutex, + /// Notified whenever a read lease is released. + read_released: Notify, #[cfg(test)] hooks: FenceTestHooks, } @@ -210,6 +310,8 @@ impl FenceController { gate, write_queues: Mutex::new(Vec::new()), write_drains: Mutex::new(Vec::new()), + read_leases: Mutex::new(ReadLeaseSet::default()), + read_released: Notify::new(), #[cfg(test)] hooks: FenceTestHooks::default(), }) @@ -281,6 +383,82 @@ impl FenceController { live } + /// Admit read work of `class` and hold a read lease of `kind` for it (section 9). The gate + /// is checked under the lease lock, so a read fence that closed admission before this call + /// is seen here, and one that closes after it counts this lease and waits for it. `cancel` + /// is how the read drain stops the work at its deadline: it must make the work end and + /// drop the lease, never wait for it to end. The lease is released when dropped. + pub fn acquire_read_lease( + self: &Arc, + class: OperationClass, + kind: LeaseKind, + cancel: impl Fn() + Send + Sync + 'static, + ) -> Result { + let cancelled = Arc::new(AtomicBool::new(false)); + let id = { + let mut leases = self.read_leases.lock(); + self.permits(class)?; + let id = leases.next_id; + leases.next_id += 1; + leases.live.insert( + id, + LeaseEntry { + kind, + cancel: Box::new(cancel), + cancelled: cancelled.clone(), + }, + ); + id + }; + Ok(ReadLease { + controller: self.clone(), + id, + kind, + cancelled, + }) + } + + /// The read leases currently held, by kind. + pub fn read_lease_counts(&self) -> ReadLeaseCounts { + let leases = self.read_leases.lock(); + let mut counts = ReadLeaseCounts::default(); + for entry in leases.live.values() { + match entry.kind { + LeaseKind::Sql => counts.sql += 1, + LeaseKind::Dump => counts.dump += 1, + LeaseKind::Replication => counts.replication += 1, + } + } + counts + } + + /// Cancel every read lease held now (the read drain's deadline). Each lease's work is asked + /// to stop once; the leases stay counted until they are actually released. Returns how many + /// were asked. + pub(crate) fn cancel_read_leases(&self) -> usize { + let leases = self.read_leases.lock(); + let mut asked = 0; + for entry in leases.live.values() { + if !entry.cancelled.swap(true, Ordering::AcqRel) { + (entry.cancel)(); + asked += 1; + } + } + asked + } + + /// Notified on every read-lease release. Enable the notification before checking + /// [`read_lease_counts`](Self::read_lease_counts), so a release in between is not missed. + pub(crate) fn read_released(&self) -> &Notify { + &self.read_released + } + + fn release_read_lease(&self, id: u64) { + let removed = self.read_leases.lock().live.remove(&id); + drop(removed); + self.read_released.notify_waiters(); + } + /// Take the namespace's transition lock. Every fence command on the namespace runs while /// holding it, from its first check to its response. pub async fn begin_transition(self: &Arc) -> Transition { @@ -330,7 +508,8 @@ impl FenceController { /// Publish a new gate. `fence: None` keeps the published fence. The write generation moves /// whenever the state, the owning operation, the indeterminate flag or the installing gate - /// changes. + /// changes. Every publication removes the read-closing gate: the commit that follows it + /// either persists the read fence or proves that nothing changed. fn publish( &self, fence: Option, @@ -348,6 +527,7 @@ impl FenceController { gate.fence = fence; gate.indeterminate = indeterminate; gate.installing = installing; + gate.closing_reads = None; if changed { gate.write_generation += 1; generation_changed = true; @@ -406,6 +586,35 @@ impl Transition { } } + /// Publish the in-memory read-closing gate for `SetSourceReadFence` `key` (section 9, + /// step 2): new SQL programs, dumps, replication calls and ATTACHes of the namespace are + /// refused with `MIGRATION_READ_FENCED`. A read lease is only ever taken after checking the + /// gate under the lease lock, so every lease taken once this returns was refused, and the + /// drain only has to wait for the leases already held. It is replaced by whatever the + /// command's commit publishes, or removed with + /// [`reopen_read_admission`](Self::reopen_read_admission) when the command is proven not to + /// have committed. + pub fn close_read_admission(&mut self, key: CommandKey) { + self.controller + .gate + .send_modify(|gate| gate.closing_reads = Some(key)); + tracing::debug!( + namespace = %self.controller.namespace, + operation_id = %key.0, + command_id = %key.1, + "closed namespace read admission" + ); + } + + /// Remove the read-closing gate of a command that was proven not to have committed. + pub fn reopen_read_admission(&mut self) { + self.controller.gate.send_if_modified(|gate| { + let was_closing = gate.closing_reads.is_some(); + gate.closing_reads = None; + was_closing + }); + } + /// Commit `request` in the metastore and publish the result. pub async fn apply( &mut self, @@ -535,6 +744,34 @@ fn pending_indeterminate((operation_id, command_id): CommandKey) -> FenceError { .with_detail(FenceDetail::IndeterminateCommit) } +/// The read leases of one running program ([`FenceConnState::begin_read_program`]), released +/// when dropped. +#[derive(Debug)] +pub struct ProgramReadLease { + conn: Arc, + lease: ReadLease, +} + +impl ProgramReadLease { + /// Whether the read drain cancelled the program at its deadline. + pub fn cancelled_by_fence(&self) -> bool { + self.lease.cancelled_by_fence() + || self + .conn + .attached_leases + .lock() + .iter() + .any(ReadLease::cancelled_by_fence) + } +} + +impl Drop for ProgramReadLease { + fn drop(&mut self) { + let attached = std::mem::take(&mut *self.conn.attached_leases.lock()); + drop(attached); + } +} + /// The fence state of one connection, shared by its WAL wrapper and its `CoreConnection` /// (section 7.4). /// @@ -554,6 +791,14 @@ pub struct FenceConnState { txn_generation: AtomicU64, /// The typed outcome of the last refusal at the WAL. denial: Mutex>, + /// The connection's cancel flag (its progress handler interrupts the running statement + /// while it is set), through which the read drain cancels a program at its deadline. + cancel: std::sync::OnceLock>, + /// The namespaces attached on this connection, by schema alias. An attachment outlives the + /// program that made it, so every later program takes a read lease on each of them too. + attached: Mutex)>>, + /// The read leases the running program holds on attached namespaces. + attached_leases: Mutex>, } impl FenceConnState { @@ -565,9 +810,90 @@ impl FenceConnState { program_generation: AtomicU64::new(generation), txn_generation: AtomicU64::new(generation), denial: Mutex::new(None), + cancel: std::sync::OnceLock::new(), + attached: Mutex::new(Vec::new()), + attached_leases: Mutex::new(Vec::new()), }) } + /// The flag that cancels the statement running on this connection. Set once, by the + /// connection that owns this state. + pub fn set_cancel_flag(&self, flag: Arc) { + let _ = self.cancel.set(flag); + } + + /// The class this connection's reads are admitted as: a normal connection reads as + /// `NormalRead`; a capability connection reads under its capability. + pub fn read_class(&self) -> OperationClass { + match self.class { + OperationClass::NormalWrite => OperationClass::NormalRead, + class => class, + } + } + + fn lease_canceller(&self) -> impl Fn() + Send + Sync + 'static { + let flag = self.cancel.get().cloned(); + move || { + if let Some(flag) = &flag { + flag.store(true, Ordering::Relaxed); + } + } + } + + /// Admit a program for reading and hold its read leases until the returned guard is + /// dropped (section 9): one on this connection's namespace and one on each namespace + /// attached on the connection. `still_attached` lists the aliases attached on the + /// connection now (it is only asked when [`attach`](Self::attach) recorded one); `None` + /// keeps every recorded attachment. + pub fn begin_read_program( + self: &Arc, + still_attached: impl FnOnce() -> Option>, + ) -> Result { + let lease = self.controller.acquire_read_lease( + self.read_class(), + LeaseKind::Sql, + self.lease_canceller(), + )?; + let guard = ProgramReadLease { + conn: self.clone(), + lease, + }; + let attached = { + let mut attached = self.attached.lock(); + if !attached.is_empty() { + if let Some(live) = still_attached() { + attached.retain(|(alias, _)| live.iter().any(|a| a == alias)); + } + } + attached.clone() + }; + for (_, controller) in attached { + let lease = controller.acquire_read_lease( + OperationClass::NormalRead, + LeaseKind::Sql, + self.lease_canceller(), + )?; + self.attached_leases.lock().push(lease); + } + Ok(guard) + } + + /// The running program is attaching `controller`'s namespace as `alias`: admit it as a + /// normal read of that namespace, hold a lease on it for the rest of the program, and + /// remember the attachment for the connection's later programs. + pub fn attach(&self, alias: &str, controller: Arc) -> Result<(), FenceError> { + let lease = controller.acquire_read_lease( + OperationClass::NormalRead, + LeaseKind::Sql, + self.lease_canceller(), + )?; + self.attached_leases.lock().push(lease); + let mut attached = self.attached.lock(); + attached.retain(|(a, _)| a != alias); + attached.push((alias.to_owned(), controller)); + Ok(()) + } + pub fn controller(&self) -> &Arc { &self.controller } diff --git a/libsql-server/src/namespace/fence/drain.rs b/libsql-server/src/namespace/fence/drain.rs index 6165784f60..5cc70563b8 100644 --- a/libsql-server/src/namespace/fence/drain.rs +++ b/libsql-server/src/namespace/fence/drain.rs @@ -56,6 +56,9 @@ impl FenceController { FenceCommand::AcquireSourceWriteFence { .. } => { acquire_source_write_fence(&mut transition, &meta, request, ctx).await } + FenceCommand::SetSourceReadFence { .. } => { + super::read::set_source_read_fence(&mut transition, &meta, request, ctx).await + } _ => transition.apply(&meta, request, ctx).await, } }) @@ -277,14 +280,14 @@ fn capture_boundary( Ok(FrozenBoundary { log_id, frame_no }) } -fn now_ms() -> i64 { +pub(super) fn now_ms() -> i64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map_or(0, |d| i64::try_from(d.as_millis()).unwrap_or(i64::MAX)) } #[cfg(test)] -mod tests { +pub(crate) mod tests { use std::sync::Arc; use std::time::Duration; @@ -306,27 +309,27 @@ mod tests { use crate::namespace::RestoreOption; use crate::replication::primary::logger::ReplicationLogger; - const OP: Uuid = Uuid::from_u128(0xa); + pub(crate) const OP: Uuid = Uuid::from_u128(0xa); const OTHER_OP: Uuid = Uuid::from_u128(0xb); /// Long enough that no test ever reaches it: a drain must never finish because of time. - const LONG: DrainPolicy = DrainPolicy { + pub(crate) const LONG: DrainPolicy = DrainPolicy { deadline_ms: 600_000, on_deadline: OnDeadline::Fail, }; - const PROMPT: Duration = Duration::from_secs(30); + pub(crate) const PROMPT: Duration = Duration::from_secs(30); /// A primary namespace `ns` with a table `t`, served by a real `NamespaceStore`, so that the /// drain goes through the connection manager and replication logger the configurator /// registered. - struct Source { + pub(crate) struct Source { _dir: TempDir, - store: NamespaceStore, - fence: Arc, + pub(crate) store: NamespaceStore, + pub(crate) fence: Arc, logger: Arc, } impl Source { - async fn new() -> Self { + pub(crate) async fn new() -> Self { let dir = tempdir().unwrap(); let store = open_store(dir.path()).await; store @@ -362,7 +365,7 @@ mod tests { this } - async fn conn(&self) -> Arc { + pub(crate) async fn conn(&self) -> Arc { let maker = self .store .with("ns".into(), |ns| ns.db.connection_maker()) @@ -380,7 +383,12 @@ mod tests { *self.logger.new_frame_notifier.borrow() } - fn acquire(&self, op: Uuid, command_id: u128, policy: DrainPolicy) -> FenceRequest { + pub(crate) fn acquire( + &self, + op: Uuid, + command_id: u128, + policy: DrainPolicy, + ) -> FenceRequest { FenceRequest { namespace: "ns".into(), operation_id: op, @@ -395,7 +403,7 @@ mod tests { } /// Run `request` through the store, on a task of its own. - fn execute( + pub(crate) fn execute( &self, request: FenceRequest, ) -> tokio::task::JoinHandle> { @@ -414,7 +422,7 @@ mod tests { } /// Wait until the published gate is in `state`. - async fn until_state(&self, state: FenceState) { + pub(crate) async fn until_state(&self, state: FenceState) { let mut rx = self.fence.subscribe(); tokio::time::timeout(PROMPT, rx.wait_for(|g| g.state() == state)) .await @@ -422,7 +430,7 @@ mod tests { .unwrap(); } - async fn count(&self) -> i64 { + pub(crate) async fn count(&self) -> i64 { let conn = self.conn().await; tokio::task::spawn_blocking(move || { conn.with_raw(|c| c.query_row("select count(*) from t", (), |r| r.get(0))) @@ -435,14 +443,14 @@ mod tests { /// Run `sql` as one raw program on `conn`, off the async runtime (it can block on the write /// slot). - async fn raw(conn: &Arc, sql: &'static str) -> rusqlite::Result<()> { + pub(crate) async fn raw(conn: &Arc, sql: &'static str) -> rusqlite::Result<()> { let conn = conn.clone(); tokio::task::spawn_blocking(move || conn.with_raw(|c| c.execute_batch(sql))) .await .unwrap() } - fn assert_fenced(result: rusqlite::Result<()>) { + pub(crate) fn assert_fenced(result: rusqlite::Result<()>) { match result { Err(rusqlite::Error::SqliteFailure(e, _)) => { assert_eq!(e.code, ErrorCode::AuthorizationForStatementDenied, "{e}") @@ -451,7 +459,7 @@ mod tests { } } - fn fence_outcome(result: &crate::Result) -> FenceOutcome { + pub(crate) fn fence_outcome(result: &crate::Result) -> FenceOutcome { match result { Ok(c) => c.receipt.outcome, Err(Error::NamespaceFence(e)) => e.outcome(), diff --git a/libsql-server/src/namespace/fence/hooks.rs b/libsql-server/src/namespace/fence/hooks.rs index 28a00c673e..75631e122c 100644 --- a/libsql-server/src/namespace/fence/hooks.rs +++ b/libsql-server/src/namespace/fence/hooks.rs @@ -27,6 +27,10 @@ pub enum HookPoint { AfterInstallingGate, /// Under the transition lock, immediately before the metastore transaction runs. BeforeMetastoreCommit, + /// The in-memory read-closing gate of `SetSourceReadFence` has been published. + AfterClosingReads, + /// The read drain's deadline passed; the leases still held are about to be cancelled. + BeforeReadLeaseCancel, /// The metastore transaction returned a committed result. AfterMetastoreCommit, /// The committed result is about to be published to the gate. diff --git a/libsql-server/src/namespace/fence/mod.rs b/libsql-server/src/namespace/fence/mod.rs index 7e64e0cd7d..5c6f5d4382 100644 --- a/libsql-server/src/namespace/fence/mod.rs +++ b/libsql-server/src/namespace/fence/mod.rs @@ -8,9 +8,9 @@ //! markers with their strict durable encoding ([`record`]), the pure transition function //! ([`transition`]), the metastore tables, compare-and-swap and marker file that persist them //! ([`store`], driven by `MetaStore::apply_fence_command`), and the in-memory authority built -//! on them: the per-namespace [`controller`] with its gate, the positive write [`drain`], the -//! [`registry`] that holds the controllers outside the namespace cache, and the test [`hooks`] -//! on their paths. +//! on them: the per-namespace [`controller`] with its gate and read leases, the positive write +//! [`drain`], the source [`read`] fence, the [`registry`] that holds the controllers outside the +//! namespace cache, and the test [`hooks`] on their paths. // The persistence, controller and protocol layers that consume these types land in the // following commits of this series; until then most of the module is unused by the rest of @@ -22,6 +22,7 @@ pub mod controller; pub mod drain; pub mod hooks; pub mod outcome; +pub mod read; pub mod record; pub mod registry; pub mod state; diff --git a/libsql-server/src/namespace/fence/read.rs b/libsql-server/src/namespace/fence/read.rs new file mode 100644 index 0000000000..cc4a60b906 --- /dev/null +++ b/libsql-server/src/namespace/fence/read.rs @@ -0,0 +1,588 @@ +//! The source read fence (`docs/NAMESPACE_FENCE.md` section 9). +//! +//! `SetSourceReadFence` closes read admission in memory, persists `SOURCE_READ_DRAINING`, and +//! then waits for every read lease that was already held to be released: running SQL programs +//! (including cursors still producing rows and ATTACHes of the namespace from other +//! namespaces), dumps and replication streams. Leases are taken only after checking the gate +//! under the controller's lease lock, so once admission is closed the set can only shrink. At the +//! deadline the remaining work is cancelled and the drain keeps waiting for the actual +//! releases; it never takes elapsed time as proof that a reader has finished. Once no lease is +//! held it persists `SOURCE_READ_FENCED`. + +use std::time::Duration; + +use tokio::time::Instant; + +use crate::namespace::meta_store::{FenceCommit, FenceContext, MetaStore}; + +use super::command::{DrainPolicy, FenceCommand, FenceRequest}; +use super::controller::{FenceController, Transition}; +use super::drain::{now_ms, FORCED_ROLLBACK_GRACE}; +use super::hooks::HookPoint; +use super::outcome::FenceOutcome; +use super::state::{FenceState, OperationClass}; +use super::transition::DrainCompletion; + +/// How long `SetSourceReadFence` waits for running reads and streams before it cancels them, +/// when neither the request nor `--namespace-fence-default-read-drain-ms` names a deadline. +pub const DEFAULT_READ_DRAIN: Duration = Duration::from_secs(30); + +/// `SetSourceReadFence`, steps 1 to 5 of section 9, under `transition`. +/// +/// Returns the `APPLIED` commit of `SOURCE_READ_FENCED` once every read lease is released, the +/// `DRAINING` commit of `SOURCE_READ_DRAINING` when leases are still held after the deadline and +/// the cancellation that follows it (read admission stays closed and a replay of the same +/// command resumes the drain), or the stored result of a replay. +pub async fn set_source_read_fence( + transition: &mut Transition, + meta: &MetaStore, + request: FenceRequest, + mut ctx: FenceContext, +) -> crate::Result { + let controller = transition.controller().clone(); + let key = (request.operation_id, request.command_id); + let policy = match &request.command { + FenceCommand::SetSourceReadFence { drain_policy } => { + drain_policy.unwrap_or_else(|| meta.fence_default_read_drain()) + } + _ => return transition.apply(meta, request, ctx).await, + }; + + // Step 2: close read admission in memory. From here on every new read lease is refused, so + // the drain only waits for the leases already held. Where reads are already closed (a + // resumed drain, a command being reconciled, a read-fenced namespace) nothing is closed. + // A command the checks refuse (step 1, the metastore's) reopens it. + if controller.permits(OperationClass::NormalRead).is_ok() + && controller.gate().state() == FenceState::SourceWriteFenced + { + transition.close_read_admission(key); + let _ = controller.hook(HookPoint::AfterClosingReads).await; + } + + // Step 3: persist SOURCE_READ_DRAINING; its publication replaces the in-memory gate. + let commit = match transition.apply(meta, request, ctx.clone()).await { + Ok(commit) => commit, + Err(e) => { + transition.reopen_read_admission(); + return Err(e); + } + }; + if commit.receipt.outcome != FenceOutcome::Draining { + return Ok(commit); + } + let drain_key = (commit.receipt.operation_id, commit.receipt.command_id); + + // Step 4. + if !drain_readers(&controller, policy).await { + return Ok(commit); + } + + // Step 5. + ctx.now_ms = now_ms(); + transition + .complete_drain(meta, drain_key, DrainCompletion::SourceReads, ctx) + .await +} + +/// Wait until every read lease of the namespace is released. At the deadline the leases still +/// held are cancelled (SQL programs through their connection's cancel flag, dumps and streams +/// through their own), and the drain keeps waiting for the actual releases for one more +/// deadline, but at least [`FORCED_ROLLBACK_GRACE`]. `false` when leases are still held then: +/// the answer is `DRAINING`, never a guess. +async fn drain_readers(controller: &FenceController, policy: DrainPolicy) -> bool { + let namespace = controller.namespace().clone(); + let deadline_after = Duration::from_millis(policy.deadline_ms); + let deadline = Instant::now() + deadline_after; + if wait_for_read_leases(controller, deadline).await { + return true; + } + let _ = controller.hook(HookPoint::BeforeReadLeaseCancel).await; + let held = controller.read_lease_counts(); + let asked = controller.cancel_read_leases(); + tracing::info!( + %namespace, + deadline_ms = policy.deadline_ms, + sql = held.sql, + dump = held.dump, + replication = held.replication, + cancelled = asked, + "read drain deadline passed; cancelling the reads and streams still running" + ); + let grace = Instant::now() + deadline_after.max(FORCED_ROLLBACK_GRACE); + if wait_for_read_leases(controller, grace).await { + return true; + } + let held = controller.read_lease_counts(); + tracing::warn!( + %namespace, + sql = held.sql, + dump = held.dump, + replication = held.replication, + "read leases still held after cancellation; answering DRAINING" + ); + false +} + +/// Wait, on the controller's release notification, until no read lease is held. `false` when +/// `deadline` passes first. +async fn wait_for_read_leases(controller: &FenceController, deadline: Instant) -> bool { + loop { + let released = controller.read_released().notified(); + tokio::pin!(released); + // Registered before the check, so a release in between is not missed. + released.as_mut().enable(); + if controller.read_lease_counts().total() == 0 { + return true; + } + tokio::select! { + _ = &mut released => {} + _ = tokio::time::sleep_until(deadline) => return false, + } + } +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + use std::time::Duration; + + use rusqlite::functions::FunctionFlags; + use uuid::Uuid; + + use super::*; + use crate::auth::Authenticated; + use crate::connection::config::DatabaseConfig; + use crate::connection::program::Program; + use crate::connection::{Connection as _, RequestContext}; + use crate::database::Connection; + use crate::error::Error; + use crate::namespace::fence::command::OnDeadline; + use crate::namespace::fence::drain::tests::{fence_outcome, raw, Source, LONG, OP, PROMPT}; + use crate::namespace::fence::outcome::FenceError; + use crate::namespace::RestoreOption; + use crate::query_result_builder::test::{StepResult, TestBuilder}; + use crate::query_result_builder::QueryResultBuilder as _; + + const OTHER_OP: Uuid = Uuid::from_u128(0xb); + /// A deadline that has already passed: the drain cancels at once. + const NOW: DrainPolicy = DrainPolicy { + deadline_ms: 0, + on_deadline: OnDeadline::Fail, + }; + + /// A write-fenced source. + async fn fenced_source() -> Source { + let s = Source::new().await; + raw(&s.conn().await, "insert into t values (1), (2)") + .await + .unwrap(); + let acquired = s.execute(s.acquire(OP, 1, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&acquired), FenceOutcome::Applied); + s + } + + fn request(s: &Source, op: Uuid, command_id: u128, command: FenceCommand) -> FenceRequest { + let gate = s.fence.gate(); + FenceRequest { + namespace: "ns".into(), + operation_id: op, + command_id: Uuid::from_u128(command_id), + expected_state: gate.state(), + expected_revision: gate.revision(), + command, + } + } + + fn read_fence(s: &Source, command_id: u128, policy: DrainPolicy) -> FenceRequest { + request( + s, + OP, + command_id, + FenceCommand::SetSourceReadFence { + drain_policy: Some(policy), + }, + ) + } + + fn clear_read_fence(s: &Source, command_id: u128) -> FenceRequest { + request(s, OP, command_id, FenceCommand::ClearSourceReadFence) + } + + async fn conn_to(s: &Source, ns: &'static str) -> Arc { + let maker = s + .store + .with(ns.into(), |ns| ns.db.connection_maker()) + .await + .unwrap(); + Arc::new(maker.create().await.unwrap()) + } + + /// Run `stmts` as one program on `conn`, the way a SQL request runs. + async fn program( + s: &Source, + ns: &'static str, + conn: &Arc, + stmts: &'static [&'static str], + ) -> crate::Result> { + let ctx = RequestContext::new( + Authenticated::FullAccess, + ns.into(), + s.store.meta_store().clone(), + ); + conn.execute_program(Program::seq(stmts), ctx, TestBuilder::default(), None) + .await + .map(|b| b.into_ret()) + } + + /// [`program`], on a task of its own. + fn spawn_program( + s: &Source, + ns: &'static str, + conn: &Arc, + stmts: &'static [&'static str], + ) -> tokio::task::JoinHandle>> { + let ctx = RequestContext::new( + Authenticated::FullAccess, + ns.into(), + s.store.meta_store().clone(), + ); + let conn = conn.clone(); + tokio::spawn(async move { + conn.execute_program(Program::seq(stmts), ctx, TestBuilder::default(), None) + .await + .map(|b| b.into_ret()) + }) + } + + fn read_fenced(e: &Error) -> &FenceError { + match e { + Error::NamespaceFence(f) if f.outcome() == FenceOutcome::MigrationReadFenced => f, + other => panic!("expected MIGRATION_READ_FENCED, got {other:?}"), + } + } + + fn step_read_fenced(step: &StepResult) { + match step { + Err(e) => { + read_fenced(e); + } + Ok(rows) => panic!("expected MIGRATION_READ_FENCED, got rows {rows:?}"), + } + } + + fn assert_ok(steps: &[StepResult]) { + for (i, step) in steps.iter().enumerate() { + assert!(step.is_ok(), "step {i} failed: {step:?}"); + } + } + + /// A `park()` SQL function on `conn`: it signals the returned receiver when a program + /// reaches it, and returns once the returned sender is used. SQLite cannot interrupt it, + /// so a program parked in it holds its read lease until the test lets it go. + fn park( + conn: &Connection, + ) -> ( + tokio::sync::mpsc::UnboundedReceiver<()>, + std::sync::mpsc::Sender<()>, + ) { + let (reached_tx, reached_rx) = tokio::sync::mpsc::unbounded_channel::<()>(); + let (resume_tx, resume_rx) = std::sync::mpsc::channel::<()>(); + let parked = std::panic::AssertUnwindSafe((reached_tx, std::sync::Mutex::new(resume_rx))); + conn.with_raw(move |c| { + c.create_scalar_function("park", 0, FunctionFlags::SQLITE_UTF8, move |_| { + let (reached, resume) = &*parked; + reached.send(()).unwrap(); + resume.lock().unwrap().recv().unwrap(); + Ok(1) + }) + }) + .unwrap(); + (reached_rx, resume_tx) + } + + /// A `reached()` SQL function on `conn` that signals the returned receiver and returns. + fn signal(conn: &Connection) -> tokio::sync::mpsc::UnboundedReceiver<()> { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<()>(); + let tx = std::panic::AssertUnwindSafe(tx); + conn.with_raw(move |c| { + c.create_scalar_function("reached", 0, FunctionFlags::SQLITE_UTF8, move |_| { + tx.send(()).unwrap(); + Ok(1) + }) + }) + .unwrap(); + rx + } + + /// A program that is running when the read fence arrives keeps its lease: the drain waits + /// for it and acknowledges only after it finished. Programs that start meanwhile are + /// refused. + #[tokio::test(flavor = "multi_thread")] + async fn read_fence_waits_for_running_program() { + let s = fenced_source().await; + let conn = s.conn().await; + let (mut reached, resume) = park(&conn); + let running = spawn_program( + &s, + "ns", + &conn, + &["select park()", "select count(*) from t"], + ); + reached.recv().await.unwrap(); + assert_eq!(s.fence.read_lease_counts().sql, 1); + + let fence = s.execute(read_fence(&s, 2, LONG)); + s.until_state(FenceState::SourceReadDraining).await; + // New work is refused while the old program still runs. + let refused = program(&s, "ns", &s.conn().await, &["select count(*) from t"]).await; + read_fenced(&refused.unwrap_err()); + assert!(!fence.is_finished()); + + resume.send(()).unwrap(); + let steps = running.await.unwrap().unwrap(); + assert_ok(&steps); + let fenced = fence.await.unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + assert_eq!(s.fence.gate().state(), FenceState::SourceReadFenced); + assert_eq!(s.fence.read_lease_counts().total(), 0); + } + + /// A program admitted after the read-closing gate is published, but before the state is + /// persisted, is refused: the drain never has to wait for work that started after it closed + /// admission. + #[tokio::test(flavor = "multi_thread")] + async fn program_after_closing_gate_is_refused() { + let s = fenced_source().await; + let closed = s.fence.hooks().pause_at(HookPoint::AfterClosingReads); + let fence = s.execute(read_fence(&s, 2, LONG)); + closed.reached().await; + assert_eq!(s.fence.gate().state(), FenceState::SourceWriteFenced); + let refused = program(&s, "ns", &s.conn().await, &["select count(*) from t"]).await; + read_fenced(&refused.unwrap_err()); + assert_eq!(s.fence.read_lease_counts().total(), 0); + closed.resume(); + assert_eq!(fence_outcome(&fence.await.unwrap()), FenceOutcome::Applied); + } + + /// At the deadline a running program is cancelled through its connection's cancel flag; + /// the drain waits for the actual release and then acknowledges. The program reports the + /// read fence. + #[tokio::test(flavor = "multi_thread")] + async fn read_fence_cancels_at_deadline() { + let s = fenced_source().await; + let conn = s.conn().await; + let mut reached = signal(&conn); + let running = spawn_program( + &s, + "ns", + &conn, + &[ + "select reached()", + "with recursive c(x) as (select 1 union all select x + 1 from c) \ + select count(*) from c", + ], + ); + reached.recv().await.unwrap(); + assert_eq!(s.fence.read_lease_counts().sql, 1); + + let fenced = tokio::time::timeout(PROMPT, s.execute(read_fence(&s, 2, NOW))) + .await + .expect("the cancelled program never released its lease") + .unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + let result = tokio::time::timeout(PROMPT, running) + .await + .unwrap() + .unwrap(); + read_fenced(&result.unwrap_err()); + assert!(conn.is_autocommit().await.unwrap()); + } + + /// A lease that is not released after cancellation keeps the drain from acknowledging: the + /// answer is `DRAINING`, reads stay closed, and a replay of the same command completes the + /// drain once the lease is gone. + #[tokio::test(flavor = "multi_thread")] + async fn unreleased_lease_answers_draining_and_replay_completes() { + let s = fenced_source().await; + let conn = s.conn().await; + let (mut reached, resume) = park(&conn); + let running = spawn_program(&s, "ns", &conn, &["select park()"]); + reached.recv().await.unwrap(); + + let request = read_fence(&s, 2, NOW); + let draining = s.execute(request.clone()).await.unwrap(); + assert_eq!(fence_outcome(&draining), FenceOutcome::Draining); + assert_eq!(s.fence.gate().state(), FenceState::SourceReadDraining); + let refused = program(&s, "ns", &s.conn().await, &["select 1"]).await; + read_fenced(&refused.unwrap_err()); + + resume.send(()).unwrap(); + // The program was cancelled by the fence; it reports the fence, not its rows. + read_fenced(&running.await.unwrap().unwrap_err()); + let replayed = s.execute(request).await.unwrap(); + assert_eq!(fence_outcome(&replayed), FenceOutcome::Applied); + assert_eq!(s.fence.gate().state(), FenceState::SourceReadFenced); + } + + /// A connection left idle inside a transaction holds no lease, so the drain does not wait + /// for it; its next program is refused and its transaction rolled back. + #[tokio::test(flavor = "multi_thread")] + async fn idle_txn_fails_on_next_program() { + let s = fenced_source().await; + let conn = s.conn().await; + assert_ok( + &program(&s, "ns", &conn, &["begin", "select count(*) from t"]) + .await + .unwrap(), + ); + assert!(!conn.is_autocommit().await.unwrap()); + assert_eq!(s.fence.read_lease_counts().total(), 0); + + let fenced = s.execute(read_fence(&s, 2, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + + let refused = program(&s, "ns", &conn, &["select count(*) from t", "commit"]).await; + read_fenced(&refused.unwrap_err()); + assert!(conn.is_autocommit().await.unwrap()); + } + + /// Clearing the read fence reopens reads and leaves writes fenced. + #[tokio::test(flavor = "multi_thread")] + async fn clear_read_fence_reopens_reads_not_writes() { + let s = fenced_source().await; + let fenced = s.execute(read_fence(&s, 2, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + let conn = s.conn().await; + read_fenced(&program(&s, "ns", &conn, &["select 1"]).await.unwrap_err()); + + let revision = s.fence.gate().revision(); + let cleared = s.execute(clear_read_fence(&s, 3)).await.unwrap(); + assert_eq!(fence_outcome(&cleared), FenceOutcome::Applied); + let gate = s.fence.gate(); + assert_eq!(gate.state(), FenceState::SourceWriteFenced); + assert!(gate.revision() > revision); + + let steps = program( + &s, + "ns", + &conn, + &["select count(*) from t", "insert into t values (3)"], + ) + .await + .unwrap(); + assert!(matches!( + steps[0].as_ref().unwrap().as_slice(), + [row] if matches!(row.as_slice(), [crate::query::Value::Integer(2)]) + )); + match &steps[1] { + Err(Error::NamespaceFence(e)) => { + assert_eq!(e.outcome(), FenceOutcome::MigrationWriteFenced) + } + other => panic!("expected MIGRATION_WRITE_FENCED, got {other:?}"), + } + assert_eq!(s.count().await, 2); + } + + /// A read fence the checks refuse (another operation owns the source) reopens read + /// admission it had closed, and reads are served again. + #[tokio::test(flavor = "multi_thread")] + async fn refused_read_fence_reopens_reads() { + let s = fenced_source().await; + let request = request( + &s, + OTHER_OP, + 2, + FenceCommand::SetSourceReadFence { + drain_policy: Some(LONG), + }, + ); + let refused = s.execute(request).await.unwrap(); + assert_eq!( + fence_outcome(&refused), + FenceOutcome::FenceOwnedByAnotherOperation + ); + let gate = s.fence.gate(); + assert_eq!(gate.state(), FenceState::SourceWriteFenced); + assert!(gate.closing_reads.is_none()); + assert_ok( + &program(&s, "ns", &s.conn().await, &["select count(*) from t"]) + .await + .unwrap(), + ); + } + + /// ATTACH of a namespace is a read of that namespace: a program on another namespace that + /// has it attached holds a lease on it, so its read fence waits for that program; once it is + /// read-fenced, a new ATTACH and any program on a connection that still has it attached are + /// refused, and a connection that detached it is served again. + #[tokio::test(flavor = "multi_thread")] + async fn attach_of_read_fenced_namespace_denied() { + let s = Source::new().await; + raw(&s.conn().await, "insert into t values (1)") + .await + .unwrap(); + let config = s + .store + .with("ns".into(), |ns| ns.db_config_store.clone()) + .await + .unwrap(); + config + .store(DatabaseConfig { + allow_attach: true, + txn_timeout: Some(Duration::from_secs(600)), + ..Default::default() + }) + .await + .unwrap(); + s.store + .create( + "other".into(), + RestoreOption::Latest, + DatabaseConfig::default(), + ) + .await + .unwrap(); + let acquired = s.execute(s.acquire(OP, 1, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&acquired), FenceOutcome::Applied); + + let other = conn_to(&s, "other").await; + let steps = program( + &s, + "other", + &other, + &["attach ns as a", "select count(*) from a.t"], + ) + .await + .unwrap(); + assert_ok(&steps); + assert_eq!(s.fence.read_lease_counts().total(), 0); + + // A later program on the connection that has `ns` attached holds a lease on it. + let (mut reached, resume) = park(&other); + let running = spawn_program(&s, "other", &other, &["select park()"]); + reached.recv().await.unwrap(); + assert_eq!(s.fence.read_lease_counts().sql, 1); + let fence = s.execute(read_fence(&s, 2, LONG)); + s.until_state(FenceState::SourceReadDraining).await; + assert!(!fence.is_finished()); + resume.send(()).unwrap(); + assert_ok(&running.await.unwrap().unwrap()); + assert_eq!(fence_outcome(&fence.await.unwrap()), FenceOutcome::Applied); + + // The connection that still has it attached is refused outright... + read_fenced( + &program(&s, "other", &other, &["select 1"]) + .await + .unwrap_err(), + ); + // ...a fresh ATTACH is refused as a step... + let fresh = conn_to(&s, "other").await; + let steps = program(&s, "other", &fresh, &["attach ns as b"]) + .await + .unwrap(); + step_read_fenced(&steps[0]); + // ...and a connection that detached it is served again. + let steps = program(&s, "other", &fresh, &["select 1"]).await.unwrap(); + assert_ok(&steps); + } +} diff --git a/libsql-server/src/namespace/meta_store.rs b/libsql-server/src/namespace/meta_store.rs index 1a16994c3f..ebd1ba2a64 100644 --- a/libsql-server/src/namespace/meta_store.rs +++ b/libsql-server/src/namespace/meta_store.rs @@ -112,6 +112,7 @@ struct FenceSettings { receipt_retention: Duration, /// The write drain deadline of an `AcquireSourceWriteFence` that names no drain policy. default_write_drain: Duration, + default_read_drain: Duration, } fn setup_connection(conn: &rusqlite::Connection) -> Result<()> { @@ -240,6 +241,9 @@ impl MetaStoreInner { default_write_drain: config .namespace_fence_default_write_drain .unwrap_or(crate::namespace::fence::drain::DEFAULT_WRITE_DRAIN), + default_read_drain: config + .namespace_fence_default_read_drain + .unwrap_or(crate::namespace::fence::read::DEFAULT_READ_DRAIN), }; let mut this = MetaStoreInner { @@ -1315,6 +1319,17 @@ impl MetaStore { } } + /// The drain policy of a `SetSourceReadFence` that names none: the configured deadline, + /// after which running reads are cancelled and streams terminated (`on_deadline` does not + /// apply to reads). + pub fn fence_default_read_drain(&self) -> DrainPolicy { + DrainPolicy { + deadline_ms: u64::try_from(self.inner.fence.default_read_drain.as_millis()) + .unwrap_or(u64::MAX), + on_deadline: OnDeadline::Fail, + } + } + /// Whether this metastore holds fence state, so fences are loaded and enforced. pub fn fence_enforced(&self) -> bool { self.inner.fence.tables diff --git a/libsql-server/src/namespace/mod.rs b/libsql-server/src/namespace/mod.rs index 28ba60ea26..f75dbd700f 100644 --- a/libsql-server/src/namespace/mod.rs +++ b/libsql-server/src/namespace/mod.rs @@ -29,8 +29,16 @@ mod schema_lock; mod store; pub type ResetCb = Box; +/// Resolves a namespace that a program ATTACHes: its directory, and its fence controller, which +/// admits the attachment as a read of that namespace (`docs/NAMESPACE_FENCE.md` section 9). pub type ResolveNamespacePathFn = - Arc crate::Result> + Sync + Send + 'static>; + Arc crate::Result + Sync + Send + 'static>; + +/// A namespace resolved for ATTACH. +pub struct AttachTarget { + pub path: Arc, + pub fence: Arc, +} pub enum ResetOp { Reset(NamespaceName), @@ -103,9 +111,6 @@ impl Namespace { Ok(()) } - // Read by the protocol layers that consult the gate outside a connection (dump, - // replication, lifecycle), which land later in this series. - #[allow(dead_code)] pub(crate) fn fence(&self) -> &Arc { &self.fence } diff --git a/libsql-server/src/namespace/store.rs b/libsql-server/src/namespace/store.rs index af406dfd96..a472fc39c3 100644 --- a/libsql-server/src/namespace/store.rs +++ b/libsql-server/src/namespace/store.rs @@ -22,6 +22,7 @@ use crate::stats::Stats; use super::broadcasters::{BroadcasterHandle, BroadcasterRegistry}; use super::configurator::{DynConfigurator, NamespaceConfigurators}; use super::fence::command::{FenceCommand, FenceRequest}; +use super::fence::controller::FenceController; use super::fence::record::ServerIdentity; use super::fence::registry::FenceRegistry; use super::meta_store::{FenceCommit, FenceContext, MetaStore, MetaStoreHandle}; @@ -376,8 +377,12 @@ impl NamespaceStore { Arc::new({ let store = self.clone(); move |ns: &NamespaceName| { - tokio::runtime::Handle::current() - .block_on(store.with(ns.clone(), |ns| ns.path.clone())) + tokio::runtime::Handle::current().block_on(store.with(ns.clone(), |ns| { + super::AttachTarget { + path: ns.path.clone(), + fence: ns.fence().clone(), + } + })) } }) }) @@ -560,6 +565,23 @@ impl NamespaceStore { .await } + /// The fence controller that admits reads of `namespace` without loading it: `None` when + /// the namespace does not exist (and has no fence state). A namespace whose fence state is + /// unavailable is refused. + pub(crate) async fn fence_gate( + &self, + namespace: &NamespaceName, + ) -> crate::Result>> { + self.inner.fences.check_available(namespace)?; + if let Some(controller) = self.inner.fences.get(namespace) { + return Ok(Some(controller)); + } + if self.inner.metadata.exists(namespace).await { + return Ok(Some(self.inner.fences.controller(namespace))); + } + Ok(None) + } + pub(crate) fn schema_locks(&self) -> &SchemaLocksRegistry { &self.inner.schema_locks } From 1a21f11943db6e8e167af61ac89031e81520aea6 Mon Sep 17 00:00:00 2001 From: River Date: Tue, 29 Sep 2026 18:02:32 +0000 Subject: [PATCH 2/3] libsql-server: end dump and replication streams on the read fence Dump and replication now hold read leases on the namespace fence, so the source read fence drains them positively: - /dump is admitted by the fence gate before a connection is created (a failure to create one is an error, not a panic) and the export holds a dump lease. At the read drain's deadline the export is cancelled before its next row, and a write blocked on a peer that stopped reading fails at once, so the lease is released without the peer. A cancelled dump ends its body with the fence error and never reaches its final COMMIT;. - Replication hello, log_entries, batch_log_entries and snapshot are refused at their start with FAILED_PRECONDITION and x-libsql-fence-code while streams are denied; refusals are counted and logged at most once a minute per namespace. Streams (and the frames of a batch) are served through FencedStream, whose watcher ends the stream as soon as the gate closes or the drain cancels it, dropping the inner stream and the lease itself; the next poll yields the typed terminal status. - With --enable-namespace-fence, the RPC server and the user HTTP server send HTTP/2 keepalive pings (--namespace-fence-keepalive-interval-s, default 30 s, 20 s timeout) so dead peers are detected. Co-authored-by: Tomasz Szymczyszyn --- libsql-server/src/connection/dump/exporter.rs | 35 +- libsql-server/src/h2c.rs | 38 +- libsql-server/src/http/user/dump.rs | 111 +++- libsql-server/src/http/user/mod.rs | 21 +- libsql-server/src/lib.rs | 9 + libsql-server/src/main.rs | 10 + libsql-server/src/namespace/fence/drain.rs | 12 +- libsql-server/src/namespace/fence/mod.rs | 4 +- libsql-server/src/namespace/fence/read.rs | 8 +- libsql-server/src/namespace/fence/stream.rs | 594 ++++++++++++++++++ libsql-server/src/namespace/store.rs | 9 +- libsql-server/src/rpc/mod.rs | 19 +- .../src/rpc/replication/replication_log.rs | 136 +++- 13 files changed, 948 insertions(+), 58 deletions(-) create mode 100644 libsql-server/src/namespace/fence/stream.rs diff --git a/libsql-server/src/connection/dump/exporter.rs b/libsql-server/src/connection/dump/exporter.rs index 1a1c69728b..66ae566edc 100644 --- a/libsql-server/src/connection/dump/exporter.rs +++ b/libsql-server/src/connection/dump/exporter.rs @@ -7,15 +7,30 @@ use anyhow::bail; use rusqlite::types::ValueRef; use rusqlite::OptionalExtension; -struct DumpState { +struct DumpState<'a, W: Write> { /// true if db is in writable_schema mode writable_schema: bool, writer: W, + /// Checked before every row: once it returns `true` the export stops with [`DumpCancelled`]. + cancelled: &'a dyn Fn() -> bool, } +/// The export was stopped by its caller before it completed (for example by a namespace read +/// fence). The output is incomplete: it never reaches its final `COMMIT;`. +#[derive(Debug, thiserror::Error)] +#[error("the dump was cancelled before it completed")] +pub struct DumpCancelled; + use rusqlite::ffi::{sqlite3_keyword_check, sqlite3_table_column_metadata, SQLITE_OK}; -impl DumpState { +impl DumpState<'_, W> { + fn check_cancelled(&self) -> anyhow::Result<()> { + if (self.cancelled)() { + return Err(DumpCancelled.into()); + } + Ok(()) + } + fn run_schema_dump_query( &mut self, txn: &rusqlite::Connection, @@ -25,6 +40,7 @@ impl DumpState { let mut stmt = txn.prepare(stmt)?; let mut rows = stmt.query(())?; while let Some(row) = rows.next()? { + self.check_cancelled()?; let ValueRef::Text(table) = row.get_ref(0)? else { bail!("invalid schema table") }; @@ -103,6 +119,7 @@ impl DumpState { let mut stmt = txn.prepare(&select)?; let mut rows = stmt.query(())?; while let Some(row) = rows.next()? { + self.check_cancelled()?; write!(self.writer, "{insert}")?; if row_id_col.is_some() { write_value_ref(&mut self.writer, row.get_ref(0)?)?; @@ -128,6 +145,7 @@ impl DumpState { let col_count = stmt.column_count(); let mut rows = stmt.query(())?; while let Some(row) = rows.next()? { + self.check_cancelled()?; let ValueRef::Text(sql) = row.get_ref(0)? else { bail!("the first row in a table dump query should be of type text") }; @@ -437,6 +455,17 @@ pub fn export_dump( db: &mut rusqlite::Connection, writer: impl Write, preserve_rowids: bool, +) -> anyhow::Result<()> { + export_dump_cancellable(db, writer, preserve_rowids, &|| false) +} + +/// [`export_dump`], stopped with a [`DumpCancelled`] error as soon as `cancelled` returns `true` +/// (it is checked before every row and before the final `COMMIT;`). +pub fn export_dump_cancellable( + db: &mut rusqlite::Connection, + writer: impl Write, + preserve_rowids: bool, + cancelled: &dyn Fn() -> bool, ) -> anyhow::Result<()> { let mut txn = db.transaction()?; txn.execute("PRAGMA writable_schema=ON", ())?; @@ -444,6 +473,7 @@ pub fn export_dump( let mut state = DumpState { writable_schema: false, writer, + cancelled, }; writeln!(state.writer, "PRAGMA foreign_keys=OFF;")?; @@ -469,6 +499,7 @@ AND type IN ('index','trigger','view')"; writeln!(state.writer, "PRAGMA writable_schema=OFF;")?; } + state.check_cancelled()?; writeln!(state.writer, "COMMIT;")?; let _ = savepoint.execute("PRAGMA writable_schema = OFF;", ()); diff --git a/libsql-server/src/h2c.rs b/libsql-server/src/h2c.rs index 93d4999543..2f85f6f256 100644 --- a/libsql-server/src/h2c.rs +++ b/libsql-server/src/h2c.rs @@ -40,6 +40,7 @@ use std::marker::PhantomData; use std::pin::Pin; +use std::time::Duration; use axum::{body::BoxBody, http::HeaderValue}; use bytes::Bytes; @@ -56,6 +57,7 @@ type BoxError = Box; #[derive(Debug, Clone)] pub struct H2cMaker { s: S, + http2_keepalive_interval: Option, _pd: PhantomData, } @@ -63,9 +65,32 @@ impl H2cMaker { pub fn new(s: S) -> Self { Self { s, + http2_keepalive_interval: None, _pd: PhantomData, } } + + /// Send HTTP/2 keepalive pings on upgraded `h2c` connections at this interval. + pub fn with_http2_keepalive(mut self, interval: Option) -> Self { + self.http2_keepalive_interval = interval; + self + } +} + +/// How long a keepalive ping may go unanswered before the connection is closed. +pub const HTTP2_KEEPALIVE_TIMEOUT: Duration = Duration::from_secs(20); + +/// Configure HTTP/2 keepalive on a server builder; `None` leaves it unchanged. +pub fn with_http2_keepalive( + builder: hyper::server::Builder, + interval: Option, +) -> hyper::server::Builder { + match interval { + Some(interval) => builder + .http2_keep_alive_interval(interval) + .http2_keep_alive_timeout(HTTP2_KEEPALIVE_TIMEOUT), + None => builder, + } } impl Service<&C> for H2cMaker @@ -95,10 +120,12 @@ where fn call(&mut self, conn: &C) -> Self::Future { let connect_info = conn.connect_info(); let s = self.s.clone(); + let http2_keepalive_interval = self.http2_keepalive_interval; Box::pin(async move { Ok(H2c { s, connect_info, + http2_keepalive_interval, _pd: PhantomData, }) }) @@ -112,6 +139,7 @@ where pub struct H2c { s: S, connect_info: TcpConnectInfo, + http2_keepalive_interval: Option, _pd: PhantomData, } @@ -139,6 +167,7 @@ where fn call(&mut self, mut req: hyper::Request) -> Self::Future { let mut svc = self.s.clone(); let connect_info = self.connect_info.clone(); + let http2_keepalive_interval = self.http2_keepalive_interval; Box::pin(async move { req.extensions_mut().insert(connect_info.clone()); @@ -169,8 +198,13 @@ where tracing::debug!("Successfully upgraded the connection, speaking h2 now"); - if let Err(e) = hyper::server::conn::Http::new() - .http2_only(true) + let mut http = hyper::server::conn::Http::new(); + http.http2_only(true); + if let Some(interval) = http2_keepalive_interval { + http.http2_keep_alive_interval(interval) + .http2_keep_alive_timeout(HTTP2_KEEPALIVE_TIMEOUT); + } + if let Err(e) = http .serve_connection( upgraded_io, tower::service_fn(move |mut r: hyper::Request| { diff --git a/libsql-server/src/http/user/dump.rs b/libsql-server/src/http/user/dump.rs index 16efcc52a7..ae0260482e 100644 --- a/libsql-server/src/http/user/dump.rs +++ b/libsql-server/src/http/user/dump.rs @@ -1,5 +1,7 @@ use std::future::Future; +use std::io::Write; use std::pin::Pin; +use std::sync::Arc; use std::task; use axum::extract::{Query, State as AxumState}; @@ -9,9 +11,14 @@ use pin_project_lite::pin_project; use serde::Deserialize; use crate::auth::Authenticated; -use crate::connection::dump::exporter::export_dump; -use crate::connection::Connection as _; +use crate::connection::dump::exporter::export_dump_cancellable; +use crate::connection::{Connection as _, MakeConnection}; +use crate::database::Connection; use crate::error::Error; +use crate::namespace::fence::controller::{FenceController, LeaseKind}; +use crate::namespace::fence::stream::{ + acquire_stream_lease, cancelled_by_read_fence, StreamCancel, +}; use crate::BLOCKING_RT; use super::db_factory::namespace_from_headers; @@ -95,36 +102,114 @@ pub(super) async fn handle_dump( return Err(Error::NamespaceDoesntExist(namespace.to_string())); } - let conn_maker = state + let (conn_maker, fence) = state .namespaces .with(namespace, |ns| { if !ns.db.is_primary() { return Err(Error::NotAPrimary); } - Ok::<_, crate::Error>(ns.db.connection_maker()) + Ok::<_, crate::Error>((ns.db.connection_maker(), ns.fence().clone())) }) .await??; - let conn = conn_maker.create().await.unwrap(); + let stream = dump_stream(&fence, conn_maker, query.preserve_row_ids.unwrap_or(false)).await?; + + Ok(axum::body::StreamBody::new(stream)) +} + +/// The dump of one namespace as a byte stream (`docs/NAMESPACE_FENCE.md` section 9). +/// +/// The dump is admitted by the namespace's fence gate before any connection is created, and +/// holds a `Dump` read lease until the export has stopped. When the read drain cancels it, the +/// export stops before its next row, and also if it is blocked because the peer is not reading, +/// so the lease is released without the peer's help; the stream then ends with the fence error +/// instead of ending cleanly, so the response body is aborted and the client never receives a +/// dump that looks complete. +pub(crate) async fn dump_stream( + fence: &Arc, + conn_maker: Arc>, + preserve_row_ids: bool, +) -> crate::Result>> { + let (lease, cancel) = + acquire_stream_lease(fence, LeaseKind::Dump).map_err(Error::NamespaceFence)?; + + let conn = conn_maker.create().await?; let (reader, writer) = tokio::io::duplex(8 * 1024); + let writer = CancellableWriter { + inner: writer, + cancel: cancel.clone(), + handle: tokio::runtime::Handle::current(), + }; let join_handle = BLOCKING_RT.spawn_blocking(move || { - let writer = tokio_util::io::SyncIoBridge::new(writer); - conn.with_raw(|conn| { - export_dump(conn, writer, query.preserve_row_ids.unwrap_or(false)).map_err(Into::into) - }) + // Released once the export has stopped and its read transaction is gone. + let _lease = lease; + let result = conn.with_raw(|conn| { + export_dump_cancellable(conn, writer, preserve_row_ids, &|| cancel.is_cancelled()) + }); + match result { + Ok(()) => Ok(()), + Err(_) if cancel.is_cancelled() => Err(Error::NamespaceFence(cancelled_by_read_fence( + LeaseKind::Dump, + ))), + Err(e) => Err(e.into()), + } }); let stream = tokio_util::io::ReaderStream::new(reader); - let stream = DumpStream { + Ok(DumpStream { stream: stream.fuse(), join_handle: Some(join_handle), - }; + }) +} - let stream = axum::body::StreamBody::new(stream); +/// The export's side of the dump pipe. A write waits for the reader (the HTTP body), unless +/// the dump is cancelled, in which case it fails at once. +struct CancellableWriter { + inner: tokio::io::DuplexStream, + cancel: Arc, + handle: tokio::runtime::Handle, +} - Ok(stream) +impl CancellableWriter { + fn cancelled() -> std::io::Error { + std::io::Error::new(std::io::ErrorKind::Other, "dump cancelled") + } +} + +impl Write for CancellableWriter { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + use tokio::io::AsyncWriteExt as _; + let Self { + inner, + cancel, + handle, + } = self; + handle.block_on(async { + tokio::select! { + biased; + _ = cancel.cancelled() => Err(Self::cancelled()), + r = inner.write(buf) => r, + } + }) + } + + fn flush(&mut self) -> std::io::Result<()> { + use tokio::io::AsyncWriteExt as _; + let Self { + inner, + cancel, + handle, + } = self; + handle.block_on(async { + tokio::select! { + biased; + _ = cancel.cancelled() => Err(Self::cancelled()), + r = inner.flush() => r, + } + }) + } } diff --git a/libsql-server/src/http/user/mod.rs b/libsql-server/src/http/user/mod.rs index 1575a3574b..7d46dce8cd 100644 --- a/libsql-server/src/http/user/mod.rs +++ b/libsql-server/src/http/user/mod.rs @@ -1,5 +1,5 @@ pub mod db_factory; -mod dump; +pub(crate) mod dump; mod extract; mod hrana_over_http_1; mod listen; @@ -10,6 +10,7 @@ mod types; pub mod timing; use std::sync::Arc; +use std::time::Duration; use anyhow::Context; use axum::extract::{FromRef, FromRequest, FromRequestParts, Path as AxumPath, State as AxumState}; @@ -256,6 +257,8 @@ pub struct UserApi { pub enable_console: bool, pub self_url: Option, pub primary_url: Option, + /// HTTP/2 keepalive interval (see `Server::http2_keepalive_interval`). + pub http2_keepalive_interval: Option, } impl UserApi @@ -443,14 +446,18 @@ where ); let router = router.fallback(handle_fallback); - let h2c = crate::h2c::H2cMaker::new(router); + let keepalive = self.http2_keepalive_interval; + let h2c = crate::h2c::H2cMaker::new(router).with_http2_keepalive(keepalive); task_manager.spawn_with_shutdown_notify(|shutdown| async move { - hyper::server::Server::builder(acceptor) - .serve(h2c) - .with_graceful_shutdown(shutdown.notified()) - .await - .context("http server")?; + crate::h2c::with_http2_keepalive( + hyper::server::Server::builder(acceptor), + keepalive, + ) + .serve(h2c) + .with_graceful_shutdown(shutdown.notified()) + .await + .context("http server")?; Ok(()) }); } diff --git a/libsql-server/src/lib.rs b/libsql-server/src/lib.rs index 1642ad951a..66fbbcf762 100644 --- a/libsql-server/src/lib.rs +++ b/libsql-server/src/lib.rs @@ -151,6 +151,10 @@ pub struct Server anyhow::Result<()> + Send + Sync + 'static>>, + /// HTTP/2 keepalive interval of the RPC server and the user-port gRPC services, set when + /// namespace fences are enabled so that the streams of dead peers are detected + /// (`docs/NAMESPACE_FENCE.md` section 9). `None` keeps hyper's default (no keepalive). + pub http2_keepalive_interval: Option, } impl Default for Server { @@ -180,6 +184,7 @@ impl Default for Server { force_load_wals: false, sync_conccurency: 8, set_log_level: None, + http2_keepalive_interval: None, } } } @@ -196,6 +201,7 @@ struct Services { db_config: DbConfig, user_auth_strategy: Auth, pub set_log_level: Option anyhow::Result<()> + Send + Sync + 'static>>, + http2_keepalive_interval: Option, } struct TaskManager { @@ -290,6 +296,7 @@ where enable_console: self.user_api_config.enable_http_console, self_url: self.user_api_config.self_url, primary_url: self.user_api_config.primary_url, + http2_keepalive_interval: self.http2_keepalive_interval, }; let user_http_service = user_http.configure(task_manager); @@ -529,6 +536,7 @@ where db_config: self.db_config, user_auth_strategy, set_log_level: self.set_log_level.take(), + http2_keepalive_interval: self.http2_keepalive_interval, } } @@ -677,6 +685,7 @@ where config.tls_config, idle_shutdown_kicker.clone(), replication_service, // internal replicaton service + self.http2_keepalive_interval, )); } diff --git a/libsql-server/src/main.rs b/libsql-server/src/main.rs index f8f56e7494..02a8011c0f 100644 --- a/libsql-server/src/main.rs +++ b/libsql-server/src/main.rs @@ -280,6 +280,13 @@ struct Cli { #[clap(long, env = "SQLD_NAMESPACE_FENCE_DEFAULT_READ_DRAIN_MS")] namespace_fence_default_read_drain_ms: Option, + /// HTTP/2 keepalive interval, in seconds, of the RPC server and the user-port gRPC + /// services when namespace fences are enabled, so that replication streams of dead peers + /// are detected (a ping unanswered for 20 seconds closes the connection). Defaults to 30 + /// seconds; ignored unless `--enable-namespace-fence` is set. + #[clap(long, env = "SQLD_NAMESPACE_FENCE_KEEPALIVE_INTERVAL_S")] + namespace_fence_keepalive_interval_s: Option, + /// Shutdown timeout duration in seconds, defaults to 30 seconds. #[clap(long, env = "SQLD_SHUTDOWN_TIMEOUT")] shutdown_timeout: Option, @@ -756,6 +763,9 @@ async fn build_server( force_load_wals: config.force_load_wals, sync_conccurency: config.sync_conccurency, set_log_level: Some(Box::new(set_log_level)), + http2_keepalive_interval: config.enable_namespace_fence.then(|| { + Duration::from_secs(config.namespace_fence_keepalive_interval_s.unwrap_or(30)) + }), }) } diff --git a/libsql-server/src/namespace/fence/drain.rs b/libsql-server/src/namespace/fence/drain.rs index 5cc70563b8..3e5c8ba5b8 100644 --- a/libsql-server/src/namespace/fence/drain.rs +++ b/libsql-server/src/namespace/fence/drain.rs @@ -304,7 +304,7 @@ pub(crate) mod tests { use crate::namespace::fence::record::ServerIdentity; use crate::namespace::fence::state::FenceState; use crate::namespace::meta_store::FenceCommitKind; - use crate::namespace::store::fence_tests::open_store; + use crate::namespace::store::fence_tests::open_store_with_max_log_size; use crate::namespace::store::NamespaceStore; use crate::namespace::RestoreOption; use crate::replication::primary::logger::ReplicationLogger; @@ -325,13 +325,19 @@ pub(crate) mod tests { _dir: TempDir, pub(crate) store: NamespaceStore, pub(crate) fence: Arc, - logger: Arc, + pub(crate) logger: Arc, } impl Source { pub(crate) async fn new() -> Self { + Self::with_max_log_size(1_000_000_000).await + } + + /// A source whose replication log is compacted into a snapshot once it holds more than + /// `max_log_size` MB (`0`: at the next compaction). + pub(crate) async fn with_max_log_size(max_log_size: u64) -> Self { let dir = tempdir().unwrap(); - let store = open_store(dir.path()).await; + let store = open_store_with_max_log_size(dir.path(), max_log_size).await; store .create( "ns".into(), diff --git a/libsql-server/src/namespace/fence/mod.rs b/libsql-server/src/namespace/fence/mod.rs index 5c6f5d4382..74f227cf23 100644 --- a/libsql-server/src/namespace/fence/mod.rs +++ b/libsql-server/src/namespace/fence/mod.rs @@ -9,7 +9,8 @@ //! ([`transition`]), the metastore tables, compare-and-swap and marker file that persist them //! ([`store`], driven by `MetaStore::apply_fence_command`), and the in-memory authority built //! on them: the per-namespace [`controller`] with its gate and read leases, the positive write -//! [`drain`], the source [`read`] fence, the [`registry`] that holds the controllers outside the +//! [`drain`], the source [`read`] fence and its +//! [`stream`] leases for dump and replication, the [`registry`] that holds the controllers outside the //! namespace cache, and the test [`hooks`] on their paths. // The persistence, controller and protocol layers that consume these types land in the @@ -27,6 +28,7 @@ pub mod record; pub mod registry; pub mod state; pub mod store; +pub mod stream; pub mod transition; #[cfg(test)] diff --git a/libsql-server/src/namespace/fence/read.rs b/libsql-server/src/namespace/fence/read.rs index cc4a60b906..50feb0b3b1 100644 --- a/libsql-server/src/namespace/fence/read.rs +++ b/libsql-server/src/namespace/fence/read.rs @@ -142,7 +142,7 @@ async fn wait_for_read_leases(controller: &FenceController, deadline: Instant) - } #[cfg(test)] -mod tests { +pub(crate) mod tests { use std::sync::Arc; use std::time::Duration; @@ -165,13 +165,13 @@ mod tests { const OTHER_OP: Uuid = Uuid::from_u128(0xb); /// A deadline that has already passed: the drain cancels at once. - const NOW: DrainPolicy = DrainPolicy { + pub(crate) const NOW: DrainPolicy = DrainPolicy { deadline_ms: 0, on_deadline: OnDeadline::Fail, }; /// A write-fenced source. - async fn fenced_source() -> Source { + pub(crate) async fn fenced_source() -> Source { let s = Source::new().await; raw(&s.conn().await, "insert into t values (1), (2)") .await @@ -193,7 +193,7 @@ mod tests { } } - fn read_fence(s: &Source, command_id: u128, policy: DrainPolicy) -> FenceRequest { + pub(crate) fn read_fence(s: &Source, command_id: u128, policy: DrainPolicy) -> FenceRequest { request( s, OP, diff --git a/libsql-server/src/namespace/fence/stream.rs b/libsql-server/src/namespace/fence/stream.rs new file mode 100644 index 0000000000..e57ba2f181 --- /dev/null +++ b/libsql-server/src/namespace/fence/stream.rs @@ -0,0 +1,594 @@ +//! Read leases for streams that serve namespace data outside SQL programs: `/dump` and the +//! replication `log_entries`, `batch_log_entries` and `snapshot` calls (`docs/NAMESPACE_FENCE.md` +//! section 9). +//! +//! A [`FencedStream`] holds a read lease of the stream's kind for as long as it can produce +//! items. A watcher task ends it when the gate stops admitting streams, or when the read drain +//! cancels it at its deadline: the watcher drops the inner stream and the lease itself, so the +//! lease is released even if the peer never polls the stream again, and the next poll yields +//! one terminal error carrying the typed fence code and then ends. +//! +//! A dump is not a stream of this kind (its data is produced by a blocking export on a +//! connection), so it holds its lease in the export itself and uses a [`StreamCancel`] to stop +//! it ([`crate::http::user::dump`]). + +use std::pin::Pin; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; +use std::task::{Context, Poll}; + +use futures::task::AtomicWaker; +use futures::Stream; +use parking_lot::Mutex; +use tokio::sync::{oneshot, Notify}; + +use super::controller::{FenceController, LeaseKind, ReadLease}; +use super::outcome::{FenceError, FenceOutcome}; +use super::state::OperationClass; + +/// How the read drain stops a stream at its deadline: a flag for code that polls it, and a +/// notification for code that waits. Setting it never blocks. +#[derive(Debug, Default)] +pub struct StreamCancel { + cancelled: AtomicBool, + notify: Notify, +} + +impl StreamCancel { + pub fn cancel(&self) { + self.cancelled.store(true, Ordering::Release); + // `notify_one` keeps a permit when nobody waits yet, so a waiter that arrives later + // still wakes. + self.notify.notify_one(); + } + + pub fn is_cancelled(&self) -> bool { + self.cancelled.load(Ordering::Acquire) + } + + /// Resolves once [`cancel`](Self::cancel) has been called. + pub async fn cancelled(&self) { + while !self.is_cancelled() { + self.notify.notified().await; + } + } +} + +/// The error a stream ends with when the read drain cancels it at its deadline. +pub fn cancelled_by_read_fence(kind: LeaseKind) -> FenceError { + FenceError::new( + FenceOutcome::MigrationReadFenced, + format!("the {kind:?} stream was ended by the namespace read fence"), + ) +} + +/// Admit a stream of `kind` and hold a read lease for it; the returned cancel is what the read +/// drain uses at its deadline. +pub fn acquire_stream_lease( + fence: &Arc, + kind: LeaseKind, +) -> Result<(ReadLease, Arc), FenceError> { + let cancel = Arc::new(StreamCancel::default()); + let lease = fence.acquire_read_lease(OperationClass::Stream, kind, { + let cancel = cancel.clone(); + move || cancel.cancel() + })?; + Ok((lease, cancel)) +} + +/// Resolves with the gate's refusal once the gate stops admitting streams. +async fn gate_denies_streams(fence: &FenceController) -> FenceError { + let mut gate = fence.subscribe(); + loop { + if let Err(e) = gate.borrow_and_update().permits(OperationClass::Stream) { + return e; + } + if gate.changed().await.is_err() { + // The controller is gone, which only happens when the namespace is destroyed: end + // the stream rather than serve it without a gate. + return FenceError::new( + FenceOutcome::FenceStateUnavailable, + "the namespace fence controller is gone", + ); + } + } +} + +enum Slot { + Live { + stream: S, + _lease: ReadLease, + /// Dropped with the stream, which stops the watcher. + _stop_watcher: oneshot::Sender<()>, + }, + /// Ended by the fence; the error is yielded once. + Terminated(Option), + /// Ended on its own. + Finished, +} + +struct Shared { + slot: Mutex>, + waker: AtomicWaker, +} + +impl Shared { + /// End the stream with `error`: drop the inner stream and the lease now, and wake the + /// consumer so that it sees the error when it next polls. + fn terminate(&self, error: FenceError) { + let ended = { + let mut slot = self.slot.lock(); + match &*slot { + Slot::Live { .. } => std::mem::replace(&mut *slot, Slot::Terminated(Some(error))), + _ => return, + } + }; + // Drop the inner stream and release the lease outside the slot lock. + drop(ended); + self.waker.wake(); + } +} + +/// A stream that holds a read lease and ends with a typed fence error when the gate stops +/// admitting streams or the read drain cancels it (see the module documentation). +pub struct FencedStream { + shared: Arc>, + map_err: F, +} + +impl FencedStream +where + S: Stream> + Unpin + Send + 'static, + F: Fn(FenceError) -> E + Unpin, +{ + /// Admit `stream` as a stream of `kind` on `fence`. Refused, with the gate's error, when the + /// gate does not admit streams. `map_err` turns the terminal fence error into the stream's + /// error type. Must be called within a Tokio runtime (the watcher is a task). + pub fn new( + fence: &Arc, + kind: LeaseKind, + stream: S, + map_err: F, + ) -> Result { + let (lease, cancel) = acquire_stream_lease(fence, kind)?; + let (stop_tx, stop_rx) = oneshot::channel(); + let shared = Arc::new(Shared { + slot: Mutex::new(Slot::Live { + stream, + _lease: lease, + _stop_watcher: stop_tx, + }), + waker: AtomicWaker::new(), + }); + let watched = Arc::downgrade(&shared); + let fence = fence.clone(); + tokio::spawn(async move { + let error = tokio::select! { + e = gate_denies_streams(&fence) => e, + _ = cancel.cancelled() => cancelled_by_read_fence(kind), + _ = stop_rx => return, + }; + if let Some(shared) = watched.upgrade() { + tracing::debug!( + namespace = %fence.namespace(), + "{kind:?} stream ended by the namespace fence: {error}" + ); + shared.terminate(error); + } + }); + Ok(Self { shared, map_err }) + } +} + +impl Stream for FencedStream +where + S: Stream> + Unpin, + F: Fn(FenceError) -> E + Unpin, +{ + type Item = Result; + + fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + let this = self.get_mut(); + this.shared.waker.register(cx.waker()); + let mut slot = this.shared.slot.lock(); + match &mut *slot { + Slot::Live { stream, .. } => match Pin::new(stream).poll_next(cx) { + Poll::Ready(None) => { + // Release the lease as soon as the stream has nothing more to serve. + let ended = std::mem::replace(&mut *slot, Slot::Finished); + drop(slot); + drop(ended); + Poll::Ready(None) + } + other => other, + }, + Slot::Terminated(error) => match error.take() { + Some(error) => Poll::Ready(Some(Err((this.map_err)(error)))), + None => Poll::Ready(None), + }, + Slot::Finished => Poll::Ready(None), + } + } +} + +/// The gRPC status of a fence error on a replication call. +pub fn fence_status(error: FenceError) -> tonic::Status { + error + .to_grpc_status() + .unwrap_or_else(|| tonic::Status::failed_precondition(error.to_string())) +} + +#[cfg(test)] +mod tests { + use bytes::Bytes; + use futures::{Stream, StreamExt}; + use libsql_replication::rpc::replication::replication_log_server::ReplicationLog; + use libsql_replication::rpc::replication::{ + Frame, HelloRequest, LogOffset, NAMESPACE_METADATA_KEY, SESSION_TOKEN_KEY, + }; + use tonic::metadata::{AsciiMetadataValue, BinaryMetadataValue}; + + use super::*; + use crate::error::Error; + use crate::http::user::dump::dump_stream; + use crate::namespace::fence::drain::tests::{fence_outcome, raw, Source, LONG, OP, PROMPT}; + use crate::namespace::fence::read::tests::{fenced_source, read_fence, NOW}; + use crate::namespace::fence::state::FenceState; + use crate::rpc::replication::replication_log::ReplicationLogService; + + /// A write-fenced source whose table holds one row of a megabyte (two in the dump's hex), + /// far more than the dump pipe buffers, so a dump whose peer stops reading blocks in the + /// middle of writing that row, where only the pipe's cancel can stop it. + async fn large_fenced_source() -> Source { + let s = Source::new().await; + raw( + &s.conn().await, + "insert into t values (randomblob(1000000))", + ) + .await + .unwrap(); + let acquired = s.execute(s.acquire(OP, 1, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&acquired), FenceOutcome::Applied); + s + } + + async fn dump(s: &Source) -> crate::Result>> { + let maker = s + .store + .with("ns".into(), |ns| ns.db.connection_maker()) + .await + .unwrap(); + dump_stream(&s.fence, maker, false).await + } + + /// Read `stream` to its end: the bytes, and the error it ended with, if any. + async fn read_to_end( + mut stream: impl Stream> + Unpin, + ) -> (Vec, Option) { + let mut out = Vec::new(); + while let Some(chunk) = tokio::time::timeout(PROMPT, stream.next()) + .await + .expect("the dump stream stalled") + { + match chunk { + Ok(bytes) => out.extend_from_slice(&bytes), + Err(e) => return (out, Some(e)), + } + } + (out, None) + } + + fn assert_read_fenced(e: &Error) { + match e { + Error::NamespaceFence(f) => assert_eq!(f.outcome(), FenceOutcome::MigrationReadFenced), + other => panic!("expected MIGRATION_READ_FENCED, got {other:?}"), + } + } + + fn ends_with_commit(dump: &[u8]) -> bool { + String::from_utf8_lossy(dump) + .trim_end() + .ends_with("COMMIT;") + } + + /// A dump whose peer has stopped reading is cancelled at the read drain's deadline: its + /// lease is released without the peer reading anything more, the fence is acknowledged, + /// and the body ends with the fence error, never with `COMMIT;`. + #[tokio::test(flavor = "multi_thread")] + async fn dump_lease_released_on_cancel() { + let s = large_fenced_source().await; + let mut stream = Box::pin(dump(&s).await.unwrap()); + // Read until the row has begun: the export is then inside the write of a row that does + // not fit in the pipe, and blocked on it. + let mut head = Vec::new(); + while !String::from_utf8_lossy(&head).contains("INSERT INTO") { + head.extend_from_slice(&stream.next().await.unwrap().unwrap()); + } + assert_eq!(s.fence.read_lease_counts().dump, 1); + + let fenced = tokio::time::timeout(PROMPT, s.execute(read_fence(&s, 2, NOW))) + .await + .expect("the dump's lease was never released") + .unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + assert_eq!(s.fence.gate().state(), FenceState::SourceReadFenced); + assert_eq!(s.fence.read_lease_counts().total(), 0); + + let (rest, error) = read_to_end(stream).await; + assert_read_fenced(&error.expect("a cancelled dump must end with an error")); + let mut body = head; + body.extend_from_slice(&rest); + assert!(!ends_with_commit(&body)); + assert!(!String::from_utf8_lossy(&body).contains("COMMIT;")); + } + + /// The read drain waits for a running dump rather than for time: the fence is acknowledged + /// only once the dump has completed and released its lease. + #[tokio::test(flavor = "multi_thread")] + async fn read_fence_waits_for_dump() { + let s = large_fenced_source().await; + let mut stream = Box::pin(dump(&s).await.unwrap()); + let first = stream.next().await.unwrap().unwrap(); + + let fencing = s.execute(read_fence(&s, 2, LONG)); + s.until_state(FenceState::SourceReadDraining).await; + assert!(!fencing.is_finished()); + assert_eq!(s.fence.read_lease_counts().dump, 1); + // New dumps are refused while the drain waits for this one. + assert_read_fenced(&dump(&s).await.err().unwrap()); + + let (rest, error) = read_to_end(stream).await; + assert!(error.is_none(), "{error:?}"); + let mut body = first.to_vec(); + body.extend_from_slice(&rest); + assert!(ends_with_commit(&body)); + + let fenced = tokio::time::timeout(PROMPT, fencing) + .await + .unwrap() + .unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + assert_eq!(s.fence.read_lease_counts().total(), 0); + } + + /// A dump requested while reads are fenced is refused before any connection is created. + #[tokio::test] + async fn dump_refused_while_read_fenced() { + let s = fenced_source().await; + let fenced = s.execute(read_fence(&s, 2, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + assert_read_fenced(&dump(&s).await.err().unwrap()); + assert_eq!(s.fence.read_lease_counts().total(), 0); + } + + struct Replication { + service: ReplicationLogService, + token: Option, + } + + impl Replication { + /// The internal replication service on `s`, after a `hello`. + async fn new(s: &Source) -> Self { + let mut this = Self { + service: ReplicationLogService::new( + s.store.clone(), + None, + None, + false, + false, + true, + ), + token: None, + }; + let hello = this + .service + .hello(this.request(HelloRequest { + handshake_version: Some(1), + })) + .await + .unwrap() + .into_inner(); + this.token = Some(AsciiMetadataValue::try_from(&hello.session_token[..]).unwrap()); + this + } + + fn request(&self, msg: T) -> tonic::Request { + let mut req = tonic::Request::new(msg); + req.metadata_mut().insert_bin( + NAMESPACE_METADATA_KEY, + BinaryMetadataValue::from_bytes(b"ns"), + ); + if let Some(token) = &self.token { + req.metadata_mut().insert(SESSION_TOKEN_KEY, token.clone()); + } + req + } + + fn offset(&self, next_offset: u64) -> tonic::Request { + self.request(LogOffset { + next_offset, + wal_flavor: None, + }) + } + } + + fn assert_read_fenced_status(status: &tonic::Status) { + assert_eq!(status.code(), tonic::Code::FailedPrecondition, "{status:?}"); + assert_eq!( + FenceError::outcome_from_grpc_status(status), + Some(FenceOutcome::MigrationReadFenced), + "{status:?}" + ); + } + + async fn next_frame( + stream: &mut (impl Stream> + Unpin), + ) -> Option> { + tokio::time::timeout(PROMPT, stream.next()) + .await + .expect("the replication stream stalled") + } + + async fn until_released(fence: &FenceController) { + tokio::time::timeout(PROMPT, async { + loop { + let released = fence.read_released().notified(); + if fence.read_lease_counts().total() == 0 { + break; + } + released.await; + } + }) + .await + .expect("the stream's lease was never released") + } + + /// A tailing `log_entries` stream on a write-fenced source is served, holds a replication + /// lease, and is ended by the read fence with a typed terminal status, without the drain + /// having to wait for its deadline. + #[tokio::test(flavor = "multi_thread")] + async fn log_entries_stream_ends_typed() { + let s = fenced_source().await; + let r = Replication::new(&s).await; + let mut stream = r + .service + .log_entries(r.offset(1)) + .await + .unwrap() + .into_inner(); + assert!(next_frame(&mut stream).await.unwrap().is_ok()); + assert_eq!(s.fence.read_lease_counts().replication, 1); + + let fenced = tokio::time::timeout(PROMPT, s.execute(read_fence(&s, 2, LONG))) + .await + .expect("the tailing stream kept the drain waiting") + .unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + assert_eq!(s.fence.read_lease_counts().total(), 0); + + let last = next_frame(&mut stream).await.unwrap(); + assert_read_fenced_status(&last.unwrap_err()); + assert!(next_frame(&mut stream).await.is_none()); + } + + /// A stream whose peer never reads again (a dead peer) releases its lease anyway: the + /// watcher drops the inner stream and the lease without the stream being polled. + #[tokio::test(flavor = "multi_thread")] + async fn stream_lease_released_without_peer_read() { + let s = fenced_source().await; + let r = Replication::new(&s).await; + let stream = r + .service + .log_entries(r.offset(1)) + .await + .unwrap() + .into_inner(); + assert_eq!(s.fence.read_lease_counts().replication, 1); + + let fenced = tokio::time::timeout(PROMPT, s.execute(read_fence(&s, 2, LONG))) + .await + .expect("the unread stream kept the drain waiting") + .unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + until_released(&s.fence).await; + + // The frames the stream had not yet produced are never served. + let mut stream = stream; + assert_read_fenced_status(&next_frame(&mut stream).await.unwrap().unwrap_err()); + assert!(next_frame(&mut stream).await.is_none()); + } + + /// A `snapshot` stream is ended by the read fence the same way. + #[tokio::test(flavor = "multi_thread")] + async fn snapshot_stream_ends_typed() { + // Compact at once, so the log's frames are in a snapshot. + let s = Source::with_max_log_size(0).await; + raw(&s.conn().await, "insert into t values (1), (2)") + .await + .unwrap(); + // The periodic compaction may already have run; either way there is nothing left to + // compact once this returns. + let logger = s.logger.clone(); + tokio::task::spawn_blocking(move || logger.maybe_compact()) + .await + .unwrap() + .unwrap(); + tokio::time::timeout(PROMPT, async { + // The snapshot is written by the compactor's own thread. + while s.logger.get_snapshot_file(1).await.unwrap().is_none() { + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + }) + .await + .expect("the snapshot was never written"); + let acquired = s.execute(s.acquire(OP, 1, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&acquired), FenceOutcome::Applied); + + let r = Replication::new(&s).await; + let stream = r.service.snapshot(r.offset(1)).await.unwrap().into_inner(); + assert_eq!(s.fence.read_lease_counts().replication, 1); + + let fenced = tokio::time::timeout(PROMPT, s.execute(read_fence(&s, 2, LONG))) + .await + .expect("the snapshot stream kept the drain waiting") + .unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + until_released(&s.fence).await; + + let mut stream = stream; + assert_read_fenced_status(&next_frame(&mut stream).await.unwrap().unwrap_err()); + assert!(next_frame(&mut stream).await.is_none()); + } + + /// While reads are fenced every replication call is refused at its start with the typed + /// status (never `UNAVAILABLE`), and no lease is taken. + #[tokio::test] + async fn replication_calls_denied_while_read_fenced() { + let s = fenced_source().await; + let r = Replication::new(&s).await; + let fenced = s.execute(read_fence(&s, 2, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + + let hello = r + .service + .hello(r.request(HelloRequest { + handshake_version: Some(1), + })) + .await; + assert_read_fenced_status(&hello.unwrap_err()); + assert_read_fenced_status(&r.service.log_entries(r.offset(1)).await.err().unwrap()); + assert_read_fenced_status( + &r.service + .batch_log_entries(r.offset(1)) + .await + .err() + .unwrap(), + ); + assert_read_fenced_status(&r.service.snapshot(r.offset(1)).await.err().unwrap()); + assert_eq!(s.fence.read_lease_counts().total(), 0); + } + + /// The deadline cancel of a stream (the backstop when the gate has not ended it) also ends + /// it with the typed status and releases the lease without the stream being polled. + #[tokio::test] + async fn read_fence_forced_termination() { + let s = fenced_source().await; + let (_tx, rx) = tokio::sync::mpsc::channel::>(1); + let mut stream = FencedStream::new( + &s.fence, + LeaseKind::Replication, + tokio_stream::wrappers::ReceiverStream::new(rx), + fence_status as fn(FenceError) -> tonic::Status, + ) + .unwrap(); + assert_eq!(s.fence.cancel_read_leases(), 1); + until_released(&s.fence).await; + let status = tokio::time::timeout(PROMPT, stream.next()) + .await + .unwrap() + .unwrap() + .unwrap_err(); + assert_read_fenced_status(&status); + assert!(stream.next().await.is_none()); + } +} diff --git a/libsql-server/src/namespace/store.rs b/libsql-server/src/namespace/store.rs index a472fc39c3..ca8e7031be 100644 --- a/libsql-server/src/namespace/store.rs +++ b/libsql-server/src/namespace/store.rs @@ -632,6 +632,13 @@ pub(crate) mod fence_tests { const OP: Uuid = Uuid::from_u128(0xa); pub(crate) async fn open_store(dir: &Path) -> NamespaceStore { + open_store_with_max_log_size(dir, 1_000_000_000).await + } + + pub(crate) async fn open_store_with_max_log_size( + dir: &Path, + max_log_size: u64, + ) -> NamespaceStore { let (maker, manager) = metastore_connection_maker(None, dir).await.unwrap(); let meta = MetaStore::new( MetaStoreConfig { @@ -660,7 +667,7 @@ pub(crate) mod fence_tests { disable_intelligent_throttling: false, }, PrimaryConfig { - max_log_size: 1_000_000_000, + max_log_size, max_log_duration: None, bottomless_replication: None, scripted_backup: None, diff --git a/libsql-server/src/rpc/mod.rs b/libsql-server/src/rpc/mod.rs index 2936a8742c..af74fa4b4c 100644 --- a/libsql-server/src/rpc/mod.rs +++ b/libsql-server/src/rpc/mod.rs @@ -29,6 +29,7 @@ pub async fn run_rpc_server( maybe_tls: Option, idle_shutdown_layer: Option, service: BoxReplicationService, + http2_keepalive_interval: Option, ) -> anyhow::Result<()> { if let Some(tls_config) = maybe_tls { let cert_pem = tokio::fs::read_to_string(&tls_config.cert).await?; @@ -81,8 +82,13 @@ pub async fn run_rpc_server( .service(router); tracing::info!("serving internal rpc server with tls"); - let h2c = crate::h2c::H2cMaker::new(svc); - hyper::server::Server::builder(acceptor).serve(h2c).await?; + let h2c = crate::h2c::H2cMaker::new(svc).with_http2_keepalive(http2_keepalive_interval); + crate::h2c::with_http2_keepalive( + hyper::server::Server::builder(acceptor), + http2_keepalive_interval, + ) + .serve(h2c) + .await?; } else { let proxy = ProxyServer::new(proxy_service); let replication = ReplicationLogServer::new(service); @@ -105,11 +111,16 @@ pub async fn run_rpc_server( ) .service(router); - let h2c = crate::h2c::H2cMaker::new(svc); + let h2c = crate::h2c::H2cMaker::new(svc).with_http2_keepalive(http2_keepalive_interval); tracing::info!("serving internal rpc server without tls"); - hyper::server::Server::builder(acceptor).serve(h2c).await?; + crate::h2c::with_http2_keepalive( + hyper::server::Server::builder(acceptor), + http2_keepalive_interval, + ) + .serve(h2c) + .await?; } Ok(()) } diff --git a/libsql-server/src/rpc/replication/replication_log.rs b/libsql-server/src/rpc/replication/replication_log.rs index 1f7585f34f..4457c49e7a 100644 --- a/libsql-server/src/rpc/replication/replication_log.rs +++ b/libsql-server/src/rpc/replication/replication_log.rs @@ -1,7 +1,8 @@ -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; use std::net::SocketAddr; use std::pin::Pin; use std::sync::{Arc, RwLock}; +use std::time::{Duration, Instant}; use bytes::Bytes; use chrono::{DateTime, Utc}; @@ -22,6 +23,10 @@ use uuid::Uuid; use crate::auth::Auth; use crate::connection::config::DatabaseConfig; +use crate::namespace::fence::controller::{FenceController, LeaseKind}; +use crate::namespace::fence::outcome::FenceError; +use crate::namespace::fence::state::OperationClass; +use crate::namespace::fence::stream::{fence_status, FencedStream}; use crate::namespace::{NamespaceName, NamespaceStore}; use crate::replication::primary::frame_stream::FrameStream; use crate::replication::{LogReadError, ReplicationLogger}; @@ -44,10 +49,16 @@ pub struct ReplicationLogService { //deprecated: generation_id: Uuid, replicas_with_hello: RwLock>, + /// When a fence denial of a replication call was last logged, per namespace. + fence_denials_logged: parking_lot::Mutex>, } pub const MAX_FRAMES_PER_BATCH: usize = 1024; +/// Replication calls denied by a namespace fence are logged at most this often per namespace +/// (they are all counted): replicas that do not understand the typed code reconnect in a loop. +const FENCE_DENIAL_LOG_INTERVAL: Duration = Duration::from_secs(60); + impl ReplicationLogService { pub fn new( namespaces: NamespaceStore, @@ -68,7 +79,60 @@ impl ReplicationLogService { generation_id: Uuid::new_v4(), replicas_with_hello: Default::default(), service_internal, + fence_denials_logged: Default::default(), + } + } + + /// The status of a replication call (or stream) that the namespace fence refuses. Counted + /// every time, and logged at most once per namespace per [`FENCE_DENIAL_LOG_INTERVAL`]. + fn fence_denied(&self, namespace: &NamespaceName, call: &str, error: FenceError) -> Status { + metrics::increment_counter!( + "libsql_server_fence_denials_total", + "code" => error.outcome().as_str(), + "surface" => "replication", + ); + let now = Instant::now(); + let log = { + let mut logged = self.fence_denials_logged.lock(); + match logged.get(namespace) { + Some(at) if now.duration_since(*at) < FENCE_DENIAL_LOG_INTERVAL => false, + _ => { + logged.insert(namespace.clone(), now); + true + } + } + }; + if log { + tracing::warn!( + namespace = %namespace, + internal = self.service_internal, + "replication {call} refused by the namespace fence: {error} \ + (further refusals for this namespace are not logged for a minute)" + ); } + fence_status(error) + } + + /// Admit `stream`, a replication stream on `namespace`, under the namespace's fence + /// (`docs/NAMESPACE_FENCE.md` section 9): it holds a replication read lease and ends with a + /// typed `FAILED_PRECONDITION` when the gate stops admitting streams. + fn fenced_stream( + &self, + namespace: &NamespaceName, + call: &str, + fence: &Arc, + stream: S, + ) -> Result Status>, Status> + where + S: futures::Stream> + Unpin + Send + 'static, + { + FencedStream::new( + fence, + LeaseKind::Replication, + stream, + fence_status as fn(FenceError) -> Status, + ) + .map_err(|e| self.fence_denied(namespace, call, e)) } async fn authenticate( @@ -124,9 +188,12 @@ impl ReplicationLogService { Ok(()) } + /// The namespace's replication log and what goes with it, for a replication `call`. Refused + /// with the typed fence status when the namespace's fence does not admit streams. async fn logger_from_namespace( &self, namespace: NamespaceName, + call: &str, req: &tonic::Request, verify_session: bool, ) -> Result< @@ -136,12 +203,13 @@ impl ReplicationLogService { usize, Arc, impl Future, + Arc, ), Status, > { - let (logger, config, version, stats, config_changed) = self + let (logger, config, version, stats, config_changed, fence) = self .namespaces - .with(namespace, |ns| -> Result<_, Status> { + .with(namespace.clone(), |ns| -> Result<_, Status> { let logger = ns .db .logger() @@ -150,23 +218,30 @@ impl ReplicationLogService { let config = ns.config(); let version = ns.config_version(); let stats = ns.stats(); + let fence = ns.fence().clone(); - Ok((logger, config, version, stats, config_changed)) + Ok((logger, config, version, stats, config_changed, fence)) }) .await - .map_err(|e| { - if let crate::error::Error::NamespaceDoesntExist(_) = e.as_ref() { + .map_err(|e| match e.as_ref() { + crate::error::Error::NamespaceDoesntExist(_) => { Status::failed_precondition(NAMESPACE_DOESNT_EXIST) - } else { - Status::internal(e.to_string()) } + crate::error::Error::NamespaceFence(f) => { + self.fence_denied(&namespace, call, f.clone()) + } + _ => Status::internal(e.to_string()), })??; + fence + .permits(OperationClass::Stream) + .map_err(|e| self.fence_denied(&namespace, call, e))?; + if verify_session { self.verify_session_token(req, version)?; } - Ok((logger, config, version, stats, config_changed)) + Ok((logger, config, version, stats, config_changed, fence)) } fn encode_session_token(&self, version: usize) -> Uuid { @@ -254,8 +329,9 @@ impl ReplicationLog for ReplicationLogService { self.authenticate(&req, namespace.clone()).await?; - let (logger, _, _, stats, config_changed) = - self.logger_from_namespace(namespace, &req, true).await?; + let (logger, _, _, stats, config_changed, fence) = self + .logger_from_namespace(namespace.clone(), "log_entries", &req, true) + .await?; let stats = if self.collect_stats { Some(stats) @@ -288,6 +364,8 @@ impl ReplicationLog for ReplicationLogService { } }; + let stream = self.fenced_stream(&namespace, "log_entries", &fence, Box::pin(stream))?; + Ok(tonic::Response::new(Box::pin(stream))) } @@ -301,7 +379,9 @@ impl ReplicationLog for ReplicationLogService { let namespace = super::super::extract_namespace(self.disable_namespaces, &req)?; self.authenticate(&req, namespace.clone()).await?; - let (logger, _, _, stats, _) = self.logger_from_namespace(namespace, &req, true).await?; + let (logger, _, _, stats, _, fence) = self + .logger_from_namespace(namespace.clone(), "batch_log_entries", &req, true) + .await?; let stats = if self.collect_stats { Some(stats) @@ -322,9 +402,12 @@ impl ReplicationLog for ReplicationLogService { .map_err(|e| Status::internal(e.to_string()))?, self.idle_shutdown_layer.clone(), ) - .map(map_frame_stream_output) - .collect::, _>>() - .await?; + .map(map_frame_stream_output); + // The batch is read under a replication read lease, so a read drain waits for it. + let frames = self + .fenced_stream(&namespace, "batch_log_entries", &fence, frames)? + .collect::, _>>() + .await?; Ok(tonic::Response::new(Frames { frames })) } @@ -348,8 +431,9 @@ impl ReplicationLog for ReplicationLogService { guard.insert((replica_addr, namespace.clone())); } } - let (logger, config, version, _, _) = - self.logger_from_namespace(namespace, &req, false).await?; + let (logger, config, version, _, _, _) = self + .logger_from_namespace(namespace, "hello", &req, false) + .await?; // If we are a shared schema and serving externally (aka to embedded replica's) then // return an error. @@ -383,7 +467,9 @@ impl ReplicationLog for ReplicationLogService { let namespace = super::super::extract_namespace(self.disable_namespaces, &req)?; self.authenticate(&req, namespace.clone()).await?; - let (logger, _, _, stats, _) = self.logger_from_namespace(namespace, &req, true).await?; + let (logger, _, _, stats, _, fence) = self + .logger_from_namespace(namespace.clone(), "snapshot", &req, true) + .await?; let stats = if self.collect_stats { Some(stats) @@ -395,9 +481,17 @@ impl ReplicationLog for ReplicationLogService { let offset = req.next_offset; match logger.get_snapshot_file(offset).await { - Ok(Some(snapshot)) => Ok(tonic::Response::new(Box::pin( - snapshot_stream::make_snapshot_stream(snapshot, offset, stats), - ))), + Ok(Some(snapshot)) => { + let stream = self.fenced_stream( + &namespace, + "snapshot", + &fence, + Box::pin(snapshot_stream::make_snapshot_stream( + snapshot, offset, stats, + )), + )?; + Ok(tonic::Response::new(Box::pin(stream))) + } Ok(None) => Err(Status::new(tonic::Code::Unavailable, "snapshot not found")), Err(e) => Err(Status::new(tonic::Code::Internal, e.to_string())), } From bd0052eec0f7d14267d135a4df820cd4207e8770 Mon Sep 17 00:00:00 2001 From: River Date: Tue, 29 Sep 2026 19:48:01 +0000 Subject: [PATCH 3/3] libsql-server: make the fenced snapshot stream test independent of compaction timing `snapshot_stream_ends_typed` polled `get_snapshot_file(1)` and unwrapped its result while the compactor was still writing the first snapshot and the snapshot merger could be replacing merged files. The lookup lists the snapshot directory and then opens the chosen file, so under load it could fail with `NotFound` (directory not created yet, or a file removed by a merge between the listing and the open) and the test panicked. Open the stream through the `snapshot` RPC itself in a bounded loop that treats only "snapshot not found" and the vanished-file error as "not yet", and asserts that no read lease is held after a failed attempt. An opened stream holds its file open, so a later merge cannot affect it. Co-authored-by: Tomasz Szymczyszyn --- libsql-server/src/namespace/fence/stream.rs | 46 +++++++++++++++++---- 1 file changed, 37 insertions(+), 9 deletions(-) diff --git a/libsql-server/src/namespace/fence/stream.rs b/libsql-server/src/namespace/fence/stream.rs index e57ba2f181..2a043d89fa 100644 --- a/libsql-server/src/namespace/fence/stream.rs +++ b/libsql-server/src/namespace/fence/stream.rs @@ -221,6 +221,7 @@ pub fn fence_status(error: FenceError) -> tonic::Status { #[cfg(test)] mod tests { use bytes::Bytes; + use futures::stream::BoxStream; use futures::{Stream, StreamExt}; use libsql_replication::rpc::replication::replication_log_server::ReplicationLog; use libsql_replication::rpc::replication::{ @@ -513,19 +514,11 @@ mod tests { .await .unwrap() .unwrap(); - tokio::time::timeout(PROMPT, async { - // The snapshot is written by the compactor's own thread. - while s.logger.get_snapshot_file(1).await.unwrap().is_none() { - tokio::time::sleep(std::time::Duration::from_millis(10)).await; - } - }) - .await - .expect("the snapshot was never written"); let acquired = s.execute(s.acquire(OP, 1, LONG)).await.unwrap(); assert_eq!(fence_outcome(&acquired), FenceOutcome::Applied); let r = Replication::new(&s).await; - let stream = r.service.snapshot(r.offset(1)).await.unwrap().into_inner(); + let stream = open_snapshot_stream(&s, &r).await; assert_eq!(s.fence.read_lease_counts().replication, 1); let fenced = tokio::time::timeout(PROMPT, s.execute(read_fence(&s, 2, LONG))) @@ -540,6 +533,41 @@ mod tests { assert!(next_frame(&mut stream).await.is_none()); } + /// Opens a `snapshot` stream from frame 1 once the compactor has written a snapshot. + /// + /// The snapshot is written, and possibly merged with the previous one, by the compactor's + /// own tasks, concurrently with this call. The server looks a snapshot up by listing the + /// snapshot directory and then opening the file it chose, so while the directory does not + /// exist yet, or while a merge removes the files it replaced, the call can fail with + /// "snapshot not found" or with the `NotFound` of the vanished file. Those two failures, + /// and only those, mean "not yet": any other status (a fence refusal in particular) fails + /// the test. A stream that was opened holds its file open, so a later merge cannot affect it. + async fn open_snapshot_stream( + s: &Source, + r: &Replication, + ) -> BoxStream<'static, Result> { + tokio::time::timeout(PROMPT, async { + loop { + match r.service.snapshot(r.offset(1)).await { + Ok(stream) => break stream.into_inner(), + Err(status) => { + let not_yet = match status.code() { + tonic::Code::Unavailable => status.message() == "snapshot not found", + tonic::Code::Internal => status.message().contains("os error 2"), + _ => false, + }; + assert!(not_yet, "unexpected snapshot status: {status:?}"); + assert_eq!(s.fence.read_lease_counts().total(), 0); + // A polling interval, not evidence: the loop ends on the condition. + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + } + } + }) + .await + .expect("the snapshot was never written") + } + /// While reads are fenced every replication call is refused at its start with the typed /// status (never `UNAVAILABLE`), and no lease is taken. #[tokio::test]