diff --git a/libsql-server/src/admin_shell.rs b/libsql-server/src/admin_shell.rs index 83f6a40e27..9f072993ae 100644 --- a/libsql-server/src/admin_shell.rs +++ b/libsql-server/src/admin_shell.rs @@ -87,11 +87,15 @@ fn run_admitted( }) { Ok(lease) => lease, Err(e) => { + crate::namespace::fence::audit::denied( + &e, + crate::namespace::fence::audit::DenialSurface::AdminShell, + ); return Ok(rpc::Response { resp: Some(Resp::Error(rpc::Error { error: e.to_string(), })), - }) + }); } }; let res = run_one(conn, q); diff --git a/libsql-server/src/config.rs b/libsql-server/src/config.rs index 48dee38dfa..ceaae19b3a 100644 --- a/libsql-server/src/config.rs +++ b/libsql-server/src/config.rs @@ -199,6 +199,50 @@ pub struct MetaStoreConfig { /// How long `SetSourceReadFence` waits for running reads and streams before it cancels /// them, when the request names no drain policy. `None` is the default of 30 seconds. pub namespace_fence_default_read_drain: Option, + /// The separate secret that authorises namespace fence adoption (`AdoptFence`), presented + /// in the `x-libsql-fence-adoption-key` header beside the admin credential. `None` + /// disables adoption. + pub namespace_fence_adoption_key: Option, +} + +/// The secret that authorises namespace fence adoption (`docs/NAMESPACE_FENCE.md` section 12). +/// +/// Only its SHA-256 digest is kept, and a presented key is compared digest to digest without +/// an early exit, so the comparison takes the same time whatever the presented key is. `Debug` +/// never prints it. +#[derive(Clone)] +pub struct FenceAdoptionKey(Arc<[u8; 32]>); + +impl FenceAdoptionKey { + /// `None` for an empty key, which would authorise nothing. + pub fn new(key: &str) -> Option { + if key.is_empty() { + return None; + } + Some(Self(Arc::new(Self::digest(key.as_bytes())))) + } + + /// Whether `presented` is the configured key. + pub fn matches(&self, presented: &[u8]) -> bool { + let presented = Self::digest(presented); + let difference = self + .0 + .iter() + .zip(presented.iter()) + .fold(0u8, |acc, (a, b)| acc | (a ^ b)); + difference == 0 + } + + fn digest(bytes: &[u8]) -> [u8; 32] { + use sha2::Digest as _; + sha2::Sha256::digest(bytes).into() + } +} + +impl std::fmt::Debug for FenceAdoptionKey { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str("FenceAdoptionKey()") + } } #[derive(Debug, Clone)] diff --git a/libsql-server/src/error.rs b/libsql-server/src/error.rs index b6903c4a82..55627a1d64 100644 --- a/libsql-server/src/error.rs +++ b/libsql-server/src/error.rs @@ -174,7 +174,13 @@ impl Error { crate::namespace::fence::outcome::FenceError::from_proxy_stable_code(code, &e.message) }); match fence { - Some(fence) => Error::NamespaceFence(fence), + Some(fence) => { + crate::namespace::fence::audit::denied( + &fence, + crate::namespace::fence::audit::DenialSurface::Proxy, + ); + Error::NamespaceFence(fence) + } None => Error::RpcQueryError(e), } } @@ -184,7 +190,13 @@ impl Error { /// anything else unchanged. pub(crate) fn from_proxy_status(status: tonic::Status) -> Self { match crate::namespace::fence::outcome::FenceError::from_grpc_status(&status) { - Some(fence) => Error::NamespaceFence(fence), + Some(fence) => { + crate::namespace::fence::audit::denied( + &fence, + crate::namespace::fence::audit::DenialSurface::Proxy, + ); + Error::NamespaceFence(fence) + } None => Error::RpcQueryExecutionError(status), } } @@ -195,6 +207,7 @@ impl Error { pub(crate) fn fence_error_response( e: &crate::namespace::fence::outcome::FenceError, ) -> axum::response::Response { + crate::namespace::fence::audit::denied(e, crate::namespace::fence::audit::DenialSurface::Http); let status = e.http_status(); tracing::debug!("HTTP API: {status}, {e}"); (status, axum::Json(e.http_error_body())).into_response() diff --git a/libsql-server/src/hrana/batch.rs b/libsql-server/src/hrana/batch.rs index 4292660088..baaf96184f 100644 --- a/libsql-server/src/hrana/batch.rs +++ b/libsql-server/src/hrana/batch.rs @@ -187,7 +187,13 @@ pub fn batch_error_from_sqld_error(sqld_error: SqldError) -> Result { BatchError::ResponseTooLarge } - SqldError::NamespaceFence(e) => BatchError::Fence(e), + SqldError::NamespaceFence(e) => { + crate::namespace::fence::audit::denied( + &e, + crate::namespace::fence::audit::DenialSurface::Hrana, + ); + BatchError::Fence(e) + } sqld_error => return Err(sqld_error), }) } diff --git a/libsql-server/src/hrana/stmt.rs b/libsql-server/src/hrana/stmt.rs index 49bb65d2f7..88824ced5c 100644 --- a/libsql-server/src/hrana/stmt.rs +++ b/libsql-server/src/hrana/stmt.rs @@ -220,7 +220,13 @@ pub fn stmt_error_from_sqld_error(sqld_error: SqldError) -> Result Ok(StmtError::Blocked { reason }), SqldError::RpcQueryError(e) => Ok(StmtError::Proxy(e.message)), - SqldError::NamespaceFence(e) => Ok(StmtError::Fence(e)), + SqldError::NamespaceFence(e) => { + crate::namespace::fence::audit::denied( + &e, + crate::namespace::fence::audit::DenialSurface::Hrana, + ); + Ok(StmtError::Fence(e)) + } SqldError::RusqliteError(rusqlite_error) | SqldError::RusqliteErrorExtended(rusqlite_error, _) => match rusqlite_error { rusqlite::Error::SqliteFailure(sqlite_error, Some(message)) => { diff --git a/libsql-server/src/http/admin/fence.rs b/libsql-server/src/http/admin/fence.rs index e7587d2781..45622fe782 100644 --- a/libsql-server/src/http/admin/fence.rs +++ b/libsql-server/src/http/admin/fence.rs @@ -4,7 +4,8 @@ //! 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. +//! [`NamespaceStore::execute_fence_command`], which owns replay, the drains and target creation +//! (`AdoptFence` through `execute_fence_command_authorised`, with the adoption key check). use std::sync::Arc; @@ -13,7 +14,7 @@ use axum::response::{IntoResponse, Response}; use axum::routing::{get, post}; use axum::Json; use bytes::Bytes; -use hyper::StatusCode; +use hyper::{HeaderMap, StatusCode}; use serde::de::DeserializeOwned; use serde::Deserialize; use serde_json::{json, Map, Value}; @@ -23,16 +24,18 @@ 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, + AdoptArgs, 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::record::{ + Adoption, 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::meta_store::{FenceCommit, MetastoreProvenance}; use crate::namespace::NamespaceName; use crate::net::Connector; @@ -43,9 +46,11 @@ use super::AppState; /// looks at part of a result. pub const MAX_VALIDATION_QUERY_ROWS: usize = 10_000; +/// The header that carries the adoption key (section 12). +pub const ADOPTION_KEY_HEADER: &str = "x-libsql-fence-adoption-key"; + /// 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] = [ +const SERVED_COMMANDS: [&str; 12] = [ "InspectFence", CommandKind::AcquireSourceWriteFence.as_str(), CommandKind::SetSourceReadFence.as_str(), @@ -57,6 +62,7 @@ const SERVED_COMMANDS: [&str; 11] = [ CommandKind::PublishTargetReadableWriteFenced.as_str(), CommandKind::EnableTargetWrites.as_str(), CommandKind::AbortQuarantinedTarget.as_str(), + CommandKind::AdoptFence.as_str(), ]; /// The fence routes, added to the admin router. @@ -66,7 +72,7 @@ pub(super) fn routes() -> axum::Router>> { move |State(state): State>>, Path(namespace): Path, body: Bytes| async move { - handle_command(state, namespace, kind, body).await + handle_command(state, namespace, kind, body, false).await }, ) }; @@ -117,6 +123,7 @@ pub(super) fn routes() -> axum::Router>> { "/v1/namespaces/:namespace/fence/target/validation-query", post(handle_validation_query), ) + .route("/v1/namespaces/:namespace/fence/adopt", post(handle_adopt)) } // --------------------------------------------------------------------------------------------- @@ -132,8 +139,7 @@ async fn handle_capabilities(State(state): State>>) -> Json( let body = json!({ "outcome": FenceOutcome::Applied.as_str(), "replayed": false, - "fence": fence_json(&namespace, &inspection.fence, controller.as_deref()), + "fence": fence_json( + &namespace, + &inspection.fence, + controller.as_deref(), + &meta.restore_provenance(), + ), "receipts": receipts, "drain": drain_json(controller.as_deref()), }); (StatusCode::OK, Json(body)).into_response() } +/// `AdoptFence` (section 12): the admin credential (checked by the middleware) and the separate +/// adoption key in [`ADOPTION_KEY_HEADER`]. Whether the key matched is handed to the transition, +/// which refuses an unauthorised adoption with `adoption_not_authorised` after replay handling, +/// like every other check. +async fn handle_adopt( + State(state): State>>, + Path(namespace): Path, + headers: HeaderMap, + body: Bytes, +) -> Response { + let presented = headers.get(ADOPTION_KEY_HEADER).map(|v| v.as_bytes()); + let authorised = state + .namespaces + .meta_store() + .fence_adoption_authorised(presented); + handle_command(state, namespace, CommandKind::AdoptFence, body, authorised).await +} + async fn handle_command( state: Arc>, namespace: String, kind: CommandKind, body: Bytes, + adoption_authorised: bool, ) -> Response { if !state.namespaces.meta_store().fence_enabled() { return StatusCode::NOT_FOUND.into_response(); @@ -231,11 +261,18 @@ async fn handle_command( Ok(request) => request, Err(e) => return error_reply(&state, &namespace, e).await, }; - match state - .namespaces - .execute_fence_command(request, server_identity()) - .await - { + let result = if kind == CommandKind::AdoptFence { + state + .namespaces + .execute_fence_command_authorised(request, server_identity(), adoption_authorised) + .await + } else { + state + .namespaces + .execute_fence_command(request, server_identity()) + .await + }; + match result { Ok(commit) => success_reply(&state, &namespace, commit), Err(e) => fence_or_error(&state, &namespace, e).await, } @@ -307,9 +344,10 @@ async fn handle_validation_query( drop(session); let controller = state.namespaces.existing_fence_controller(&namespace); + let provenance = state.namespaces.meta_store().restore_provenance(); let fence = controller .as_ref() - .map(|c| fence_json(&namespace, &c.gate().fence, Some(c))) + .map(|c| fence_json(&namespace, &c.gate().fence, Some(c), &provenance)) .unwrap_or(Value::Null); let body = json!({ "results": results, @@ -396,14 +434,20 @@ async fn error_reply( error: FenceError, ) -> Response { let mut reply = ErrorReply::new(error); + let provenance = state.namespaces.meta_store().restore_provenance(); match state.namespaces.existing_fence_controller(namespace) { Some(controller) => { - reply.fence = fence_json(namespace, &controller.gate().fence, Some(&controller)); + reply.fence = fence_json( + namespace, + &controller.gate().fence, + Some(&controller), + &provenance, + ); 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.fence = fence_json(namespace, &inspection.fence, None, &provenance); } } } @@ -424,13 +468,15 @@ fn success_reply( commit: FenceCommit, ) -> Response { let controller = state.namespaces.existing_fence_controller(namespace); + let provenance = state.namespaces.meta_store().restore_provenance(); let fence = match (&commit.record, &controller) { (Some(record), _) => fence_json( namespace, &StoredFence::Record(record.clone()), controller.as_deref(), + &provenance, ), - (None, Some(c)) => fence_json(namespace, &c.gate().fence, Some(c)), + (None, Some(c)) => fence_json(namespace, &c.gate().fence, Some(c), &provenance), (None, None) => Value::Null, }; let outcome = commit.receipt.outcome; @@ -637,9 +683,12 @@ fn parse_command( } CommandKind::EnableTargetWrites => FenceCommand::EnableTargetWrites, CommandKind::AbortQuarantinedTarget => FenceCommand::AbortQuarantinedTarget, - CommandKind::AdoptFence => { - return Err(invalid("adoption is not served by this route")); - } + CommandKind::AdoptFence => FenceCommand::AdoptFence(AdoptArgs { + current_operation_id: body.uuid("current_operation_id")?, + approvers: body.req("approvers")?, + incident_ref: body.req("incident_ref")?, + reason: body.req("reason")?, + }), }; body.finish()?; Ok(FenceRequest { @@ -894,33 +943,18 @@ fn record_fields(record: &NamespaceFenceRecord, out: &mut Map) { 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(), - ), + Value::Array(record.adoptions.iter().map(adoption_json).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. +/// when the namespace has one, supplies the live admission and the live log id; `provenance` +/// says whether the metastore that holds `fence` was restored from its backup at startup. fn fence_json( namespace: &NamespaceName, fence: &StoredFence, controller: Option<&FenceController>, + provenance: &MetastoreProvenance, ) -> Value { let gate = controller.map(|c| c.gate()); let mut out = Map::new(); @@ -990,11 +1024,36 @@ fn fence_json( out.insert("server".into(), server_json(&server_identity())); out.insert( "provenance".into(), - json!({ "metastore_restored_from_backup": false, "marker": marker }), + json!({ + "metastore_restored_from_backup": provenance.restored_from_backup, + "metastore_restored_generation": provenance.restored_generation.map(|g| g.to_string()), + "marker": marker, + }), ); Value::Object(out) } +/// The `metastore` object of the capability endpoint (section 4.4). +fn metastore_provenance_json(provenance: &MetastoreProvenance) -> Value { + json!({ + "restored_from_backup": provenance.restored_from_backup, + "restored_generation": provenance.restored_generation.map(|g| g.to_string()), + }) +} + +fn adoption_json(a: &Adoption) -> Value { + 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, + }) +} + fn receipt_json(receipt: &CommandReceipt) -> Value { json!({ "operation_id": receipt.operation_id.to_string(), @@ -1007,6 +1066,7 @@ fn receipt_json(receipt: &CommandReceipt) -> Value { "state_after": receipt.state_after.as_str(), "applied_at": timestamp(receipt.applied_at_ms), "instance_id": receipt.instance_id.to_string(), + "adoption": receipt.adoption.as_ref().map(adoption_json), }) } @@ -1022,3 +1082,57 @@ fn stored_receipt_json(stored: &StoredReceipt) -> Value { }), } } + +#[cfg(test)] +mod tests { + use super::*; + + /// A metastore restored from its backup is reported in the fence view (section 4.3) and + /// in the capability endpoint's `metastore` object (section 4.4), with the generation. + #[test] + fn restore_provenance_is_reported() { + let generation = Uuid::from_u128(0x77); + let restored = MetastoreProvenance { + restored_from_backup: true, + restored_generation: Some(generation), + }; + let namespace = NamespaceName::from_string("db".into()).unwrap(); + let unavailable = StoredFence::Unavailable { + detail: FenceDetail::MetastoreBehindMarker, + reason: "the metastore is behind the marker".into(), + marker: None, + }; + + let view = fence_json(&namespace, &unavailable, None, &restored); + assert_eq!( + view["provenance"], + json!({ + "metastore_restored_from_backup": true, + "metastore_restored_generation": generation.to_string(), + "marker": "metastore_behind_marker", + }) + ); + assert_eq!(view["state"], "UNKNOWN_UNAVAILABLE"); + assert_eq!( + metastore_provenance_json(&restored), + json!({ "restored_from_backup": true, "restored_generation": generation.to_string() }) + ); + + let plain = StoredFence::None { + namespace_exists: true, + }; + let view = fence_json(&namespace, &plain, None, &MetastoreProvenance::default()); + assert_eq!( + view["provenance"], + json!({ + "metastore_restored_from_backup": false, + "metastore_restored_generation": null, + "marker": null, + }) + ); + assert_eq!( + metastore_provenance_json(&MetastoreProvenance::default()), + json!({ "restored_from_backup": false, "restored_generation": null }) + ); + } +} diff --git a/libsql-server/src/http/admin/mod.rs b/libsql-server/src/http/admin/mod.rs index 42d64d4837..f0e86f6a7d 100644 --- a/libsql-server/src/http/admin/mod.rs +++ b/libsql-server/src/http/admin/mod.rs @@ -248,8 +248,10 @@ async fn handle_get_index() -> &'static str { "Welcome to the sqld admin API" } -async fn handle_metrics(State(metrics): State) -> String { - metrics.render() +async fn handle_metrics(State(app_state): State>>) -> String { + // The fence gauges are computed from the registry when they are read. + app_state.namespaces.update_fence_gauges(); + app_state.metrics.render() } async fn handle_get_config( diff --git a/libsql-server/src/http/user/dump.rs b/libsql-server/src/http/user/dump.rs index ae0260482e..f278fb9a48 100644 --- a/libsql-server/src/http/user/dump.rs +++ b/libsql-server/src/http/user/dump.rs @@ -131,8 +131,13 @@ pub(crate) async fn dump_stream( conn_maker: Arc>, preserve_row_ids: bool, ) -> crate::Result>> { - let (lease, cancel) = - acquire_stream_lease(fence, LeaseKind::Dump).map_err(Error::NamespaceFence)?; + let (lease, cancel) = acquire_stream_lease(fence, LeaseKind::Dump).map_err(|e| { + crate::namespace::fence::audit::denied( + &e, + crate::namespace::fence::audit::DenialSurface::Dump, + ); + Error::NamespaceFence(e) + })?; let conn = conn_maker.create().await?; @@ -151,9 +156,14 @@ pub(crate) async fn dump_stream( }); match result { Ok(()) => Ok(()), - Err(_) if cancel.is_cancelled() => Err(Error::NamespaceFence(cancelled_by_read_fence( - LeaseKind::Dump, - ))), + Err(_) if cancel.is_cancelled() => { + let e = cancelled_by_read_fence(LeaseKind::Dump); + crate::namespace::fence::audit::denied( + &e, + crate::namespace::fence::audit::DenialSurface::Dump, + ); + Err(Error::NamespaceFence(e)) + } Err(e) => Err(e.into()), } }); diff --git a/libsql-server/src/lib.rs b/libsql-server/src/lib.rs index 66fbbcf762..a3480eae32 100644 --- a/libsql-server/src/lib.rs +++ b/libsql-server/src/lib.rs @@ -11,7 +11,7 @@ use crate::connection::{Connection, MakeConnection}; use crate::database::DatabaseKind; use crate::error::Error; use crate::migration::maybe_migrate; -use crate::namespace::meta_store::{metastore_connection_maker, MetaStore}; +use crate::namespace::meta_store::{metastore_connection_maker_with_provenance, MetaStore}; use crate::net::Accept; use crate::pager::{make_pager, PAGER_CACHE_SIZE}; use crate::rpc::proxy::rpc::proxy_server::Proxy; @@ -611,9 +611,12 @@ where connection_creation_timeout: self.db_config.connection_creation_timeout, }; - let (metastore_conn_maker, meta_store_wal_manager) = - metastore_connection_maker(self.meta_store_config.bottomless.clone(), &self.path) - .await?; + let (metastore_conn_maker, meta_store_wal_manager, metastore_provenance) = + metastore_connection_maker_with_provenance( + self.meta_store_config.bottomless.clone(), + &self.path, + ) + .await?; let meta_conn = metastore_conn_maker()?; let meta_store = MetaStore::new( self.meta_store_config.clone(), @@ -623,6 +626,7 @@ where db_kind, ) .await?; + meta_store.record_restore_provenance(metastore_provenance); let (configurators, make_replication_svc) = self .make_configurators_and_replication_svc( diff --git a/libsql-server/src/main.rs b/libsql-server/src/main.rs index 02a8011c0f..7358b3482f 100644 --- a/libsql-server/src/main.rs +++ b/libsql-server/src/main.rs @@ -16,8 +16,8 @@ use tracing_subscriber::Layer; use tracing_subscriber::{prelude::*, EnvFilter}; use libsql_server::config::{ - AdminApiConfig, BottomlessConfig, DbConfig, HeartbeatConfig, MetaStoreConfig, RpcClientConfig, - RpcServerConfig, TlsConfig, UserApiConfig, + AdminApiConfig, BottomlessConfig, DbConfig, FenceAdoptionKey, HeartbeatConfig, MetaStoreConfig, + RpcClientConfig, RpcServerConfig, TlsConfig, UserApiConfig, }; use libsql_server::net::AddrIncoming; use libsql_server::version::Version; @@ -287,6 +287,17 @@ struct Cli { #[clap(long, env = "SQLD_NAMESPACE_FENCE_KEEPALIVE_INTERVAL_S")] namespace_fence_keepalive_interval_s: Option, + /// The separate secret that authorises namespace fence adoption (incident recovery of an + /// unfinished operation whose owner was lost), presented in the + /// `x-libsql-fence-adoption-key` header beside the admin credential. Adoption is disabled + /// when this is not set. + #[clap( + long, + env = "SQLD_NAMESPACE_FENCE_ADOPTION_KEY", + hide_env_values = true + )] + namespace_fence_adoption_key: Option, + /// Shutdown timeout duration in seconds, defaults to 30 seconds. #[clap(long, env = "SQLD_SHUTDOWN_TIMEOUT")] shutdown_timeout: Option, @@ -689,6 +700,13 @@ fn make_meta_store_config(config: &Cli) -> anyhow::Result { namespace_fence_default_read_drain: config .namespace_fence_default_read_drain_ms .map(Duration::from_millis), + namespace_fence_adoption_key: match config.namespace_fence_adoption_key.as_deref() { + None => None, + Some(key) => Some( + FenceAdoptionKey::new(key) + .context("--namespace-fence-adoption-key must not be empty")?, + ), + }, }) } diff --git a/libsql-server/src/metrics.rs b/libsql-server/src/metrics.rs index e5de216b79..69cd49bb32 100644 --- a/libsql-server/src/metrics.rs +++ b/libsql-server/src/metrics.rs @@ -32,6 +32,14 @@ pub static TOTAL_RESPONSE_SIZE_HIST: Lazy = Lazy::new(|| { describe_histogram!(NAME, "total response size value before connection lock"); register_histogram!(NAME) }); +pub static METASTORE_RESTORED_FROM_BACKUP: Lazy = Lazy::new(|| { + const NAME: &str = "libsql_server_metastore_restored_from_backup"; + describe_gauge!( + NAME, + "1 when the metastore was restored from its backup at startup, 0 otherwise" + ); + register_gauge!(NAME) +}); pub static STREAM_HANDLES_COUNT: Lazy = Lazy::new(|| { const NAME: &str = "libsql_server_stream_handles"; describe_gauge!(NAME, "amount of in-memory stream handles"); diff --git a/libsql-server/src/namespace/fence/audit.rs b/libsql-server/src/namespace/fence/audit.rs new file mode 100644 index 0000000000..009dee94f5 --- /dev/null +++ b/libsql-server/src/namespace/fence/audit.rs @@ -0,0 +1,901 @@ +//! Fence observability (`docs/NAMESPACE_FENCE.md` sections 12 and 15): the fence metrics and +//! the audit log. +//! +//! Every fence command a server answers emits one structured event under the tracing target +//! [`AUDIT_TARGET`], so that a log pipeline can route them apart from the server's operational +//! logs, and is counted in the metrics below. Metric labels are bounded: a namespace, an +//! operation id, a command id, a revision or a caller never appears as a label value; those +//! are in the audit event. + +use std::time::Duration; + +use uuid::Uuid; + +use super::command::CommandKind; +use super::outcome::FenceError; +use super::registry::FenceRegistry; +use super::state::FenceState; +use crate::namespace::meta_store::{FenceCommit, FenceCommitKind}; +use crate::namespace::NamespaceName; + +/// The tracing target of every fence audit event. +pub const AUDIT_TARGET: &str = "libsql_server::fence::audit"; + +pub const TRANSITIONS_TOTAL: &str = "libsql_server_fence_transitions_total"; +pub const DRAIN_DURATION_SECONDS: &str = "libsql_server_fence_drain_duration_seconds"; +pub const FORCED_TOTAL: &str = "libsql_server_fence_forced_total"; +pub const REPLAYS_TOTAL: &str = "libsql_server_fence_replays_total"; +pub const DENIALS_TOTAL: &str = "libsql_server_fence_denials_total"; +pub const NAMESPACES: &str = "libsql_server_fence_namespaces"; +pub const OLDEST_ACTIVE_AGE_SECONDS: &str = "libsql_server_fence_oldest_active_age_seconds"; +pub const ADOPTIONS_TOTAL: &str = "libsql_server_fence_adoptions_total"; + +/// The label value of a command answered with an error that is not a fence outcome (an I/O or +/// metastore failure), in place of an outcome code. +pub const OTHER_ERROR: &str = "ERROR"; + +/// Describe the fence metrics to the recorder, once, so that `/metrics` carries their help +/// text. Called when the namespace store starts. +pub fn describe_metrics() { + static ONCE: std::sync::Once = std::sync::Once::new(); + ONCE.call_once(|| { + metrics::describe_counter!( + TRANSITIONS_TOTAL, + "fence commands answered, by command and outcome code (replays and refusals included)" + ); + metrics::describe_histogram!( + DRAIN_DURATION_SECONDS, + metrics::Unit::Seconds, + "time from the start of a fence drain to its proof, by kind (write, read, import)" + ); + metrics::describe_counter!( + FORCED_TOTAL, + "work ended by a fence drain at its deadline, by kind (rollback, sql_cancel, \ + dump_cancel, stream_termination)" + ); + metrics::describe_counter!( + REPLAYS_TOTAL, + "fence commands that were replays of a recorded command, or reused its command id \ + for a different request (conflict)" + ); + metrics::describe_counter!( + DENIALS_TOTAL, + "requests refused by a namespace fence, by outcome code and surface" + ); + metrics::describe_gauge!( + NAMESPACES, + "namespaces with fence state on this server, by role and state" + ); + metrics::describe_gauge!( + OLDEST_ACTIVE_AGE_SECONDS, + metrics::Unit::Seconds, + "age of the oldest active fence record on this server, 0 when there is none" + ); + metrics::describe_counter!(ADOPTIONS_TOTAL, "committed fence adoptions"); + }); +} + +/// The drain a command waited on. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DrainKind { + Write, + Read, + Import, +} + +impl DrainKind { + pub const fn as_str(self) -> &'static str { + match self { + DrainKind::Write => "write", + DrainKind::Read => "read", + DrainKind::Import => "import", + } + } +} + +/// Work a drain ended at its deadline instead of waiting for it. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +pub enum ForcedKind { + /// A write or import transaction rolled back (`on_deadline: force_rollback`). + Rollback, + /// A running SQL program cancelled by the read drain. + SqlCancel, + /// A dump cancelled by the read drain. + DumpCancel, + /// A replication stream terminated by the read drain. + StreamTermination, +} + +impl ForcedKind { + pub const fn as_str(self) -> &'static str { + match self { + ForcedKind::Rollback => "rollback", + ForcedKind::SqlCancel => "sql_cancel", + ForcedKind::DumpCancel => "dump_cancel", + ForcedKind::StreamTermination => "stream_termination", + } + } +} + +/// Where a request refused by a fence was refused. +/// +/// A denial is counted once at each surface it crosses: a write sent to a replica and refused +/// by the primary is counted under `rpc` on the primary, and under `proxy` and the replica's +/// user protocol on the replica; a lifecycle operation refused over the admin API is counted +/// under `lifecycle` and `http`. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DenialSurface { + /// An HTTP error response: legacy `/`, the admin API's lifecycle routes, schema errors. + Http, + /// A Hrana error: `/v1`, `/v2`, `/v3`, cursors and WebSocket streams. + Hrana, + /// The primary's proxy service answering a replica. + Rpc, + /// A replica mapping a denial its primary returned through the write proxy. + Proxy, + /// `/dump`. + Dump, + /// The replication service (`hello`, `log_entries`, `batch_log_entries`, `snapshot`). + Replication, + /// The admin shell. + AdminShell, + /// Delete, reset, fork, restore, config and schema operations. + Lifecycle, +} + +impl DenialSurface { + #[cfg(test)] + pub const ALL: [DenialSurface; 8] = [ + DenialSurface::Http, + DenialSurface::Hrana, + DenialSurface::Rpc, + DenialSurface::Proxy, + DenialSurface::Dump, + DenialSurface::Replication, + DenialSurface::AdminShell, + DenialSurface::Lifecycle, + ]; + + pub const fn as_str(self) -> &'static str { + match self { + DenialSurface::Http => "http", + DenialSurface::Hrana => "hrana", + DenialSurface::Rpc => "rpc", + DenialSurface::Proxy => "proxy", + DenialSurface::Dump => "dump", + DenialSurface::Replication => "replication", + DenialSurface::AdminShell => "admin_shell", + DenialSurface::Lifecycle => "lifecycle", + } + } +} + +/// Count a request refused by a fence at `surface`. +pub fn denied(error: &FenceError, surface: DenialSurface) { + metrics::increment_counter!( + DENIALS_TOTAL, + "code" => error.outcome().as_str(), + "surface" => surface.as_str(), + ); +} + +/// Count `count` pieces of work of `kind` ended by a drain. +pub fn forced(kind: ForcedKind, count: usize) { + if count > 0 { + metrics::counter!(FORCED_TOTAL, count as u64, "kind" => kind.as_str()); + } +} + +/// What a command's drain did, filled in by the drain while it runs under the transition and +/// reported with the command's audit event. +#[derive(Debug, Clone, Default)] +pub struct CommandReport { + drain: Option<(DrainKind, Duration)>, + forced: Vec, +} + +impl CommandReport { + /// The drain of `kind` was proven after `duration`. Observed in the drain histogram. + pub fn drained(&mut self, kind: DrainKind, duration: Duration) { + metrics::histogram!( + DRAIN_DURATION_SECONDS, + duration.as_secs_f64(), + "kind" => kind.as_str(), + ); + self.drain = Some((kind, duration)); + } + + /// The drain ended `count` pieces of work of `kind` at its deadline. Counted. + pub fn forced(&mut self, kind: ForcedKind, count: usize) { + if count == 0 { + return; + } + forced(kind, count); + if !self.forced.contains(&kind) { + self.forced.push(kind); + self.forced.sort(); + } + } + + #[cfg(test)] + pub fn forced_kinds(&self) -> &[ForcedKind] { + &self.forced + } + + fn drain_ms(&self) -> Option { + self.drain.map(|(_, d)| d.as_millis()) + } + + fn forced_list(&self) -> String { + self.forced + .iter() + .map(|k| k.as_str()) + .collect::>() + .join(",") + } +} + +/// The request fields of a command, taken before the command runs, for its audit event. +#[derive(Debug, Clone)] +pub struct CommandAudit { + pub namespace: NamespaceName, + pub operation_id: Uuid, + pub command_id: Uuid, + pub command: CommandKind, + pub expected_revision: u64, + /// The published state when the command took the transition lock. + pub state_before: FenceState, +} + +impl CommandAudit { + pub fn new(request: &super::command::FenceRequest, state_before: FenceState) -> Self { + Self { + namespace: request.namespace.clone(), + operation_id: request.operation_id, + command_id: request.command_id, + command: request.command.kind(), + expected_revision: request.expected_revision, + state_before, + } + } +} + +/// Count a command's answer and emit its audit event: one per command answered, whether it +/// committed, replayed a recorded answer or was refused. +pub fn command_finished( + audit: &CommandAudit, + result: &crate::Result, + report: &CommandReport, +) { + let command = audit.command.as_str(); + match result { + Ok(commit) => { + let receipt = &commit.receipt; + let replay = match commit.kind { + FenceCommitKind::Committed => None, + FenceCommitKind::Replayed => Some("replay"), + FenceCommitKind::Resumed => Some("resume"), + }; + metrics::increment_counter!( + TRANSITIONS_TOTAL, + "command" => command, + "outcome" => receipt.outcome.as_str(), + ); + if replay.is_some() { + metrics::increment_counter!(REPLAYS_TOTAL, "result" => "replay"); + } + let adoption = receipt + .adoption + .as_ref() + .filter(|_| commit.kind == FenceCommitKind::Committed); + if let Some(adoption) = adoption { + metrics::increment_counter!(ADOPTIONS_TOTAL); + tracing::info!( + target: AUDIT_TARGET, + event = "namespace_fence_adopted", + namespace = %audit.namespace, + command, + outcome = receipt.outcome.as_str(), + state_before = audit.state_before.as_str(), + state = receipt.state_after.as_str(), + previous_operation_id = %adoption.previous_operation_id, + operation_id = %adoption.new_operation_id, + command_id = %adoption.command_id, + approvers = ?adoption.approvers, + incident_ref = %adoption.incident_ref, + reason = %adoption.reason, + revision_before = receipt.revision_before, + revision_after = receipt.revision_after, + server_instance = %receipt.instance_id, + "namespace fence adopted" + ); + return; + } + tracing::info!( + target: AUDIT_TARGET, + event = "namespace_fence_command", + namespace = %audit.namespace, + operation_id = %audit.operation_id, + command_id = %audit.command_id, + command, + outcome = receipt.outcome.as_str(), + replay = replay.unwrap_or("none"), + revision_before = receipt.revision_before, + revision_after = receipt.revision_after, + state_before = audit.state_before.as_str(), + state_after = receipt.state_after.as_str(), + drain = report.drain.map(|(k, _)| k.as_str()).unwrap_or("none"), + drain_ms = ?report.drain_ms(), + forced = %report.forced_list(), + server_instance = %receipt.instance_id, + "namespace fence command answered" + ); + } + Err(error) => { + let fence = error.fence_error(); + let outcome = fence.map(|e| e.outcome().as_str()).unwrap_or(OTHER_ERROR); + metrics::increment_counter!( + TRANSITIONS_TOTAL, + "command" => command, + "outcome" => outcome, + ); + let conflict = fence + .is_some_and(|e| e.outcome() == super::outcome::FenceOutcome::FenceCommandConflict); + if conflict { + metrics::increment_counter!(REPLAYS_TOTAL, "result" => "conflict"); + } + tracing::info!( + target: AUDIT_TARGET, + event = "namespace_fence_command", + namespace = %audit.namespace, + operation_id = %audit.operation_id, + command_id = %audit.command_id, + command, + outcome, + detail = fence.and_then(|e| e.detail()).map(|d| d.as_str()).unwrap_or("none"), + replay = if conflict { "conflict" } else { "none" }, + expected_revision = audit.expected_revision, + state_before = audit.state_before.as_str(), + drain = report.drain.map(|(k, _)| k.as_str()).unwrap_or("none"), + drain_ms = ?report.drain_ms(), + forced = %report.forced_list(), + server_instance = %super::server_identity().instance_id, + error = %error, + "namespace fence command refused" + ); + } + } +} + +/// Set the fence gauges from the registry: how many namespaces are in each state, by role, and +/// the age of the oldest active record. Every (role, state) pair is written, so a state that +/// emptied reads 0. Called before `/metrics` is rendered. +pub fn update_gauges(registry: &FenceRegistry, now_ms: i64) { + let census = registry.census(); + for state in FenceState::ALL { + let Some(role) = gauge_role(state) else { + continue; + }; + let count = census.iter().filter(|(s, _)| *s == state).count(); + metrics::gauge!( + NAMESPACES, + count as f64, + "role" => role, + "state" => state.as_str(), + ); + } + let oldest = census + .iter() + .filter_map(|(state, created_at_ms)| created_at_ms.filter(|_| state.is_active())) + .min(); + let age = oldest + .map(|created| (now_ms.saturating_sub(created)).max(0) as f64 / 1000.0) + .unwrap_or(0.0); + metrics::gauge!(OLDEST_ACTIVE_AGE_SECONDS, age); +} + +/// The `role` label of a state counted in [`NAMESPACES`]: the record's role, `unknown` for a +/// namespace whose fence state cannot be established, none for an ordinary namespace. +fn gauge_role(state: FenceState) -> Option<&'static str> { + match state.role() { + Some(role) => Some(match role { + super::state::Role::Source => "source", + super::state::Role::Target => "target", + }), + None if state == FenceState::UnknownUnavailable => Some("unknown"), + None => None, + } +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + use std::io::Write; + use std::sync::{Arc, Mutex}; + + use metrics_util::debugging::{DebugValue, DebuggingRecorder, Snapshotter}; + use metrics_util::MetricKind; + use uuid::Uuid; + + use super::*; + use crate::namespace::fence::command::{AdoptArgs, FenceCommand, FenceRequest}; + use crate::namespace::fence::outcome::FenceOutcome; + use crate::namespace::fence::record::tests::sample_record; + use crate::namespace::fence::record::{Adoption, CommandReceipt}; + use crate::namespace::fence::state::FenceState; + use crate::namespace::fence::store::StoredFence; + + #[derive(Clone, Default)] + struct Captured(Arc>>); + + impl Write for Captured { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.lock().unwrap().extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + pub(crate) fn capture(f: impl FnOnce()) -> String { + let captured = Captured::default(); + let writer = captured.clone(); + let subscriber = tracing_subscriber::fmt() + .with_ansi(false) + .with_max_level(tracing::Level::INFO) + .with_writer(move || writer.clone()) + .finish(); + tracing::subscriber::with_default(subscriber, f); + let bytes = captured.0.lock().unwrap().clone(); + String::from_utf8(bytes).unwrap() + } + + fn adopt_request() -> FenceRequest { + FenceRequest { + namespace: "ns".into(), + operation_id: Uuid::from_u128(0xb), + command_id: Uuid::from_u128(3), + expected_state: FenceState::SourceWriteFenced, + expected_revision: 2, + command: FenceCommand::AdoptFence(AdoptArgs { + current_operation_id: Uuid::from_u128(0xa), + approvers: vec!["alice".into(), "bob".into()], + incident_ref: "INC-1".into(), + reason: "control record lost".into(), + }), + } + } + + fn acquire_request() -> FenceRequest { + FenceRequest { + namespace: "ns".into(), + operation_id: Uuid::from_u128(0xc), + command_id: Uuid::from_u128(4), + expected_state: FenceState::Unfenced, + expected_revision: 0, + command: FenceCommand::AcquireSourceWriteFence { + expected_log_id: Uuid::from_u128(0x10), + drain_policy: None, + }, + } + } + + fn commit( + request: &FenceRequest, + kind: FenceCommitKind, + outcome: FenceOutcome, + (revision_before, revision_after): (u64, u64), + adoption: Option, + ) -> FenceCommit { + FenceCommit { + kind, + receipt: CommandReceipt { + namespace: "ns".into(), + operation_id: request.operation_id, + command_id: request.command_id, + command: request.command.kind(), + fingerprint: request.fingerprint(), + outcome, + revision_before, + revision_after, + state_after: FenceState::SourceWriteFenced, + applied_at_ms: 1_000, + instance_id: Uuid::from_u128(0x99), + adoption, + }, + record: None, + created_config: None, + } + } + + fn adoption_entry() -> Adoption { + Adoption { + previous_operation_id: Uuid::from_u128(0xa), + new_operation_id: Uuid::from_u128(0xb), + command_id: Uuid::from_u128(3), + approvers: vec!["alice".into(), "bob".into()], + incident_ref: "INC-1".into(), + reason: "control record lost".into(), + at_ms: 1_000, + revision: 3, + } + } + + #[track_caller] + fn assert_contains_all(out: &str, expected: &[&str]) { + for expected in expected { + assert!(out.contains(expected), "`{expected}` missing from {out}"); + } + } + + /// A committed adoption is one event under the audit target with every field section 12 + /// asks for; a replay of it is an ordinary command event without them. + #[test] + fn adoption_event_fields() { + let request = adopt_request(); + let audit = CommandAudit::new(&request, FenceState::SourceWriteFenced); + let committed = commit( + &request, + FenceCommitKind::Committed, + FenceOutcome::Applied, + (2, 3), + Some(adoption_entry()), + ); + let out = capture(|| command_finished(&audit, &Ok(committed), &CommandReport::default())); + assert_eq!(out.lines().count(), 1, "{out}"); + assert_contains_all( + &out, + &[ + AUDIT_TARGET, + "namespace fence adopted", + "event=\"namespace_fence_adopted\"", + "namespace=ns", + "command=\"AdoptFence\"", + "outcome=\"APPLIED\"", + "state_before=\"SOURCE_WRITE_FENCED\"", + "state=\"SOURCE_WRITE_FENCED\"", + &format!("previous_operation_id={}", Uuid::from_u128(0xa)), + &format!("operation_id={}", Uuid::from_u128(0xb)), + &format!("command_id={}", Uuid::from_u128(3)), + "approvers=[\"alice\", \"bob\"]", + "incident_ref=INC-1", + "reason=control record lost", + "revision_before=2", + "revision_after=3", + &format!("server_instance={}", Uuid::from_u128(0x99)), + ], + ); + + let replayed = commit( + &request, + FenceCommitKind::Replayed, + FenceOutcome::Applied, + (2, 3), + Some(adoption_entry()), + ); + let out = capture(|| command_finished(&audit, &Ok(replayed), &CommandReport::default())); + assert_eq!(out.lines().count(), 1, "{out}"); + assert!(!out.contains("namespace_fence_adopted"), "{out}"); + assert_contains_all( + &out, + &["event=\"namespace_fence_command\"", "replay=\"replay\""], + ); + } + + /// Every answer is one event under the audit target: a committed drain with its duration + /// and forced actions, a replay, and a refusal (a command-id conflict) with its code. + #[test] + fn audit_event_fields() { + let request = acquire_request(); + let audit = CommandAudit::new(&request, FenceState::Unfenced); + let mut report = CommandReport::default(); + report.forced(ForcedKind::Rollback, 1); + report.forced(ForcedKind::Rollback, 1); + report.forced(ForcedKind::SqlCancel, 0); + report.drained(DrainKind::Write, Duration::from_millis(1_500)); + assert_eq!(report.forced_kinds(), &[ForcedKind::Rollback]); + + let committed = commit( + &request, + FenceCommitKind::Committed, + FenceOutcome::Applied, + (0, 2), + None, + ); + let out = capture(|| command_finished(&audit, &Ok(committed), &report)); + assert_eq!(out.lines().count(), 1, "{out}"); + assert_contains_all( + &out, + &[ + AUDIT_TARGET, + "namespace fence command answered", + "event=\"namespace_fence_command\"", + "namespace=ns", + &format!("operation_id={}", Uuid::from_u128(0xc)), + &format!("command_id={}", Uuid::from_u128(4)), + "command=\"AcquireSourceWriteFence\"", + "outcome=\"APPLIED\"", + "replay=\"none\"", + "revision_before=0", + "revision_after=2", + "state_before=\"UNFENCED\"", + "state_after=\"SOURCE_WRITE_FENCED\"", + "drain=\"write\"", + "drain_ms=Some(1500)", + "forced=rollback", + &format!("server_instance={}", Uuid::from_u128(0x99)), + ], + ); + + let out = capture(|| { + let replayed = commit( + &request, + FenceCommitKind::Replayed, + FenceOutcome::Applied, + (0, 2), + None, + ); + command_finished(&audit, &Ok(replayed), &CommandReport::default()) + }); + assert_contains_all(&out, &["replay=\"replay\"", "drain=\"none\"", "forced= "]); + + let conflict = FenceError::new(FenceOutcome::FenceCommandConflict, "reused command id"); + let out = capture(|| { + command_finished( + &audit, + &Err(crate::Error::NamespaceFence(conflict)), + &CommandReport::default(), + ) + }); + assert_eq!(out.lines().count(), 1, "{out}"); + assert_contains_all( + &out, + &[ + AUDIT_TARGET, + "namespace fence command refused", + "outcome=\"FENCE_COMMAND_CONFLICT\"", + "replay=\"conflict\"", + "expected_revision=0", + "state_before=\"UNFENCED\"", + &format!( + "server_instance={}", + super::super::server_identity().instance_id + ), + ], + ); + + let out = capture(|| { + command_finished( + &audit, + &Err(crate::Error::NamespaceStoreShutdown), + &CommandReport::default(), + ) + }); + assert_contains_all(&out, &["outcome=\"ERROR\"", "detail=\"none\""]); + } + + type Snapshot = HashMap<(MetricKind, String, Vec<(String, String)>), DebugValue>; + + fn snapshot() -> Snapshot { + Snapshotter::current_thread_snapshot() + .expect("per-thread recorder installed") + .into_vec() + .into_iter() + .map(|(key, _, _, value)| { + let (kind, key) = key.into_parts(); + let mut labels: Vec<_> = key + .labels() + .map(|l| (l.key().to_string(), l.value().to_string())) + .collect(); + labels.sort(); + ((kind, key.name().to_string(), labels), value) + }) + .collect() + } + + fn labels(pairs: &[(&str, &str)]) -> Vec<(String, String)> { + let mut labels: Vec<_> = pairs + .iter() + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect(); + labels.sort(); + labels + } + + #[track_caller] + fn counter(s: &Snapshot, name: &str, pairs: &[(&str, &str)]) -> u64 { + match s.get(&(MetricKind::Counter, name.to_string(), labels(pairs))) { + Some(DebugValue::Counter(v)) => *v, + other => panic!("counter {name} {pairs:?}: {other:?} in {s:?}"), + } + } + + #[track_caller] + fn gauge(s: &Snapshot, name: &str, pairs: &[(&str, &str)]) -> f64 { + match s.get(&(MetricKind::Gauge, name.to_string(), labels(pairs))) { + Some(DebugValue::Gauge(v)) => v.0, + other => panic!("gauge {name} {pairs:?}: {other:?} in {s:?}"), + } + } + + /// Each metric of section 15 is recorded with its bounded labels, and no label value is a + /// namespace, an operation id or a command id. + #[test] + fn metrics_and_bounded_labels() { + let _ = DebuggingRecorder::per_thread().install(); + let acquire = acquire_request(); + let audit = CommandAudit::new(&acquire, FenceState::Unfenced); + let mut report = CommandReport::default(); + report.forced(ForcedKind::Rollback, 2); + report.forced(ForcedKind::StreamTermination, 1); + report.drained(DrainKind::Write, Duration::from_millis(20)); + let applied = commit( + &acquire, + FenceCommitKind::Committed, + FenceOutcome::Applied, + (0, 2), + None, + ); + let replayed = FenceCommit { + kind: FenceCommitKind::Replayed, + ..applied.clone() + }; + let conflict = FenceError::new(FenceOutcome::FenceCommandConflict, "reused command id"); + let _ = capture(|| { + command_finished(&audit, &Ok(applied), &report); + command_finished(&audit, &Ok(replayed), &CommandReport::default()); + command_finished( + &audit, + &Err(crate::Error::NamespaceFence(conflict)), + &CommandReport::default(), + ); + let adopt = adopt_request(); + let adopted = commit( + &adopt, + FenceCommitKind::Committed, + FenceOutcome::Applied, + (2, 3), + Some(adoption_entry()), + ); + command_finished( + &CommandAudit::new(&adopt, FenceState::SourceWriteFenced), + &Ok(adopted), + &CommandReport::default(), + ); + }); + let fenced = FenceError::new(FenceOutcome::MigrationWriteFenced, "fenced"); + for surface in DenialSurface::ALL { + denied(&fenced, surface); + } + + let mut quarantined = sample_record(); + quarantined.namespace = "t1".into(); + quarantined.role = super::super::state::Role::Target; + quarantined.state = FenceState::TargetQuarantined; + quarantined.created_at_ms = 10_000; + let mut released = sample_record(); + released.namespace = "s2".into(); + released.state = FenceState::Released; + released.created_at_ms = 1_000; + let registry = FenceRegistry::seeded([ + ("s1".into(), StoredFence::Record(sample_record())), + ("s2".into(), StoredFence::Record(released)), + ("t1".into(), StoredFence::Record(quarantined)), + ( + "u1".into(), + StoredFence::Unavailable { + detail: super::super::outcome::FenceDetail::CorruptRecord, + reason: "test".into(), + marker: None, + }, + ), + ( + "plain".into(), + StoredFence::None { + namespace_exists: true, + }, + ), + ]); + // s1 (active) was created at 100 ms, s2 (released, not active) at 1 s. + update_gauges(®istry, 60_100); + + let s = snapshot(); + let cmd = [("command", "AcquireSourceWriteFence")]; + assert_eq!( + counter(&s, TRANSITIONS_TOTAL, &[cmd[0], ("outcome", "APPLIED")]), + 2 + ); + assert_eq!( + counter( + &s, + TRANSITIONS_TOTAL, + &[cmd[0], ("outcome", "FENCE_COMMAND_CONFLICT")] + ), + 1 + ); + assert_eq!(counter(&s, REPLAYS_TOTAL, &[("result", "replay")]), 1); + assert_eq!(counter(&s, REPLAYS_TOTAL, &[("result", "conflict")]), 1); + assert_eq!(counter(&s, FORCED_TOTAL, &[("kind", "rollback")]), 2); + assert_eq!( + counter(&s, FORCED_TOTAL, &[("kind", "stream_termination")]), + 1 + ); + assert_eq!(counter(&s, ADOPTIONS_TOTAL, &[]), 1); + match s.get(&( + MetricKind::Histogram, + DRAIN_DURATION_SECONDS.to_string(), + labels(&[("kind", "write")]), + )) { + Some(DebugValue::Histogram(v)) => assert_eq!(v.len(), 1), + other => panic!("drain histogram: {other:?}"), + } + for surface in DenialSurface::ALL { + assert_eq!( + counter( + &s, + DENIALS_TOTAL, + &[ + ("code", "MIGRATION_WRITE_FENCED"), + ("surface", surface.as_str()) + ] + ), + 1 + ); + } + let ns = |role, state| gauge(&s, NAMESPACES, &[("role", role), ("state", state)]); + assert_eq!(ns("source", "SOURCE_WRITE_FENCED"), 1.0); + assert_eq!(ns("source", "RELEASED"), 1.0); + assert_eq!(ns("target", "TARGET_QUARANTINED"), 1.0); + assert_eq!(ns("unknown", "UNKNOWN_UNAVAILABLE"), 1.0); + assert_eq!(ns("source", "SOURCE_DRAINING"), 0.0); + assert_eq!(ns("target", "TARGET_WRITABLE"), 0.0); + assert_eq!(gauge(&s, OLDEST_ACTIVE_AGE_SECONDS, &[]), 60.0); + + let forbidden = [ + "ns".to_string(), + "s1".to_string(), + "t1".to_string(), + Uuid::from_u128(0xb).to_string(), + Uuid::from_u128(0xc).to_string(), + Uuid::from_u128(3).to_string(), + Uuid::from_u128(4).to_string(), + ]; + for ((_, name, labels), _) in &s { + if !name.starts_with("libsql_server_fence_") { + continue; + } + for (key, value) in labels { + assert!( + matches!( + key.as_str(), + "command" + | "outcome" + | "kind" + | "result" + | "code" + | "surface" + | "role" + | "state" + ), + "{name}: unexpected label {key}" + ); + assert!(!forbidden.contains(value), "{name}: label {key}={value}"); + } + } + + // An emptied registry sets every gauge back to zero. + update_gauges(&FenceRegistry::seeded([]), 60_100); + let s = snapshot(); + assert_eq!( + gauge( + &s, + NAMESPACES, + &[("role", "source"), ("state", "SOURCE_WRITE_FENCED")] + ), + 0.0 + ); + assert_eq!(gauge(&s, OLDEST_ACTIVE_AGE_SECONDS, &[]), 0.0); + } +} diff --git a/libsql-server/src/namespace/fence/controller.rs b/libsql-server/src/namespace/fence/controller.rs index bc358818a4..4b05a249db 100644 --- a/libsql-server/src/namespace/fence/controller.rs +++ b/libsql-server/src/namespace/fence/controller.rs @@ -500,14 +500,18 @@ impl FenceController { /// Cancel every read lease held now (the read drain's deadline). Each lease's work is asked /// to stop once; the leases stay counted until they are actually released. Returns how many - /// were asked. - pub(crate) fn cancel_read_leases(&self) -> usize { + /// were asked, counted by the kind of work asked to stop. + pub(crate) fn cancel_read_leases_by_kind(&self) -> ReadLeaseCounts { let leases = self.read_leases.lock(); - let mut asked = 0; + let mut asked = ReadLeaseCounts::default(); for entry in leases.live.values() { if !entry.cancelled.swap(true, Ordering::AcqRel) { (entry.cancel)(); - asked += 1; + match entry.kind { + LeaseKind::Sql => asked.sql += 1, + LeaseKind::Dump => asked.dump += 1, + LeaseKind::Replication => asked.replication += 1, + } } } asked @@ -757,6 +761,7 @@ impl FenceController { let guard = self.transition_lock.clone().lock_owned().await; Transition { controller: self.clone(), + report: Default::default(), _guard: guard, } } @@ -879,6 +884,8 @@ impl FenceController { /// dropped. pub struct Transition { controller: Arc, + /// What the command's drain did, for its audit event (section 15). + pub(crate) report: super::audit::CommandReport, _guard: OwnedMutexGuard<()>, } @@ -1019,9 +1026,11 @@ impl Transition { let result = match controller.hook(HookPoint::BeforeMetastoreCommit).await { HookOutcome::Continue => run.await, + #[cfg(test)] HookOutcome::Fail(e) => return Err(e.into()), // A commit that failed without applying, but whose outcome the controller cannot // know (test hook). + #[cfg(test)] HookOutcome::Indeterminate => { Err(indeterminate(key, "the commit was not acknowledged (test hook)").into()) } @@ -1030,6 +1039,7 @@ impl Transition { let result = match result { Ok(commit) => match controller.hook(HookPoint::AfterMetastoreCommit).await { HookOutcome::Continue => Ok(commit), + #[cfg(test)] HookOutcome::Indeterminate | HookOutcome::Fail(_) => Err(indeterminate( key, "the commit was not acknowledged (test hook)", diff --git a/libsql-server/src/namespace/fence/drain.rs b/libsql-server/src/namespace/fence/drain.rs index 4c93334b10..a0aef2f2cb 100644 --- a/libsql-server/src/namespace/fence/drain.rs +++ b/libsql-server/src/namespace/fence/drain.rs @@ -16,6 +16,7 @@ use tokio::time::Instant; use crate::namespace::meta_store::{FenceCommit, FenceCommitKind, FenceContext, MetaStore}; +use super::audit::{self, CommandAudit, CommandReport, DrainKind, ForcedKind}; use super::command::{DrainPolicy, FenceCommand, FenceRequest, OnDeadline}; use super::controller::{FenceController, LiveWriteDrain, Transition}; use super::hooks::HookPoint; @@ -52,7 +53,8 @@ impl FenceController { let meta = meta.clone(); tokio::spawn(async move { let mut transition = this.begin_transition().await; - match request.command { + let audit = CommandAudit::new(&request, this.gate().state()); + let result = match request.command { FenceCommand::AcquireSourceWriteFence { .. } => { acquire_source_write_fence(&mut transition, &meta, request, ctx).await } @@ -63,7 +65,9 @@ impl FenceController { super::import::seal_target_import(&mut transition, &meta, request, ctx).await } _ => transition.apply(&meta, request, ctx).await, - } + }; + audit::command_finished(&audit, &result, &transition.report); + result }) .await? } @@ -124,10 +128,14 @@ pub async fn acquire_source_write_fence( let drain_key = (commit.receipt.operation_id, commit.receipt.command_id); // Steps 5 and 6. - let boundary = match drain_writers(&controller, policy).await { + let started = Instant::now(); + let boundary = match drain_writers(&controller, policy, &mut transition.report).await { Some(boundary) => boundary, None => return Ok(commit), }; + transition + .report + .drained(DrainKind::Write, started.elapsed()); // Step 7. let acquired_on = commit.record.as_ref().and_then(|r| r.identity.log_id); @@ -163,6 +171,7 @@ pub async fn acquire_source_write_fence( async fn drain_writers( controller: &FenceController, policy: DrainPolicy, + report: &mut CommandReport, ) -> Option { let namespace = controller.namespace().clone(); let deadline_after = Duration::from_millis(policy.deadline_ms); @@ -185,6 +194,10 @@ async fn drain_writers( match policy.on_deadline { OnDeadline::ForceRollback if !forced => { forced = true; + report.forced( + ForcedKind::Rollback, + sources.iter().filter(|s| s.manager.has_writer()).count(), + ); for source in &sources { let manager = source.manager.clone(); // The rollback takes the connection's lock, which a running program @@ -290,7 +303,7 @@ fn capture_boundary( Ok(FrozenBoundary { log_id, frame_no }) } -pub(super) fn now_ms() -> i64 { +pub(crate) fn now_ms() -> i64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map_or(0, |d| i64::try_from(d.as_millis()).unwrap_or(i64::MAX)) diff --git a/libsql-server/src/namespace/fence/hooks.rs b/libsql-server/src/namespace/fence/hooks.rs index 75631e122c..23f05d3336 100644 --- a/libsql-server/src/namespace/fence/hooks.rs +++ b/libsql-server/src/namespace/fence/hooks.rs @@ -18,6 +18,7 @@ use parking_lot::Mutex; #[cfg(test)] use tokio::sync::Notify; +#[cfg(test)] use super::outcome::FenceError; /// A named point on a fence transition or a gated path. @@ -35,10 +36,6 @@ pub enum HookPoint { AfterMetastoreCommit, /// The committed result is about to be published to the gate. BeforeGatePublish, - /// In `begin_write_txn`, after the gate check admitted the transaction. - InBeginWriteTxnAfterCheck, - /// The connection manager released the write slot. - AfterManagerRelease, /// A drain is about to read the frozen boundary. BeforeBoundaryCapture, /// The rows of a quarantined target were committed; the target is not published yet. @@ -70,7 +67,9 @@ pub enum HookAction { #[derive(Debug, Clone, PartialEq, Eq)] pub enum HookOutcome { Continue, + #[cfg(test)] Fail(FenceError), + #[cfg(test)] Indeterminate, } diff --git a/libsql-server/src/namespace/fence/import.rs b/libsql-server/src/namespace/fence/import.rs index 54fea99606..a8a20a3a31 100644 --- a/libsql-server/src/namespace/fence/import.rs +++ b/libsql-server/src/namespace/fence/import.rs @@ -28,6 +28,7 @@ use crate::namespace::configurator::{load_dump_sql, read_dump}; use crate::namespace::meta_store::{FenceCommit, FenceCommitKind, FenceContext, MetaStore}; use crate::namespace::replication_wal::ReplicationWalWrapper; +use super::audit::{CommandReport, DrainKind, ForcedKind}; use super::capability::MigrationCapability; use super::command::{DrainPolicy, FenceCommand, FenceRequest, OnDeadline}; use super::controller::{FenceController, LiveWriteDrain, Transition}; @@ -174,10 +175,14 @@ pub async fn seal_target_import( let resumed = commit.kind == FenceCommitKind::Resumed; let drain_key = (commit.receipt.operation_id, commit.receipt.command_id); - if !drain_import_writers(&controller, policy).await { + let started = Instant::now(); + if !drain_import_writers(&controller, policy, &mut transition.report).await { return Ok(commit); } + transition + .report + .drained(DrainKind::Import, started.elapsed()); ctx.now_ms = now_ms(); let mut completed = transition .complete_drain(meta, drain_key, DrainCompletion::TargetImport, ctx) @@ -193,7 +198,11 @@ pub async fn seal_target_import( /// still holding the slot (an idle session's open transaction; a running call ends its own) /// and waits again for the same deadline, but at least [`FORCED_ROLLBACK_GRACE`]. `false` when /// the drain could not be proven within the policy. -async fn drain_import_writers(controller: &FenceController, policy: DrainPolicy) -> bool { +async fn drain_import_writers( + controller: &FenceController, + policy: DrainPolicy, + report: &mut CommandReport, +) -> bool { let namespace = controller.namespace().clone(); let deadline_after = Duration::from_millis(policy.deadline_ms); let mut deadline = Instant::now() + deadline_after; @@ -212,7 +221,12 @@ async fn drain_import_writers(controller: &FenceController, policy: DrainPolicy) match policy.on_deadline { OnDeadline::ForceRollback if !forced => { forced = true; - for source in controller.live_write_drains() { + let sources = controller.live_write_drains(); + report.forced( + ForcedKind::Rollback, + sources.iter().filter(|s| s.manager.has_writer()).count(), + ); + for source in sources { let manager = source.manager.clone(); // The rollback takes the connection's lock, which a running import call // holds; the release it causes is what the drain keeps waiting for. diff --git a/libsql-server/src/namespace/fence/mod.rs b/libsql-server/src/namespace/fence/mod.rs index 8c89d79f5e..a4f5ba9130 100644 --- a/libsql-server/src/namespace/fence/mod.rs +++ b/libsql-server/src/namespace/fence/mod.rs @@ -9,16 +9,13 @@ //! ([`transition`]), the metastore tables, compare-and-swap and marker file that persist them //! ([`store`], driven by `MetaStore::apply_fence_command`), and the in-memory authority built //! on them: the per-namespace [`controller`] with its gate and read leases, the positive write -//! [`drain`], the source [`read`] fence and its -//! [`stream`] leases for dump and replication, quarantined migration [`target`]s with their -//! [`capability`]-scoped [`import`] sessions and seal drain, the [`registry`] that holds the controllers outside the namespace cache, the -//! [`replica`]-server view of a primary's fence, and the test [`hooks`] on their paths. - -// The persistence, controller and protocol layers that consume these types land in the -// following commits of this series; until then most of the module is unused by the rest of -// the crate. This attribute is removed once they are wired. -#![allow(dead_code)] +//! [`drain`], the source [`read`] fence and its [`stream`] leases for dump and replication, +//! quarantined migration [`target`]s with their [`capability`]-scoped [`import`] sessions and +//! seal drain, the [`registry`] that holds the controllers outside the namespace cache, the +//! [`replica`]-server view of a primary's fence, the [`audit`] log, and the test [`hooks`] on +//! their paths. +pub mod audit; pub mod capability; pub mod command; pub mod controller; diff --git a/libsql-server/src/namespace/fence/outcome.rs b/libsql-server/src/namespace/fence/outcome.rs index 723b7b5ab6..5b4a18a7fd 100644 --- a/libsql-server/src/namespace/fence/outcome.rs +++ b/libsql-server/src/namespace/fence/outcome.rs @@ -191,6 +191,9 @@ pub enum FenceDetail { ValidationReceiptRequired, RestoreNotAllowed, AdoptionNotAuthorised, + /// Adoption would have to re-establish a record for a namespace the metastore holds no + /// configuration for (section 12). + NamespaceConfigMissing, InvalidArgument, // INVALID_FENCE_TRANSITION RoleMismatch, @@ -217,6 +220,7 @@ impl FenceDetail { FenceDetail::ValidationReceiptRequired => "validation_receipt_required", FenceDetail::RestoreNotAllowed => "restore_not_allowed", FenceDetail::AdoptionNotAuthorised => "adoption_not_authorised", + FenceDetail::NamespaceConfigMissing => "namespace_config_missing", FenceDetail::InvalidArgument => "invalid_argument", FenceDetail::RoleMismatch => "role_mismatch", FenceDetail::OperationFinished => "operation_finished", diff --git a/libsql-server/src/namespace/fence/read.rs b/libsql-server/src/namespace/fence/read.rs index 95d7beb877..0603e962b8 100644 --- a/libsql-server/src/namespace/fence/read.rs +++ b/libsql-server/src/namespace/fence/read.rs @@ -15,6 +15,7 @@ use tokio::time::Instant; use crate::namespace::meta_store::{FenceCommit, FenceCommitKind, FenceContext, MetaStore}; +use super::audit::{CommandReport, DrainKind, ForcedKind}; use super::command::{DrainPolicy, FenceCommand, FenceRequest}; use super::controller::{FenceController, Transition}; use super::drain::{now_ms, FORCED_ROLLBACK_GRACE}; @@ -75,10 +76,14 @@ pub async fn set_source_read_fence( let drain_key = (commit.receipt.operation_id, commit.receipt.command_id); // Step 4. - if !drain_readers(&controller, policy).await { + let started = Instant::now(); + if !drain_readers(&controller, policy, &mut transition.report).await { return Ok(commit); } + transition + .report + .drained(DrainKind::Read, started.elapsed()); // Step 5. ctx.now_ms = now_ms(); let mut completed = transition @@ -95,7 +100,11 @@ pub async fn set_source_read_fence( /// through their own), and the drain keeps waiting for the actual releases for one more /// deadline, but at least [`FORCED_ROLLBACK_GRACE`]. `false` when leases are still held then: /// the answer is `DRAINING`, never a guess. -async fn drain_readers(controller: &FenceController, policy: DrainPolicy) -> bool { +async fn drain_readers( + controller: &FenceController, + policy: DrainPolicy, + report: &mut CommandReport, +) -> bool { let namespace = controller.namespace().clone(); let deadline_after = Duration::from_millis(policy.deadline_ms); let deadline = Instant::now() + deadline_after; @@ -104,7 +113,11 @@ async fn drain_readers(controller: &FenceController, policy: DrainPolicy) -> boo } let _ = controller.hook(HookPoint::BeforeReadLeaseCancel).await; let held = controller.read_lease_counts(); - let asked = controller.cancel_read_leases(); + let asked = controller.cancel_read_leases_by_kind(); + report.forced(ForcedKind::SqlCancel, asked.sql); + report.forced(ForcedKind::DumpCancel, asked.dump); + report.forced(ForcedKind::StreamTermination, asked.replication); + let asked = asked.total(); tracing::info!( %namespace, deadline_ms = policy.deadline_ms, diff --git a/libsql-server/src/namespace/fence/registry.rs b/libsql-server/src/namespace/fence/registry.rs index 7dd52425ba..2c2e508652 100644 --- a/libsql-server/src/namespace/fence/registry.rs +++ b/libsql-server/src/namespace/fence/registry.rs @@ -110,10 +110,14 @@ impl FenceRegistry { /// 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) { + let refused = match self.get(namespace) { Some(controller) => controller.gate().permits(OperationClass::Lifecycle), None => Ok(()), + }; + if let Err(e) = &refused { + super::audit::denied(e, super::audit::DenialSurface::Lifecycle); } + refused } /// How many namespaces have an active fence (`docs/NAMESPACE_FENCE.md` section 4.4, @@ -132,6 +136,19 @@ impl FenceRegistry { .count() } + /// The published state of every namespace the registry holds, with the creation time of + /// its record when it has one (for the fence gauges, section 15). + pub fn census(&self) -> Vec<(super::state::FenceState, Option)> { + let controllers: Vec<_> = self.controllers.lock().values().cloned().collect(); + controllers + .iter() + .map(|controller| { + let gate = controller.gate(); + (gate.state(), gate.fence.record().map(|r| r.created_at_ms)) + }) + .collect() + } + pub fn len(&self) -> usize { self.controllers.lock().len() } diff --git a/libsql-server/src/namespace/fence/stream.rs b/libsql-server/src/namespace/fence/stream.rs index 2711105a2a..dd763abacf 100644 --- a/libsql-server/src/namespace/fence/stream.rs +++ b/libsql-server/src/namespace/fence/stream.rs @@ -730,7 +730,7 @@ mod tests { fence_status as fn(FenceError) -> tonic::Status, ) .unwrap(); - assert_eq!(s.fence.cancel_read_leases(), 1); + assert_eq!(s.fence.cancel_read_leases_by_kind().total(), 1); until_released(&s.fence).await; let status = tokio::time::timeout(PROMPT, stream.next()) .await diff --git a/libsql-server/src/namespace/fence/target.rs b/libsql-server/src/namespace/fence/target.rs index 44006e40bd..05901d8e25 100644 --- a/libsql-server/src/namespace/fence/target.rs +++ b/libsql-server/src/namespace/fence/target.rs @@ -235,7 +235,7 @@ pub(crate) mod tests { } } - fn execute( + pub(crate) fn execute( store: &NamespaceStore, request: FenceRequest, ) -> tokio::task::JoinHandle> { @@ -303,7 +303,7 @@ pub(crate) mod tests { } /// A target with successful validation, published readable and write-fenced at revision 5. - async fn write_fenced_target() -> (TempDir, NamespaceStore, Arc) { + pub(crate) async fn write_fenced_target() -> (TempDir, NamespaceStore, Arc) { let (dir, store, fence) = validating_target().await; record_validation(&store, 20, 3, ValidationResult::Ok).await; let publish = target_command( diff --git a/libsql-server/src/namespace/fence/tests.rs b/libsql-server/src/namespace/fence/tests.rs index 6d297a0231..5f2c5fec85 100644 --- a/libsql-server/src/namespace/fence/tests.rs +++ b/libsql-server/src/namespace/fence/tests.rs @@ -659,6 +659,46 @@ fn restart_at(case: &Boundary) { server.crash(); } +/// A dirty restart rebuilds the replication log under a new id. The stored record keeps the +/// log the source was acquired on and the boundary frozen in it; the controller's +/// `current_log_id`, which `InspectFence` and every admin response report as +/// `incarnation.current_log_id` (section 4.5), names the rebuilt log, so a caller can see the +/// rebuild without a separate probe. +#[test] +fn current_log_id_is_the_rebuilt_log_after_restart() { + let dir = tempdir().unwrap(); + + let server = Server::boot(dir.path()); + server.create_source(); + let (acquired_on, _) = server.run(server.log()); + let commit = server.fence_source(); + assert_eq!( + commit.record.as_ref().unwrap().state, + FenceState::SourceWriteFenced + ); + server.run(async { + assert_eq!(server.fence().await.current_log_id(), Some(acquired_on)); + }); + server.crash(); + + let server = Server::boot(dir.path()); + server.run(async { + let fence = server.fence().await; + let (live, _) = server.log().await; + assert_ne!(live, acquired_on, "a dirty restart rebuilds the log"); + assert_eq!(fence.current_log_id(), Some(live)); + + let inspected = server.inspect().await; + let record = inspected.fence.record().unwrap(); + assert_eq!(record.state, FenceState::SourceWriteFenced); + assert_eq!(record.identity.log_id, Some(acquired_on)); + assert_eq!(boundary(&commit).log_id, acquired_on); + // The data is the same; only the log was rebuilt. + assert_eq!(server.count().await, ROWS); + }); + server.crash(); +} + /// After a restart in `SOURCE_DRAINING` with a writer that was active at the crash, nothing /// advances by itself: the namespace stays closed, other commands cannot move it on, and only /// the replay of the same acquisition completes the drain, at once. @@ -2004,3 +2044,612 @@ mod legacy_mirror { server.crash(); } } + +// --------------------------------------------------------------------------------------------- +// Incident adoption (section 12; section 17 row 22) + +mod adoption { + use super::*; + use crate::config::FenceAdoptionKey; + use crate::namespace::fence::command::AdoptArgs; + use crate::namespace::fence::target::tests::{ + create, create_request, enable_request, execute as execute_target, target_command, + write_fenced_target, + }; + use crate::namespace::meta_store::metastore_connection_maker; + + fn adopt_args(approvers: &[&str]) -> AdoptArgs { + AdoptArgs { + current_operation_id: OP, + approvers: approvers.iter().map(|a| a.to_string()).collect(), + incident_ref: "INC-1".into(), + reason: "the operation's control record was lost".into(), + } + } + + fn adopt( + ns: &'static str, + command_id: u128, + expected_state: FenceState, + expected_revision: u64, + args: AdoptArgs, + ) -> FenceRequest { + FenceRequest { + namespace: ns.into(), + operation_id: OTHER_OP, + command_id: Uuid::from_u128(command_id), + expected_state, + expected_revision, + command: FenceCommand::AdoptFence(args), + } + } + + async fn run_adopt( + store: &NamespaceStore, + request: FenceRequest, + authorised: bool, + ) -> crate::Result { + let store = store.clone(); + tokio::spawn(async move { + store + .execute_fence_command_authorised(request, server_identity(), authorised) + .await + }) + .await + .unwrap() + } + + fn release_by(operation_id: Uuid, command_id: u128, revision: u64) -> FenceRequest { + FenceRequest { + operation_id, + ..release(command_id, revision) + } + } + + /// The adoption key is kept as a digest and compared whole; an empty key authorises + /// nothing; `Debug` does not print it; a metastore without a key refuses every presented + /// value. + #[test] + fn adoption_key_matching() { + let key = FenceAdoptionKey::new("s3cret").unwrap(); + assert!(key.matches(b"s3cret")); + for wrong in [&b""[..], b"s3cre", b"s3cret ", b"S3CRET"] { + assert!(!key.matches(wrong)); + } + assert!(FenceAdoptionKey::new("").is_none()); + assert!(!format!("{key:?}").contains("s3cret")); + + let dir = tempdir().unwrap(); + let server = Server::boot(dir.path()); + let meta = server.store.meta_store(); + assert!(!meta.fence_adoption_authorised(None)); + assert!(!meta.fence_adoption_authorised(Some(b"s3cret"))); + server.crash(); + } + + /// Adoption without the key, with one approver, with the same approver twice (also after + /// trimming) or without an incident reference or reason is refused with + /// `adoption_not_authorised`, and changes nothing. + #[test] + fn adopt_requires_key_and_two_approvers() { + let dir = tempdir().unwrap(); + let server = Server::boot(dir.path()); + server.create_source(); + let fenced = server.fence_source(); + let revision = fenced.record.as_ref().unwrap().revision; + server.run(async { + let state = FenceState::SourceWriteFenced; + let refusals = [ + (adopt_args(&["alice", "bob"]), false), + (adopt_args(&["alice"]), true), + (adopt_args(&["alice", "alice"]), true), + (adopt_args(&["alice", " alice "]), true), + (adopt_args(&["alice", ""]), true), + (adopt_args(&["alice", "bob", "carol"]), true), + ( + AdoptArgs { + incident_ref: " ".into(), + ..adopt_args(&["alice", "bob"]) + }, + true, + ), + ( + AdoptArgs { + reason: String::new(), + ..adopt_args(&["alice", "bob"]) + }, + true, + ), + ]; + for (i, (args, authorised)) in refusals.into_iter().enumerate() { + let request = adopt("ns", 10 + i as u128, state, revision, args.clone()); + let result = run_adopt(&server.store, request, authorised).await; + let e = fence_error(&result); + assert_eq!( + e.outcome(), + FenceOutcome::FencePreconditionFailed, + "{args:?}" + ); + assert_eq!( + e.detail(), + Some(FenceDetail::AdoptionNotAuthorised), + "{args:?}" + ); + } + // Nothing changed: owner, revision, receipts. + let inspection = server.inspect().await; + let record = inspection.fence.record().unwrap(); + assert_eq!((record.operation_id, record.revision), (OP, revision)); + assert!(record.adoptions.is_empty()); + assert!(inspection + .receipts + .iter() + .all(|r| r.operation_id == OP.to_string())); + }); + server.crash(); + } + + /// Adopting a write-fenced source changes the owner and the revision and nothing else: the + /// state and every admission stay as they were (the write generation moves, as on every + /// change of owner, section 7.2), SQL writes are still refused, the old owner is refused + /// and the new owner finishes the operation. A replay returns the stored receipt. + #[test] + fn adopt_keeps_gates_closed() { + let dir = tempdir().unwrap(); + let server = Server::boot(dir.path()); + server.create_source(); + let fenced = server.fence_source(); + let before = fenced.record.clone().unwrap(); + server.run(async { + let fence = server.fence().await; + let generation = fence.write_generation(); + let request = adopt( + "ns", + 10, + FenceState::SourceWriteFenced, + before.revision, + adopt_args(&["alice", "bob"]), + ); + let commit = run_adopt(&server.store, request.clone(), true) + .await + .unwrap(); + assert_eq!(commit.kind, FenceCommitKind::Committed); + assert_eq!(commit.receipt.outcome, FenceOutcome::Applied); + let after = commit.record.clone().unwrap(); + assert_eq!(after.operation_id, OTHER_OP); + assert_eq!(after.revision, before.revision + 1); + assert_eq!(after.state, before.state); + assert_eq!(after.frozen_boundary, before.frozen_boundary); + assert_eq!(after.identity, before.identity); + assert_eq!(after.legacy_blocks, before.legacy_blocks); + assert_eq!(after.adoptions.len(), 1); + let adoption = commit.receipt.adoption.as_ref().unwrap(); + assert_eq!(adoption.previous_operation_id, OP); + assert_eq!(adoption.approvers, vec!["alice", "bob"]); + assert_eq!(&after.adoptions[0], adoption); + + // The gate: same state and admissions, the new owner and revision. + let gate = fence.gate(); + assert_eq!(gate.state(), FenceState::SourceWriteFenced); + assert_eq!(gate.operation_id(), Some(OTHER_OP)); + assert_eq!(gate.revision(), before.revision + 1); + assert!(!gate.write().is_open()); + assert!(gate.read().is_open()); + assert_eq!(fence.write_generation(), generation + 1); + assert!(!server.writes_admitted().await); + assert_eq!(server.count().await, ROWS); + // The marker follows the record. + let marker = fence_store::read_marker(&dir.path().join("dbs"), &"ns".into()) + .unwrap() + .unwrap(); + assert_eq!(marker.unwrap().record, after); + + // Replay, with or without the key: the stored receipt, nothing moves. + for authorised in [true, false] { + let replay = run_adopt(&server.store, request.clone(), authorised) + .await + .unwrap(); + assert_eq!(replay.kind, FenceCommitKind::Replayed); + assert_eq!(replay.receipt, commit.receipt); + } + assert_eq!(fence.write_generation(), generation + 1); + + // The old owner is refused. + let e = server + .execute(release(20, after.revision)) + .await + .unwrap() + .unwrap_err(); + assert_eq!( + fence_error(&Err(e)).outcome(), + FenceOutcome::FenceOwnedByAnotherOperation + ); + assert!(!server.writes_admitted().await); + // The new owner finishes the operation. + let released = server + .execute(release_by(OTHER_OP, 21, after.revision)) + .await + .unwrap() + .unwrap(); + assert_eq!(released.receipt.outcome, FenceOutcome::Applied); + assert_eq!(fence.gate().state(), FenceState::Released); + assert!(server.writes_admitted().await); + }); + server.crash(); + } + + /// Adopting a quarantined target moves its import capability to the new owner: the old + /// owner can no longer open an import session, the new owner can and imports, and normal + /// SQL is still refused with `MIGRATION_TARGET_QUARANTINED`. + #[tokio::test(flavor = "multi_thread")] + async fn adopt_quarantined_target() { + let dir = tempdir().unwrap(); + let store = open_store(dir.path()).await; + create(&store, create_request("tgt", 1)) + .await + .unwrap() + .unwrap(); + let fence = store.fence_controller(&"tgt".into()); + let request = adopt( + "tgt", + 10, + FenceState::TargetQuarantined, + 1, + adopt_args(&["alice", "bob"]), + ); + let commit = run_adopt(&store, request, true).await.unwrap(); + let after = commit.record.unwrap(); + assert_eq!( + (after.state, after.revision, after.operation_id), + (FenceState::TargetQuarantined, 2, OTHER_OP) + ); + assert_eq!(fence.gate().state(), FenceState::TargetQuarantined); + for class in [OperationClass::NormalRead, OperationClass::NormalWrite] { + assert_eq!( + fence.permits(class).unwrap_err().outcome(), + FenceOutcome::MigrationTargetQuarantined + ); + } + + let e = store + .open_import_session("tgt".into(), OP, 2) + .await + .err() + .expect("the old owner has no import capability"); + assert_eq!( + fence_error(&Err(e)).outcome(), + FenceOutcome::FenceOwnedByAnotherOperation + ); + let mut session = store + .open_import_session("tgt".into(), OTHER_OP, 2) + .await + .unwrap(); + session + .with_raw(|conn| conn.execute_batch("create table t (x); insert into t values (1)")) + .await + .unwrap() + .unwrap(); + } + + /// Finished operations cannot be adopted: `TARGET_WRITABLE`, `TARGET_ABORTED` and a + /// `RELEASED` source are refused with `operation_finished`, and nothing changes. + #[tokio::test(flavor = "multi_thread")] + async fn adopt_cannot_touch_writable() { + // TARGET_WRITABLE. + let (_dir, store, fence) = write_fenced_target().await; + execute_target(&store, enable_request(30)) + .await + .unwrap() + .unwrap(); + assert_eq!(fence.gate().state(), FenceState::TargetWritable); + let revision = fence.gate().revision(); + let request = adopt( + "tgt", + 40, + FenceState::TargetWritable, + revision, + adopt_args(&["alice", "bob"]), + ); + let e = fence_error(&run_adopt(&store, request, true).await).clone(); + assert_eq!(e.outcome(), FenceOutcome::InvalidFenceTransition); + assert_eq!(e.detail(), Some(FenceDetail::OperationFinished)); + assert_eq!(fence.gate().operation_id(), Some(OP)); + assert_eq!(fence.gate().revision(), revision); + + // TARGET_ABORTED. + let dir = tempdir().unwrap(); + let store = open_store(dir.path()).await; + create(&store, create_request("tgt", 1)) + .await + .unwrap() + .unwrap(); + let abort = target_command( + 2, + FenceState::TargetQuarantined, + 1, + FenceCommand::AbortQuarantinedTarget, + ); + execute_target(&store, abort).await.unwrap().unwrap(); + let fence = store.fence_controller(&"tgt".into()); + assert_eq!(fence.gate().state(), FenceState::TargetAborted); + let request = adopt( + "tgt", + 3, + FenceState::TargetAborted, + 2, + adopt_args(&["alice", "bob"]), + ); + let e = fence_error(&run_adopt(&store, request, true).await).clone(); + assert_eq!(e.outcome(), FenceOutcome::InvalidFenceTransition); + assert_eq!(e.detail(), Some(FenceDetail::OperationFinished)); + assert_eq!(fence.gate().revision(), 2); + } + + /// A released source cannot be adopted either. + #[test] + fn adopt_cannot_touch_released() { + let dir = tempdir().unwrap(); + let server = Server::boot(dir.path()); + server.create_source(); + let revision = server.fence_source().record.unwrap().revision; + server.run(async { + server.execute(release(2, revision)).await.unwrap().unwrap(); + let request = adopt( + "ns", + 3, + FenceState::Released, + revision + 1, + adopt_args(&["alice", "bob"]), + ); + let e = fence_error(&run_adopt(&server.store, request, true).await).clone(); + assert_eq!(e.detail(), Some(FenceDetail::OperationFinished)); + }); + server.crash(); + } + + async fn raw_meta(dir: &Path) -> crate::namespace::meta_store::MetaStoreConnection { + let (maker, _) = metastore_connection_maker(None, dir).await.unwrap(); + maker().unwrap() + } + + fn unavailable_detail(result: crate::Result<()>) -> FenceDetail { + match result { + Err(Error::NamespaceFence(e)) => { + assert_eq!(e.outcome(), FenceOutcome::FenceStateUnavailable, "{e}"); + e.detail().unwrap() + } + other => panic!("expected FENCE_STATE_UNAVAILABLE, got {other:?}"), + } + } + + /// After a metastore rollback (a restore from a backup older than the fence), the + /// namespace is `UNKNOWN_UNAVAILABLE` with `metastore_behind_marker`. Adoption re-establishes + /// the record the marker holds, under the adopting operation, at the marker's revision + 1, + /// both when the fence row is gone (the backup predates the fence) and when it is at an + /// older revision (the backup predates the last transition). The namespace is then served + /// again behind the same gate, with its own `block_*` values in memory, and the new owner + /// finishes the operation. + #[test] + fn adopt_recovers_metastore_rollback() { + for row_behind in [false, true] { + let dir = tempdir().unwrap(); + let server = Server::boot(dir.path()); + server.create_source(); + let fenced = server.fence_source().record.unwrap(); + let saved_row: (i64, i64, Vec) = server.run(async { + raw_meta(dir.path()) + .await + .query_row( + "SELECT format_version, revision, record FROM namespace_fences", + (), + |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)), + ) + .unwrap() + }); + // With the row kept behind, one more transition moves the marker on. + let marker_record = if row_behind { + let read = FenceRequest { + namespace: "ns".into(), + operation_id: OP, + command_id: Uuid::from_u128(2), + expected_state: FenceState::SourceWriteFenced, + expected_revision: fenced.revision, + command: FenceCommand::SetSourceReadFence { + drain_policy: Some(LONG), + }, + }; + let commit = server.run(async { server.execute(read).await.unwrap().unwrap() }); + assert_eq!(commit.receipt.outcome, FenceOutcome::Applied); + commit.record.unwrap() + } else { + fenced.clone() + }; + server.crash(); + + // What restoring the metastore from an older backup leaves: no receipts, and the + // fence row gone or at the older revision. The config row keeps the legacy + // mirror the acquisition wrote into it (block_writes on). + server_run_blocking(|| async { + let conn = raw_meta(dir.path()).await; + conn.execute("DELETE FROM namespace_fence_receipts", ()) + .unwrap(); + if row_behind { + conn.execute( + "UPDATE namespace_fences SET format_version = ?1, revision = ?2, \ + record = ?3", + (saved_row.0, saved_row.1, &saved_row.2), + ) + .unwrap(); + } else { + conn.execute("DELETE FROM namespace_fences", ()).unwrap(); + } + }); + + let server = Server::boot(dir.path()); + server.run(async { + let store = &server.store; + assert_eq!( + unavailable_detail(store.with("ns".into(), |_| ()).await), + FenceDetail::MetastoreBehindMarker + ); + assert!(store.meta_store().lookup(&"ns".into()).await.is_err()); + + // The adoption has to name the marker's state, revision and owner. + let wrong_revision = adopt( + "ns", + 10, + marker_record.state, + saved_row.1 as u64, + adopt_args(&["alice", "bob"]), + ); + if row_behind { + let e = fence_error(&run_adopt(store, wrong_revision, true).await).clone(); + assert_eq!(e.outcome(), FenceOutcome::FenceRevisionMismatch); + } + let request = adopt( + "ns", + 11, + marker_record.state, + marker_record.revision, + adopt_args(&["alice", "bob"]), + ); + let e = fence_error(&run_adopt(store, request.clone(), false).await).clone(); + assert_eq!(e.detail(), Some(FenceDetail::AdoptionNotAuthorised)); + + let commit = run_adopt(store, request, true).await.unwrap(); + assert_eq!(commit.receipt.outcome, FenceOutcome::Applied); + let adopted = commit.record.unwrap(); + assert_eq!(adopted.operation_id, OTHER_OP); + assert_eq!(adopted.revision, marker_record.revision + 1); + assert_eq!(adopted.state, marker_record.state); + assert_eq!(adopted.frozen_boundary, marker_record.frozen_boundary); + assert_eq!(adopted.adoptions.len(), 1); + + // Established again: durable, served, and behind the same gate. + let inspection = server.inspect().await; + assert_eq!(inspection.fence.record(), Some(&adopted)); + let fence = server.fence().await; + assert_eq!(fence.gate().state(), marker_record.state); + assert_eq!(fence.gate().operation_id(), Some(OTHER_OP)); + assert!(!server.writes_admitted().await); + let config = store.meta_store().lookup(&"ns".into()).await.unwrap(); + let config = config.unwrap().get(); + assert!(!config.block_writes, "the namespace's own values are back"); + assert!(!config.block_reads); + let marker = fence_store::read_marker(&dir.path().join("dbs"), &"ns".into()) + .unwrap() + .unwrap(); + assert_eq!(marker.unwrap().record, adopted); + + // The new owner finishes the operation. + let mut revision = adopted.revision; + if row_behind { + let clear = FenceRequest { + namespace: "ns".into(), + operation_id: OTHER_OP, + command_id: Uuid::from_u128(20), + expected_state: FenceState::SourceReadFenced, + expected_revision: revision, + command: FenceCommand::ClearSourceReadFence, + }; + revision = server + .execute(clear) + .await + .unwrap() + .unwrap() + .record + .unwrap() + .revision; + } + server + .execute(release_by(OTHER_OP, 21, revision)) + .await + .unwrap() + .unwrap(); + assert!(server.writes_admitted().await); + assert_eq!(server.count().await, ROWS); + }); + server.crash(); + } + } + + /// Run `f` on a throwaway runtime (between two server lifetimes). + fn server_run_blocking(f: F) + where + F: FnOnce() -> Fut, + Fut: Future, + { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap() + .block_on(f()); + } + + /// A metastore restored from a backup that predates the namespace itself holds neither a + /// fence row nor a config row for it. Adoption is refused with + /// `FENCE_PRECONDITION_FAILED` / `namespace_config_missing` (after the authorisation + /// checks): the marker holds the fence record, not the namespace's configuration, and + /// adoption does not invent one. Nothing is written and the name stays unavailable. + #[test] + fn adopt_recovered_name_without_config_row() { + let dir = tempdir().unwrap(); + let server = Server::boot(dir.path()); + server.create_source(); + let fenced = server.fence_source().record.unwrap(); + server.crash(); + let marker_before = read_marker_bytes(&dir.path().join("dbs")); + server_run_blocking(|| async { + let conn = raw_meta(dir.path()).await; + for sql in [ + "DELETE FROM namespace_fence_receipts", + "DELETE FROM namespace_fences", + "DELETE FROM namespace_configs", + ] { + conn.execute(sql, ()).unwrap(); + } + }); + + let server = Server::boot(dir.path()); + server.run(async { + let store = &server.store; + assert_eq!( + unavailable_detail(store.with("ns".into(), |_| ()).await), + FenceDetail::MetastoreBehindMarker + ); + let request = adopt( + "ns", + 10, + fenced.state, + fenced.revision, + adopt_args(&["alice", "bob"]), + ); + let e = fence_error(&run_adopt(store, request.clone(), false).await).clone(); + assert_eq!(e.detail(), Some(FenceDetail::AdoptionNotAuthorised)); + let e = fence_error(&run_adopt(store, request, true).await).clone(); + assert_eq!(e.outcome(), FenceOutcome::FencePreconditionFailed); + assert_eq!(e.detail(), Some(FenceDetail::NamespaceConfigMissing)); + + // Nothing was written, and the name is still unavailable. + let conn = raw_meta(dir.path()).await; + for table in [ + "namespace_fences", + "namespace_fence_receipts", + "namespace_configs", + ] { + let n: i64 = conn + .query_row(&format!("SELECT count(*) FROM {table}"), (), |r| r.get(0)) + .unwrap(); + assert_eq!(n, 0, "{table}"); + } + assert_eq!(read_marker_bytes(&dir.path().join("dbs")), marker_before); + assert_eq!( + unavailable_detail(store.with("ns".into(), |_| ()).await), + FenceDetail::MetastoreBehindMarker + ); + assert!(store.meta_store().handle("ns".into()).await.is_err()); + assert!(store.fence_controller(&"ns".into()).gate().is_unavailable()); + }); + server.crash(); + } +} diff --git a/libsql-server/src/namespace/meta_store.rs b/libsql-server/src/namespace/meta_store.rs index 8665745526..4a63018ac3 100644 --- a/libsql-server/src/namespace/meta_store.rs +++ b/libsql-server/src/namespace/meta_store.rs @@ -23,7 +23,7 @@ use tokio::sync::{ }; use uuid::Uuid; -use crate::config::BottomlessConfig; +use crate::config::{BottomlessConfig, FenceAdoptionKey}; use crate::connection::config::DatabaseConfig; use crate::database::DatabaseKind; use crate::schema::{MigrationDetails, MigrationSummary}; @@ -97,6 +97,35 @@ struct MetaStoreInner { /// record for. They are refused by lookups and by every config or lifecycle change, and /// never default-created. A fence command that commits for the name takes it out. recovered: Mutex>, + /// Where this metastore's contents came from at startup, recorded once the server has + /// opened it (section 13.3). + restore_provenance: std::sync::OnceLock, + /// The secret that authorises `AdoptFence` (section 12); `None` disables adoption. + fence_adoption_key: Option, +} + +/// Whether the metastore was restored from its bottomless backup when the server started +/// (`docs/NAMESPACE_FENCE.md` sections 4.3, 4.4 and 13.3). A restored metastore can hold fence +/// records older than the namespace markers; the marker comparison makes those namespaces +/// `UNKNOWN_UNAVAILABLE`, and this is what the admin API, the startup log and the +/// `libsql_server_metastore_restored_from_backup` gauge report about the restore itself. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct MetastoreProvenance { + /// The metastore database was restored from its backup at startup. + pub restored_from_backup: bool, + /// The backup generation it was restored from, when known. + pub restored_generation: Option, +} + +impl MetastoreProvenance { + /// The provenance of a bottomless restore that reported `did_recover`, where `generation` + /// is the replicator's generation right after the restore: the generation restored from. + pub fn from_restore(did_recover: bool, generation: Option) -> Self { + Self { + restored_from_backup: did_recover, + restored_generation: generation.filter(|_| did_recover), + } + } } /// How this metastore treats namespace fences (`docs/NAMESPACE_FENCE.md` section 13.1). @@ -140,15 +169,33 @@ fn setup_connection(conn: &rusqlite::Connection) -> Result<()> { Ok(()) } +/// [`metastore_connection_maker_with_provenance`] without the provenance, for tests. +#[cfg(test)] pub async fn metastore_connection_maker( config: Option, base_path: &Path, ) -> crate::Result<( impl Fn() -> crate::Result, MetaStoreWalManager, +)> { + let (maker, wal_manager, _) = + metastore_connection_maker_with_provenance(config, base_path).await?; + Ok((maker, wal_manager)) +} + +/// [`metastore_connection_maker`], also returning whether the bottomless restore recovered the +/// metastore from its backup, for [`MetaStore::record_restore_provenance`]. +pub async fn metastore_connection_maker_with_provenance( + config: Option, + base_path: &Path, +) -> crate::Result<( + impl Fn() -> crate::Result, + MetaStoreWalManager, + MetastoreProvenance, )> { let db_path = base_path.join("metastore"); tokio::fs::create_dir_all(&db_path).await?; + let mut provenance = MetastoreProvenance::default(); let replicator = match config { Some(config) => { let options = bottomless::replicator::Options { @@ -175,7 +222,11 @@ pub async fn metastore_connection_maker( options, ) .await?; - let (action, _did_recover) = replicator.restore(None, None).await?; + let (action, did_recover) = replicator.restore(None, None).await?; + // A restore that recovered the database leaves the replicator on the generation it + // restored from; a new generation, if any, is only started below. + provenance = + MetastoreProvenance::from_restore(did_recover, replicator.generation().ok()); // TODO: this logic should probably be moved to bottomless. match action { bottomless::replicator::RestoreAction::SnapshotMainDbFile => { @@ -213,7 +264,7 @@ pub async fn metastore_connection_maker( } }; - Ok((maker, wal_manager)) + Ok((maker, wal_manager, provenance)) } impl MetaStoreInner { @@ -254,6 +305,8 @@ impl MetaStoreInner { dbs_path, fence, recovered: Default::default(), + restore_provenance: Default::default(), + fence_adoption_key: config.namespace_fence_adoption_key.clone(), }; if config.allow_recover_from_fs { @@ -600,6 +653,17 @@ fn unavailable_error(stored: &StoredFence) -> FenceError { }) } +/// A config write or a delete under the stored fence: refused (and counted as a lifecycle +/// denial) unless its state permits lifecycle operations. +fn lifecycle_permitted(stored: &StoredFence) -> std::result::Result<(), FenceError> { + stored.permits(OperationClass::Lifecycle).inspect_err(|e| { + crate::namespace::fence::audit::denied( + e, + crate::namespace::fence::audit::DenialSurface::Lifecycle, + ) + }) +} + /// Why a name that has no config must not be created: its directory holds a marker, so it is a /// target being created or a namespace the metastore lost (section 13.3). fn marker_denial(dbs_path: &Path, namespace: &NamespaceName) -> Result> { @@ -692,7 +756,7 @@ fn try_process( if inner.fence.tables { let (stored, _) = fence_store::read_fence(&tx, &inner.dbs_path, namespace).map_err(fence_store_error)?; - stored.permits(OperationClass::Lifecycle)?; + lifecycle_permitted(&stored)?; } if let Some(schema) = config.shared_schema_name.as_ref() { if inner.db_kind.is_primary() { @@ -923,6 +987,19 @@ fn apply_fence_command( created_config = Some(Arc::new(logical)); } else { let Some(config) = &config else { + if let FenceCommand::AdoptFence(_) = &request.command { + // Section 12: adoption changes the owner and nothing else. The marker holds + // the fence record, not the namespace's configuration (its JWT key, size + // limit, durability, backup id), so re-establishing the record would mean + // inventing a configuration. The namespace stays unavailable. + return Err(FenceError::new( + FenceOutcome::FencePreconditionFailed, + "the metastore holds no configuration for this namespace; adoption \ + re-establishes a fence, not a namespace configuration", + ) + .with_detail(FenceDetail::NamespaceConfigMissing) + .into()); + } return Err(FenceError::new( FenceOutcome::FenceStateUnavailable, "the fenced namespace has no config row", @@ -948,6 +1025,13 @@ fn apply_fence_command( // The command established the fence from the durable state; whatever startup could not // recover about this name is settled. inner.recovered.lock().remove(ns); + if let (StoredFence::Unavailable { .. }, Some(next)) = (&stored, &record) { + // An adoption re-established the record the marker held (section 12). While the name + // was unavailable its in-memory config kept whatever the stale config row held; from + // now on it carries the namespace's own values, as `restore_fences` gives every + // established record at startup. + restore_own_blocks(inner, ns, next); + } let current = record.or_else(|| stored.record().cloned()); after_fence_commit( @@ -965,6 +1049,18 @@ fn apply_fence_command( }) } +/// Put the namespace's own `block_*` values ([`fence_store::own_config`]) into its in-memory +/// config, if it has one. The caller holds the connection lock, which is taken before the +/// config map's everywhere. +fn restore_own_blocks(inner: &MetaStoreInner, ns: &NamespaceName, record: &NamespaceFenceRecord) { + let configs = inner.configs.blocking_lock(); + if let Some(sender) = configs.get(ns) { + let config = sender.borrow().config.clone(); + let config = fence_store::own_config(&config, record); + sender.send_modify(|c| c.config = Arc::new(config)); + } +} + /// Whether the marker has to be written after a commit: the record changed, or it had fallen /// behind. fn record_changed( @@ -1143,8 +1239,13 @@ impl MetaStore { ); } - let (maker, wal) = - metastore_connection_maker(config.bottomless.clone(), base_path).await?; + // The rebuilt metastore is restored from the backup again, so what the + // server reports is this restore, not the one of the broken metastore. + let (maker, wal, provenance) = metastore_connection_maker_with_provenance( + config.bottomless.clone(), + base_path, + ) + .await?; let conn = maker()?; @@ -1156,6 +1257,7 @@ impl MetaStore { }) .await .unwrap()?; + let _ = inner.restore_provenance.set(provenance); tracing::info!("metastore destroy on error successful"); @@ -1256,7 +1358,7 @@ impl MetaStore { if self.inner.fence.tables { let (stored, _) = fence_store::read_fence(&tx, &self.inner.dbs_path, &namespace) .map_err(fence_store_error)?; - stored.permits(OperationClass::Lifecycle)?; + lifecycle_permitted(&stored)?; if !matches!(stored, StoredFence::None { .. }) { // The marker goes before the commit: a crash in between leaves a record // without a marker, which is repaired on load, rather than a marker @@ -1350,6 +1452,51 @@ impl MetaStore { self.inner.fence.enabled } + /// Whether `presented`, the value of a request's `x-libsql-fence-adoption-key` header, + /// authorises `AdoptFence` (section 12): an adoption key is configured and `presented` is + /// that key. Compared in constant time. + pub fn fence_adoption_authorised(&self, presented: Option<&[u8]>) -> bool { + match (&self.inner.fence_adoption_key, presented) { + (Some(key), Some(presented)) => key.matches(presented), + _ => false, + } + } + + /// Records where this metastore's contents came from at startup, as reported by + /// [`metastore_connection_maker_with_provenance`], then sets the + /// `libsql_server_metastore_restored_from_backup` gauge and, after a restore from backup, + /// logs a startup warning. What was recorded first wins: when `destroy_on_error` rebuilt + /// the metastore while opening it, the restore of the rebuilt metastore is already recorded + /// and is what is reported. The server calls this once, right after opening the metastore. + pub fn record_restore_provenance(&self, provenance: MetastoreProvenance) { + let provenance = *self.inner.restore_provenance.get_or_init(|| provenance); + crate::metrics::METASTORE_RESTORED_FROM_BACKUP.set(if provenance.restored_from_backup { + 1.0 + } else { + 0.0 + }); + if provenance.restored_from_backup { + tracing::warn!( + restored_generation = ?provenance.restored_generation, + fence_tables = self.inner.fence.tables, + unavailable_namespaces = self.inner.recovered.lock().len(), + "the metastore was restored from its backup at startup; fence records the backup \ + does not hold are detected through the namespace markers and reported as \ + UNKNOWN_UNAVAILABLE" + ); + } + } + + /// Where this metastore's contents came from at startup; the default (not restored) until + /// [`record_restore_provenance`](Self::record_restore_provenance) is called. + pub fn restore_provenance(&self) -> MetastoreProvenance { + self.inner + .restore_provenance + .get() + .copied() + .unwrap_or_default() + } + /// The drain policy of an `AcquireSourceWriteFence` that names none: the configured /// deadline, then `DRAINING`. pub fn fence_default_write_drain(&self) -> DrainPolicy { @@ -2958,4 +3105,127 @@ mod fence_tests { assert_eq!(e.detail(), Some(FenceDetail::MetastoreBehindMarker)); } } + + /// Metastore restore provenance (`docs/NAMESPACE_FENCE.md` sections 4.3, 4.4 and 13.3). + mod provenance { + use super::*; + + #[test] + fn generation_only_after_a_recovery() { + let generation = Uuid::from_u128(0x77); + assert_eq!( + MetastoreProvenance::from_restore(true, Some(generation)), + MetastoreProvenance { + restored_from_backup: true, + restored_generation: Some(generation), + } + ); + // A restore that found the local database up to date (or nothing to restore) did + // not recover anything, whatever generation the replicator is on. + assert_eq!( + MetastoreProvenance::from_restore(false, Some(generation)), + MetastoreProvenance::default() + ); + assert_eq!( + MetastoreProvenance::from_restore(true, None), + MetastoreProvenance { + restored_from_backup: true, + restored_generation: None, + } + ); + } + + #[tokio::test] + async fn not_restored_until_recorded_and_first_record_wins() { + let dir = tempdir().unwrap(); + let store = open(dir.path(), true).await; + assert_eq!(store.restore_provenance(), MetastoreProvenance::default()); + + let restored = MetastoreProvenance { + restored_from_backup: true, + restored_generation: Some(Uuid::from_u128(0x77)), + }; + store.record_restore_provenance(restored); + assert_eq!(store.restore_provenance(), restored); + // Every clone of the store reports the same provenance. + assert_eq!(store.clone().restore_provenance(), restored); + + store.record_restore_provenance(MetastoreProvenance::default()); + assert_eq!(store.restore_provenance(), restored); + } + + /// An S3 endpoint backed by a temporary directory, on a port of its own. + async fn mock_s3() -> (tempfile::TempDir, String) { + use s3s::auth::SimpleAuth; + use s3s::service::S3ServiceBuilder; + + let root = tempdir().unwrap(); + let mut s3 = S3ServiceBuilder::new(s3s_fs::FileSystem::new(root.path()).unwrap()); + s3.set_auth(SimpleAuth::from_single("key", "secret")); + let service = s3.build().into_shared().into_make_service(); + let server = hyper::Server::bind(&([127, 0, 0, 1], 0).into()).serve(service); + let endpoint = format!("http://{}", server.local_addr()); + tokio::spawn(server); + (root, endpoint) + } + + fn bottomless(endpoint: String) -> BottomlessConfig { + BottomlessConfig { + access_key_id: "key".into(), + secret_access_key: "secret".into(), + session_token: None, + region: "us-east-1".into(), + backup_id: "metastore-provenance".into(), + bucket_name: "provenance".into(), + backup_interval: Duration::from_millis(100), + bucket_endpoint: endpoint, + } + } + + async fn open_bottomless( + dir: &Path, + config: &BottomlessConfig, + ) -> (MetaStore, MetastoreProvenance) { + let (maker, manager, provenance) = + metastore_connection_maker_with_provenance(Some(config.clone()), dir) + .await + .unwrap(); + let store = MetaStore::new( + MetaStoreConfig::default(), + dir, + maker().unwrap(), + manager, + DatabaseKind::Primary, + ) + .await + .unwrap(); + store.record_restore_provenance(provenance); + (store, provenance) + } + + /// A metastore opened on an empty directory from a backup that holds one reports the + /// restore and the generation it came from; one opened with nothing to restore does + /// not. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn bottomless_restore_is_reported() { + let (_s3, endpoint) = mock_s3().await; + let config = bottomless(endpoint); + + let first = tempdir().unwrap(); + let (store, provenance) = open_bottomless(first.path(), &config).await; + assert_eq!(provenance, MetastoreProvenance::default()); + assert_eq!(store.restore_provenance(), MetastoreProvenance::default()); + let _handle = create_namespace(&store, "db").await; + // Uploads everything the backup does not hold yet. + store.shutdown().await.unwrap(); + + let second = tempdir().unwrap(); + let (store, provenance) = open_bottomless(second.path(), &config).await; + assert!(provenance.restored_from_backup, "{provenance:?}"); + assert!(provenance.restored_generation.is_some(), "{provenance:?}"); + assert_eq!(store.restore_provenance(), provenance); + // The restored metastore holds what the first one backed up. + assert!(store.exists(&"db".into()).await); + } + } } diff --git a/libsql-server/src/namespace/store.rs b/libsql-server/src/namespace/store.rs index aa22de9e26..a4c09790fe 100644 --- a/libsql-server/src/namespace/store.rs +++ b/libsql-server/src/namespace/store.rs @@ -21,6 +21,7 @@ use crate::stats::Stats; use super::broadcasters::{BroadcasterHandle, BroadcasterRegistry}; use super::configurator::{DynConfigurator, NamespaceConfigurators}; +use super::fence::audit::CommandAudit; use super::fence::capability::CapabilityPurpose; use super::fence::command::{FenceCommand, FenceRequest}; use super::fence::controller::{FenceController, Transition}; @@ -100,6 +101,7 @@ impl NamespaceStore { // Every namespace with fence state gets its controller before anything is served // (section 8.5). + super::fence::audit::describe_metrics(); let fences = FenceRegistry::seeded(metadata.load_fences().await?); if fences.len() > 0 { tracing::info!("loaded {} namespace fence controllers", fences.len()); @@ -633,6 +635,28 @@ impl NamespaceStore { request: FenceRequest, server: ServerIdentity, ) -> crate::Result { + self.execute_fence_command_authorised(request, server, false) + .await + } + + /// [`execute_fence_command`](Self::execute_fence_command) for a request that may carry + /// the adoption key: `adoption_authorised` says whether it did (section 12). Only + /// `AdoptFence` looks at it. Like every command, a committed adoption is written to the + /// audit log (target `libsql_server::fence::audit`), with its approvers, incident reference + /// and reason. + pub(crate) async fn execute_fence_command_authorised( + &self, + request: FenceRequest, + server: ServerIdentity, + adoption_authorised: bool, + ) -> crate::Result { + if let FenceCommand::AdoptFence(_) = &request.command { + let namespace = request.namespace.clone(); + let controller = self.inner.fences.controller(&namespace); + let mut ctx = FenceContext::now(server, None); + ctx.adoption_authorised = adoption_authorised; + return controller.execute(&self.inner.metadata, request, ctx).await; + } let controller = match request.command { FenceCommand::CreateTargetQuarantined { .. } => { return self.run_create_target(request, server).await @@ -726,9 +750,13 @@ impl NamespaceStore { let this = self.clone(); tokio::spawn(async move { let mut transition = controller.begin_transition().await; + let audit = CommandAudit::new(&request, controller.gate().state()); let ctx = FenceContext::now(server, None); - this.create_target_under(&mut transition, request, ctx) - .await + let result = this + .create_target_under(&mut transition, request, ctx) + .await; + super::fence::audit::command_finished(&audit, &result, &transition.report); + result }) .await? } @@ -981,6 +1009,12 @@ impl NamespaceStore { self.inner.fences.active_count() } + /// Set the fence gauges (`docs/NAMESPACE_FENCE.md` section 15) from the registry, before + /// `/metrics` is rendered. + pub(crate) fn update_fence_gauges(&self) { + super::fence::audit::update_gauges(&self.inner.fences, super::fence::drain::now_ms()); + } + pub(crate) fn schema_locks(&self) -> &SchemaLocksRegistry { &self.inner.schema_locks } diff --git a/libsql-server/src/rpc/proxy.rs b/libsql-server/src/rpc/proxy.rs index b189b1143e..51b46c0dec 100644 --- a/libsql-server/src/rpc/proxy.rs +++ b/libsql-server/src/rpc/proxy.rs @@ -46,6 +46,12 @@ pub mod rpc { let stable_code = other .fence_error() .and_then(|e| e.outcome().proxy_stable_code()); + if let Some(fence) = other.fence_error().filter(|_| stable_code.is_some()) { + crate::namespace::fence::audit::denied( + fence, + crate::namespace::fence::audit::DenialSurface::Rpc, + ); + } let code = match other { _ if stable_code.is_some() => ErrorCode::SqlError, SqldError::LibSqlInvalidQueryParams(_) => ErrorCode::SqlError, @@ -326,11 +332,12 @@ impl ProxyService { crate::error::Error::NamespaceDoesntExist(_) => None, // A namespace the fence refuses is refused with the typed status, never // retried by the write proxy (`docs/NAMESPACE_FENCE.md` section 6.1). - e if fence_status(e).is_some() => Err(fence_status(e).unwrap())?, - _ => Err(tonic::Status::internal(format!( - "Error fetching jwt key for a namespace: {}", - e - )))?, + e => Err(fence_status(e).unwrap_or_else(|| { + tonic::Status::internal(format!( + "Error fetching jwt key for a namespace: {}", + e + )) + }))?, }, Ok(Err(e)) => Err(tonic::Status::internal(format!( "Error fetching jwt key for a namespace: {}", @@ -582,8 +589,15 @@ pub async fn garbage_collect(clients: &mut HashMap> /// The typed status of a fence denial on the proxy service (`docs/NAMESPACE_FENCE.md` section /// 6.1): `FAILED_PRECONDITION` with the stable code, never `UNAVAILABLE`, which the write /// proxy retries without bound. +/// Counted as a denial on the `rpc` surface when it is one. fn fence_status(e: &crate::error::Error) -> Option { - e.fence_error()?.to_grpc_status() + let fence = e.fence_error()?; + let status = fence.to_grpc_status()?; + crate::namespace::fence::audit::denied( + fence, + crate::namespace::fence::audit::DenialSurface::Rpc, + ); + Some(status) } /// The status for an error looking up the namespace a proxy request names. diff --git a/libsql-server/src/rpc/replication/replication_log.rs b/libsql-server/src/rpc/replication/replication_log.rs index 7b6c5d51d7..1502205f7c 100644 --- a/libsql-server/src/rpc/replication/replication_log.rs +++ b/libsql-server/src/rpc/replication/replication_log.rs @@ -87,10 +87,9 @@ impl ReplicationLogService { /// The status of a replication call (or stream) that the namespace fence refuses. Counted /// every time, and logged at most once per namespace per [`FENCE_DENIAL_LOG_INTERVAL`]. fn fence_denied(&self, namespace: &NamespaceName, call: &str, error: FenceError) -> Status { - metrics::increment_counter!( - "libsql_server_fence_denials_total", - "code" => error.outcome().as_str(), - "surface" => "replication", + crate::namespace::fence::audit::denied( + &error, + crate::namespace::fence::audit::DenialSurface::Replication, ); let now = Instant::now(); let log = { diff --git a/libsql-server/tests/fence/admin.rs b/libsql-server/tests/fence/admin.rs index 5a5a2aec03..365ecb1fae 100644 --- a/libsql-server/tests/fence/admin.rs +++ b/libsql-server/tests/fence/admin.rs @@ -45,6 +45,7 @@ fn capabilities() { "PublishTargetReadableWriteFenced", "EnableTargetWrites", "AbortQuarantinedTarget", + "AdoptFence", ] { assert!(commands.contains(&command), "{command} missing: {body}"); } @@ -56,11 +57,34 @@ fn capabilities() { .unwrap() .starts_with("sqld ")); Uuid::parse_str(body["server"]["instance_id"].as_str().unwrap())?; - assert_eq!(body["metastore"]["restored_from_backup"], false); + // Metastore restore provenance: this server's metastore was not restored from a + // backup, which the capability endpoint, every fence view and the gauge all report. + assert_eq!(body["metastore"]["restored_from_backup"], false, "{body}"); + assert_eq!( + body["metastore"]["restored_generation"], + json!(null), + "{body}" + ); + assert_eq!( + crate::common::snapshot_metrics() + .get_gauge("libsql_server_metastore_restored_from_backup"), + Some(0.0) + ); // An active fence is counted. admin.create_namespace("src").await?; let log_id = load_and_log_id(&admin, "src").await?; + let (status, body) = admin.inspect("src").await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!( + body["fence"]["provenance"], + json!({ + "metastore_restored_from_backup": false, + "metastore_restored_generation": null, + "marker": null, + }), + "{body}" + ); let (status, body) = admin .command( "src", @@ -69,6 +93,14 @@ fn capabilities() { ) .await?; assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!( + body["fence"]["provenance"]["metastore_restored_from_backup"], false, + "{body}" + ); + assert_eq!( + body["fence"]["provenance"]["marker"], "consistent", + "{body}" + ); let (_, body) = admin.get("/v1/fence/capabilities").await?; assert_eq!(body["active_fences"], 1, "{body}"); @@ -704,3 +736,185 @@ fn inspect_reports_drain_counters() { }); sim.run().unwrap(); } + +const ADOPTION_KEY: &str = "fence-adoption-key"; + +fn adopt_body( + new_op: Uuid, + cmd: Uuid, + state: &str, + rev: u64, + current_op: Uuid, +) -> serde_json::Value { + command_body( + new_op, + cmd, + state, + rev, + json!({ + "current_operation_id": current_op.to_string(), + "approvers": ["alice", "bob"], + "incident_ref": "INC-1", + "reason": "the operation's control record was lost", + }), + ) +} + +/// `AdoptFence` over HTTP (section 12): the admin credential and the adoption key header are +/// both required, two distinct approvers are required, and an adoption moves the owner and +/// nothing else. The old owner is refused afterwards and the new one finishes the operation. +#[test] +fn adopt_over_http() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary( + &mut sim, + tmp.path().to_path_buf(), + Primary { + adoption_key: Some(ADOPTION_KEY), + ..Default::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 (old, new) = (uuid(0xa), uuid(0xb)); + let conn = connect("src")?; + let (status, acquired) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(old, uuid(1), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{acquired}"); + assert_eq!(state_of(&acquired), ("SOURCE_WRITE_FENCED", 2)); + + let body = adopt_body(new, uuid(2), "SOURCE_WRITE_FENCED", 2, old); + // No key, a wrong key, and the admin credential missing. + for key in [None, Some("not-the-key")] { + let (status, refused) = admin.adopt("src", body.clone(), key).await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{refused}"); + assert_eq!(refused["outcome"], "FENCE_PRECONDITION_FAILED"); + assert_eq!(refused["detail"], "adoption_not_authorised", "{refused}"); + assert_eq!(state_of(&refused), ("SOURCE_WRITE_FENCED", 2), "{refused}"); + } + let (status, _) = Admin::new(None) + .adopt("src", body.clone(), Some(ADOPTION_KEY)) + .await?; + assert_eq!(status, StatusCode::UNAUTHORIZED); + // One approver, the same approver twice, and an unknown field. + for approvers in [json!(["alice"]), json!(["alice", "alice"])] { + let mut bad = adopt_body(new, uuid(3), "SOURCE_WRITE_FENCED", 2, old); + bad["approvers"] = approvers; + let (status, refused) = admin.adopt("src", bad, Some(ADOPTION_KEY)).await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{refused}"); + assert_eq!(refused["detail"], "adoption_not_authorised", "{refused}"); + } + let mut bad = adopt_body(new, uuid(3), "SOURCE_WRITE_FENCED", 2, old); + bad["gate"] = json!("open"); + let (status, refused) = admin.adopt("src", bad, Some(ADOPTION_KEY)).await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{refused}"); + assert_eq!(refused["detail"], "invalid_argument", "{refused}"); + + let (status, adopted) = admin.adopt("src", body.clone(), Some(ADOPTION_KEY)).await?; + assert_eq!(status, StatusCode::OK, "{adopted}"); + assert_eq!(adopted["outcome"], "APPLIED"); + assert_eq!(adopted["replayed"], false); + assert_eq!(state_of(&adopted), ("SOURCE_WRITE_FENCED", 3)); + assert_eq!(adopted["fence"]["operation_id"], new.to_string()); + assert_eq!(adopted["fence"]["admission"]["write"], "closed"); + assert_eq!(adopted["fence"]["admission"]["read"], "open"); + assert_eq!( + adopted["fence"]["frozen_boundary"], + acquired["fence"]["frozen_boundary"] + ); + assert_eq!(adopted["receipt"]["command"], "AdoptFence"); + let adoption = &adopted["receipt"]["adoption"]; + assert_eq!(adoption["previous_operation_id"], old.to_string()); + assert_eq!(adoption["new_operation_id"], new.to_string()); + assert_eq!(adoption["approvers"], json!(["alice", "bob"])); + assert_eq!(adoption["incident_ref"], "INC-1"); + assert_eq!(adopted["fence"]["adoptions"][0], *adoption); + // Writes stay closed; reads stay open. + assert!(conn.execute("insert into t values (2)", ()).await.is_err()); + conn.query("select * from t", ()).await?; + + // A replay returns the stored receipt, even without the key. + let (status, replay) = admin.adopt("src", body.clone(), None).await?; + assert_eq!(status, StatusCode::OK, "{replay}"); + assert_eq!(replay["replayed"], true); + assert_eq!(replay["receipt"], adopted["receipt"]); + + // The old owner is refused; the new owner finishes the operation. + let (status, refused) = admin + .command( + "src", + "source/release-write-fence", + command_body(old, uuid(4), "SOURCE_WRITE_FENCED", 3, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::CONFLICT, "{refused}"); + assert_eq!(refused["outcome"], "FENCE_OWNED_BY_ANOTHER_OPERATION"); + let (status, released) = admin + .command( + "src", + "source/release-write-fence", + command_body(new, uuid(5), "SOURCE_WRITE_FENCED", 3, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{released}"); + assert_eq!(state_of(&released), ("RELEASED", 4)); + conn.execute("insert into t values (3)", ()).await?; + + // A finished operation cannot be adopted. + let (status, refused) = admin + .adopt( + "src", + adopt_body(uuid(0xc), uuid(6), "RELEASED", 4, new), + Some(ADOPTION_KEY), + ) + .await?; + assert_eq!(status, StatusCode::CONFLICT, "{refused}"); + assert_eq!(refused["outcome"], "INVALID_FENCE_TRANSITION"); + assert_eq!(refused["detail"], "operation_finished", "{refused}"); + Ok(()) + }); + sim.run().unwrap(); +} + +/// Without `--namespace-fence-adoption-key`, adoption is disabled whatever the request presents. +#[test] +fn adopt_disabled_without_key() { + 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 (status, acquired) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(uuid(0xa), uuid(1), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{acquired}"); + for key in [None, Some(""), Some(ADOPTION_KEY)] { + let (status, refused) = admin + .adopt( + "src", + adopt_body(uuid(0xb), uuid(2), "SOURCE_WRITE_FENCED", 2, uuid(0xa)), + key, + ) + .await?; + assert_eq!(status, StatusCode::PRECONDITION_FAILED, "{refused}"); + assert_eq!(refused["detail"], "adoption_not_authorised", "{refused}"); + assert_eq!(refused["fence"]["operation_id"], uuid(0xa).to_string()); + } + Ok(()) + }); + sim.run().unwrap(); +} diff --git a/libsql-server/tests/fence/lifecycle.rs b/libsql-server/tests/fence/lifecycle.rs index 2ec5f77e54..5d57915753 100644 --- a/libsql-server/tests/fence/lifecycle.rs +++ b/libsql-server/tests/fence/lifecycle.rs @@ -1,15 +1,18 @@ //! 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 std::sync::Arc; + use hyper::StatusCode; use libsql::Value as SqlValue; use serde_json::{json, Value}; use tempfile::tempdir; +use tokio::sync::Notify; use uuid::Uuid; use super::{ - acquire_body, command_body, connect, load_and_log_id, make_primary, sim, state_of, Admin, - Primary, ADMIN_KEY, + acquire_body, command_body, connect, load_and_log_id, make_primary, make_restartable_primary, + sim, state_of, user_execute, Admin, Primary, ADMIN_KEY, }; fn uuid(n: u128) -> Uuid { @@ -200,3 +203,382 @@ fn lifecycle_rejected_while_fenced() { }); sim.run().unwrap(); } + +#[track_caller] +fn assert_fence(body: &Value, state: &str, revision: u64, write: &str, read: &str) { + assert_eq!(state_of(body), (state, revision), "{body}"); + assert_eq!(body["fence"]["admission"]["write"], write, "{body}"); + assert_eq!(body["fence"]["admission"]["read"], read, "{body}"); +} + +async fn assert_user_locked(ns: &str, sql: &str, code: &str) -> anyhow::Result<()> { + let (status, body) = user_execute(ns, sql).await?; + assert_eq!(status, StatusCode::LOCKED, "{ns}: {body}"); + assert_eq!(body["code"], code, "{ns}: {body}"); + Ok(()) +} + +async fn assert_user_ok(ns: &str, sql: &str) -> anyhow::Result<()> { + let (status, body) = user_execute(ns, sql).await?; + assert_eq!(status, StatusCode::OK, "{ns}: {body}"); + Ok(()) +} + +/// Create a target and move it through validation into `TARGET_WRITE_FENCED`. Returns the exact +/// publish request and response so a restart test can replay the command and compare its receipt. +async fn target_write_fenced( + admin: &Admin, + ns: &str, + op: Uuid, + command_base: u128, +) -> anyhow::Result<(Value, Value)> { + let (status, created) = admin + .command( + ns, + "target/create-quarantined", + command_body(op, uuid(command_base), "ABSENT", 0, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{created}"); + let (_, revision) = state_of(&created); + + let (status, sealed) = admin + .command( + ns, + "target/seal-import", + command_body( + op, + uuid(command_base + 1), + "TARGET_QUARANTINED", + revision, + json!({}), + ), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{sealed}"); + let (_, revision) = state_of(&sealed); + + let (status, validated) = admin + .command( + ns, + "target/validation-receipt", + command_body( + op, + uuid(command_base + 2), + "TARGET_VALIDATING", + revision, + json!({ "result": "ok", "summary": "restart integration" }), + ), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{validated}"); + let (_, revision) = state_of(&validated); + + let publish = command_body( + op, + uuid(command_base + 3), + "TARGET_VALIDATING", + revision, + json!({}), + ); + let (status, published) = admin + .command(ns, "target/publish-readable", publish.clone()) + .await?; + assert_eq!(status, StatusCode::OK, "{published}"); + assert_eq!(state_of(&published).0, "TARGET_WRITE_FENCED", "{published}"); + Ok((publish, published)) +} + +/// Capacity eviction removes the live namespace but not its controller. A reload through the user +/// protocol therefore has the same revision/generation and still rejects writes. +#[test] +fn evicted_namespace_reloads_same_gate() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary( + &mut sim, + tmp.path().to_path_buf(), + Primary { + max_active_namespaces: 1, + ..Primary::default() + }, + ); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + admin.create_namespace("evicted").await?; + let log_id = load_and_log_id(&admin, "evicted").await?; + let operation = uuid(0xac00); + let (status, fenced) = admin + .command( + "evicted", + "source/acquire-write-fence", + acquire_body(operation, uuid(0xac01), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{fenced}"); + let (state, revision) = state_of(&fenced); + let generation = fenced["fence"]["admission"]["generation"].clone(); + assert_fence(&fenced, "SOURCE_WRITE_FENCED", revision, "closed", "open"); + assert_eq!(state, "SOURCE_WRITE_FENCED"); + + // With capacity one, creating each distinct namespace loads it and forces the prior live + // namespace out. More than one successor also lets the asynchronous eviction listener + // finish before `evicted` is accessed again. + for i in 0..4 { + admin.create_namespace(&format!("pressure-{i}")).await?; + } + + // This read reloads `evicted`; the controller is outside the capacity-limited cache. + assert_user_ok("evicted", "select count(*) from t").await?; + assert_user_locked( + "evicted", + "insert into t values (2)", + "MIGRATION_WRITE_FENCED", + ) + .await?; + let (status, reloaded) = admin.inspect("evicted").await?; + assert_eq!(status, StatusCode::OK, "{reloaded}"); + assert_fence(&reloaded, "SOURCE_WRITE_FENCED", revision, "closed", "open"); + assert_eq!( + reloaded["fence"]["admission"]["generation"], generation, + "{reloaded}" + ); + Ok(()) + }); + sim.run().unwrap(); +} + +/// Every active source/target gate is reconstructed before traffic after a clean server restart, +/// and the exact command which produced it still replays its durable receipt. +#[test] +fn restart_keeps_fence() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + let restart = Arc::new(Notify::new()); + let restarted = Arc::new(Notify::new()); + make_restartable_primary( + &mut sim, + tmp.path().to_path_buf(), + Primary::default(), + restart.clone(), + restarted.clone(), + ); + sim.client("client", async move { + let admin = Admin::new(Some(ADMIN_KEY)); + + admin.create_namespace("source-write").await?; + let log_id = load_and_log_id(&admin, "source-write").await?; + let source_write_op = uuid(0xac10); + let source_write_request = acquire_body(source_write_op, uuid(0xac11), &log_id); + let (status, source_write) = admin + .command( + "source-write", + "source/acquire-write-fence", + source_write_request.clone(), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{source_write}"); + + admin.create_namespace("source-read").await?; + let log_id = load_and_log_id(&admin, "source-read").await?; + let source_read_op = uuid(0xac20); + let (status, acquired) = admin + .command( + "source-read", + "source/acquire-write-fence", + acquire_body(source_read_op, uuid(0xac21), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{acquired}"); + let source_read_request = command_body( + source_read_op, + uuid(0xac22), + "SOURCE_WRITE_FENCED", + state_of(&acquired).1, + json!({ "drain_policy": { "deadline_ms": 5000 } }), + ); + let (status, source_read) = admin + .command( + "source-read", + "source/set-read-fence", + source_read_request.clone(), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{source_read}"); + + let target_quarantined_op = uuid(0xac30); + let target_quarantined_request = + command_body(target_quarantined_op, uuid(0xac31), "ABSENT", 0, json!({})); + let (status, target_quarantined) = admin + .command( + "target-quarantined", + "target/create-quarantined", + target_quarantined_request.clone(), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{target_quarantined}"); + + let (target_write_request, target_write) = + target_write_fenced(&admin, "target-write-fenced", uuid(0xac40), 0xac41).await?; + + let expected = [ + ( + "source-write", + "SOURCE_WRITE_FENCED", + state_of(&source_write).1, + "closed", + "open", + ), + ( + "source-read", + "SOURCE_READ_FENCED", + state_of(&source_read).1, + "closed", + "closed", + ), + ( + "target-quarantined", + "TARGET_QUARANTINED", + state_of(&target_quarantined).1, + "closed", + "closed", + ), + ( + "target-write-fenced", + "TARGET_WRITE_FENCED", + state_of(&target_write).1, + "closed", + "open", + ), + ]; + for (ns, state, revision, write, read) in expected { + let (status, body) = admin.inspect(ns).await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_fence(&body, state, revision, write, read); + } + + restart.notify_waiters(); + restarted.notified().await; + // A response proves the restarted admin server is serving from the rebuilt store. + let admin = Admin::new(Some(ADMIN_KEY)); + assert_eq!(admin.get("/v1/fence/capabilities").await?.0, StatusCode::OK); + + for (ns, state, revision, write, read) in expected { + let (status, body) = admin.inspect(ns).await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_fence(&body, state, revision, write, read); + } + + for (ns, route, request, original) in [ + ( + "source-write", + "source/acquire-write-fence", + source_write_request, + source_write, + ), + ( + "source-read", + "source/set-read-fence", + source_read_request, + source_read, + ), + ( + "target-quarantined", + "target/create-quarantined", + target_quarantined_request, + target_quarantined, + ), + ( + "target-write-fenced", + "target/publish-readable", + target_write_request, + target_write, + ), + ] { + let (status, replay) = admin.command(ns, route, request).await?; + assert_eq!(status, StatusCode::OK, "{ns}: {replay}"); + assert_eq!(replay["replayed"], true, "{ns}: {replay}"); + assert_eq!(replay["receipt"], original["receipt"], "{ns}: {replay}"); + } + + assert_user_ok("source-write", "select count(*) from t").await?; + assert_user_locked( + "source-write", + "insert into t values (2)", + "MIGRATION_WRITE_FENCED", + ) + .await?; + assert_user_locked("source-read", "select * from t", "MIGRATION_READ_FENCED").await?; + assert_user_locked( + "target-quarantined", + "select 1", + "MIGRATION_TARGET_QUARANTINED", + ) + .await?; + assert_user_ok("target-write-fenced", "select 1").await?; + assert_user_locked( + "target-write-fenced", + "create table denied (x)", + "MIGRATION_WRITE_FENCED", + ) + .await?; + Ok(()) + }); + sim.run().unwrap(); +} + +/// `TARGET_WRITABLE` is durable too: a restart cannot put the target back in quarantine or close +/// writes, and a lost enable-writes response remains inspectable through exact replay. +#[test] +fn restart_after_enable_writes_stays_writable() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + let restart = Arc::new(Notify::new()); + let restarted = Arc::new(Notify::new()); + make_restartable_primary( + &mut sim, + tmp.path().to_path_buf(), + Primary::default(), + restart.clone(), + restarted.clone(), + ); + sim.client("client", async move { + let admin = Admin::new(Some(ADMIN_KEY)); + let operation = uuid(0xac50); + let (_, published) = + target_write_fenced(&admin, "target-writable", operation, 0xac51).await?; + let enable = command_body( + operation, + uuid(0xac55), + "TARGET_WRITE_FENCED", + state_of(&published).1, + json!({}), + ); + let (status, enabled) = admin + .command("target-writable", "target/enable-writes", enable.clone()) + .await?; + assert_eq!(status, StatusCode::OK, "{enabled}"); + let revision = state_of(&enabled).1; + assert_fence(&enabled, "TARGET_WRITABLE", revision, "open", "open"); + assert_user_ok("target-writable", "create table before_restart (x)").await?; + + restart.notify_waiters(); + restarted.notified().await; + let admin = Admin::new(Some(ADMIN_KEY)); + let (status, body) = admin.inspect("target-writable").await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_fence(&body, "TARGET_WRITABLE", revision, "open", "open"); + assert_user_ok("target-writable", "insert into before_restart values (1)").await?; + assert_user_ok("target-writable", "create table after_restart (x)").await?; + + let (status, replay) = admin + .command("target-writable", "target/enable-writes", enable) + .await?; + assert_eq!(status, StatusCode::OK, "{replay}"); + assert_eq!(replay["replayed"], true, "{replay}"); + assert_eq!(replay["receipt"], enabled["receipt"], "{replay}"); + assert_fence(&replay, "TARGET_WRITABLE", revision, "open", "open"); + Ok(()) + }); + sim.run().unwrap(); +} diff --git a/libsql-server/tests/fence/mod.rs b/libsql-server/tests/fence/mod.rs index 2e9cbae183..e7dddf1785 100644 --- a/libsql-server/tests/fence/mod.rs +++ b/libsql-server/tests/fence/mod.rs @@ -4,19 +4,23 @@ mod admin; mod lifecycle; +mod observability; mod protocol; use std::path::PathBuf; +use std::sync::Arc; use std::time::Duration; use hyper::StatusCode; use libsql_server::auth::user_auth_strategies::http_basic::HttpBasic; use libsql_server::auth::Auth; use libsql_server::config::{ - AdminApiConfig, MetaStoreConfig, RpcClientConfig, RpcServerConfig, UserApiConfig, + AdminApiConfig, FenceAdoptionKey, MetaStoreConfig, RpcClientConfig, RpcServerConfig, + UserApiConfig, }; use s3s::header::AUTHORIZATION; use serde_json::{json, Value}; +use tokio::sync::Notify; use turmoil::{Builder, Sim}; use uuid::Uuid; @@ -27,12 +31,17 @@ use crate::common::net::{ pub const ADMIN_KEY: &str = "fence-admin-key"; +#[derive(Clone, Copy)] pub struct Primary { /// `None` starts the admin API without an auth key. pub admin_key: Option<&'static str>, pub fence_enabled: bool, /// A basic-auth credential the user API requires; `None` leaves it unauthenticated. pub user_credential: Option<&'static str>, + /// The fence adoption key; `None` leaves adoption disabled. + pub adoption_key: Option<&'static str>, + /// Capacity of the live namespace cache. Fence controllers live outside this cache. + pub max_active_namespaces: usize, } impl Default for Primary { @@ -41,6 +50,8 @@ impl Default for Primary { admin_key: Some(ADMIN_KEY), fence_enabled: true, user_credential: None, + adoption_key: None, + max_active_namespaces: 100, } } } @@ -51,44 +62,79 @@ pub fn sim() -> Sim<'static> { .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(); +async fn primary_server(path: PathBuf, primary: Primary) -> anyhow::Result { let Primary { admin_key, fence_enabled, user_credential, + adoption_key, + max_active_namespaces, } = primary; + Ok(TestServer { + path: path.into(), + user_api_config: UserApiConfig { + auth_strategy: match user_credential { + Some(credential) => Auth::new(HttpBasic::new(credential.into())), + None => UserApiConfig::::default().auth_strategy, + }, + ..Default::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, + namespace_fence_adoption_key: adoption_key.and_then(FenceAdoptionKey::new), + ..Default::default() + }, + disable_namespaces: false, + disable_default_namespace: true, + max_active_namespaces, + ..Default::default() + }) +} + +/// 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(); sim.host("primary", move || { let path = path.clone(); async move { - let server = TestServer { - path: path.into(), - user_api_config: UserApiConfig { - auth_strategy: match user_credential { - Some(credential) => Auth::new(HttpBasic::new(credential.into())), - None => UserApiConfig::::default().auth_strategy, - }, - ..Default::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() - }; + primary_server(path, primary).await?.start_sim(8080).await?; + Ok(()) + } + }); +} + +/// A primary that shuts down when `restart` is notified, then starts again on the same path. +/// `restarted` is notified after the second server has rebound the admin and RPC listeners; an +/// admin request made after that notification is the readiness barrier for the restarted server. +pub fn make_restartable_primary( + sim: &mut Sim, + path: PathBuf, + primary: Primary, + restart: Arc, + restarted: Arc, +) { + init_tracing(); + sim.host("primary", move || { + let path = path.clone(); + let restart = restart.clone(); + let restarted = restarted.clone(); + async move { + let mut server = primary_server(path.clone(), primary).await?; + server.shutdown = restart; + server.start_sim(8080).await?; + + let server = primary_server(path, primary).await?; + restarted.notify_one(); server.start_sim(8080).await?; Ok(()) } @@ -195,6 +241,22 @@ impl Admin { self.get(&format!("/v1/namespaces/{ns}/fence")).await } + /// `AdoptFence`, presenting `adoption_key` (if any) in the adoption key header. + pub async fn adopt( + &self, + ns: &str, + body: Value, + adoption_key: Option<&str>, + ) -> anyhow::Result<(StatusCode, Value)> { + let url = format!("http://primary:9090/v1/namespaces/{ns}/fence/adopt"); + let mut headers = self.headers(); + let name = hyper::header::HeaderName::from_static("x-libsql-fence-adoption-key"); + if let Some(key) = adoption_key { + headers.push((name, key)); + } + Self::json(self.client.post_with_headers(&url, &headers, body).await?).await + } + /// A fence command: `route` is the part after `/fence/`. pub async fn command( &self, @@ -244,6 +306,18 @@ pub fn connect(ns: &str) -> anyhow::Result { Ok(db.connect()?) } +/// Execute one statement through the Hrana v1 user endpoint, returning its typed body. +pub async fn user_execute(ns: &str, sql: &str) -> anyhow::Result<(StatusCode, Value)> { + let response = Client::new() + .post( + &format!("http://{ns}.primary:8080/v1/execute"), + json!({ "stmt": { "sql": sql } }), + ) + .await?; + let status = response.status(); + Ok((status, response.json_value().await?)) +} + /// 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 { diff --git a/libsql-server/tests/fence/observability.rs b/libsql-server/tests/fence/observability.rs new file mode 100644 index 0000000000..8b8b3c28db --- /dev/null +++ b/libsql-server/tests/fence/observability.rs @@ -0,0 +1,275 @@ +//! Fence metrics over a running server (`docs/NAMESPACE_FENCE.md` section 15). + +use hyper::StatusCode; +use metrics_util::debugging::DebugValue; +use metrics_util::MetricKind; +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, + user_execute, Admin, Primary, ADMIN_KEY, +}; + +fn uuid(n: u128) -> Uuid { + Uuid::from_u128(n) +} + +/// The value of the metric `name` of `kind` whose labels are exactly `labels`. +fn metric(kind: MetricKind, name: &str, labels: &[(&str, &str)]) -> Option { + let snapshot = crate::common::snapshot_metrics(); + let mut wanted: Vec<(String, String)> = labels + .iter() + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect(); + wanted.sort(); + snapshot + .snapshot() + .iter() + .find(|(key, _)| { + let mut have: Vec<(String, String)> = key + .key() + .labels() + .map(|l| (l.key().to_string(), l.value().to_string())) + .collect(); + have.sort(); + key.kind() == kind && key.key().name() == name && have == wanted + }) + .map(|(_, (_, _, value))| match value { + DebugValue::Counter(v) => DebugValue::Counter(*v), + DebugValue::Gauge(v) => DebugValue::Gauge(*v), + DebugValue::Histogram(v) => DebugValue::Histogram(v.clone()), + }) +} + +#[track_caller] +fn counter(name: &str, labels: &[(&str, &str)]) -> u64 { + match metric(MetricKind::Counter, name, labels) { + Some(DebugValue::Counter(v)) => v, + other => panic!("counter {name} {labels:?}: {other:?}"), + } +} + +#[track_caller] +fn gauge(name: &str, labels: &[(&str, &str)]) -> f64 { + match metric(MetricKind::Gauge, name, labels) { + Some(DebugValue::Gauge(v)) => v.0, + other => panic!("gauge {name} {labels:?}: {other:?}"), + } +} + +/// After a source walk with a forced rollback, a replay, a command-id conflict and denials on +/// the user protocol and a lifecycle route, every fence metric of section 15 is present with +/// its bounded labels, and no fence metric has a label naming the namespace, the operation or +/// a command. +#[test] +fn metrics_and_labels() { + 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); + + // A transaction admitted before the fence is still open at the deadline: the drain + // rolls it back and is then proven. + let conn = connect("src")?; + let tx = conn.transaction().await?; + tx.execute("insert into t values (2)", ()).await?; + let acquire = command_body( + op, + uuid(1), + "UNFENCED", + 0, + json!({ + "expected_namespace_identity": { "log_id": log_id }, + "drain_policy": { "deadline_ms": 100, "on_deadline": "force_rollback" }, + }), + ); + let (status, acquired) = admin + .command("src", "source/acquire-write-fence", acquire.clone()) + .await?; + assert_eq!(status, StatusCode::OK, "{acquired}"); + assert_eq!(state_of(&acquired).0, "SOURCE_WRITE_FENCED", "{acquired}"); + assert!(tx.commit().await.is_err()); + + // A replay, and a conflict on the same command id. + let (status, replay) = admin + .command("src", "source/acquire-write-fence", acquire) + .await?; + assert_eq!(status, StatusCode::OK, "{replay}"); + assert_eq!(replay["replayed"], true); + let (status, body) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(op, uuid(1), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::CONFLICT, "{body}"); + + // Denials: a write over Hrana, and a config change (a lifecycle operation) over the + // admin API. (A delete is refused on a blocking thread, which this test's per-thread + // metrics recorder does not see.) + let (status, body) = user_execute("src", "insert into t values (3)").await?; + assert_eq!(status, StatusCode::LOCKED, "{body}"); + assert_eq!(body["code"], "MIGRATION_WRITE_FENCED", "{body}"); + let (status, body) = admin + .post( + "/v1/namespaces/src/config", + json!({ "block_reads": false, "block_writes": true, "block_reason": null }), + ) + .await?; + assert_eq!(status, StatusCode::LOCKED, "{body}"); + + // The gauges are computed when `/metrics` is read. + let (status, _) = admin.get("/metrics").await?; + assert_eq!(status, StatusCode::OK); + + // First: a snapshot drains the recorded histogram values. + match metric( + MetricKind::Histogram, + "libsql_server_fence_drain_duration_seconds", + &[("kind", "write")], + ) { + Some(DebugValue::Histogram(v)) => assert_eq!(v.len(), 1), + other => panic!("drain histogram: {other:?}"), + } + let acquire = ("command", "AcquireSourceWriteFence"); + assert_eq!( + counter( + "libsql_server_fence_transitions_total", + &[acquire, ("outcome", "APPLIED")] + ), + 2 + ); + assert_eq!( + counter( + "libsql_server_fence_transitions_total", + &[acquire, ("outcome", "FENCE_COMMAND_CONFLICT")] + ), + 1 + ); + assert_eq!( + counter("libsql_server_fence_replays_total", &[("result", "replay")]), + 1 + ); + assert_eq!( + counter( + "libsql_server_fence_replays_total", + &[("result", "conflict")] + ), + 1 + ); + assert_eq!( + counter("libsql_server_fence_forced_total", &[("kind", "rollback")]), + 1 + ); + for surface in ["hrana", "lifecycle", "http"] { + assert!( + counter( + "libsql_server_fence_denials_total", + &[("code", "MIGRATION_WRITE_FENCED"), ("surface", surface)] + ) >= 1, + "{surface}" + ); + } + assert_eq!( + gauge( + "libsql_server_fence_namespaces", + &[("role", "source"), ("state", "SOURCE_WRITE_FENCED")] + ), + 1.0 + ); + assert_eq!( + gauge( + "libsql_server_fence_namespaces", + &[("role", "target"), ("state", "TARGET_QUARANTINED")] + ), + 0.0 + ); + let age = gauge("libsql_server_fence_oldest_active_age_seconds", &[]); + assert!(age >= 0.0, "{age}"); + + // No fence metric carries a namespace, operation or command id as a label value. + let forbidden = [ + "src".to_string(), + op.to_string(), + uuid(1).to_string(), + log_id.clone(), + ]; + let snapshot = crate::common::snapshot_metrics(); + let mut seen = 0; + for (key, _) in snapshot.snapshot() { + let name = key.key().name(); + if !name.starts_with("libsql_server_fence_") { + continue; + } + seen += 1; + for label in key.key().labels() { + assert!( + matches!( + label.key(), + "command" + | "outcome" + | "kind" + | "result" + | "code" + | "surface" + | "role" + | "state" + ), + "{name}: unexpected label {}", + label.key() + ); + assert!( + !forbidden.iter().any(|f| f == label.value()), + "{name}: {}={}", + label.key(), + label.value() + ); + } + } + assert!(seen > 10, "{seen}"); + + // Releasing moves the source out of the write-fenced count. + let (status, body) = admin + .command( + "src", + "source/release-write-fence", + command_body( + op, + uuid(2), + "SOURCE_WRITE_FENCED", + state_of(&acquired).1, + json!({}), + ), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + admin.get("/metrics").await?; + assert_eq!( + gauge( + "libsql_server_fence_namespaces", + &[("role", "source"), ("state", "SOURCE_WRITE_FENCED")] + ), + 0.0 + ); + assert_eq!( + gauge( + "libsql_server_fence_namespaces", + &[("role", "source"), ("state", "RELEASED")] + ), + 1.0 + ); + assert_eq!( + gauge("libsql_server_fence_oldest_active_age_seconds", &[]), + 0.0 + ); + Ok(()) + }); + sim.run().unwrap(); +}