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/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/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/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/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/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 a1ab51c4c4..02a8011c0f 100644 --- a/libsql-server/src/main.rs +++ b/libsql-server/src/main.rs @@ -274,6 +274,19 @@ 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, + + /// 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, @@ -673,6 +686,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), }) } @@ -747,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/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..3e5c8ba5b8 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; @@ -301,34 +304,40 @@ 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; - 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, - logger: Arc, + pub(crate) store: NamespaceStore, + pub(crate) fence: Arc, + pub(crate) logger: Arc, } impl Source { - async fn new() -> Self { + 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(), @@ -362,7 +371,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 +389,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 +409,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 +428,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 +436,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 +449,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 +465,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..74f227cf23 100644 --- a/libsql-server/src/namespace/fence/mod.rs +++ b/libsql-server/src/namespace/fence/mod.rs @@ -8,9 +8,10 @@ //! 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 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 // following commits of this series; until then most of the module is unused by the rest of @@ -22,10 +23,12 @@ pub mod controller; pub mod drain; pub mod hooks; pub mod outcome; +pub mod read; 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 new file mode 100644 index 0000000000..50feb0b3b1 --- /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)] +pub(crate) 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. + pub(crate) const NOW: DrainPolicy = DrainPolicy { + deadline_ms: 0, + on_deadline: OnDeadline::Fail, + }; + + /// A write-fenced source. + pub(crate) 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, + } + } + + pub(crate) 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/fence/stream.rs b/libsql-server/src/namespace/fence/stream.rs new file mode 100644 index 0000000000..2a043d89fa --- /dev/null +++ b/libsql-server/src/namespace/fence/stream.rs @@ -0,0 +1,622 @@ +//! 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::BoxStream; + 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(); + 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 = 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))) + .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()); + } + + /// 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] + 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/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..ca8e7031be 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 } @@ -610,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 { @@ -638,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())), }