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
91 changes: 91 additions & 0 deletions libsql-server/src/admin_shell.rs
Original file line number Diff line number Diff line change
Expand Up @@ -343,4 +343,95 @@ mod fence_tests {
);
assert_eq!(s.fence.read_lease_counts().total(), 0);
}

/// Section 17 row 12: a quarantined target is written only through its operation's import
/// capability. The admin shell, which reaches the namespace with admin authority and runs raw
/// SQL, can neither read nor write it: every query is refused by the read admission with the
/// quarantine code, and a raw write that skipped the admission is still refused at the WAL.
/// The import session keeps working, and nothing the shell sent changed the data.
#[tokio::test(flavor = "multi_thread")]
async fn admin_shell_cannot_write_quarantined() {
use crate::namespace::fence::target::tests::{create, create_request, OP as TARGET_OP};
use crate::namespace::open_test_store as open_store;

let dir = tempfile::tempdir().unwrap();
let store = open_store(dir.path()).await;
create(&store, create_request("tgt", 1))
.await
.unwrap()
.unwrap();
let mut session = store
.open_import_session("tgt".into(), TARGET_OP, 1)
.await
.unwrap();
session
.with_raw(|c| c.execute_batch("create table t (x); insert into t values (1)"))
.await
.unwrap()
.unwrap();

let shell = AdminShell::new(store.clone());
let sql = [
"insert into t values (2)",
"delete from t",
"create table u (y)",
"pragma user_version = 7",
"begin immediate",
"select count(*) from t",
];
let queries = tokio_stream::iter(sql.map(|q| Ok(rpc::Query { query: q.into() })));
let responses: Vec<_> = shell
.with_namespace(Bytes::from_static(b"tgt"), queries)
.await
.unwrap()
.collect()
.await;
assert_eq!(responses.len(), sql.len());
for (q, resp) in sql.iter().zip(&responses) {
let resp = resp.as_ref().unwrap();
assert!(
error(resp).starts_with("MIGRATION_TARGET_QUARANTINED"),
"{q}: {}",
error(resp)
);
}

// Without the shell's read admission, the raw write is refused by the WAL gate: the
// connection holds no capability.
let (fence, maker) = store
.with("tgt".into(), |ns| {
(ns.fence().clone(), ns.db.connection_maker())
})
.await
.unwrap();
let conn = maker.create().await.unwrap();
for q in ["insert into t values (3)", "create table v (z)"] {
let resp = conn.with_raw(|c| run_one(c, q.into())).unwrap();
assert!(error(&resp).contains("authoriz"), "{q}: {}", error(&resp));
}
assert_eq!(fence.read_lease_counts().total(), 0);

// The capability still writes, and it sees only its own rows.
let rows: i64 = session
.with_raw(|c| {
c.execute("insert into t values (4)", ())?;
c.query_row("select count(*) from t", (), |r| r.get(0))
})
.await
.unwrap()
.unwrap();
assert_eq!(rows, 2);
let tables: i64 = session
.with_raw(|c| {
c.query_row(
"select count(*) from sqlite_schema where type = 'table'",
(),
|r| r.get(0),
)
})
.await
.unwrap()
.unwrap();
assert_eq!(tables, 1);
}
}
37 changes: 36 additions & 1 deletion libsql-server/src/connection/legacy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ use tokio::time::Duration;
use crate::error::Error;
use crate::metrics::DESCRIBE_COUNT;
use crate::namespace::broadcasters::BroadcasterHandle;
use crate::namespace::fence::capability::MigrationCapability;
use crate::namespace::fence::controller::{FenceConnState, FenceController};
use crate::namespace::fence::state::OperationClass;
use crate::namespace::meta_store::MetaStoreHandle;
Expand Down Expand Up @@ -143,6 +144,33 @@ where

#[tracing::instrument(skip(self))]
pub(super) async fn make_connection(&self) -> Result<LegacyConnection<W>> {
self.make_connection_with(FenceConnState::new(
self.fence.clone(),
OperationClass::NormalWrite,
))
.await
}

/// Open a connection that works under `capability` (an import or validation session,
/// `docs/NAMESPACE_FENCE.md` section 11). It shares the maker's write slot, WAL and
/// replication log with every other connection, and is admitted as the capability's class
/// only while the capability is valid. It is not counted by the connection throttle: it is
/// operation-owned work, and the fence, not the throttle, bounds how much of it runs.
pub(crate) async fn make_capability_connection(
&self,
capability: MigrationCapability,
) -> Result<LegacyConnection<W>> {
self.make_connection_with(FenceConnState::with_capability(
self.fence.clone(),
capability,
))
.await
}

async fn make_connection_with(
&self,
fence: Arc<FenceConnState>,
) -> Result<LegacyConnection<W>> {
LegacyConnection::new(
self.db_path.clone(),
self.extensions.clone(),
Expand All @@ -161,7 +189,7 @@ where
self.resolve_attach_path.clone(),
self.connection_manager.clone(),
self.make_wal_manager.clone(),
FenceConnState::new(self.fence.clone(), OperationClass::NormalWrite),
fence,
)
.await
}
Expand Down Expand Up @@ -213,6 +241,13 @@ impl LegacyConnection<libsql_sys::wal::wrapper::PassthroughWalWrapper> {
}
}

