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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
114 changes: 108 additions & 6 deletions libsql-server/src/admin_shell.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -41,34 +43,68 @@ impl AdminShell {
queries: impl Stream<Item = Result<rpc::Query, tonic::Status>>,
) -> anyhow::Result<impl Stream<Item = Result<rpc::Response, tonic::Status>>> {
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<FenceController>,
queries: impl Stream<Item = Result<rpc::Query, tonic::Status>>,
) -> impl Stream<Item = Result<rpc::Response, tonic::Status>> {
async_stream::stream! {
tokio::pin!(queries);
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
}
}
}

/// 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<FenceController>,
conn: &mut rusqlite::Connection,
q: String,
) -> Result<rpc::Response, tonic::Status> {
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<rpc::Response, tonic::Status> {
match try_run_one(conn, q) {
Ok(resp) => Ok(resp),
Expand Down Expand Up @@ -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);
}
}
3 changes: 3 additions & 0 deletions libsql-server/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Duration>,
/// 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<Duration>,
}

#[derive(Debug, Clone)]
Expand Down
67 changes: 66 additions & 1 deletion libsql-server/src/connection/connection_core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -102,6 +103,8 @@ impl<W: Wal + Send + 'static> CoreConnection<W> {
);

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();
Expand Down Expand Up @@ -218,6 +221,22 @@ impl<W: Wal + Send + 'static> CoreConnection<W> {
// 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(
Expand Down Expand Up @@ -272,16 +291,40 @@ impl<W: Wal + Send + 'static> CoreConnection<W> {
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<DescribeResponse> {
// 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}");
Expand Down Expand Up @@ -425,6 +468,28 @@ impl<W: Wal + Send + 'static> CoreConnection<W> {
}
}

/// The schema aliases attached on `conn` (other than `main` and `temp`).
/// `None` when they cannot be listed.
fn attached_schemas<W: Wal>(conn: &libsql_sys::Connection<W>) -> Option<Vec<String>> {
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;
Expand Down
35 changes: 33 additions & 2 deletions libsql-server/src/connection/dump/exporter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,15 +7,30 @@ use anyhow::bail;
use rusqlite::types::ValueRef;
use rusqlite::OptionalExtension;

struct DumpState<W: Write> {
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<W: Write> DumpState<W> {
impl<W: Write> 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,
Expand All @@ -25,6 +40,7 @@ impl<W: Write> DumpState<W> {
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")
};
Expand Down Expand Up @@ -103,6 +119,7 @@ impl<W: Write> DumpState<W> {
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)?)?;
Expand All @@ -128,6 +145,7 @@ impl<W: Write> DumpState<W> {
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")
};
Expand Down Expand Up @@ -437,13 +455,25 @@ 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", ())?;
let savepoint = txn.savepoint_with_name("dump")?;
let mut state = DumpState {
writable_schema: false,
writer,
cancelled,
};

writeln!(state.writer, "PRAGMA foreign_keys=OFF;")?;
Expand All @@ -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;", ());
Expand Down
2 changes: 1 addition & 1 deletion libsql-server/src/connection/legacy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand Down
9 changes: 7 additions & 2 deletions libsql-server/src/connection/program.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading
Loading