From dd3c5a8a6006413ef4494fa9e8759af4459a8f26 Mon Sep 17 00:00:00 2001 From: River Date: Wed, 30 Sep 2026 08:21:03 +0000 Subject: [PATCH 1/2] libsql-server: namespace fence admin API and capability discovery Serve the fence contract of docs/NAMESPACE_FENCE.md section 4 on the admin listener (new http/admin/fence.rs): - GET /v1/fence/capabilities, always served: protocol version, whether fences are enabled, served commands, states, proxy stable_code support, server build and instance id, and the number of active fences counted from the fence registry. - GET /v1/namespaces/:ns/fence (InspectFence): the durable record, the owning operation's receipts (or ?receipts=all), live admission, drain counters and the live replication log id; read-only, and also served while fence tables exist with the flag off. - One POST route per source and target command, all through NamespaceStore::execute_fence_command, and validation-query, which runs one read-only program under the operation's validation capability with a 10 000-row bound. Command routes answer 404 unless --enable-namespace-fence is on, refuse to run without an admin auth key (admin_auth_required) or on a replica (not_primary), parse bodies strictly (invalid_argument), and refuse restore options on target creation (restore_not_allowed). Responses carry the outcome, replay flag, fence view, receipt and drain counters with the outcome's admin status. Co-authored-by: Tomasz Szymczyszyn --- libsql-server/src/http/admin/fence.rs | 1024 +++++++++++++++++ libsql-server/src/http/admin/mod.rs | 6 + .../src/namespace/fence/controller.rs | 31 + libsql-server/src/namespace/fence/mod.rs | 17 + libsql-server/src/namespace/fence/registry.rs | 16 + libsql-server/src/namespace/store.rs | 34 +- libsql-server/tests/common/http.rs | 14 + libsql-server/tests/fence/admin.rs | 731 ++++++++++++ libsql-server/tests/fence/mod.rs | 188 +++ libsql-server/tests/tests.rs | 1 + 10 files changed, 2059 insertions(+), 3 deletions(-) create mode 100644 libsql-server/src/http/admin/fence.rs create mode 100644 libsql-server/tests/fence/admin.rs create mode 100644 libsql-server/tests/fence/mod.rs diff --git a/libsql-server/src/http/admin/fence.rs b/libsql-server/src/http/admin/fence.rs new file mode 100644 index 0000000000..e7587d2781 --- /dev/null +++ b/libsql-server/src/http/admin/fence.rs @@ -0,0 +1,1024 @@ +//! The namespace fence admin API (`docs/NAMESPACE_FENCE.md` section 4). +//! +//! `GET /v1/fence/capabilities` is always served. The other routes answer `404` unless the +//! server was started with `--enable-namespace-fence`, except `InspectFence`, which is also +//! served while fence state exists in the metastore with the flag off (fences are enforced +//! either way, section 13.1). Every mutating route runs through +//! [`NamespaceStore::execute_fence_command`], which owns replay, the drains and target creation. + +use std::sync::Arc; + +use axum::extract::{Path, Query, State}; +use axum::response::{IntoResponse, Response}; +use axum::routing::{get, post}; +use axum::Json; +use bytes::Bytes; +use hyper::StatusCode; +use serde::de::DeserializeOwned; +use serde::Deserialize; +use serde_json::{json, Map, Value}; +use uuid::Uuid; + +use crate::auth::parse_jwt_keys; +use crate::error::Error; +use crate::hrana::proto; +use crate::namespace::fence::command::{ + CommandKind, DrainPolicy, FenceCommand, FenceRequest, OnDeadline, TargetConfig, + ValidationResult, +}; +use crate::namespace::fence::controller::{DrainCounters, FenceController}; +use crate::namespace::fence::outcome::{FenceDetail, FenceError, FenceOutcome}; +use crate::namespace::fence::record::{CommandReceipt, NamespaceFenceRecord, ServerIdentity}; +use crate::namespace::fence::state::FenceState; +use crate::namespace::fence::store::{StoredFence, StoredReceipt}; +use crate::namespace::fence::{server_identity, FENCE_PROTOCOL_VERSION, PROXY_STABLE_CODE}; +use crate::namespace::meta_store::FenceCommit; +use crate::namespace::NamespaceName; +use crate::net::Connector; + +use super::AppState; + +/// The most rows one `validation-query` request returns, over all of its statements. A query +/// that would return more is refused rather than truncated, so a validation never silently +/// looks at part of a result. +pub const MAX_VALIDATION_QUERY_ROWS: usize = 10_000; + +/// The commands this server serves over the admin API, reported by capability discovery. +/// Adoption is served once its route exists. +const SERVED_COMMANDS: [&str; 11] = [ + "InspectFence", + CommandKind::AcquireSourceWriteFence.as_str(), + CommandKind::SetSourceReadFence.as_str(), + CommandKind::ClearSourceReadFence.as_str(), + CommandKind::ReleaseSourceWriteFence.as_str(), + CommandKind::CreateTargetQuarantined.as_str(), + CommandKind::SealTargetImport.as_str(), + CommandKind::RecordTargetValidation.as_str(), + CommandKind::PublishTargetReadableWriteFenced.as_str(), + CommandKind::EnableTargetWrites.as_str(), + CommandKind::AbortQuarantinedTarget.as_str(), +]; + +/// The fence routes, added to the admin router. +pub(super) fn routes() -> axum::Router>> { + let command = |kind: CommandKind| { + post( + move |State(state): State>>, + Path(namespace): Path, + body: Bytes| async move { + handle_command(state, namespace, kind, body).await + }, + ) + }; + axum::Router::new() + .route("/v1/fence/capabilities", get(handle_capabilities)) + .route("/v1/namespaces/:namespace/fence", get(handle_inspect)) + .route( + "/v1/namespaces/:namespace/fence/source/acquire-write-fence", + command(CommandKind::AcquireSourceWriteFence), + ) + .route( + "/v1/namespaces/:namespace/fence/source/set-read-fence", + command(CommandKind::SetSourceReadFence), + ) + .route( + "/v1/namespaces/:namespace/fence/source/clear-read-fence", + command(CommandKind::ClearSourceReadFence), + ) + .route( + "/v1/namespaces/:namespace/fence/source/release-write-fence", + command(CommandKind::ReleaseSourceWriteFence), + ) + .route( + "/v1/namespaces/:namespace/fence/target/create-quarantined", + command(CommandKind::CreateTargetQuarantined), + ) + .route( + "/v1/namespaces/:namespace/fence/target/seal-import", + command(CommandKind::SealTargetImport), + ) + .route( + "/v1/namespaces/:namespace/fence/target/validation-receipt", + command(CommandKind::RecordTargetValidation), + ) + .route( + "/v1/namespaces/:namespace/fence/target/publish-readable", + command(CommandKind::PublishTargetReadableWriteFenced), + ) + .route( + "/v1/namespaces/:namespace/fence/target/enable-writes", + command(CommandKind::EnableTargetWrites), + ) + .route( + "/v1/namespaces/:namespace/fence/target/abort", + command(CommandKind::AbortQuarantinedTarget), + ) + .route( + "/v1/namespaces/:namespace/fence/target/validation-query", + post(handle_validation_query), + ) +} + +// --------------------------------------------------------------------------------------------- +// Handlers + +async fn handle_capabilities(State(state): State>>) -> Json { + let meta = state.namespaces.meta_store(); + Json(json!({ + "fence_protocol_version": FENCE_PROTOCOL_VERSION, + "enabled": meta.fence_enabled(), + "commands": SERVED_COMMANDS, + "states": FenceState::ALL.iter().map(|s| s.as_str()).collect::>(), + "proxy_stable_code": PROXY_STABLE_CODE, + "server": server_json(&server_identity()), + "active_fences": state.namespaces.active_fences(), + // Metastore restore provenance is not tracked yet. + "metastore": { "restored_from_backup": false, "restored_generation": null }, + })) +} + +#[derive(Debug, Default, Deserialize)] +struct InspectQuery { + #[serde(default)] + receipts: Option, +} + +async fn handle_inspect( + State(state): State>>, + Path(namespace): Path, + Query(query): Query, +) -> Response { + let meta = state.namespaces.meta_store(); + if !meta.fence_enabled() && !meta.fence_enforced() { + return StatusCode::NOT_FOUND.into_response(); + } + let namespace = match NamespaceName::from_string(namespace) { + Ok(ns) => ns, + Err(e) => return invalid_argument(e.to_string()).into_response(), + }; + if !state.namespaces.is_primary() { + return ErrorReply::new(not_primary()).into_response(); + } + let all = match query.receipts.as_deref() { + None => false, + Some("all") => true, + Some(other) => { + return invalid_argument(format!("unknown `receipts` value `{other}`")).into_response() + } + }; + let (inspection, controller) = match state.namespaces.inspect_fence(&namespace).await { + Ok(found) => found, + Err(e) => return fence_or_error(&state, &namespace, e).await, + }; + if matches!( + inspection.fence, + StoredFence::None { + namespace_exists: false + } + ) && controller.as_ref().map_or(true, |c| { + let gate = c.gate(); + matches!( + gate.fence, + StoredFence::None { + namespace_exists: false + } + ) && gate.creating_target.is_none() + }) { + return ( + StatusCode::NOT_FOUND, + Json(json!({ "error": format!("namespace `{namespace}` does not exist") })), + ) + .into_response(); + } + + let owner = inspection + .fence + .record() + .map(|r| r.operation_id.to_string()); + let receipts: Vec = inspection + .receipts + .iter() + .filter(|r| all || owner.as_deref() == Some(r.operation_id.as_str())) + .map(stored_receipt_json) + .collect(); + let body = json!({ + "outcome": FenceOutcome::Applied.as_str(), + "replayed": false, + "fence": fence_json(&namespace, &inspection.fence, controller.as_deref()), + "receipts": receipts, + "drain": drain_json(controller.as_deref()), + }); + (StatusCode::OK, Json(body)).into_response() +} + +async fn handle_command( + state: Arc>, + namespace: String, + kind: CommandKind, + body: Bytes, +) -> Response { + if !state.namespaces.meta_store().fence_enabled() { + return StatusCode::NOT_FOUND.into_response(); + } + let namespace = match NamespaceName::from_string(namespace) { + Ok(ns) => ns, + Err(e) => return invalid_argument(e.to_string()).into_response(), + }; + if let Err(e) = mutating_preconditions(&state) { + return error_reply(&state, &namespace, e).await; + } + let request = match parse_command(namespace.clone(), kind, &body) { + Ok(request) => request, + Err(e) => return error_reply(&state, &namespace, e).await, + }; + match state + .namespaces + .execute_fence_command(request, server_identity()) + .await + { + Ok(commit) => success_reply(&state, &namespace, commit), + Err(e) => fence_or_error(&state, &namespace, e).await, + } +} + +async fn handle_validation_query( + State(state): State>>, + Path(namespace): Path, + body: Bytes, +) -> Response { + if !state.namespaces.meta_store().fence_enabled() { + return StatusCode::NOT_FOUND.into_response(); + } + let namespace = match NamespaceName::from_string(namespace) { + Ok(ns) => ns, + Err(e) => return invalid_argument(e.to_string()).into_response(), + }; + if let Err(e) = mutating_preconditions(&state) { + return error_reply(&state, &namespace, e).await; + } + let query = match parse_validation_query(&body) { + Ok(query) => query, + Err(e) => return error_reply(&state, &namespace, e).await, + }; + if let Some(expected) = query.expected_state { + let current = state + .namespaces + .existing_fence_controller(&namespace) + .map(|c| c.gate().state()); + if let Some(current) = current.filter(|s| *s != expected) { + let e = FenceError::new( + FenceOutcome::FenceRevisionMismatch, + format!("the namespace is {current}, not {expected}"), + ); + return error_reply(&state, &namespace, e).await; + } + } + let mut session = match state + .namespaces + .open_validation_session( + namespace.clone(), + query.operation_id, + query.expected_revision, + ) + .await + { + Ok(session) => session, + Err(e) => return fence_or_error(&state, &namespace, e).await, + }; + + let mut results = Vec::with_capacity(query.stmts.len()); + let mut remaining = MAX_VALIDATION_QUERY_ROWS; + for (index, stmt) in query.stmts.into_iter().enumerate() { + let budget = remaining; + let ran = session + .with_raw(move |conn| run_validation_stmt(conn, &stmt, budget)) + .await; + match ran { + Ok(Ok(result)) => { + remaining -= result.rows.len(); + results.push(result); + } + Ok(Err(e)) => { + return error_reply(&state, &namespace, e.into_fence_error(index)).await; + } + Err(e) => return error_reply(&state, &namespace, e).await, + } + } + drop(session); + + let controller = state.namespaces.existing_fence_controller(&namespace); + let fence = controller + .as_ref() + .map(|c| fence_json(&namespace, &c.gate().fence, Some(c))) + .unwrap_or(Value::Null); + let body = json!({ + "results": results, + "fence": fence, + "drain": drain_json(controller.as_deref()), + }); + (StatusCode::OK, Json(body)).into_response() +} + +// --------------------------------------------------------------------------------------------- +// Preconditions and replies + +/// Section 4.1, for every route that changes fence state or works under a capability. The admin +/// authentication itself is the admin router's middleware and has already run. +fn mutating_preconditions(state: &AppState) -> Result<(), FenceError> { + if !state.namespaces.is_primary() { + return Err(not_primary()); + } + if !state.admin_auth_configured { + return Err(FenceError::new( + FenceOutcome::FencePreconditionFailed, + "namespace fence commands need an admin auth key: without one the admin API is \ + unauthenticated", + ) + .with_detail(FenceDetail::AdminAuthRequired)); + } + Ok(()) +} + +fn not_primary() -> FenceError { + FenceError::new( + FenceOutcome::FencePreconditionFailed, + "namespace fences live on the primary; this server is a replica", + ) + .with_detail(FenceDetail::NotPrimary) +} + +fn invalid_argument(message: impl Into) -> ErrorReply { + ErrorReply::new( + FenceError::new(FenceOutcome::FencePreconditionFailed, message) + .with_detail(FenceDetail::InvalidArgument), + ) +} + +/// An error reply without a fence view (the namespace name itself could not be used). +struct ErrorReply { + error: FenceError, + fence: Value, + drain: Value, +} + +impl ErrorReply { + fn new(error: FenceError) -> Self { + Self { + error, + fence: Value::Null, + drain: Value::Null, + } + } +} + +impl IntoResponse for ErrorReply { + fn into_response(self) -> Response { + let outcome = self.error.outcome(); + let mut body = json!({ + "outcome": outcome.as_str(), + "replayed": false, + "error": self.error.message(), + "fence": self.fence, + "drain": self.drain, + }); + if let Some(detail) = self.error.detail() { + body["detail"] = json!(detail.as_str()); + } + (outcome.admin_http_status(), Json(body)).into_response() + } +} + +/// An error reply carrying the namespace's current fence view: the live gate if the namespace +/// has a controller, otherwise what the metastore holds. +async fn error_reply( + state: &AppState, + namespace: &NamespaceName, + error: FenceError, +) -> Response { + let mut reply = ErrorReply::new(error); + match state.namespaces.existing_fence_controller(namespace) { + Some(controller) => { + reply.fence = fence_json(namespace, &controller.gate().fence, Some(&controller)); + reply.drain = drain_json(Some(&controller)); + } + None => { + if let Ok((inspection, _)) = state.namespaces.inspect_fence(namespace).await { + reply.fence = fence_json(namespace, &inspection.fence, None); + } + } + } + reply.into_response() +} + +/// A fence refusal in the fence response shape; any other error as the admin API reports it. +async fn fence_or_error(state: &AppState, namespace: &NamespaceName, e: Error) -> Response { + match e { + Error::NamespaceFence(e) => error_reply(state, namespace, e).await, + e => e.into_response(), + } +} + +fn success_reply( + state: &AppState, + namespace: &NamespaceName, + commit: FenceCommit, +) -> Response { + let controller = state.namespaces.existing_fence_controller(namespace); + let fence = match (&commit.record, &controller) { + (Some(record), _) => fence_json( + namespace, + &StoredFence::Record(record.clone()), + controller.as_deref(), + ), + (None, Some(c)) => fence_json(namespace, &c.gate().fence, Some(c)), + (None, None) => Value::Null, + }; + let outcome = commit.receipt.outcome; + let body = json!({ + "outcome": outcome.as_str(), + "replayed": commit.kind != crate::namespace::meta_store::FenceCommitKind::Committed, + "fence": fence, + "receipt": receipt_json(&commit.receipt), + "drain": drain_json(controller.as_deref()), + }); + (outcome.admin_http_status(), Json(body)).into_response() +} + +// --------------------------------------------------------------------------------------------- +// Request parsing + +/// A JSON object request body whose fields are taken one by one; whatever is left at the end is +/// an unknown field and refused. +struct Body(Map); + +impl Body { + fn parse(bytes: &[u8]) -> Result { + if bytes.iter().all(u8::is_ascii_whitespace) { + return Err(invalid("the request body must be a JSON object")); + } + match serde_json::from_slice::(bytes) { + Ok(Value::Object(map)) => Ok(Self(map)), + Ok(_) => Err(invalid("the request body must be a JSON object")), + Err(e) => Err(invalid(format!("the request body is not valid JSON: {e}"))), + } + } + + fn has(&self, key: &str) -> bool { + self.0.contains_key(key) + } + + fn opt(&mut self, key: &str) -> Result, FenceError> { + match self.0.remove(key) { + None | Some(Value::Null) => Ok(None), + // Through text rather than `from_value`: some protocol types (the Hrana values of + // a statement) only deserialize from borrowed strings. + Some(v) => serde_json::from_str(&v.to_string()) + .map(Some) + .map_err(|e| invalid(format!("invalid `{key}`: {e}"))), + } + } + + fn req(&mut self, key: &str) -> Result { + self.opt(key)? + .ok_or_else(|| invalid(format!("missing `{key}`"))) + } + + fn uuid(&mut self, key: &str) -> Result { + let s: String = self.req(key)?; + Uuid::parse_str(&s).map_err(|e| invalid(format!("invalid `{key}`: {e}"))) + } + + fn state(&mut self, key: &str) -> Result { + let s: String = self.req(key)?; + s.parse() + .map_err(|e: crate::namespace::fence::state::UnknownFenceState| invalid(e.to_string())) + } + + fn finish(self) -> Result<(), FenceError> { + match self.0.keys().next() { + None => Ok(()), + Some(key) => Err(invalid(format!("unknown field `{key}`"))), + } + } +} + +fn invalid(message: impl Into) -> FenceError { + FenceError::new(FenceOutcome::FencePreconditionFailed, message) + .with_detail(FenceDetail::InvalidArgument) +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct DrainPolicyBody { + deadline_ms: u64, + #[serde(default)] + on_deadline: Option, +} + +fn drain_policy(body: &mut Body) -> Result, FenceError> { + let Some(p) = body.opt::("drain_policy")? else { + return Ok(None); + }; + let on_deadline = match p.on_deadline.as_deref() { + None | Some("fail") => OnDeadline::Fail, + Some("force_rollback") => OnDeadline::ForceRollback, + Some(other) => { + return Err(invalid(format!( + "invalid `drain_policy.on_deadline` `{other}`: expected `fail` or `force_rollback`" + ))) + } + }; + Ok(Some(DrainPolicy { + deadline_ms: p.deadline_ms, + on_deadline, + })) +} + +/// Restore and dump options a target creation refuses: import goes through the migration +/// capability (section 4.4). +const RESTORE_FIELDS: [&str; 5] = [ + "dump_url", + "restore", + "restore_option", + "timestamp", + "from_backup", +]; + +fn parse_command( + namespace: NamespaceName, + kind: CommandKind, + bytes: &[u8], +) -> Result { + let mut body = Body::parse(bytes)?; + if kind == CommandKind::CreateTargetQuarantined { + if let Some(field) = RESTORE_FIELDS.iter().find(|f| body.has(f)) { + return Err(FenceError::new( + FenceOutcome::FencePreconditionFailed, + format!( + "`{field}` is not accepted: a migration target is created empty and filled \ + through the operation's import capability" + ), + ) + .with_detail(FenceDetail::RestoreNotAllowed)); + } + if body.has("shared_schema") || body.has("shared_schema_name") { + return Err(FenceError::new( + FenceOutcome::FencePreconditionFailed, + "namespace fences do not support shared schemas", + ) + .with_detail(FenceDetail::SharedSchemaUnsupported)); + } + } + + let operation_id = body.uuid("operation_id")?; + let command_id = body.uuid("command_id")?; + let expected_state = body.state("expected_state")?; + let expected_revision: u64 = body.req("expected_revision")?; + + let command = match kind { + CommandKind::AcquireSourceWriteFence => { + #[derive(Deserialize)] + #[serde(deny_unknown_fields)] + struct Identity { + log_id: String, + } + let identity: Identity = body.req("expected_namespace_identity")?; + let expected_log_id = Uuid::parse_str(&identity.log_id).map_err(|e| { + invalid(format!("invalid `expected_namespace_identity.log_id`: {e}")) + })?; + FenceCommand::AcquireSourceWriteFence { + expected_log_id, + drain_policy: drain_policy(&mut body)?, + } + } + CommandKind::SetSourceReadFence => FenceCommand::SetSourceReadFence { + drain_policy: drain_policy(&mut body)?, + }, + CommandKind::ClearSourceReadFence => FenceCommand::ClearSourceReadFence, + CommandKind::ReleaseSourceWriteFence => FenceCommand::ReleaseSourceWriteFence, + CommandKind::CreateTargetQuarantined => { + let jwt_key: Option = body.opt("jwt_key")?; + if let Some(key) = jwt_key.as_deref() { + parse_jwt_keys(key).map_err(|e| invalid(format!("invalid `jwt_key`: {e}")))?; + } + let max_db_size: Option = body.opt("max_db_size")?; + FenceCommand::CreateTargetQuarantined { + config: TargetConfig { + max_db_size: max_db_size.map(|s| s.as_u64()), + jwt_key, + txn_timeout_s: body.opt("txn_timeout_s")?, + allow_attach: body.opt("allow_attach")?.unwrap_or(false), + durability_mode: body.opt("durability_mode")?, + bottomless_db_id: body.opt("bottomless_db_id")?, + }, + } + } + CommandKind::SealTargetImport => FenceCommand::SealTargetImport { + drain_policy: drain_policy(&mut body)?, + }, + CommandKind::RecordTargetValidation => { + let result: String = body.req("result")?; + let result = match result.as_str() { + "ok" => ValidationResult::Ok, + "failed" => ValidationResult::Failed, + other => { + return Err(invalid(format!( + "invalid `result` `{other}`: expected `ok` or `failed`" + ))) + } + }; + FenceCommand::RecordTargetValidation { + result, + summary: body.opt("summary")?.unwrap_or_default(), + } + } + CommandKind::PublishTargetReadableWriteFenced => { + FenceCommand::PublishTargetReadableWriteFenced + } + CommandKind::EnableTargetWrites => FenceCommand::EnableTargetWrites, + CommandKind::AbortQuarantinedTarget => FenceCommand::AbortQuarantinedTarget, + CommandKind::AdoptFence => { + return Err(invalid("adoption is not served by this route")); + } + }; + body.finish()?; + Ok(FenceRequest { + namespace, + operation_id, + command_id, + expected_state, + expected_revision, + command, + }) +} + +struct ValidationQuery { + operation_id: Uuid, + expected_state: Option, + expected_revision: u64, + stmts: Vec, +} + +fn parse_validation_query(bytes: &[u8]) -> Result { + let mut body = Body::parse(bytes)?; + let operation_id = body.uuid("operation_id")?; + let expected_state = if body.has("expected_state") { + Some(body.state("expected_state")?) + } else { + None + }; + let expected_revision = body.req("expected_revision")?; + let stmts: Vec = body.req("stmts")?; + body.finish()?; + if stmts.is_empty() { + return Err(invalid("`stmts` is empty")); + } + Ok(ValidationQuery { + operation_id, + expected_state, + expected_revision, + stmts, + }) +} + +// --------------------------------------------------------------------------------------------- +// Validation queries + +enum StmtFailure { + Invalid(String), + Sqlite(rusqlite::Error), + TooManyRows, +} + +impl StmtFailure { + fn into_fence_error(self, index: usize) -> FenceError { + match self { + StmtFailure::Invalid(message) => invalid(format!("statement {index}: {message}")), + StmtFailure::Sqlite(rusqlite::Error::SqliteFailure(e, message)) + if e.code == rusqlite::ErrorCode::ReadOnly => + { + FenceError::new( + FenceOutcome::OperationCapabilityRequired, + format!( + "statement {index}: the validation capability is read-only: {}", + message.unwrap_or_else(|| e.to_string()) + ), + ) + } + StmtFailure::Sqlite(e) => invalid(format!("statement {index}: {e}")), + StmtFailure::TooManyRows => invalid(format!( + "statement {index}: the request returns more than {MAX_VALIDATION_QUERY_ROWS} rows" + )), + } + } +} + +impl From for StmtFailure { + fn from(e: rusqlite::Error) -> Self { + StmtFailure::Sqlite(e) + } +} + +fn to_sql_value(value: &proto::Value) -> Result { + use rusqlite::types::Value as V; + Ok(match value { + proto::Value::None => return Err(StmtFailure::Invalid("an argument has no value".into())), + proto::Value::Null => V::Null, + proto::Value::Integer { value } => V::Integer(*value), + proto::Value::Float { value } => V::Real(*value), + proto::Value::Text { value } => V::Text(value.to_string()), + proto::Value::Blob { value } => V::Blob(value.to_vec()), + }) +} + +fn from_sql_value(value: rusqlite::types::ValueRef<'_>) -> proto::Value { + use rusqlite::types::ValueRef as V; + match value { + V::Null => proto::Value::Null, + V::Integer(value) => proto::Value::Integer { value }, + V::Real(value) => proto::Value::Float { value }, + V::Text(bytes) => proto::Value::Text { + value: String::from_utf8_lossy(bytes).into(), + }, + V::Blob(bytes) => proto::Value::Blob { + value: Bytes::copy_from_slice(bytes), + }, + } +} + +/// Run one statement of a `validation-query` and collect at most `budget` rows. +fn run_validation_stmt( + conn: &mut rusqlite::Connection, + stmt: &proto::Stmt, + budget: usize, +) -> Result { + let sql = stmt.sql.as_deref().ok_or_else(|| { + StmtFailure::Invalid("`sql` is required (`sql_id` is not supported)".into()) + })?; + let mut prepared = conn.prepare(sql)?; + if !stmt.args.is_empty() && !stmt.named_args.is_empty() { + return Err(StmtFailure::Invalid( + "`args` and `named_args` cannot be combined".into(), + )); + } + for (i, arg) in stmt.args.iter().enumerate() { + prepared.raw_bind_parameter(i + 1, to_sql_value(arg)?)?; + } + for arg in &stmt.named_args { + let index = prepared + .parameter_index(&arg.name)? + .ok_or_else(|| StmtFailure::Invalid(format!("unknown parameter `{}`", arg.name)))?; + prepared.raw_bind_parameter(index, to_sql_value(&arg.value)?)?; + } + let cols: Vec = prepared + .columns() + .iter() + .map(|c| proto::Col { + name: Some(c.name().to_string()), + decltype: c.decl_type().map(str::to_string), + }) + .collect(); + let want_rows = stmt.want_rows.unwrap_or(true); + let column_count = cols.len(); + let mut rows = Vec::new(); + let mut raw = prepared.raw_query(); + while let Some(row) = raw.next()? { + if !want_rows { + continue; + } + if rows.len() == budget { + return Err(StmtFailure::TooManyRows); + } + let mut values = Vec::with_capacity(column_count); + for i in 0..column_count { + values.push(from_sql_value(row.get_ref(i)?)); + } + rows.push(proto::Row { values }); + } + Ok(proto::StmtResult { + cols, + rows, + ..Default::default() + }) +} + +// --------------------------------------------------------------------------------------------- +// Response views + +fn timestamp(ms: i64) -> Value { + chrono::DateTime::::from_timestamp_millis(ms) + .map(|t| json!(t.to_rfc3339_opts(chrono::SecondsFormat::Millis, true))) + .unwrap_or(Value::Null) +} + +fn server_json(server: &ServerIdentity) -> Value { + json!({ "build": server.build, "instance_id": server.instance_id.to_string() }) +} + +fn drain_json(controller: Option<&FenceController>) -> Value { + let counters = controller.map(|c| c.drain_counters()).unwrap_or_default(); + let DrainCounters { + active_writers, + read_leases, + import_writers, + } = counters; + json!({ + "active_writers": active_writers, + "read_leases": { + "sql": read_leases.sql, + "dump": read_leases.dump, + "replication": read_leases.replication, + }, + "import_writers": import_writers, + }) +} + +fn record_fields(record: &NamespaceFenceRecord, out: &mut Map) { + out.insert("role".into(), json!(record.role.as_str())); + out.insert("revision".into(), json!(record.revision)); + out.insert( + "operation_id".into(), + json!(record.operation_id.to_string()), + ); + out.insert( + "frozen_boundary".into(), + record + .frozen_boundary + .map(|b| json!({ "log_id": b.log_id.to_string(), "frame_no": b.frame_no })) + .unwrap_or(Value::Null), + ); + out.insert( + "drain_policy".into(), + record + .drain_policy + .map(|p| json!({ "deadline_ms": p.deadline_ms, "on_deadline": p.on_deadline.as_str() })) + .unwrap_or(Value::Null), + ); + out.insert( + "drain_started_at".into(), + record + .drain_started_at_ms + .map(timestamp) + .unwrap_or(Value::Null), + ); + out.insert( + "validation".into(), + record + .validation + .as_ref() + .map(|v| { + json!({ + "operation_id": v.operation_id.to_string(), + "command_id": v.command_id.to_string(), + "result": v.result.as_str(), + "summary": v.summary, + "snapshot": v.snapshot.map(|s| json!({ + "log_id": s.log_id.to_string(), + "frame_no": s.frame_no, + "page_count": s.page_count, + })), + "recorded_at": timestamp(v.recorded_at_ms), + }) + }) + .unwrap_or(Value::Null), + ); + out.insert("created_at".into(), timestamp(record.created_at_ms)); + out.insert( + "last_transition_at".into(), + timestamp(record.last_transition_at_ms), + ); + out.insert( + "last_command_id".into(), + json!(record.last_command_id.to_string()), + ); + out.insert("written_by".into(), server_json(&record.written_by)); + out.insert( + "adoptions".into(), + Value::Array( + record + .adoptions + .iter() + .map(|a| { + json!({ + "previous_operation_id": a.previous_operation_id.to_string(), + "new_operation_id": a.new_operation_id.to_string(), + "command_id": a.command_id.to_string(), + "approvers": a.approvers, + "incident_ref": a.incident_ref, + "reason": a.reason, + "at": timestamp(a.at_ms), + "revision": a.revision, + }) + }) + .collect(), + ), + ); +} + +/// The fence view of section 4.3. `fence` is the durable state being reported; `controller`, +/// when the namespace has one, supplies the live admission and the live log id. +fn fence_json( + namespace: &NamespaceName, + fence: &StoredFence, + controller: Option<&FenceController>, +) -> Value { + let gate = controller.map(|c| c.gate()); + let mut out = Map::new(); + out.insert("namespace".into(), json!(namespace.as_str())); + let state = match &gate { + Some(g) if g.is_creating_target() => FenceState::TargetQuarantined, + _ => fence.state(), + }; + out.insert("state".into(), json!(state.as_str())); + out.insert("role".into(), Value::Null); + out.insert("revision".into(), json!(fence.revision())); + out.insert("operation_id".into(), Value::Null); + let current_log_id = controller.and_then(|c| c.current_log_id()); + let (log_id, incarnation_id) = fence + .record() + .map(|r| (r.identity.log_id, r.identity.target_incarnation_id)) + .unwrap_or((None, None)); + out.insert( + "incarnation".into(), + json!({ + "log_id": log_id.map(|id| id.to_string()), + "target_incarnation_id": incarnation_id.map(|id| id.to_string()), + "current_log_id": current_log_id.map(|id| id.to_string()), + }), + ); + let (write, read, generation) = match &gate { + Some(g) => (g.write(), g.read(), g.write_generation), + None => (state.write_admission(), state.read_admission(), 0), + }; + out.insert( + "admission".into(), + json!({ + "write": write.as_str(), + "read": read.as_str(), + "generation": generation, + "indeterminate": gate.as_ref().is_some_and(|g| g.indeterminate.is_some()), + }), + ); + let marker = match fence { + StoredFence::None { .. } => Value::Null, + StoredFence::Record(record) => { + record_fields(record, &mut out); + json!("consistent") + } + StoredFence::Unavailable { + detail, + reason, + marker, + } => { + out.insert("detail".into(), json!(detail.as_str())); + out.insert("reason".into(), json!(reason)); + out.insert( + "marker_record".into(), + marker + .as_ref() + .map(|m| { + let mut inner = Map::new(); + inner.insert("state".into(), json!(m.state.as_str())); + record_fields(m, &mut inner); + Value::Object(inner) + }) + .unwrap_or(Value::Null), + ); + json!(detail.as_str()) + } + }; + out.insert("server".into(), server_json(&server_identity())); + out.insert( + "provenance".into(), + json!({ "metastore_restored_from_backup": false, "marker": marker }), + ); + Value::Object(out) +} + +fn receipt_json(receipt: &CommandReceipt) -> Value { + json!({ + "operation_id": receipt.operation_id.to_string(), + "command_id": receipt.command_id.to_string(), + "command": receipt.command.as_str(), + "fingerprint": receipt.fingerprint.to_string(), + "outcome": receipt.outcome.as_str(), + "revision_before": receipt.revision_before, + "revision_after": receipt.revision_after, + "state_after": receipt.state_after.as_str(), + "applied_at": timestamp(receipt.applied_at_ms), + "instance_id": receipt.instance_id.to_string(), + }) +} + +fn stored_receipt_json(stored: &StoredReceipt) -> Value { + match &stored.receipt { + Ok(receipt) => receipt_json(receipt), + Err(e) => json!({ + "operation_id": stored.operation_id, + "command_id": stored.command_id, + "revision_after": stored.revision_after, + "applied_at": timestamp(stored.applied_at_ms), + "error": e.to_string(), + }), + } +} diff --git a/libsql-server/src/http/admin/mod.rs b/libsql-server/src/http/admin/mod.rs index 2d8de1cdd1..be99413913 100644 --- a/libsql-server/src/http/admin/mod.rs +++ b/libsql-server/src/http/admin/mod.rs @@ -30,6 +30,7 @@ use crate::namespace::{DumpStream, NamespaceName, NamespaceStore, RestoreOption} use crate::net::Connector; use crate::LIBSQL_PAGE_SIZE; +pub mod fence; pub mod stats; #[derive(Clone)] @@ -49,6 +50,9 @@ struct AppState { connector: C, metrics: Metrics, set_env_filter: Option anyhow::Result<()> + Sync + Send + 'static>>, + /// Whether an admin auth key is configured. Namespace fence commands refuse to run without + /// one (`docs/NAMESPACE_FENCE.md` section 4.1). + admin_auth_configured: bool, } impl FromRef>> for Metrics { @@ -170,12 +174,14 @@ where .route("/profile/heap/disable/:id", post(disable_profile_heap)) .route("/profile/heap/:id", delete(delete_profile_heap)) .route("/log-filter", post(handle_set_log_filter)) + .merge(fence::routes()) .with_state(Arc::new(AppState { namespaces: namespaces.clone(), connector, user_http_server, metrics, set_env_filter, + admin_auth_configured: auth.is_some(), })) .layer( tower_http::trace::TraceLayer::new_for_http() diff --git a/libsql-server/src/namespace/fence/controller.rs b/libsql-server/src/namespace/fence/controller.rs index 4dfd45743b..621576bacb 100644 --- a/libsql-server/src/namespace/fence/controller.rs +++ b/libsql-server/src/namespace/fence/controller.rs @@ -233,6 +233,14 @@ pub enum LeaseKind { Replication, } +/// The live drain counters of a namespace (see [`FenceController::drain_counters`]). +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct DrainCounters { + pub active_writers: usize, + pub read_leases: ReadLeaseCounts, + pub import_writers: usize, +} + /// The number of read leases held on a namespace, by kind. #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] pub struct ReadLeaseCounts { @@ -652,6 +660,29 @@ impl FenceController { self.capabilities.lock().import_writers } + /// The live drain counters reported by `InspectFence` and every admin response + /// (`docs/NAMESPACE_FENCE.md` section 4.3): connections holding a write slot for a write + /// transaction, read leases by kind, and running import calls. A snapshot; never waits. + pub fn drain_counters(&self) -> DrainCounters { + let active_writers = self + .live_write_drains() + .iter() + .filter(|source| source.manager.has_writer()) + .count(); + DrainCounters { + active_writers, + read_leases: self.read_lease_counts(), + import_writers: self.import_writers(), + } + } + + /// The replication log id of the namespace as it is loaded now, if it is loaded on this + /// server as a primary. After a dirty restart this can differ from the log id a source was + /// acquired on (section 8.5). + pub fn current_log_id(&self) -> Option { + self.live_write_drains().last().map(|source| source.log_id) + } + /// The capabilities issued and still live. pub fn live_capabilities(&self) -> usize { self.capabilities.lock().live.len() diff --git a/libsql-server/src/namespace/fence/mod.rs b/libsql-server/src/namespace/fence/mod.rs index a8d80da59c..af5b76182e 100644 --- a/libsql-server/src/namespace/fence/mod.rs +++ b/libsql-server/src/namespace/fence/mod.rs @@ -45,3 +45,20 @@ pub(crate) mod proto { /// Version of the fence admin protocol reported by capability discovery. pub const FENCE_PROTOCOL_VERSION: u32 = 1; + +/// Whether this server fills the proxy protocol's additive `Error.stable_code` field and maps +/// it on the replica side (`docs/NAMESPACE_FENCE.md` section 6.1). Reported by capability +/// discovery so that deployment tooling can check every server before fences are used. +pub const PROXY_STABLE_CODE: bool = false; + +/// The identity of this server process: its build and an id generated once per process. It is +/// written into records and receipts, and reported by the admin API. +pub fn server_identity() -> record::ServerIdentity { + static IDENTITY: std::sync::OnceLock = std::sync::OnceLock::new(); + IDENTITY + .get_or_init(|| record::ServerIdentity { + build: crate::version::version(), + instance_id: uuid::Uuid::new_v4(), + }) + .clone() +} diff --git a/libsql-server/src/namespace/fence/registry.rs b/libsql-server/src/namespace/fence/registry.rs index 7c03102cd1..176cfddd4e 100644 --- a/libsql-server/src/namespace/fence/registry.rs +++ b/libsql-server/src/namespace/fence/registry.rs @@ -82,6 +82,22 @@ impl FenceRegistry { } } + /// How many namespaces have an active fence (`docs/NAMESPACE_FENCE.md` section 4.4, + /// `active_fences`): a record in any state but `RELEASED` or `TARGET_WRITABLE`, an + /// unavailable state, a target being created, or a commit whose outcome is not known yet. + pub fn active_count(&self) -> usize { + let controllers: Vec<_> = self.controllers.lock().values().cloned().collect(); + controllers + .iter() + .filter(|controller| { + let gate = controller.gate(); + gate.state().is_active() + || gate.indeterminate.is_some() + || gate.is_creating_target() + }) + .count() + } + pub fn len(&self) -> usize { self.controllers.lock().len() } diff --git a/libsql-server/src/namespace/store.rs b/libsql-server/src/namespace/store.rs index 889370ad79..09d27a5695 100644 --- a/libsql-server/src/namespace/store.rs +++ b/libsql-server/src/namespace/store.rs @@ -32,7 +32,7 @@ use super::fence::registry::FenceRegistry; use super::fence::state::{FenceState, Role}; use super::fence::store::StoredFence; use super::fence::target::{self, CreateTargetRequest, ValidationSession}; -use super::meta_store::{FenceCommit, FenceContext, MetaStore, MetaStoreHandle}; +use super::meta_store::{FenceCommit, FenceContext, FenceInspection, MetaStore, MetaStoreHandle}; use super::schema_lock::SchemaLocksRegistry; use super::{Namespace, ResetCb, ResetOp, ResolveNamespacePathFn, RestoreOption}; @@ -556,8 +556,6 @@ impl NamespaceStore { /// (`docs/NAMESPACE_FENCE.md` sections 5.3 and 8). `AcquireSourceWriteFence` loads the /// namespace first, so that its connection manager and replication log are registered with /// the namespace's controller before the drain needs them. - // The admin routes that call this are not part of the server yet. - #[cfg_attr(not(test), allow(dead_code))] pub(crate) async fn execute_fence_command( &self, request: FenceRequest, @@ -881,6 +879,36 @@ impl NamespaceStore { Ok(None) } + /// Whether this store serves primary namespaces. Fences live on the primary that owns the + /// WAL; a replica-kind server refuses every fence route (`docs/NAMESPACE_FENCE.md` 4.1). + pub(crate) fn is_primary(&self) -> bool { + !self.inner.db_kind.is_replica() + } + + /// The fence controller `namespace` already has, without creating one. + pub(crate) fn existing_fence_controller( + &self, + namespace: &NamespaceName, + ) -> Option> { + self.inner.fences.get(namespace) + } + + /// `InspectFence`: the durable fence and receipts of `namespace` as the metastore holds + /// them, and the namespace's controller if it has one (for the live gate and drain + /// counters). Read-only: it neither loads the namespace nor creates a controller. + pub(crate) async fn inspect_fence( + &self, + namespace: &NamespaceName, + ) -> crate::Result<(FenceInspection, Option>)> { + let inspection = self.inner.metadata.inspect_fence(namespace.clone()).await?; + Ok((inspection, self.inner.fences.get(namespace))) + } + + /// How many namespaces on this server have an active fence (capability discovery). + pub(crate) fn active_fences(&self) -> usize { + self.inner.fences.active_count() + } + pub(crate) fn schema_locks(&self) -> &SchemaLocksRegistry { &self.inner.schema_locks } diff --git a/libsql-server/tests/common/http.rs b/libsql-server/tests/common/http.rs index 8716a60503..a08a928478 100644 --- a/libsql-server/tests/common/http.rs +++ b/libsql-server/tests/common/http.rs @@ -41,6 +41,20 @@ impl Client { Ok(Response(self.0.get(s.parse()?).await?)) } + pub(crate) async fn get_with_headers( + &self, + url: &str, + headers: &[(HeaderName, &str)], + ) -> anyhow::Result { + let mut request = hyper::Request::get(url).body(Body::empty())?; + for (key, val) in headers { + request + .headers_mut() + .insert(key.clone(), val.parse().unwrap()); + } + Ok(Response(self.0.request(request).await?)) + } + pub(crate) async fn post(&self, url: &str, body: T) -> anyhow::Result { self.post_with_headers(url, &[], body).await } diff --git a/libsql-server/tests/fence/admin.rs b/libsql-server/tests/fence/admin.rs new file mode 100644 index 0000000000..d1ac8413c4 --- /dev/null +++ b/libsql-server/tests/fence/admin.rs @@ -0,0 +1,731 @@ +//! The fence admin API over HTTP (`docs/NAMESPACE_FENCE.md` section 4). + +use hyper::StatusCode; +use serde_json::{json, Value}; +use tempfile::tempdir; +use uuid::Uuid; + +use super::{command_body, connect, make_primary, sim, state_of, Admin, Primary, ADMIN_KEY}; + +fn uuid(n: u128) -> Uuid { + Uuid::from_u128(n) +} + +/// Load `ns` on the server with one write, and return the replication log id the server +/// reports for it. +async fn load_and_log_id(admin: &Admin, ns: &str) -> anyhow::Result { + let conn = connect(ns)?; + conn.execute("create table if not exists t (x)", ()).await?; + conn.execute("insert into t values (1)", ()).await?; + let (status, body) = admin.inspect(ns).await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body), ("UNFENCED", 0), "{body}"); + Ok(body["fence"]["incarnation"]["current_log_id"] + .as_str() + .unwrap_or_else(|| panic!("no current_log_id: {body}")) + .to_string()) +} + +fn acquire_body(op: Uuid, cmd: Uuid, log_id: &str) -> Value { + command_body( + op, + cmd, + "UNFENCED", + 0, + json!({ + "expected_namespace_identity": { "log_id": log_id }, + "drain_policy": { "deadline_ms": 5000, "on_deadline": "fail" }, + }), + ) +} + +#[test] +fn capabilities() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let (status, body) = admin.get("/v1/fence/capabilities").await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(body["fence_protocol_version"], 1); + assert_eq!(body["enabled"], true); + assert_eq!(body["active_fences"], 0); + assert_eq!(body["proxy_stable_code"], false); + let commands: Vec<&str> = body["commands"] + .as_array() + .unwrap() + .iter() + .map(|c| c.as_str().unwrap()) + .collect(); + for command in [ + "InspectFence", + "AcquireSourceWriteFence", + "SetSourceReadFence", + "ClearSourceReadFence", + "ReleaseSourceWriteFence", + "CreateTargetQuarantined", + "SealTargetImport", + "RecordTargetValidation", + "PublishTargetReadableWriteFenced", + "EnableTargetWrites", + "AbortQuarantinedTarget", + ] { + assert!(commands.contains(&command), "{command} missing: {body}"); + } + let states = body["states"].as_array().unwrap(); + assert!(states.contains(&json!("SOURCE_WRITE_FENCED")), "{body}"); + assert!(states.contains(&json!("UNKNOWN_UNAVAILABLE")), "{body}"); + assert!(body["server"]["build"] + .as_str() + .unwrap() + .starts_with("sqld ")); + Uuid::parse_str(body["server"]["instance_id"].as_str().unwrap())?; + assert_eq!(body["metastore"]["restored_from_backup"], false); + + // An active fence is counted. + admin.create_namespace("src").await?; + let log_id = load_and_log_id(&admin, "src").await?; + let (status, body) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(uuid(1), uuid(2), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let (_, body) = admin.get("/v1/fence/capabilities").await?; + assert_eq!(body["active_fences"], 1, "{body}"); + + // The admin API's own authentication still applies. + let (status, _) = Admin::new(None).get("/v1/fence/capabilities").await?; + assert_eq!(status, StatusCode::UNAUTHORIZED); + Ok(()) + }); + sim.run().unwrap(); +} + +#[test] +fn capabilities_when_disabled() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary( + &mut sim, + tmp.path().to_path_buf(), + Primary { + fence_enabled: false, + ..Default::default() + }, + ); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let (status, body) = admin.get("/v1/fence/capabilities").await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(body["enabled"], false); + assert_eq!(body["fence_protocol_version"], 1); + + admin.create_namespace("src").await?; + let (status, _) = admin.inspect("src").await?; + assert_eq!(status, StatusCode::NOT_FOUND); + for route in [ + "source/acquire-write-fence", + "source/release-write-fence", + "target/create-quarantined", + "target/validation-query", + ] { + let (status, body) = admin + .command( + "src", + route, + acquire_body(uuid(1), uuid(2), &uuid(3).to_string()), + ) + .await?; + assert_eq!(status, StatusCode::NOT_FOUND, "{route}: {body}"); + } + // The namespace is untouched. + let conn = connect("src")?; + conn.execute("create table t (x)", ()).await?; + Ok(()) + }); + sim.run().unwrap(); +} + +#[test] +fn mutating_routes_require_admin_key() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary( + &mut sim, + tmp.path().to_path_buf(), + Primary { + admin_key: None, + ..Default::default() + }, + ); + sim.client("client", async { + let admin = Admin::new(None); + admin.create_namespace("src").await?; + let log_id = load_and_log_id(&admin, "src").await?; + + let (status, body) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(uuid(1), uuid(2), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{body}"); + assert_eq!(body["outcome"], "FENCE_PRECONDITION_FAILED"); + assert_eq!(body["detail"], "admin_auth_required"); + assert_eq!(state_of(&body), ("UNFENCED", 0), "{body}"); + + let (status, body) = admin + .command( + "tgt", + "target/create-quarantined", + command_body(uuid(1), uuid(3), "ABSENT", 0, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{body}"); + assert_eq!(body["detail"], "admin_auth_required"); + + let (status, body) = admin + .command( + "tgt", + "target/validation-query", + json!({ + "operation_id": uuid(1).to_string(), + "expected_revision": 1, + "stmts": [{ "sql": "select 1" }], + }), + ) + .await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{body}"); + assert_eq!(body["detail"], "admin_auth_required"); + + // Nothing was fenced or created, and reading the state is still possible. + let (status, body) = admin.inspect("src").await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body), ("UNFENCED", 0)); + let (status, _) = admin.inspect("tgt").await?; + assert_eq!(status, StatusCode::NOT_FOUND); + connect("src")? + .execute("insert into t values (2)", ()) + .await?; + Ok(()) + }); + sim.run().unwrap(); +} + +/// Acceptance test: two operations race to acquire the same source; exactly one owns it. +#[test] +fn concurrent_acquire_one_owner() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + admin.create_namespace("src").await?; + let log_id = load_and_log_id(&admin, "src").await?; + + let other = Admin::new(Some(ADMIN_KEY)); + let (a, b) = tokio::join!( + admin.command( + "src", + "source/acquire-write-fence", + acquire_body(uuid(0xa), uuid(1), &log_id), + ), + other.command( + "src", + "source/acquire-write-fence", + acquire_body(uuid(0xb), uuid(2), &log_id), + ), + ); + let (a, b) = (a?, b?); + let mut results = [a, b]; + results.sort_by_key(|(status, _)| status.as_u16()); + let [(won_status, won), (lost_status, lost)] = results; + assert_eq!(won_status, StatusCode::OK, "{won}"); + assert_eq!(won["outcome"], "APPLIED"); + assert_eq!(state_of(&won).0, "SOURCE_WRITE_FENCED"); + assert_eq!(lost_status, StatusCode::CONFLICT, "{lost}"); + assert_eq!(lost["outcome"], "FENCE_OWNED_BY_ANOTHER_OPERATION"); + // The loser is shown who owns the namespace. + assert_eq!(lost["fence"]["operation_id"], won["fence"]["operation_id"]); + + let (_, body) = admin.inspect("src").await?; + assert_eq!(body["fence"]["operation_id"], won["fence"]["operation_id"]); + assert!(connect("src")? + .execute("insert into t values (2)", ()) + .await + .is_err()); + Ok(()) + }); + sim.run().unwrap(); +} + +#[test] +fn source_walk_over_http() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + admin.create_namespace("src").await?; + let log_id = load_and_log_id(&admin, "src").await?; + let op = uuid(0xa); + let conn = connect("src")?; + + // A wrong identity is refused before anything changes. + let (status, body) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(op, uuid(1), &uuid(0xdead).to_string()), + ) + .await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{body}"); + assert_eq!(body["detail"], "namespace_identity_mismatch"); + + let (status, acquired) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(op, uuid(2), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{acquired}"); + assert_eq!(acquired["outcome"], "APPLIED"); + assert_eq!(acquired["replayed"], false); + let (state, rev) = state_of(&acquired); + assert_eq!(state, "SOURCE_WRITE_FENCED"); + assert_eq!(acquired["fence"]["role"], "SOURCE"); + assert_eq!(acquired["fence"]["admission"]["write"], "closed"); + assert_eq!(acquired["fence"]["admission"]["read"], "open"); + assert_eq!( + acquired["fence"]["frozen_boundary"]["log_id"], + log_id.as_str() + ); + assert_eq!(acquired["receipt"]["command"], "AcquireSourceWriteFence"); + assert_eq!(acquired["drain"]["active_writers"], 0); + assert!(conn.execute("insert into t values (2)", ()).await.is_err()); + conn.query("select * from t", ()).await?; + + // Replay returns the stored receipt. + let (status, replay) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(op, uuid(2), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{replay}"); + assert_eq!(replay["replayed"], true); + assert_eq!(replay["receipt"], acquired["receipt"]); + + // The same command id with a different request is a conflict. + let (status, body) = admin + .command( + "src", + "source/acquire-write-fence", + command_body( + op, + uuid(2), + "UNFENCED", + 0, + json!({ "expected_namespace_identity": { "log_id": log_id } }), + ), + ) + .await?; + assert_eq!(status, StatusCode::CONFLICT, "{body}"); + assert_eq!(body["outcome"], "FENCE_COMMAND_CONFLICT"); + + // A stale revision is refused. + let (status, body) = admin + .command( + "src", + "source/set-read-fence", + command_body(op, uuid(3), "SOURCE_WRITE_FENCED", rev - 1, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::CONFLICT, "{body}"); + assert_eq!(body["outcome"], "FENCE_REVISION_MISMATCH"); + assert_eq!(state_of(&body), ("SOURCE_WRITE_FENCED", rev)); + + // Read fence, then clear it. + let (status, body) = admin + .command( + "src", + "source/set-read-fence", + command_body( + op, + uuid(4), + "SOURCE_WRITE_FENCED", + rev, + json!({ "drain_policy": { "deadline_ms": 5000 } }), + ), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let (state, rev) = state_of(&body); + assert_eq!(state, "SOURCE_READ_FENCED"); + assert_eq!(body["fence"]["admission"]["read"], "closed"); + assert!(conn.query("select * from t", ()).await.is_err()); + + let (status, body) = admin + .command( + "src", + "source/clear-read-fence", + command_body(op, uuid(5), "SOURCE_READ_FENCED", rev, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let (state, rev) = state_of(&body); + assert_eq!(state, "SOURCE_WRITE_FENCED"); + connect("src")?.query("select * from t", ()).await?; + + // Release reopens writes. + let release = command_body(op, uuid(6), "SOURCE_WRITE_FENCED", rev, json!({})); + let (status, released) = admin + .command("src", "source/release-write-fence", release.clone()) + .await?; + assert_eq!(status, StatusCode::OK, "{released}"); + assert_eq!(state_of(&released).0, "RELEASED"); + assert_eq!(released["fence"]["admission"]["write"], "open"); + connect("src")? + .execute("insert into t values (3)", ()) + .await?; + + let (status, replay) = admin + .command("src", "source/release-write-fence", release) + .await?; + assert_eq!(status, StatusCode::OK, "{replay}"); + assert_eq!(replay["replayed"], true); + assert_eq!(replay["receipt"], released["receipt"]); + + // Inspect shows the operation's receipts, and all of them on request. + let (status, body) = admin.inspect("src").await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body).0, "RELEASED"); + let receipts = body["receipts"].as_array().unwrap(); + assert!(receipts.len() >= 4, "{body}"); + assert!(receipts + .iter() + .all(|r| r["operation_id"] == op.to_string().as_str())); + let (_, all) = admin.get("/v1/namespaces/src/fence?receipts=all").await?; + assert!(all["receipts"].as_array().unwrap().len() >= receipts.len()); + + // Malformed requests are typed refusals. + let (status, body) = admin + .command( + "src", + "source/release-write-fence", + json!({ "operation_id": "x" }), + ) + .await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{body}"); + assert_eq!(body["detail"], "invalid_argument"); + let (status, body) = admin + .command( + "src", + "source/release-write-fence", + command_body(op, uuid(7), "RELEASED", 0, json!({ "surprise": 1 })), + ) + .await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{body}"); + assert_eq!(body["detail"], "invalid_argument"); + Ok(()) + }); + sim.run().unwrap(); +} + +#[test] +fn target_walk_over_http() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let op = uuid(0xa); + + // Restore options are refused, and nothing is created. + let (status, body) = admin + .command( + "tgt", + "target/create-quarantined", + command_body( + op, + uuid(1), + "ABSENT", + 0, + json!({ "dump_url": "file:///tmp/dump.sql" }), + ), + ) + .await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{body}"); + assert_eq!(body["detail"], "restore_not_allowed"); + assert_eq!(admin.inspect("tgt").await?.0, StatusCode::NOT_FOUND); + + let create = command_body( + op, + uuid(2), + "ABSENT", + 0, + json!({ "max_db_size": 10_000_000, "durability_mode": "strong" }), + ); + let (status, created) = admin + .command("tgt", "target/create-quarantined", create.clone()) + .await?; + assert_eq!(status, StatusCode::OK, "{created}"); + assert_eq!(created["outcome"], "APPLIED"); + let (state, rev) = state_of(&created); + assert_eq!(state, "TARGET_QUARANTINED"); + assert_eq!(created["fence"]["role"], "TARGET"); + assert!(created["fence"]["incarnation"]["target_incarnation_id"].is_string()); + let (status, replay) = admin + .command("tgt", "target/create-quarantined", create) + .await?; + assert_eq!(status, StatusCode::OK, "{replay}"); + assert_eq!(replay["replayed"], true); + + // Normal SQL is refused while the target is quarantined. + assert!(connect("tgt")?.query("select 1", ()).await.is_err()); + let (status, body) = admin + .post("/v1/namespaces/tgt/create", json!({})) + .await?; + assert!(!status.is_success(), "{status} {body}"); + + // Validation queries are refused before the import is sealed. + let (status, body) = admin + .command( + "tgt", + "target/validation-query", + json!({ + "operation_id": op.to_string(), + "expected_revision": rev, + "stmts": [{ "sql": "select 1" }], + }), + ) + .await?; + assert_eq!(status, StatusCode::FORBIDDEN, "{body}"); + assert_eq!(body["outcome"], "OPERATION_CAPABILITY_REQUIRED"); + + let (status, body) = admin + .command( + "tgt", + "target/seal-import", + command_body(op, uuid(3), "TARGET_QUARANTINED", rev, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let (state, rev) = state_of(&body); + assert_eq!(state, "TARGET_VALIDATING"); + + let (status, body) = admin + .command( + "tgt", + "target/validation-query", + json!({ + "operation_id": op.to_string(), + "expected_state": "TARGET_VALIDATING", + "expected_revision": rev, + "stmts": [ + { "sql": "select count(*) as n from sqlite_master" }, + { "sql": "select ? + 1 as v", "args": [{ "type": "integer", "value": "41" }] }, + ], + }), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(body["results"][0]["cols"][0]["name"], "n"); + assert_eq!( + body["results"][1]["rows"][0][0], + json!({ "type": "integer", "value": "42" }) + ); + + // A validation query cannot write. + let (status, body) = admin + .command( + "tgt", + "target/validation-query", + json!({ + "operation_id": op.to_string(), + "expected_revision": rev, + "stmts": [{ "sql": "create table sneaky (x)" }], + }), + ) + .await?; + assert_eq!(status, StatusCode::FORBIDDEN, "{body}"); + assert_eq!(body["outcome"], "OPERATION_CAPABILITY_REQUIRED"); + + // Another operation cannot validate. + let (status, body) = admin + .command( + "tgt", + "target/validation-query", + json!({ + "operation_id": uuid(0xb).to_string(), + "expected_revision": rev, + "stmts": [{ "sql": "select 1" }], + }), + ) + .await?; + assert_eq!(status, StatusCode::CONFLICT, "{body}"); + assert_eq!(body["outcome"], "FENCE_OWNED_BY_ANOTHER_OPERATION"); + + // Publication needs a successful validation receipt. + let (status, body) = admin + .command( + "tgt", + "target/publish-readable", + command_body(op, uuid(4), "TARGET_VALIDATING", rev, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{body}"); + assert_eq!(body["detail"], "validation_receipt_required"); + + let (status, body) = admin + .command( + "tgt", + "target/validation-receipt", + command_body( + op, + uuid(5), + "TARGET_VALIDATING", + rev, + json!({ "result": "ok", "summary": "row counts match" }), + ), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let (state, rev) = state_of(&body); + assert_eq!(state, "TARGET_VALIDATING"); + assert_eq!(body["fence"]["validation"]["result"], "ok"); + assert!(body["fence"]["validation"]["snapshot"]["page_count"].is_u64()); + + let (status, body) = admin + .command( + "tgt", + "target/publish-readable", + command_body(op, uuid(6), "TARGET_VALIDATING", rev, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let (state, rev) = state_of(&body); + assert_eq!(state, "TARGET_WRITE_FENCED"); + let conn = connect("tgt")?; + conn.query("select 1", ()).await?; + assert!(conn.execute("create table t (x)", ()).await.is_err()); + + let enable = command_body(op, uuid(7), "TARGET_WRITE_FENCED", rev, json!({})); + let (status, body) = admin + .command("tgt", "target/enable-writes", enable.clone()) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(body["outcome"], "APPLIED"); + let (state, rev_after) = state_of(&body); + assert_eq!(state, "TARGET_WRITABLE"); + connect("tgt")?.execute("create table t (x)", ()).await?; + + let (status, body) = admin + .command("tgt", "target/enable-writes", enable) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(body["replayed"], true); + let (status, body) = admin + .command( + "tgt", + "target/enable-writes", + command_body(op, uuid(8), "TARGET_WRITE_FENCED", rev, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(body["outcome"], "ALREADY_APPLIED"); + assert_eq!(state_of(&body), ("TARGET_WRITABLE", rev_after)); + + // Abort is not possible once writes are enabled. + let (status, body) = admin + .command( + "tgt", + "target/abort", + command_body(op, uuid(9), "TARGET_WRITABLE", rev_after, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::CONFLICT, "{body}"); + assert_eq!(body["outcome"], "INVALID_FENCE_TRANSITION"); + Ok(()) + }); + sim.run().unwrap(); +} + +/// A write transaction open when the fence is requested holds the drain: `InspectFence` counts +/// it, the acquisition answers `DRAINING` (202) at its deadline, and replaying the command once +/// the transaction has committed completes the fence. +#[test] +fn inspect_reports_drain_counters() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + admin.create_namespace("src").await?; + let log_id = load_and_log_id(&admin, "src").await?; + + let (_, body) = admin.inspect("src").await?; + assert_eq!( + body["drain"], + json!({ + "active_writers": 0, + "read_leases": { "sql": 0, "dump": 0, "replication": 0 }, + "import_writers": 0, + }) + ); + + let conn = connect("src")?; + let tx = conn.transaction().await?; + tx.execute("insert into t values (2)", ()).await?; + + let (_, body) = admin.inspect("src").await?; + assert_eq!(body["drain"]["active_writers"], 1, "{body}"); + + let op = uuid(0xa); + let acquire = command_body( + op, + uuid(1), + "UNFENCED", + 0, + json!({ + "expected_namespace_identity": { "log_id": log_id }, + "drain_policy": { "deadline_ms": 100, "on_deadline": "fail" }, + }), + ); + let (status, body) = admin + .command("src", "source/acquire-write-fence", acquire.clone()) + .await?; + assert_eq!(status, StatusCode::ACCEPTED, "{body}"); + assert_eq!(body["outcome"], "DRAINING"); + assert_eq!(state_of(&body).0, "SOURCE_DRAINING"); + assert_eq!(body["fence"]["admission"]["write"], "closed"); + assert_eq!(body["drain"]["active_writers"], 1, "{body}"); + + // The transaction admitted before the fence commits; new writes are refused. + tx.commit().await?; + assert!(connect("src")? + .execute("insert into t values (3)", ()) + .await + .is_err()); + + let (status, body) = admin + .command("src", "source/acquire-write-fence", acquire) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(body["outcome"], "APPLIED"); + assert_eq!(state_of(&body).0, "SOURCE_WRITE_FENCED"); + assert_eq!(body["drain"]["active_writers"], 0); + let mut rows = connect("src")?.query("select count(*) from t", ()).await?; + let n: i64 = rows.next().await?.unwrap().get(0)?; + assert_eq!(n, 2); + Ok(()) + }); + sim.run().unwrap(); +} diff --git a/libsql-server/tests/fence/mod.rs b/libsql-server/tests/fence/mod.rs new file mode 100644 index 0000000000..a7a149907f --- /dev/null +++ b/libsql-server/tests/fence/mod.rs @@ -0,0 +1,188 @@ +#![allow(deprecated)] + +//! Namespace fence integration tests (`docs/NAMESPACE_FENCE.md`), driven over the admin API. + +mod admin; + +use std::path::PathBuf; +use std::time::Duration; + +use hyper::StatusCode; +use libsql_server::config::{AdminApiConfig, MetaStoreConfig, RpcServerConfig, UserApiConfig}; +use s3s::header::AUTHORIZATION; +use serde_json::{json, Value}; +use turmoil::{Builder, Sim}; +use uuid::Uuid; + +use crate::common::http::Client; +use crate::common::net::{ + init_tracing, SimServer as _, TestServer, TurmoilAcceptor, TurmoilConnector, +}; + +pub const ADMIN_KEY: &str = "fence-admin-key"; + +pub struct Primary { + /// `None` starts the admin API without an auth key. + pub admin_key: Option<&'static str>, + pub fence_enabled: bool, +} + +impl Default for Primary { + fn default() -> Self { + Self { + admin_key: Some(ADMIN_KEY), + fence_enabled: true, + } + } +} + +pub fn sim() -> Sim<'static> { + Builder::new() + .simulation_duration(Duration::from_secs(1000)) + .build() +} + +/// A primary on host `primary`: user API on 8080, admin API on 9090. +pub fn make_primary(sim: &mut Sim, path: PathBuf, primary: Primary) { + init_tracing(); + let Primary { + admin_key, + fence_enabled, + } = primary; + sim.host("primary", move || { + let path = path.clone(); + async move { + let server = TestServer { + path: path.into(), + user_api_config: UserApiConfig::default(), + admin_api_config: Some(AdminApiConfig { + acceptor: TurmoilAcceptor::bind(([0, 0, 0, 0], 9090)).await?, + connector: TurmoilConnector, + disable_metrics: true, + auth_key: admin_key.map(Into::into), + }), + rpc_server_config: Some(RpcServerConfig { + acceptor: TurmoilAcceptor::bind(([0, 0, 0, 0], 4567)).await?, + tls_config: None, + }), + meta_store_config: MetaStoreConfig { + namespace_fence: fence_enabled, + ..Default::default() + }, + disable_namespaces: false, + disable_default_namespace: true, + ..Default::default() + }; + server.start_sim(8080).await?; + Ok(()) + } + }); +} + +/// The admin API of `primary`, authenticating with `key` when there is one. +pub struct Admin { + client: Client, + key: Option, +} + +impl Admin { + pub fn new(key: Option<&str>) -> Self { + Self { + client: Client::new(), + key: key.map(|k| format!("basic {k}")), + } + } + + fn headers(&self) -> Vec<(hyper::header::HeaderName, &str)> { + self.key + .as_deref() + .map(|k| vec![(AUTHORIZATION, k)]) + .unwrap_or_default() + } + + async fn json(resp: crate::common::http::Response) -> anyhow::Result<(StatusCode, Value)> { + let status = resp.status(); + let body = resp.body_string().await?; + let value = if body.trim().is_empty() { + Value::Null + } else { + serde_json::from_str(&body).unwrap_or(Value::String(body)) + }; + Ok((status, value)) + } + + pub async fn get(&self, path: &str) -> anyhow::Result<(StatusCode, Value)> { + let url = format!("http://primary:9090{path}"); + Self::json(self.client.get_with_headers(&url, &self.headers()).await?).await + } + + pub async fn post(&self, path: &str, body: Value) -> anyhow::Result<(StatusCode, Value)> { + let url = format!("http://primary:9090{path}"); + Self::json( + self.client + .post_with_headers(&url, &self.headers(), body) + .await?, + ) + .await + } + + pub async fn create_namespace(&self, ns: &str) -> anyhow::Result<()> { + let (status, body) = self + .post(&format!("/v1/namespaces/{ns}/create"), json!({})) + .await?; + anyhow::ensure!(status.is_success(), "create {ns}: {status} {body}"); + Ok(()) + } + + pub async fn inspect(&self, ns: &str) -> anyhow::Result<(StatusCode, Value)> { + self.get(&format!("/v1/namespaces/{ns}/fence")).await + } + + /// A fence command: `route` is the part after `/fence/`. + pub async fn command( + &self, + ns: &str, + route: &str, + body: Value, + ) -> anyhow::Result<(StatusCode, Value)> { + self.post(&format!("/v1/namespaces/{ns}/fence/{route}"), body) + .await + } +} + +/// The common fields of a fence command, with `extra` merged in. +pub fn command_body( + operation_id: Uuid, + command_id: Uuid, + expected_state: &str, + expected_revision: u64, + extra: Value, +) -> Value { + let mut body = json!({ + "operation_id": operation_id.to_string(), + "command_id": command_id.to_string(), + "expected_state": expected_state, + "expected_revision": expected_revision, + }); + if let Value::Object(extra) = extra { + body.as_object_mut().unwrap().extend(extra); + } + body +} + +pub fn state_of(body: &Value) -> (&str, u64) { + ( + body["fence"]["state"].as_str().unwrap_or(""), + body["fence"]["revision"].as_u64().unwrap_or(u64::MAX), + ) +} + +/// A connection to namespace `ns` over the user API. +pub fn connect(ns: &str) -> anyhow::Result { + let db = libsql::Database::open_remote_with_connector( + format!("http://{ns}.primary:8080"), + "", + TurmoilConnector, + )?; + Ok(db.connect()?) +} diff --git a/libsql-server/tests/tests.rs b/libsql-server/tests/tests.rs index ab475df546..55a88d2dab 100644 --- a/libsql-server/tests/tests.rs +++ b/libsql-server/tests/tests.rs @@ -6,6 +6,7 @@ mod common; mod auth; mod cluster; mod embedded_replica; +mod fence; mod hrana; mod namespaces; mod standalone; From dbfec828db71ced454ea6048fc97811e20779d74 Mon Sep 17 00:00:00 2001 From: River Date: Wed, 30 Sep 2026 09:13:46 +0000 Subject: [PATCH 2/2] libsql-server: deny lifecycle operations on fenced namespaces Refuse config changes, delete, reset, fork (as source or destination), create over an existing record (with or without a dump URL), linking to a shared schema and schema migration while a namespace's fence denies lifecycle work. A single check reads the fence registry, so it also sees in-memory gates (a closing transition, a target being created, an indeterminate commit) and never loads the namespace. Paths that persist through the metastore are refused again inside its transaction. A fork reads the source's log without a read lease, so it now holds the source's transition lock for its whole run and checks the gate under it: a write or read fence command either completes first and the fork is refused, or waits for the fork. Reset writes nothing to the metastore and is checked under the same lock. The fork destination is checked before anything is stored and again with its namespace entry held; previously a fork onto an existing, unloaded namespace would publish its config in memory and remove its directory. A dump URL is no longer fetched for a create that is refused. Registering a schema migration job checks the schema and every linked namespace, so a link written by a binary that does not know fences cannot lead to a migration step being refused halfway through a job. Co-authored-by: Tomasz Szymczyszyn --- libsql-server/src/http/admin/mod.rs | 13 +- libsql-server/src/namespace/fence/registry.rs | 11 + libsql-server/src/namespace/store.rs | 268 +++++++++++++++++- libsql-server/src/schema/db.rs | 17 ++ libsql-server/src/schema/error.rs | 3 + libsql-server/src/schema/scheduler.rs | 258 +++++++++++++++++ libsql-server/tests/common/http.rs | 16 +- libsql-server/tests/fence/admin.rs | 35 +-- libsql-server/tests/fence/lifecycle.rs | 201 +++++++++++++ libsql-server/tests/fence/mod.rs | 40 +++ 10 files changed, 817 insertions(+), 45 deletions(-) create mode 100644 libsql-server/tests/fence/lifecycle.rs diff --git a/libsql-server/src/http/admin/mod.rs b/libsql-server/src/http/admin/mod.rs index be99413913..42d64d4837 100644 --- a/libsql-server/src/http/admin/mod.rs +++ b/libsql-server/src/http/admin/mod.rs @@ -332,10 +332,11 @@ async fn handle_post_config( // Check that the jwt keys are correct parse_jwt_keys(jwt_key)?; } - let store = app_state - .namespaces - .config_store(NamespaceName::from_string(namespace.clone())?) - .await?; + let namespace_name = NamespaceName::from_string(namespace.clone())?; + // Config mutation is lifecycle work: refused while a fence denies it, before the namespace + // is loaded (and again in the metastore transaction that would store it). + app_state.namespaces.check_lifecycle(&namespace_name)?; + let store = app_state.namespaces.config_store(namespace_name).await?; let original = (*store.get()).clone(); let mut updated = original.clone(); updated.block_reads = req.block_reads; @@ -402,6 +403,10 @@ async fn handle_create_namespace( ) -> crate::Result<()> { let mut config = DatabaseConfig::default(); + // Creating over a name whose fence denies lifecycle work is refused before a dump is + // fetched or anything is stored. + app_state.namespaces.check_lifecycle(&namespace)?; + if let Some(jwt_key) = req.jwt_key { // Check that the jwt keys are correct parse_jwt_keys(&jwt_key)?; diff --git a/libsql-server/src/namespace/fence/registry.rs b/libsql-server/src/namespace/fence/registry.rs index 176cfddd4e..24dd0ea23e 100644 --- a/libsql-server/src/namespace/fence/registry.rs +++ b/libsql-server/src/namespace/fence/registry.rs @@ -82,6 +82,17 @@ impl FenceRegistry { } } + /// Refuse generic lifecycle and configuration work on `namespace` while its gate denies it + /// (section 3.3, the lifecycle column): an active fence, a closing transition being + /// installed, a target being created, an indeterminate commit or an unavailable state. A + /// name without a controller has no fence state and is not refused here. + pub fn check_lifecycle(&self, namespace: &NamespaceName) -> Result<(), FenceError> { + match self.get(namespace) { + Some(controller) => controller.gate().permits(OperationClass::Lifecycle), + None => Ok(()), + } + } + /// How many namespaces have an active fence (`docs/NAMESPACE_FENCE.md` section 4.4, /// `active_fences`): a record in any state but `RELEASED` or `TARGET_WRITABLE`, an /// unavailable state, a target being created, or a commit whose outcome is not known yet. diff --git a/libsql-server/src/namespace/store.rs b/libsql-server/src/namespace/store.rs index 09d27a5695..9a83252b53 100644 --- a/libsql-server/src/namespace/store.rs +++ b/libsql-server/src/namespace/store.rs @@ -182,6 +182,17 @@ impl NamespaceStore { namespace: NamespaceName, restore_option: RestoreOption, ) -> anyhow::Result<()> { + // Reset destroys the namespace's data and writes nothing to the metastore, so the fence + // check is made here, under the namespace's transition lock: a fence command either + // finished before this check or starts after the reset (section 3.3, lifecycle). + let _transition = self + .inner + .fences + .controller(&namespace) + .begin_transition() + .await; + self.check_lifecycle(&namespace)?; + // The process for reseting is as follow: // - get a lock on the namespace entry, if the entry exists, then it's a lock on the entry, // if it doesn't exist, insert an empty entry and take a lock on it @@ -221,18 +232,26 @@ impl NamespaceStore { Box::new(move |op| { let this = this.clone(); tokio::spawn(async move { - match op { - ResetOp::Reset(ns) => { - tracing::info!("received reset signal for: {ns}"); - if let Err(e) = this.reset(ns.clone(), RestoreOption::Latest).await { - tracing::error!("error resetting namespace `{ns}`: {e}"); - } - } - } + let _ = this.handle_reset_op(op).await; }); }) } + /// A reset requested by a replica's replicator. A namespace whose fence denies lifecycle + /// work is not reset: the refusal is logged and returned. + async fn handle_reset_op(&self, op: ResetOp) -> anyhow::Result<()> { + match op { + ResetOp::Reset(ns) => { + tracing::info!("received reset signal for: {ns}"); + let result = self.reset(ns.clone(), RestoreOption::Latest).await; + if let Err(e) = &result { + tracing::error!("error resetting namespace `{ns}`: {e}"); + } + result + } + } + } + pub async fn fork( &self, from: NamespaceName, @@ -245,14 +264,23 @@ impl NamespaceStore { } // The destination is refused before anything is stored for it when it is being created - // as a migration target or its fence state is unknown. + // as a migration target, its fence state is unknown, or its fence denies lifecycle work + // (an existing fenced namespace, whose directory the fork would otherwise replace). self.inner.fences.check_available(&to)?; + self.check_lifecycle(&to)?; // check that the source namespace exists if !self.inner.metadata.exists(&from).await { return Err(crate::error::Error::NamespaceDoesntExist(from.to_string())); } + // A fork reads the source's data without a read lease, so it runs under the source's + // transition lock and checks the source's gate under it: a fence command on the source + // (a write or read fence) either finished before this check, and the fork is refused, + // or waits for the fork to finish (section 3.3, fork as source). + let _from_transition = self.inner.fences.controller(&from).begin_transition().await; + self.check_lifecycle(&from)?; + let to_entry = self .inner .store @@ -262,6 +290,9 @@ impl NamespaceStore { if to_lock.is_some() { return Err(crate::error::Error::NamespaceAlreadyExist(to.to_string())); } + // With the destination's entry held, a fence command cannot load the destination, so + // the check cannot go stale before the fork has stored and flushed its config. + self.check_lifecycle(&to)?; // FIXME: we could potentially delete the namespace while trying to fork it if !self.inner.metadata.exists(&from).await { @@ -462,8 +493,10 @@ impl NamespaceStore { db_config: DatabaseConfig, ) -> crate::Result<()> { // A name that is being created as a migration target, or whose fence state is unknown, - // is refused before anything is stored for it. + // is refused before anything is stored for it; so is a name whose fence denies lifecycle + // work (creating over an existing record, with or without a restore). self.inner.fences.check_available(&namespace)?; + self.check_lifecycle(&namespace)?; if let Some(shared_schema_name) = &db_config.shared_schema_name { // we hold a lock for the duration of the namespace creation let _lock = self @@ -552,6 +585,16 @@ impl NamespaceStore { &self.inner.metadata } + /// Refuse generic lifecycle and configuration work on `namespace` (config mutation, + /// delete, reset, fork on either side, create over an existing record, restore, dump load, + /// shared-schema linking, schema migration) while its fence denies it + /// (`docs/NAMESPACE_FENCE.md` section 3.3), without loading the namespace. A name without + /// fence state is not refused here: the existing checks apply to it. Paths that persist + /// through the metastore are refused again inside its transaction. + pub(crate) fn check_lifecycle(&self, namespace: &NamespaceName) -> crate::Result<()> { + Ok(self.inner.fences.check_lifecycle(namespace)?) + } + /// Run one fence command on its namespace, including the drain it starts /// (`docs/NAMESPACE_FENCE.md` sections 5.3 and 8). `AcquireSourceWriteFence` loads the /// namespace first, so that its connection manager and replication log are registered with @@ -947,6 +990,7 @@ pub(crate) mod fence_tests { use super::*; use crate::config::MetaStoreConfig; + use crate::connection::Connection as _; use crate::namespace::configurator::{BaseNamespaceConfig, PrimaryConfig, PrimaryConfigurator}; use crate::namespace::fence::command::{FenceCommand, FenceRequest}; use crate::namespace::fence::outcome::{FenceDetail, FenceOutcome}; @@ -1148,4 +1192,208 @@ pub(crate) mod fence_tests { store.destroy("ns".into(), false).await.unwrap(); assert!(store.inner.fences.get(&"ns".into()).is_none()); } + + fn release(ns: &'static str, command_id: u128) -> FenceRequest { + FenceRequest { + namespace: ns.into(), + operation_id: OP, + command_id: Uuid::from_u128(command_id), + expected_state: FenceState::SourceDraining, + expected_revision: 1, + command: FenceCommand::ReleaseSourceWriteFence, + } + } + + /// Create `ns` holding a table `t` with one row. + async fn create_with_row(store: &NamespaceStore, ns: &'static str) { + store + .create(ns.into(), RestoreOption::Latest, Default::default()) + .await + .unwrap(); + let conn = store + .with(ns.into(), |ns| ns.db.connection_maker()) + .await + .unwrap() + .create() + .await + .unwrap(); + tokio::task::spawn_blocking(move || { + conn.with_raw(|c| c.execute_batch("create table t (x); insert into t values (1);")) + }) + .await + .unwrap() + .unwrap(); + } + + /// The number of rows in `ns`'s table `t`, or the error reading it. + async fn rows(store: &NamespaceStore, ns: &'static str) -> rusqlite::Result { + let conn = store + .with(ns.into(), |ns| ns.db.connection_maker()) + .await + .unwrap() + .create() + .await + .unwrap(); + tokio::task::spawn_blocking(move || { + conn.with_raw(|c| c.query_row("select count(*) from t", (), |r| r.get(0))) + }) + .await + .unwrap() + } + + #[track_caller] + fn assert_fenced(result: crate::Result<()>, outcome: FenceOutcome) { + match result { + Err(Error::NamespaceFence(e)) => assert_eq!(e.outcome(), outcome, "{e}"), + other => panic!("expected {outcome}, got {other:?}"), + } + } + + #[track_caller] + fn assert_fenced_anyhow(result: anyhow::Result<()>, outcome: FenceOutcome) { + match result { + Err(e) => match e.downcast_ref::() { + Some(Error::NamespaceFence(e)) => assert_eq!(e.outcome(), outcome, "{e}"), + _ => panic!("expected {outcome}, got {e:?}"), + }, + Ok(()) => panic!("expected {outcome}, got Ok"), + } + } + + /// Reset, which destroys the namespace's data and recreates it, is refused while the + /// namespace is fenced, both called directly and as the replicator's reset callback does. + #[tokio::test(flavor = "multi_thread")] + async fn reset_refused_while_fenced() { + let tmp = tempdir().unwrap(); + let store = open_store(tmp.path()).await; + create_with_row(&store, "ns").await; + let fence = store.inner.fences.controller(&"ns".into()); + fence + .apply_command(store.meta_store(), acquire("ns"), ctx()) + .await + .unwrap(); + + assert_fenced_anyhow( + store.reset("ns".into(), RestoreOption::Latest).await, + FenceOutcome::MigrationWriteFenced, + ); + assert_fenced_anyhow( + store.handle_reset_op(ResetOp::Reset("ns".into())).await, + FenceOutcome::MigrationWriteFenced, + ); + // The namespace was not touched and still serves reads. + assert_eq!(rows(&store, "ns").await.unwrap(), 1); + + // Once the fence is released, reset works as before, and its data is gone. + fence + .apply_command(store.meta_store(), release("ns", 2), ctx()) + .await + .unwrap(); + assert_eq!(fence.gate().state(), FenceState::Released); + store + .handle_reset_op(ResetOp::Reset("ns".into())) + .await + .unwrap(); + assert!(rows(&store, "ns").await.is_err()); + } + + /// Fork is lifecycle work on both sides: a fenced source is not copied, and a fenced + /// destination (whose directory a fork would replace) is not overwritten. Create over a + /// fenced name, delete and config mutation, including linking the namespace to a shared + /// schema, are refused too. + #[tokio::test(flavor = "multi_thread")] + async fn lifecycle_refused_while_fenced() { + let tmp = tempdir().unwrap(); + let store = open_store(tmp.path()).await; + create_with_row(&store, "src").await; + create_with_row(&store, "other").await; + let fence = store.inner.fences.controller(&"src".into()); + fence + .apply_command(store.meta_store(), acquire("src"), ctx()) + .await + .unwrap(); + let write_fenced = FenceOutcome::MigrationWriteFenced; + + // Fork with the fenced namespace as the source: nothing is created. + assert_fenced( + store + .fork("src".into(), "copy".into(), Default::default(), None) + .await, + write_fenced, + ); + assert!(!store.exists(&"copy".into()).await); + assert!(!tmp.path().join("dbs").join("copy").exists()); + // Fork onto the fenced namespace: its data and config are untouched. + let config_before = store.config_store("src".into()).await.unwrap().get(); + assert_fenced( + store + .fork( + "other".into(), + "src".into(), + DatabaseConfig { + block_reason: Some("fork".into()), + ..Default::default() + }, + None, + ) + .await, + write_fenced, + ); + assert_eq!(rows(&store, "src").await.unwrap(), 1); + let config_after = store.config_store("src".into()).await.unwrap().get(); + assert_eq!(config_after.block_reason, config_before.block_reason); + assert_eq!(config_after.block_writes, config_before.block_writes); + + // Create over it, with or without a restore, and delete. + assert_fenced( + store + .create("src".into(), RestoreOption::Latest, Default::default()) + .await, + write_fenced, + ); + assert_fenced(store.destroy("src".into(), false).await, write_fenced); + + // Config mutation, including linking the namespace to a shared schema, is refused in the + // metastore transaction that would store it. + let handle = store.config_store("src".into()).await.unwrap(); + assert_fenced( + handle + .store(DatabaseConfig { + block_reason: Some("changed".into()), + ..Default::default() + }) + .await, + write_fenced, + ); + assert_fenced( + handle + .store(DatabaseConfig { + shared_schema_name: Some("other".into()), + ..Default::default() + }) + .await, + write_fenced, + ); + assert_eq!(handle.get().block_reason, config_before.block_reason); + assert!(handle.get().shared_schema_name.is_none()); + assert_eq!(rows(&store, "src").await.unwrap(), 1); + + // After release the same operations follow the existing policy again. + fence + .apply_command(store.meta_store(), release("src", 2), ctx()) + .await + .unwrap(); + store + .fork("src".into(), "copy".into(), Default::default(), None) + .await + .unwrap(); + assert_eq!(rows(&store, "copy").await.unwrap(), 1); + assert!(matches!( + store + .create("src".into(), RestoreOption::Latest, Default::default()) + .await, + Err(Error::NamespaceAlreadyExist(_)) + )); + store.destroy("src".into(), false).await.unwrap(); + } } diff --git a/libsql-server/src/schema/db.rs b/libsql-server/src/schema/db.rs index d0bce10128..aef39b7be0 100644 --- a/libsql-server/src/schema/db.rs +++ b/libsql-server/src/schema/db.rs @@ -110,6 +110,23 @@ pub(crate) fn schema_has_linked_dbs( Ok(has_linked) } +/// The namespaces linked to `schema`. +pub(crate) fn linked_namespaces( + conn: &rusqlite::Connection, + schema: &NamespaceName, +) -> Result, Error> { + let mut stmt = + conn.prepare("SELECT namespace FROM shared_schema_links WHERE shared_schema_name = ?")?; + let names = stmt + .query_map([schema.as_str()], |row| row.get::<_, String>(0))? + .collect::>>()?; + // A link whose name does not decode cannot name a fenced namespace. + Ok(names + .into_iter() + .filter_map(|name| NamespaceName::from_string(name).ok()) + .collect()) +} + /// Create a migration job, and returns the job_id pub(super) fn register_schema_migration_job( conn: &mut rusqlite::Connection, diff --git a/libsql-server/src/schema/error.rs b/libsql-server/src/schema/error.rs index 13f21f3c15..528251a2b3 100644 --- a/libsql-server/src/schema/error.rs +++ b/libsql-server/src/schema/error.rs @@ -45,6 +45,8 @@ pub enum Error { InteractiveTxnNotAllowed, #[error("Connection left in transaction state")] ConnectionInTxnState, + #[error("{0}")] + NamespaceFence(#[from] crate::namespace::fence::outcome::FenceError), } impl ResponseError for Error {} @@ -58,6 +60,7 @@ impl IntoResponse for &Error { self.format_err(StatusCode::BAD_REQUEST) } Error::MigrationExecuteError(e) => e.as_ref().into_response(), + Error::NamespaceFence(e) => self.format_err(e.outcome().admin_http_status()), _ => self.format_err(StatusCode::INTERNAL_SERVER_ERROR), } } diff --git a/libsql-server/src/schema/scheduler.rs b/libsql-server/src/schema/scheduler.rs index d9431b2d86..b9844b7373 100644 --- a/libsql-server/src/schema/scheduler.rs +++ b/libsql-server/src/schema/scheduler.rs @@ -410,6 +410,24 @@ impl Scheduler { .schema_locks() .acquire_exlusive(schema.clone()) .await; + // Schema migration is lifecycle work on the schema and on every namespace linked to it. + // Fences refuse shared schemas and linked namespaces, and linking a fenced namespace is + // refused, so this finds nothing in normal operation; if it does (a link made by a binary + // that does not know fences), no job is registered rather than a migration step being + // refused at a fenced namespace's WAL. + self.namespace_store + .check_lifecycle(&schema) + .map_err(fence_error)?; + let linked = with_conn_async(self.migration_db.clone(), { + let schema = schema.clone(); + move |conn| super::db::linked_namespaces(conn, &schema) + }) + .await?; + for namespace in &linked { + self.namespace_store + .check_lifecycle(namespace) + .map_err(fence_error)?; + } with_conn_async(self.migration_db.clone(), move |conn| { register_schema_migration_job(conn, &schema, &migration) }) @@ -427,6 +445,14 @@ impl Scheduler { } } +/// The schema error for a fence refusal returned by `NamespaceStore::check_lifecycle`. +fn fence_error(e: crate::Error) -> Error { + match e { + crate::Error::NamespaceFence(e) => Error::NamespaceFence(e), + e => Error::Registration(e.into()), + } +} + async fn try_step_task( _permit: OwnedSemaphorePermit, namespace_store: NamespaceStore, @@ -1229,4 +1255,236 @@ mod test { .is_err()); } } + + /// Namespace fences and shared schemas (`docs/NAMESPACE_FENCE.md` section 13.4). + mod fence { + use uuid::Uuid; + + use super::*; + use crate::config::MetaStoreConfig; + use crate::namespace::fence::command::{FenceCommand, FenceRequest}; + use crate::namespace::fence::outcome::{FenceDetail, FenceOutcome}; + use crate::namespace::fence::record::ServerIdentity; + use crate::namespace::fence::state::FenceState; + use crate::namespace::meta_store::FenceContext; + + const LOG: Uuid = Uuid::from_u128(0x10); + const OP: Uuid = Uuid::from_u128(0xa); + + fn server() -> ServerIdentity { + ServerIdentity { + build: "test".into(), + instance_id: Uuid::from_u128(0x99), + } + } + + fn acquire(ns: &'static str, command_id: u128) -> FenceRequest { + FenceRequest { + namespace: ns.into(), + operation_id: OP, + command_id: Uuid::from_u128(command_id), + expected_state: FenceState::Unfenced, + expected_revision: 0, + command: FenceCommand::AcquireSourceWriteFence { + expected_log_id: LOG, + drain_policy: None, + }, + } + } + + fn release(ns: &'static str, command_id: u128) -> FenceRequest { + FenceRequest { + namespace: ns.into(), + operation_id: OP, + command_id: Uuid::from_u128(command_id), + expected_state: FenceState::SourceDraining, + expected_revision: 1, + command: FenceCommand::ReleaseSourceWriteFence, + } + } + + /// A primary store with fences enabled and a shared schema `schema` with one linked + /// namespace `linked`. + async fn setup( + path: &Path, + ) -> (NamespaceStore, Scheduler, mpsc::Receiver) { + let (maker, manager) = metastore_connection_maker(None, path).await.unwrap(); + let meta_store = MetaStore::new( + MetaStoreConfig { + namespace_fence: true, + ..Default::default() + }, + path, + maker().unwrap(), + manager, + DatabaseKind::Primary, + ) + .await + .unwrap(); + let (sender, receiver) = mpsc::channel(100); + let config = make_config(sender.into(), path); + let store = + NamespaceStore::new(false, false, 10, meta_store, config, DatabaseKind::Primary) + .await + .unwrap(); + let scheduler = Scheduler::new(store.clone(), maker().unwrap()) + .await + .unwrap(); + store + .create( + "schema".into(), + RestoreOption::Latest, + DatabaseConfig { + is_shared_schema: true, + ..Default::default() + }, + ) + .await + .unwrap(); + store + .create( + "linked".into(), + RestoreOption::Latest, + DatabaseConfig { + shared_schema_name: Some("schema".into()), + ..Default::default() + }, + ) + .await + .unwrap(); + (store, scheduler, receiver) + } + + /// Fence `ns` at the metastore (`SOURCE_DRAINING`, write admission closed). + async fn fence(store: &NamespaceStore, ns: &'static str) { + store + .fence_controller(&ns.into()) + .apply_command( + store.meta_store(), + acquire(ns, 1), + FenceContext::now(server(), Some(LOG)), + ) + .await + .unwrap(); + } + + #[track_caller] + fn assert_fence_error(result: crate::Result<()>, outcome: FenceOutcome) { + match result { + Err(crate::Error::NamespaceFence(e)) => assert_eq!(e.outcome(), outcome, "{e}"), + other => panic!("expected {outcome}, got {other:?}"), + } + } + + /// A shared schema and a namespace linked to one cannot be fenced, and a fenced + /// namespace cannot be linked to a shared schema. + #[tokio::test(flavor = "multi_thread")] + async fn acquire_rejects_shared_schema() { + let tmp = tempdir().unwrap(); + let (store, scheduler, _receiver) = setup(tmp.path()).await; + + for ns in ["schema", "linked"] { + match store.execute_fence_command(acquire(ns, 1), server()).await { + Err(crate::Error::NamespaceFence(e)) => { + assert_eq!(e.outcome(), FenceOutcome::FencePreconditionFailed, "{e}"); + assert_eq!(e.detail(), Some(FenceDetail::SharedSchemaUnsupported)); + } + other => panic!("{ns}: expected shared_schema_unsupported, got {other:?}"), + } + // The refused acquisition left the namespace unfenced and writable. + let gate = store.fence_controller(&ns.into()).gate(); + assert_eq!(gate.state(), FenceState::Unfenced); + assert!(gate.write().is_open()); + } + + // A fenced namespace is not linked to the schema, whether by creating it with a + // shared schema or by changing its config. + store + .create("plain".into(), RestoreOption::Latest, Default::default()) + .await + .unwrap(); + fence(&store, "plain").await; + let linked_config = || DatabaseConfig { + shared_schema_name: Some("schema".into()), + ..Default::default() + }; + assert_fence_error( + store + .create("plain".into(), RestoreOption::Latest, linked_config()) + .await, + FenceOutcome::MigrationWriteFenced, + ); + let handle = store.config_store("plain".into()).await.unwrap(); + assert_fence_error( + handle.store(linked_config()).await, + FenceOutcome::MigrationWriteFenced, + ); + assert!(handle.get().shared_schema_name.is_none()); + let links = super::super::super::db::linked_namespaces( + &scheduler.migration_db.lock(), + &"schema".into(), + ) + .unwrap(); + assert_eq!(links, vec![NamespaceName::from("linked")]); + } + + /// A schema migration is lifecycle work on every linked namespace: if a fenced namespace + /// is linked to the schema (here by writing the link directly, as a binary that does not + /// know fences could), no migration job is registered until the fence is released. + #[tokio::test(flavor = "multi_thread")] + async fn migration_not_registered_while_linked_namespace_fenced() { + let tmp = tempdir().unwrap(); + let (store, scheduler, _receiver) = setup(tmp.path()).await; + store + .create("plain".into(), RestoreOption::Latest, Default::default()) + .await + .unwrap(); + fence(&store, "plain").await; + scheduler + .migration_db + .lock() + .execute( + "INSERT INTO shared_schema_links (shared_schema_name, namespace) \ + VALUES ('schema', 'plain')", + (), + ) + .unwrap(); + + let migration = || Program::seq(&["create table test (c)"]).into(); + match scheduler + .register_migration_job("schema".into(), migration()) + .await + { + Err(Error::NamespaceFence(e)) => { + assert_eq!(e.outcome(), FenceOutcome::MigrationWriteFenced, "{e}") + } + other => panic!("expected MIGRATION_WRITE_FENCED, got {other:?}"), + } + assert!(!super::super::super::db::has_pending_migration_jobs( + &scheduler.migration_db.lock(), + &"schema".into(), + ) + .unwrap()); + + // Released, the namespace is ordinary again and the migration is registered. + store + .fence_controller(&"plain".into()) + .apply_command( + store.meta_store(), + release("plain", 2), + FenceContext::now(server(), Some(LOG)), + ) + .await + .unwrap(); + scheduler + .register_migration_job("schema".into(), migration()) + .await + .unwrap(); + assert!(super::super::super::db::has_pending_migration_jobs( + &scheduler.migration_db.lock(), + &"schema".into(), + ) + .unwrap()); + } + } } diff --git a/libsql-server/tests/common/http.rs b/libsql-server/tests/common/http.rs index a08a928478..ccaad17171 100644 --- a/libsql-server/tests/common/http.rs +++ b/libsql-server/tests/common/http.rs @@ -90,12 +90,26 @@ impl Client { &self, url: &str, body: T, + ) -> anyhow::Result { + self.delete_with_headers(url, &[], body).await + } + + pub(crate) async fn delete_with_headers( + &self, + url: &str, + headers: &[(HeaderName, &str)], + body: T, ) -> anyhow::Result { let bytes: Bytes = serde_json::to_vec(&body)?.into(); let body = Body::from(bytes); - let request = hyper::Request::delete(url) + let mut request = hyper::Request::delete(url) .header("Content-Type", "application/json") .body(body)?; + for (key, val) in headers { + request + .headers_mut() + .insert(key.clone(), val.parse().unwrap()); + } let resp = self.0.request(request).await?; Ok(Response(resp)) diff --git a/libsql-server/tests/fence/admin.rs b/libsql-server/tests/fence/admin.rs index d1ac8413c4..2c64935c46 100644 --- a/libsql-server/tests/fence/admin.rs +++ b/libsql-server/tests/fence/admin.rs @@ -1,44 +1,19 @@ //! The fence admin API over HTTP (`docs/NAMESPACE_FENCE.md` section 4). use hyper::StatusCode; -use serde_json::{json, Value}; +use serde_json::json; use tempfile::tempdir; use uuid::Uuid; -use super::{command_body, connect, make_primary, sim, state_of, Admin, Primary, ADMIN_KEY}; +use super::{ + acquire_body, command_body, connect, load_and_log_id, make_primary, sim, state_of, Admin, + Primary, ADMIN_KEY, +}; fn uuid(n: u128) -> Uuid { Uuid::from_u128(n) } -/// Load `ns` on the server with one write, and return the replication log id the server -/// reports for it. -async fn load_and_log_id(admin: &Admin, ns: &str) -> anyhow::Result { - let conn = connect(ns)?; - conn.execute("create table if not exists t (x)", ()).await?; - conn.execute("insert into t values (1)", ()).await?; - let (status, body) = admin.inspect(ns).await?; - assert_eq!(status, StatusCode::OK, "{body}"); - assert_eq!(state_of(&body), ("UNFENCED", 0), "{body}"); - Ok(body["fence"]["incarnation"]["current_log_id"] - .as_str() - .unwrap_or_else(|| panic!("no current_log_id: {body}")) - .to_string()) -} - -fn acquire_body(op: Uuid, cmd: Uuid, log_id: &str) -> Value { - command_body( - op, - cmd, - "UNFENCED", - 0, - json!({ - "expected_namespace_identity": { "log_id": log_id }, - "drain_policy": { "deadline_ms": 5000, "on_deadline": "fail" }, - }), - ) -} - #[test] fn capabilities() { let mut sim = sim(); diff --git a/libsql-server/tests/fence/lifecycle.rs b/libsql-server/tests/fence/lifecycle.rs new file mode 100644 index 0000000000..90ac6d9749 --- /dev/null +++ b/libsql-server/tests/fence/lifecycle.rs @@ -0,0 +1,201 @@ +//! Lifecycle and configuration operations on fenced namespaces, over the admin API +//! (`docs/NAMESPACE_FENCE.md` section 3.3, the lifecycle column; section 17 row 17). + +use hyper::StatusCode; +use libsql::Value as SqlValue; +use serde_json::{json, Value}; +use tempfile::tempdir; +use uuid::Uuid; + +use super::{ + acquire_body, command_body, connect, load_and_log_id, make_primary, sim, state_of, Admin, + Primary, ADMIN_KEY, +}; + +fn uuid(n: u128) -> Uuid { + Uuid::from_u128(n) +} + +/// Every generic lifecycle and configuration route, attempted on `ns`, with a name for the +/// failure message. `ns` must be refused by each of them. +async fn lifecycle_attempts( + admin: &Admin, + ns: &str, +) -> anyhow::Result> { + let mut results = Vec::new(); + let (status, body) = admin.delete(&format!("/v1/namespaces/{ns}")).await?; + results.push(("delete", status, body)); + let (status, body) = admin + .post(&format!("/v1/namespaces/{ns}/fork/{ns}-copy"), json!({})) + .await?; + results.push(("fork as source", status, body)); + let (status, body) = admin + .post(&format!("/v1/namespaces/other/fork/{ns}"), json!({})) + .await?; + results.push(("fork as destination", status, body)); + // The dump file does not exist: a refusal made before the dump is fetched is the fence's. + let (status, body) = admin + .post( + &format!("/v1/namespaces/{ns}/create"), + json!({ "dump_url": "file:///nonexistent/dump.sql" }), + ) + .await?; + results.push(("create with dump_url", status, body)); + let (status, body) = admin + .post(&format!("/v1/namespaces/{ns}/create"), json!({})) + .await?; + results.push(("create over the record", status, body)); + let (status, body) = admin + .post( + &format!("/v1/namespaces/{ns}/create"), + json!({ "shared_schema_name": "schema" }), + ) + .await?; + results.push(("link to a shared schema", status, body)); + let (status, body) = admin + .post( + &format!("/v1/namespaces/{ns}/config"), + json!({ "block_reads": false, "block_writes": false, "block_reason": null }), + ) + .await?; + results.push(("config", status, body)); + Ok(results) +} + +async fn count_rows(ns: &str) -> anyhow::Result { + let mut rows = connect(ns)?.query("select count(*) from t", ()).await?; + let row = rows.next().await?.expect("one row"); + match row.get_value(0)? { + SqlValue::Integer(n) => Ok(n), + other => anyhow::bail!("unexpected count {other:?}"), + } +} + +#[test] +fn lifecycle_rejected_while_fenced() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + admin.create_namespace("other").await?; + let (status, body) = admin + .post( + "/v1/namespaces/schema/create", + json!({ "shared_schema": true }), + ) + .await?; + assert!(status.is_success(), "{status} {body}"); + + // A write-fenced source. + admin.create_namespace("src").await?; + let log_id = load_and_log_id(&admin, "src").await?; + let source_op = uuid(0x100); + let (status, body) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(source_op, uuid(1), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + // SOURCE_DRAINING at revision 1, then SOURCE_WRITE_FENCED once the drain is proven. + assert_eq!(state_of(&body), ("SOURCE_WRITE_FENCED", 2), "{body}"); + + // A quarantined target. A target cannot be created with a shared schema. + let target_op = uuid(0x200); + let (status, body) = admin + .command( + "tgt", + "target/create-quarantined", + command_body( + target_op, + uuid(2), + "ABSENT", + 0, + json!({ "shared_schema_name": "schema" }), + ), + ) + .await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{body}"); + assert_eq!(body["detail"], "shared_schema_unsupported", "{body}"); + let (status, body) = admin + .command( + "tgt", + "target/create-quarantined", + command_body(target_op, uuid(3), "ABSENT", 0, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body), ("TARGET_QUARANTINED", 1), "{body}"); + + for (ns, state, revision, code) in [ + ("src", "SOURCE_WRITE_FENCED", 2, "MIGRATION_WRITE_FENCED"), + ( + "tgt", + "TARGET_QUARANTINED", + 1, + "MIGRATION_TARGET_QUARANTINED", + ), + ] { + for (what, status, body) in lifecycle_attempts(&admin, ns).await? { + assert_eq!(status, StatusCode::LOCKED, "{what} on {ns}: {body}"); + let message = body["error"].as_str().unwrap_or_default(); + assert!( + message.starts_with(code), + "{what} on {ns}: expected {code}, got {body}" + ); + } + // Nothing moved: same state and revision, and no copy was created. + let (status, body) = admin.inspect(ns).await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body), (state, revision), "{body}"); + let (status, body) = admin + .get(&format!("/v1/namespaces/{ns}-copy/config")) + .await?; + assert_eq!(status, StatusCode::NOT_FOUND, "{body}"); + } + // The source's data is intact and still served to readers. + assert_eq!(count_rows("src").await?, 1); + // Reading config and stats is still allowed. + let (status, body) = admin.get("/v1/namespaces/src/config").await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let (status, body) = admin.get("/v1/namespaces/src/stats").await?; + assert_eq!(status, StatusCode::OK, "{body}"); + + // Released, the source is an ordinary namespace again: it takes writes and lifecycle + // operations follow the existing policy. + let (status, body) = admin + .command( + "src", + "source/release-write-fence", + command_body(source_op, uuid(4), "SOURCE_WRITE_FENCED", 2, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body).0, "RELEASED", "{body}"); + connect("src")? + .execute("insert into t values (2)", ()) + .await?; + assert_eq!(count_rows("src").await?, 2); + let (status, body) = admin + .post( + "/v1/namespaces/src/config", + json!({ "block_reads": false, "block_writes": false }), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let (status, body) = admin + .post("/v1/namespaces/src/fork/src-copy", json!({})) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(count_rows("src-copy").await?, 2); + let (status, body) = admin.delete("/v1/namespaces/src-copy").await?; + assert_eq!(status, StatusCode::OK, "{body}"); + + // The target stays quarantined: it serves no SQL. + assert!(count_rows("tgt").await.is_err()); + Ok(()) + }); + sim.run().unwrap(); +} diff --git a/libsql-server/tests/fence/mod.rs b/libsql-server/tests/fence/mod.rs index a7a149907f..c8be3cde94 100644 --- a/libsql-server/tests/fence/mod.rs +++ b/libsql-server/tests/fence/mod.rs @@ -3,6 +3,7 @@ //! Namespace fence integration tests (`docs/NAMESPACE_FENCE.md`), driven over the admin API. mod admin; +mod lifecycle; use std::path::PathBuf; use std::time::Duration; @@ -126,6 +127,16 @@ impl Admin { .await } + pub async fn delete(&self, path: &str) -> anyhow::Result<(StatusCode, Value)> { + let url = format!("http://primary:9090{path}"); + Self::json( + self.client + .delete_with_headers(&url, &self.headers(), json!({})) + .await?, + ) + .await + } + pub async fn create_namespace(&self, ns: &str) -> anyhow::Result<()> { let (status, body) = self .post(&format!("/v1/namespaces/{ns}/create"), json!({})) @@ -186,3 +197,32 @@ pub fn connect(ns: &str) -> anyhow::Result { )?; Ok(db.connect()?) } + +/// Load `ns` on the server with one write, and return the replication log id the server +/// reports for it. +pub async fn load_and_log_id(admin: &Admin, ns: &str) -> anyhow::Result { + let conn = connect(ns)?; + conn.execute("create table if not exists t (x)", ()).await?; + conn.execute("insert into t values (1)", ()).await?; + let (status, body) = admin.inspect(ns).await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body), ("UNFENCED", 0), "{body}"); + Ok(body["fence"]["incarnation"]["current_log_id"] + .as_str() + .unwrap_or_else(|| panic!("no current_log_id: {body}")) + .to_string()) +} + +/// An `AcquireSourceWriteFence` body for a namespace in `UNFENCED` at revision 0. +pub fn acquire_body(op: Uuid, cmd: Uuid, log_id: &str) -> Value { + command_body( + op, + cmd, + "UNFENCED", + 0, + json!({ + "expected_namespace_identity": { "log_id": log_id }, + "drain_policy": { "deadline_ms": 5000, "on_deadline": "fail" }, + }), + ) +}