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..42d64d4837 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() @@ -326,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; @@ -396,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/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..24dd0ea23e 100644 --- a/libsql-server/src/namespace/fence/registry.rs +++ b/libsql-server/src/namespace/fence/registry.rs @@ -82,6 +82,33 @@ 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. + 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..9a83252b53 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}; @@ -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,12 +585,20 @@ 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 /// 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 +922,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 } @@ -919,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}; @@ -1120,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 8716a60503..ccaad17171 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 } @@ -76,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 new file mode 100644 index 0000000000..2c64935c46 --- /dev/null +++ b/libsql-server/tests/fence/admin.rs @@ -0,0 +1,706 @@ +//! The fence admin API over HTTP (`docs/NAMESPACE_FENCE.md` section 4). + +use hyper::StatusCode; +use serde_json::json; +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) +} + +#[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/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 new file mode 100644 index 0000000000..c8be3cde94 --- /dev/null +++ b/libsql-server/tests/fence/mod.rs @@ -0,0 +1,228 @@ +#![allow(deprecated)] + +//! 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; + +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 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!({})) + .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()?) +} + +/// 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" }, + }), + ) +} 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;