impl<T> LegacyConnection<T> {
/// The fence state shared by this connection's WAL wrapper and `CoreConnection`.
pub(crate) fn fence_state(&self) -> &Arc<FenceConnState> {
&self.fence
}
}

impl<T> Clone for LegacyConnection<T> {
fn clone(&self) -> Self {
Self {
Expand Down
5 changes: 5 additions & 0 deletions libsql-server/src/connection/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -286,6 +286,11 @@ pub struct MakeThrottledConnection<F> {
}

impl<F> MakeThrottledConnection<F> {
/// The connection maker this one throttles.
pub(crate) fn inner(&self) -> &F {
&self.connection_maker
}

fn new(
semaphore: Arc<Semaphore>,
connection_maker: F,
Expand Down
81 changes: 45 additions & 36 deletions libsql-server/src/namespace/configurator/helpers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -301,6 +301,16 @@ async fn run_periodic_compactions(logger: Arc<ReplicationLogger>) -> anyhow::Res
}

async fn load_dump<S>(dump: S, conn: PrimaryConnection) -> crate::Result<(), LoadDumpError>
where
S: Stream<Item = std::io::Result<Bytes>> + Unpin,
{
let dump_content = read_dump(dump).await?;
tokio::task::spawn_blocking(move || conn.with_raw(|conn| load_dump_sql(&dump_content, conn)))
.await?
}

/// Read a whole dump into memory.
pub(crate) async fn read_dump<S>(dump: S) -> crate::Result<String, LoadDumpError>
where
S: Stream<Item = std::io::Result<Bytes>> + Unpin,
{
Expand All @@ -310,13 +320,36 @@ where
.read_to_string(&mut dump_content)
.await
.map_err(|e| LoadDumpError::Internal(format!("Failed to read dump content: {}", e)))?;
Ok(dump_content)
}

/// Parse `dump_content` and run its statements, one at a time, on `conn`. The dump must run
/// inside one transaction that it commits itself; `ATTACH` is refused. This is the loader both
/// for a namespace created from a dump and for an import session into a quarantined migration
/// target, which runs it under its capability (`docs/NAMESPACE_FENCE.md` section 11).
pub(crate) fn load_dump_sql(
dump_content: &str,
conn: &mut rusqlite::Connection,
) -> crate::Result<(), LoadDumpError> {
if dump_content.to_lowercase().contains("attach") {
return Err(LoadDumpError::InvalidSqlInput(
"attach statements are not allowed in dumps".to_string(),
));
}

conn.authorizer(Some(|auth: AuthContext<'_>| match auth.action {
AuthAction::Attach { filename: _ } => Authorization::Deny,
_ => Authorization::Allow,
}));
let result = run_dump_statements(dump_content, conn);
conn.authorizer(None::<fn(AuthContext<'_>) -> Authorization>);
result
}

fn run_dump_statements(
dump_content: &str,
conn: &mut rusqlite::Connection,
) -> crate::Result<(), LoadDumpError> {
let mut parser = Box::new(Parser::new(dump_content.as_bytes()));
let mut skipped_wasm_table = false;
let mut n_stmt = 0;
Expand All @@ -336,37 +369,20 @@ where
}
}

if n_stmt > 2 && conn.is_autocommit().await.unwrap() {
if n_stmt > 2 && conn.is_autocommit() {
return Err(LoadDumpError::NoTxn);
}

let stmt_sql = cmd.to_string();
tokio::task::spawn_blocking({
let conn = conn.clone();
move || -> crate::Result<(), LoadDumpError> {
conn.with_raw(|conn| {
conn.authorizer(Some(|auth: AuthContext<'_>| match auth.action {
AuthAction::Attach { filename: _ } => Authorization::Deny,
_ => Authorization::Allow,
}));
conn.execute(&stmt_sql, ())
})
.map_err(|e| match e {
rusqlite::Error::SqlInputError {
msg, sql, offset, ..
} => LoadDumpError::InvalidSqlInput(format!(
"msg: {}, sql: {}, offset: {}",
msg, sql, offset
)),
e => LoadDumpError::Internal(format!(
"statement: {}, error: {}",
n_stmt, e
)),
})?;
Ok(())
}
})
.await??;
conn.execute(&stmt_sql, ()).map_err(|e| match e {
rusqlite::Error::SqlInputError {
msg, sql, offset, ..
} => LoadDumpError::InvalidSqlInput(format!(
"msg: {}, sql: {}, offset: {}",
msg, sql, offset
)),
e => LoadDumpError::Internal(format!("statement: {}, error: {}", n_stmt, e)),
})?;
}
Ok(None) => break,
Err(e) => {
Expand All @@ -389,15 +405,8 @@ where
}
}

if !conn.is_autocommit().await.unwrap() {
tokio::task::spawn_blocking({
let conn = conn.clone();
move || -> crate::Result<(), LoadDumpError> {
conn.with_raw(|conn| conn.execute("rollback", ()))?;
Ok(())
}
})
.await??;
if !conn.is_autocommit() {
conn.execute("rollback", ())?;
return Err(LoadDumpError::NoCommit);
}

Expand Down
1 change: 1 addition & 0 deletions libsql-server/src/namespace/configurator/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ mod primary;
mod replica;
mod schema;

pub(crate) use helpers::{load_dump_sql, read_dump};
pub use primary::PrimaryConfigurator;
pub use replica::ReplicaConfigurator;
pub use schema::SchemaConfigurator;
Expand Down
Loading
Loading