diff --git a/libsql-replication/proto/metadata.proto b/libsql-replication/proto/metadata.proto index 797e9da7f8..374e78ec3a 100644 --- a/libsql-replication/proto/metadata.proto +++ b/libsql-replication/proto/metadata.proto @@ -27,4 +27,17 @@ message DatabaseConfig { optional bool shared_schema = 11; optional string shared_schema_name = 12; optional DurabilityMode durability_mode = 13; + // The namespace fence as seen by the primary when it answered. Only ever filled by the + // primary's replication `Hello`, and only while a fence is active; it is never part of a + // stored configuration. Absent from older primaries; older replicas ignore it and still + // see the legacy `block_*` fields above. + optional ReplicatedFence fence = 14; +} + +// The part of a namespace fence a replica needs to apply the primary's read admission. +message ReplicatedFence { + // The fence state name, e.g. "SOURCE_READ_FENCED". + string state = 1; + // The revision of the fence record the state belongs to. + uint64 revision = 2; } diff --git a/libsql-replication/proto/proxy.proto b/libsql-replication/proto/proxy.proto index 5949cb078d..7d1f3897e8 100644 --- a/libsql-replication/proto/proxy.proto +++ b/libsql-replication/proto/proxy.proto @@ -43,6 +43,11 @@ message Error { ErrorCode code = 1; string message = 2; int32 extended_code = 3; + // Stable machine-readable outcome of the error, e.g. "MIGRATION_WRITE_FENCED" for a + // request refused by a namespace fence. Absent when the error has no typed outcome, and + // always absent from older servers: a receiver treats an absent field as "no typed + // outcome" and falls back to `code`. + optional string stable_code = 4; } message ResultRows { diff --git a/libsql-replication/src/generated/metadata.rs b/libsql-replication/src/generated/metadata.rs index dac5a62511..8dbe2f7266 100644 --- a/libsql-replication/src/generated/metadata.rs +++ b/libsql-replication/src/generated/metadata.rs @@ -32,6 +32,23 @@ pub struct DatabaseConfig { pub shared_schema_name: ::core::option::Option<::prost::alloc::string::String>, #[prost(enumeration = "DurabilityMode", optional, tag = "13")] pub durability_mode: ::core::option::Option, + /// The namespace fence as seen by the primary when it answered. Only ever filled by the + /// primary's replication `Hello`, and only while a fence is active; it is never part of a + /// stored configuration. Absent from older primaries; older replicas ignore it and still + /// see the legacy `block_*` fields above. + #[prost(message, optional, tag = "14")] + pub fence: ::core::option::Option, +} +/// The part of a namespace fence a replica needs to apply the primary's read admission. +#[allow(clippy::derive_partial_eq_without_eq)] +#[derive(Clone, PartialEq, ::prost::Message)] +pub struct ReplicatedFence { + /// The fence state name, e.g. "SOURCE_READ_FENCED". + #[prost(string, tag = "1")] + pub state: ::prost::alloc::string::String, + /// The revision of the fence record the state belongs to. + #[prost(uint64, tag = "2")] + pub revision: u64, } #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)] #[repr(i32)] diff --git a/libsql-replication/src/generated/proxy.rs b/libsql-replication/src/generated/proxy.rs index a76bb9bf65..cb63ce6ea2 100644 --- a/libsql-replication/src/generated/proxy.rs +++ b/libsql-replication/src/generated/proxy.rs @@ -77,6 +77,12 @@ pub struct Error { pub message: ::prost::alloc::string::String, #[prost(int32, tag = "3")] pub extended_code: i32, + /// Stable machine-readable outcome of the error, e.g. "MIGRATION_WRITE_FENCED" for a + /// request refused by a namespace fence. Absent when the error has no typed outcome, and + /// always absent from older servers: a receiver treats an absent field as "no typed + /// outcome" and falls back to `code`. + #[prost(string, optional, tag = "4")] + pub stable_code: ::core::option::Option<::prost::alloc::string::String>, } /// Nested message and enum types in `Error`. pub mod error { diff --git a/libsql-replication/src/rpc.rs b/libsql-replication/src/rpc.rs index 57f22ffa1a..242ae083d7 100644 --- a/libsql-replication/src/rpc.rs +++ b/libsql-replication/src/rpc.rs @@ -105,3 +105,118 @@ pub mod metadata { #![allow(clippy::all)] include!("generated/metadata.rs"); } + +#[cfg(test)] +mod test { + use prost::Message; + + use super::metadata::{DatabaseConfig, ReplicatedFence}; + use super::proxy::{error::ErrorCode, Error}; + + /// `proxy.Error` as a peer built before `stable_code` existed knows it. + #[derive(Clone, PartialEq, ::prost::Message)] + struct ErrorWithoutStableCode { + #[prost(enumeration = "ErrorCode", tag = "1")] + code: i32, + #[prost(string, tag = "2")] + message: String, + #[prost(int32, tag = "3")] + extended_code: i32, + } + + /// The legacy part of `metadata.DatabaseConfig`, as a peer built before `fence` existed + /// knows it (the fields in between are skipped the same way as `fence` is). + #[derive(Clone, PartialEq, ::prost::Message)] + struct DatabaseConfigWithoutFence { + #[prost(bool, tag = "1")] + block_reads: bool, + #[prost(bool, tag = "2")] + block_writes: bool, + #[prost(string, optional, tag = "3")] + block_reason: Option, + #[prost(uint64, tag = "4")] + max_db_pages: u64, + } + + fn error(stable_code: Option<&str>) -> Error { + Error { + code: ErrorCode::SqlError as i32, + message: "writes are fenced".into(), + extended_code: 23, + stable_code: stable_code.map(Into::into), + } + } + + #[test] + fn proxy_error_stable_code_is_additive() { + // A newer server's error, read by an older replica: the known fields are intact and + // the stable code is skipped. + let new = error(Some("MIGRATION_WRITE_FENCED")); + let old = ErrorWithoutStableCode::decode(&new.encode_to_vec()[..]).unwrap(); + assert_eq!( + old, + ErrorWithoutStableCode { + code: ErrorCode::SqlError as i32, + message: "writes are fenced".into(), + extended_code: 23, + } + ); + + // An older server's error, read by a newer replica: no typed outcome. + let decoded = Error::decode(&old.encode_to_vec()[..]).unwrap(); + assert_eq!(decoded, error(None)); + + // Without a stable code the encoding is exactly the older one, so an error that has no + // typed outcome is unchanged on the wire. + assert_eq!(error(None).encode_to_vec(), old.encode_to_vec()); + + // And the field round-trips between newer peers. + assert_eq!(Error::decode(&new.encode_to_vec()[..]).unwrap(), new); + } + + fn config(fence: Option) -> DatabaseConfig { + DatabaseConfig { + block_reads: true, + block_writes: true, + block_reason: Some("namespace fence".into()), + max_db_pages: 1024, + fence, + ..Default::default() + } + } + + #[test] + fn replicated_fence_is_additive() { + let fence = ReplicatedFence { + state: "SOURCE_READ_FENCED".into(), + revision: 3, + }; + + // A newer primary's config, read by an older replica: the legacy block fields are + // intact, so the older replica still applies the legacy mirror of the fence. + let new = config(Some(fence.clone())); + let old = DatabaseConfigWithoutFence::decode(&new.encode_to_vec()[..]).unwrap(); + assert_eq!( + old, + DatabaseConfigWithoutFence { + block_reads: true, + block_writes: true, + block_reason: Some("namespace fence".into()), + max_db_pages: 1024, + } + ); + + // An older primary's config, read by a newer replica: no fence. + let decoded = DatabaseConfig::decode(&old.encode_to_vec()[..]).unwrap(); + assert_eq!(decoded.fence, None); + assert_eq!(decoded, config(None)); + + // Without a fence the encoding is exactly the older one: a stored configuration, which + // never carries a fence, is unchanged. + assert_eq!(config(None).encode_to_vec(), old.encode_to_vec()); + + // And the fence round-trips between newer peers. + let round_trip = DatabaseConfig::decode(&new.encode_to_vec()[..]).unwrap(); + assert_eq!(round_trip.fence, Some(fence)); + } +} diff --git a/libsql-server/src/connection/config.rs b/libsql-server/src/connection/config.rs index 970014415d..1e835d395f 100644 --- a/libsql-server/src/connection/config.rs +++ b/libsql-server/src/connection/config.rs @@ -108,6 +108,9 @@ impl From<&DatabaseConfig> for metadata::DatabaseConfig { shared_schema: Some(value.is_shared_schema), shared_schema_name: value.shared_schema_name.as_ref().map(|s| s.to_string()), durability_mode: Some(metadata::DurabilityMode::from(value.durability_mode).into()), + // Never part of a configuration: the replication `hello` fills it from the live + // fence gate, and it is not stored. + fence: None, } } } diff --git a/libsql-server/src/connection/write_proxy.rs b/libsql-server/src/connection/write_proxy.rs index 030920c6ec..cc86a2170d 100644 --- a/libsql-server/src/connection/write_proxy.rs +++ b/libsql-server/src/connection/write_proxy.rs @@ -266,13 +266,15 @@ impl RemoteConnection { let response_stream = match client.stream_exec(req).await { Ok(i) => i.into_inner(), Err(e) => { + // Only `UNAVAILABLE` is retried. A fence denial is `FAILED_PRECONDITION` + // with its stable code, answered to the client as the primary's denial. if e.code() == Code::Unavailable { tracing::error!("retrying proxy connection: {}", e); tokio::time::sleep(Duration::from_millis(500) * 2u32.pow(retries)).await; retries += 1; continue; } else { - return Err(e.into()); + return Err(Error::from_proxy_status(e)); } } }; @@ -387,7 +389,7 @@ where ) } exec_resp::Response::DescribeResp(_) => Err(Error::PrimaryStreamMisuse), - exec_resp::Response::Error(e) => Err(Error::RpcQueryError(e)), + exec_resp::Response::Error(e) => Err(Error::from_proxy_error(e)), } }; @@ -430,7 +432,7 @@ where Ok(false) } - exec_resp::Response::Error(e) => Err(Error::RpcQueryError(e)), + exec_resp::Response::Error(e) => Err(Error::from_proxy_error(e)), exec_resp::Response::ProgramResp(_) => Err(Error::PrimaryStreamMisuse), }; diff --git a/libsql-server/src/error.rs b/libsql-server/src/error.rs index f0cb631769..b6903c4a82 100644 --- a/libsql-server/src/error.rs +++ b/libsql-server/src/error.rs @@ -151,6 +151,55 @@ pub trait ResponseError: std::error::Error { impl ResponseError for Error {} +impl Error { + /// The fence denial this error carries, looking through the wrappers it can arrive in. + pub(crate) fn fence_error(&self) -> Option<&crate::namespace::fence::outcome::FenceError> { + match self { + Error::NamespaceFence(e) => Some(e), + Error::Migration(crate::schema::Error::NamespaceFence(e)) => Some(e), + Error::Ref(this) => this.fence_error(), + Error::Anyhow(e) => e.downcast_ref::().and_then(Error::fence_error), + _ => None, + } + } + + /// A step or program error the primary returned through the write proxy. A fence denial + /// from a primary that fills the additive `stable_code` field + /// (`docs/NAMESPACE_FENCE.md` section 6.1) becomes the same [`Error::NamespaceFence`] a + /// local denial is, so the replica answers its client exactly as the primary would. Any + /// other error, and every error from a primary that does not fill the field, stays + /// [`Error::RpcQueryError`]. + pub(crate) fn from_proxy_error(e: crate::rpc::proxy::rpc::Error) -> Self { + let fence = e.stable_code.as_deref().and_then(|code| { + crate::namespace::fence::outcome::FenceError::from_proxy_stable_code(code, &e.message) + }); + match fence { + Some(fence) => Error::NamespaceFence(fence), + None => Error::RpcQueryError(e), + } + } + + /// A gRPC status from the primary's proxy service: its typed fence denial + /// (`FAILED_PRECONDITION` with the stable code, section 6.1) as [`Error::NamespaceFence`], + /// 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), + None => Error::RpcQueryExecutionError(status), + } + } +} + +/// The HTTP response for a fence denial (`docs/NAMESPACE_FENCE.md` section 6): the fence status +/// and the JSON error body with the additive `code` (and `detail`) fields. +pub(crate) fn fence_error_response( + e: &crate::namespace::fence::outcome::FenceError, +) -> axum::response::Response { + let status = e.http_status(); + tracing::debug!("HTTP API: {status}, {e}"); + (status, axum::Json(e.http_error_body())).into_response() +} + impl IntoResponse for Error { fn into_response(self) -> axum::response::Response { (&self).into_response() @@ -226,7 +275,7 @@ impl IntoResponse for &Error { AttachInMigration => self.format_err(StatusCode::BAD_REQUEST), RuntimeTaskJoinError(_) => self.format_err(StatusCode::INTERNAL_SERVER_ERROR), NotAPrimary => self.format_err(StatusCode::BAD_REQUEST), - NamespaceFence(e) => self.format_err(e.outcome().admin_http_status()), + NamespaceFence(e) => fence_error_response(e), } } } @@ -338,3 +387,63 @@ impl IntoResponse for &ForkError { } } } + +#[cfg(test)] +mod fence_tests { + use super::*; + use crate::namespace::fence::outcome::{FenceDetail, FenceError, FenceOutcome}; + + async fn response(e: &Error) -> (StatusCode, serde_json::Value) { + let response = e.into_response(); + let status = response.status(); + let body = hyper::body::to_bytes(response.into_body()).await.unwrap(); + (status, serde_json::from_slice(&body).unwrap()) + } + + /// Section 6: a fence denial is `423` with the additive `code` (and `detail`) field, found + /// through every wrapper the error can arrive in; other errors keep their shape. + #[tokio::test] + async fn fence_errors_carry_code() { + for outcome in [ + FenceOutcome::MigrationWriteFenced, + FenceOutcome::MigrationReadFenced, + FenceOutcome::MigrationTargetQuarantined, + FenceOutcome::FenceStateUnavailable, + ] { + let fence = FenceError::new(outcome, "denied"); + let wrapped = [ + Error::NamespaceFence(fence.clone()), + Error::Ref(std::sync::Arc::new(Error::NamespaceFence(fence.clone()))), + Error::Anyhow(anyhow::anyhow!(Error::NamespaceFence(fence.clone()))), + Error::Migration(crate::schema::Error::NamespaceFence(fence.clone())), + ]; + for e in &wrapped { + assert_eq!(e.fence_error(), Some(&fence), "{e:?}"); + let (status, body) = response(e).await; + assert_eq!(status, StatusCode::LOCKED, "{e:?}"); + assert_eq!(body["code"], outcome.as_str(), "{e:?}"); + assert_eq!(body["error"], fence.to_string(), "{e:?}"); + assert!(body.get("detail").is_none(), "{body}"); + } + } + + let unavailable = FenceError::new(FenceOutcome::FenceStateUnavailable, "corrupt") + .with_detail(FenceDetail::CorruptRecord); + let (_, body) = response(&Error::NamespaceFence(unavailable)).await; + assert_eq!(body["detail"], "corrupt_record"); + + let precondition = FenceError::new(FenceOutcome::FencePreconditionFailed, "no") + .with_detail(FenceDetail::NotPrimary); + let (status, body) = response(&Error::NamespaceFence(precondition)).await; + assert_eq!(status, StatusCode::PRECONDITION_FAILED); + assert_eq!(body["code"], "FENCE_PRECONDITION_FAILED"); + assert_eq!(body["detail"], "not_primary"); + + // The legacy `block_*` refusal keeps its mapping and has no code. + let blocked = Error::Blocked(Some("maintenance".into())); + assert_eq!(blocked.fence_error(), None); + let (status, body) = response(&blocked).await; + assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); + assert!(body.get("code").is_none(), "{body}"); + } +} diff --git a/libsql-server/src/hrana/batch.rs b/libsql-server/src/hrana/batch.rs index 9b73551774..4292660088 100644 --- a/libsql-server/src/hrana/batch.rs +++ b/libsql-server/src/hrana/batch.rs @@ -27,6 +27,10 @@ pub enum BatchError { ResponseTooLarge, #[error("Schema migration error: {message}")] SchemaError { message: String }, + /// The whole batch was refused by the namespace fence (`docs/NAMESPACE_FENCE.md` + /// section 6), for example a read under a read fence. + #[error(transparent)] + Fence(crate::namespace::fence::outcome::FenceError), } fn proto_cond_to_cond( @@ -183,6 +187,7 @@ pub fn batch_error_from_sqld_error(sqld_error: SqldError) -> Result { BatchError::ResponseTooLarge } + SqldError::NamespaceFence(e) => BatchError::Fence(e), sqld_error => return Err(sqld_error), }) } @@ -201,6 +206,7 @@ impl BatchError { Self::TransactionBusy => "TRANSACTION_BUSY", Self::ResponseTooLarge => "RESPONSE_TOO_LARGE", Self::SchemaError { message: _ } => "SCHEMA_MIGRATION_ERROR", + Self::Fence(e) => e.outcome().as_str(), } } } diff --git a/libsql-server/src/hrana/stmt.rs b/libsql-server/src/hrana/stmt.rs index 4414c72e1d..49bb65d2f7 100644 --- a/libsql-server/src/hrana/stmt.rs +++ b/libsql-server/src/hrana/stmt.rs @@ -49,6 +49,10 @@ pub enum StmtError { ResponseTooLarge, #[error("error executing a request on the primary: {0}")] Proxy(String), + /// A denial by the namespace fence (`docs/NAMESPACE_FENCE.md` section 6). Its Hrana code is + /// the fence's stable code. + #[error(transparent)] + Fence(crate::namespace::fence::outcome::FenceError), } pub async fn execute_stmt( @@ -216,6 +220,7 @@ 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::RusqliteError(rusqlite_error) | SqldError::RusqliteErrorExtended(rusqlite_error, _) => match rusqlite_error { rusqlite::Error::SqliteFailure(sqlite_error, Some(message)) => { @@ -271,6 +276,7 @@ impl StmtError { Self::Blocked { .. } => "BLOCKED", Self::ResponseTooLarge => "RESPONSE_TOO_LARGE", Self::Proxy(_) => "PROXY_ERROR", + Self::Fence(e) => e.outcome().as_str(), } } } diff --git a/libsql-server/src/http/user/hrana_over_http_1.rs b/libsql-server/src/http/user/hrana_over_http_1.rs index 57ad971d91..e3809e7211 100644 --- a/libsql-server/src/http/user/hrana_over_http_1.rs +++ b/libsql-server/src/http/user/hrana_over_http_1.rs @@ -13,6 +13,9 @@ use super::db_factory::MakeConnectionExtractor; enum ResponseError { #[error(transparent)] Stmt(hrana::stmt::StmtError), + /// A whole batch refused by the namespace fence. + #[error(transparent)] + Fence(crate::namespace::fence::outcome::FenceError), } pub async fn handle_index() -> hyper::Response { @@ -81,6 +84,7 @@ pub(crate) async fn handle_batch( hrana::batch::execute_batch(&db, ctx, pgm, req_body.batch.replication_index) .await .map(|result| RespBody { result }) + .map_err(catch_batch_fence_error) .context("Could not execute batch") }) .await?; @@ -131,6 +135,7 @@ where fn response_error_response(err: ResponseError) -> hyper::Response { use hrana::stmt::StmtError; let status = match &err { + ResponseError::Fence(err) => err.http_status(), ResponseError::Stmt(err) => match err { StmtError::SqlParse { .. } | StmtError::SqlNoStmt @@ -145,6 +150,7 @@ fn response_error_response(err: ResponseError) -> hyper::Response { hyper::StatusCode::SERVICE_UNAVAILABLE } StmtError::SqliteError { .. } => hyper::StatusCode::INTERNAL_SERVER_ERROR, + StmtError::Fence(err) => err.http_status(), }, }; @@ -184,10 +190,21 @@ fn catch_stmt_error(err: anyhow::Error) -> anyhow::Error { } } +/// A batch the fence refused as a whole is answered like a refused statement, with the fence's +/// status and code, rather than as an internal error. +fn catch_batch_fence_error(err: anyhow::Error) -> anyhow::Error { + match err.downcast::() { + Ok(hrana::batch::BatchError::Fence(e)) => anyhow!(ResponseError::Fence(e)), + Ok(batch_err) => anyhow!(batch_err), + Err(err) => err, + } +} + impl ResponseError { pub fn code(&self) -> &'static str { match self { Self::Stmt(err) => err.code(), + Self::Fence(err) => err.outcome().as_str(), } } } diff --git a/libsql-server/src/http/user/mod.rs b/libsql-server/src/http/user/mod.rs index 7d46dce8cd..eed6e751b5 100644 --- a/libsql-server/src/http/user/mod.rs +++ b/libsql-server/src/http/user/mod.rs @@ -3,7 +3,7 @@ pub(crate) mod dump; mod extract; mod hrana_over_http_1; mod listen; -mod result_builder; +pub(crate) mod result_builder; mod trace; mod types; #[macro_use] @@ -143,9 +143,12 @@ async fn handle_query( let db = connection_maker.create().await?; let builder = JsonHttpPayloadBuilder::new(); - let builder = db + let mut builder = db .execute_batch_or_rollback(batch, ctx, builder, query.replication_index) .await?; + if let Some(denial) = builder.take_fence_denial() { + return Err(Error::NamespaceFence(denial)); + } let res = ( [(header::CONTENT_TYPE, "application/json")], diff --git a/libsql-server/src/http/user/result_builder.rs b/libsql-server/src/http/user/result_builder.rs index 1364b5d956..cd341ba6cb 100644 --- a/libsql-server/src/http/user/result_builder.rs +++ b/libsql-server/src/http/user/result_builder.rs @@ -7,6 +7,7 @@ use serde::{Serialize, Serializer}; use serde_json::ser::{CompactFormatter, Formatter}; use std::sync::atomic::Ordering; +use crate::namespace::fence::outcome::FenceError; use crate::query_result_builder::{ Column, JsonFormatter, QueryBuilderConfig, QueryResultBuilder, QueryResultBuilderError, TOTAL_RESPONSE_SIZE, @@ -25,6 +26,9 @@ pub struct JsonHttpPayloadBuilder { step_row_count: usize, is_step_error: bool, is_step_empty: bool, + /// The first step error that was a fence denial. The legacy API has no per-step codes, so + /// a fenced batch is answered as a whole with the fence's status and code. + fence_denial: Option, } #[derive(Default)] @@ -112,8 +116,14 @@ impl JsonHttpPayloadBuilder { step_row_count: 0, is_step_error: false, is_step_empty: false, + fence_denial: None, } } + + /// The first fence denial reported as a step error, if any. + pub fn take_fence_denial(&mut self) -> Option { + self.fence_denial.take() + } } impl<'a> Serialize for HttpJsonValueSerializer<'a> { @@ -209,6 +219,9 @@ impl QueryResultBuilder for JsonHttpPayloadBuilder { } fn step_error(&mut self, error: crate::error::Error) -> Result<(), QueryResultBuilderError> { + if self.fence_denial.is_none() { + self.fence_denial = error.fence_error().cloned(); + } self.is_step_error = true; self.is_step_empty = false; self.buffer.truncate(self.checkpoint); diff --git a/libsql-server/src/namespace/fence/controller.rs b/libsql-server/src/namespace/fence/controller.rs index 621576bacb..e61421f341 100644 --- a/libsql-server/src/namespace/fence/controller.rs +++ b/libsql-server/src/namespace/fence/controller.rs @@ -18,6 +18,8 @@ use parking_lot::Mutex; use tokio::sync::{watch, Notify, OwnedMutexGuard}; use uuid::Uuid; +use libsql_replication::rpc::metadata::ReplicatedFence; + use crate::connection::connection_manager::{ConnectionManager, WeakConnectionManager}; use crate::error::Error; use crate::namespace::meta_store::{FenceCommit, FenceContext, MetaStore}; @@ -157,6 +159,17 @@ impl GateSnapshot { Ok(()) } + /// The fence as a replica needs it (`docs/NAMESPACE_FENCE.md` section 6.2): the state and + /// revision of an active fence, `None` while no fence is active. The primary fills it into + /// the configuration its replication `hello` returns; it is never stored. + pub fn replicated(&self) -> Option { + let state = self.state(); + state.is_active().then(|| ReplicatedFence { + state: state.as_str().into(), + revision: self.revision(), + }) + } + /// Whether a closing transition is being installed. pub fn is_installing(&self) -> bool { self.installing.is_some() diff --git a/libsql-server/src/namespace/fence/mod.rs b/libsql-server/src/namespace/fence/mod.rs index af5b76182e..43c384659a 100644 --- a/libsql-server/src/namespace/fence/mod.rs +++ b/libsql-server/src/namespace/fence/mod.rs @@ -49,7 +49,7 @@ pub const FENCE_PROTOCOL_VERSION: u32 = 1; /// Whether this server fills the proxy protocol's additive `Error.stable_code` field and maps /// it on the replica side (`docs/NAMESPACE_FENCE.md` section 6.1). Reported by capability /// discovery so that deployment tooling can check every server before fences are used. -pub const PROXY_STABLE_CODE: bool = false; +pub const PROXY_STABLE_CODE: bool = true; /// The identity of this server process: its build and an id generated once per process. It is /// written into records and receipts, and reported by the admin API. diff --git a/libsql-server/src/namespace/fence/outcome.rs b/libsql-server/src/namespace/fence/outcome.rs index cf4df0c9af..723b7b5ab6 100644 --- a/libsql-server/src/namespace/fence/outcome.rs +++ b/libsql-server/src/namespace/fence/outcome.rs @@ -275,6 +275,27 @@ impl FenceError { &self.message } + /// The status of this error on either HTTP API: the user API's mapping for a data-plane + /// denial (`423`), and the admin API's for any other outcome. + pub fn http_status(&self) -> StatusCode { + self.outcome + .user_http_status() + .unwrap_or_else(|| self.outcome.admin_http_status()) + } + + /// The JSON error body of the HTTP APIs for this error: the usual `error` message plus the + /// additive stable `code` and, when there is one, the bounded `detail`. + pub fn http_error_body(&self) -> serde_json::Value { + let mut body = serde_json::json!({ + "error": self.to_string(), + "code": self.outcome.as_str(), + }); + if let Some(detail) = self.detail { + body["detail"] = detail.as_str().into(); + } + body + } + /// A gRPC status for this error, if the outcome has a gRPC mapping. The code is in the /// [`GRPC_FENCE_CODE_METADATA`] entry and prefixes the message. pub fn to_grpc_status(&self) -> Option { @@ -297,6 +318,38 @@ impl FenceError { .parse() .ok() } + + /// The fence denial a peer reported as a gRPC status produced by + /// [`FenceError::to_grpc_status`], with the peer's message. `None` for any other status, + /// including one whose code is not the fence mapping of the outcome it names. + pub fn from_grpc_status(status: &tonic::Status) -> Option { + let outcome = Self::outcome_from_grpc_status(status)?; + if outcome.grpc_code() != Some(status.code()) { + return None; + } + Some(Self::from_peer(outcome, status.message())) + } + + /// The fence denial a peer reported in the proxy protocol's `Error.stable_code` field + /// (section 6.1), with the peer's message. `None` for a code this server does not know or + /// that is not a data-plane denial, which the peer never sends: such an error keeps its + /// untyped mapping. + pub fn from_proxy_stable_code(stable_code: &str, message: &str) -> Option { + let outcome = stable_code.parse::().ok()?; + outcome.proxy_stable_code()?; + Some(Self::from_peer(outcome, message)) + } + + /// A denial reported by a peer. The peer's message is `": "` (this type's + /// `Display`); the prefix is dropped so that it is not repeated. The bounded `detail` is + /// not carried by the peer protocols. + fn from_peer(outcome: FenceOutcome, message: &str) -> FenceError { + let message = message + .strip_prefix(outcome.as_str()) + .and_then(|m| m.strip_prefix(": ")) + .unwrap_or(message); + Self::new(outcome, message) + } } #[cfg(test)] @@ -404,6 +457,49 @@ mod tests { assert!(control.to_grpc_status().is_none()); } + /// A denial crosses the proxy and RPC protocols as the same outcome and message, and + /// nothing else is mistaken for one. + #[test] + fn peer_denials_round_trip() { + let err = FenceError::new(FenceOutcome::MigrationWriteFenced, "writes are fenced"); + let status = err.to_grpc_status().unwrap(); + assert_eq!(FenceError::from_grpc_status(&status), Some(err.clone())); + // The metadata alone is not enough: the code must be the outcome's gRPC mapping. + let mut forged = tonic::Status::unavailable("MIGRATION_WRITE_FENCED: x"); + forged.metadata_mut().insert( + GRPC_FENCE_CODE_METADATA, + tonic::metadata::MetadataValue::from_static("MIGRATION_WRITE_FENCED"), + ); + assert_eq!(FenceError::from_grpc_status(&forged), None); + assert_eq!( + FenceError::from_grpc_status(&tonic::Status::failed_precondition("x")), + None + ); + + for outcome in FenceOutcome::ALL { + let decoded = FenceError::from_proxy_stable_code(outcome.as_str(), "m"); + match outcome.proxy_stable_code() { + Some(_) => assert_eq!(decoded, Some(FenceError::new(outcome, "m"))), + None => assert_eq!(decoded, None, "{outcome}"), + } + } + assert_eq!( + FenceError::from_proxy_stable_code("MIGRATION_READ_FENCED", &err_msg()), + Some(FenceError::new( + FenceOutcome::MigrationReadFenced, + "reads are fenced" + )) + ); + assert_eq!( + FenceError::from_proxy_stable_code("SOMETHING_NEW", "m"), + None + ); + + fn err_msg() -> String { + FenceError::new(FenceOutcome::MigrationReadFenced, "reads are fenced").to_string() + } + } + #[test] #[should_panic] fn success_is_not_an_error() { diff --git a/libsql-server/src/namespace/fence/stream.rs b/libsql-server/src/namespace/fence/stream.rs index 2a043d89fa..2711105a2a 100644 --- a/libsql-server/src/namespace/fence/stream.rs +++ b/libsql-server/src/namespace/fence/stream.rs @@ -223,15 +223,18 @@ mod tests { use bytes::Bytes; use futures::stream::BoxStream; use futures::{Stream, StreamExt}; + use libsql_replication::rpc::metadata::{self, ReplicatedFence}; use libsql_replication::rpc::replication::replication_log_server::ReplicationLog; use libsql_replication::rpc::replication::{ Frame, HelloRequest, LogOffset, NAMESPACE_METADATA_KEY, SESSION_TOKEN_KEY, }; use tonic::metadata::{AsciiMetadataValue, BinaryMetadataValue}; + use uuid::Uuid; use super::*; use crate::error::Error; use crate::http::user::dump::dump_stream; + use crate::namespace::fence::command::{FenceCommand, FenceRequest}; use crate::namespace::fence::drain::tests::{fence_outcome, raw, Source, LONG, OP, PROMPT}; use crate::namespace::fence::read::tests::{fenced_source, read_fence, NOW}; use crate::namespace::fence::state::FenceState; @@ -323,6 +326,48 @@ mod tests { assert!(!String::from_utf8_lossy(&body).contains("COMMIT;")); } + /// The same cancellation seen through the HTTP response `/dump` returns: the body fails + /// instead of ending, which hyper sends as an aborted response (no final chunk), so a + /// client can never take the bytes it received for a complete dump. + #[tokio::test(flavor = "multi_thread")] + async fn dump_response_aborted_on_cancel() { + use axum::response::IntoResponse as _; + use hyper::body::HttpBody as _; + + let s = large_fenced_source().await; + let stream = dump(&s).await.unwrap(); + let response = axum::body::StreamBody::new(stream).into_response(); + assert_eq!(response.status(), hyper::StatusCode::OK); + let mut body = response.into_body(); + let mut head = Vec::new(); + while !String::from_utf8_lossy(&head).contains("INSERT INTO") { + head.extend_from_slice(&body.data().await.unwrap().unwrap()); + } + + let fenced = tokio::time::timeout(PROMPT, s.execute(read_fence(&s, 2, NOW))) + .await + .expect("the dump's lease was never released") + .unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + + let mut failed = false; + while let Some(chunk) = tokio::time::timeout(PROMPT, body.data()) + .await + .expect("the dump body stalled") + { + match chunk { + Ok(bytes) => head.extend_from_slice(&bytes), + Err(e) => { + assert!(e.to_string().contains("MIGRATION_READ_FENCED"), "{e}"); + failed = true; + break; + } + } + } + assert!(failed, "the response body ended cleanly"); + assert!(!String::from_utf8_lossy(&head).contains("COMMIT;")); + } + /// The read drain waits for a running dump rather than for time: the fence is acknowledged /// only once the dump has completed and released its lease. #[tokio::test(flavor = "multi_thread")] @@ -413,6 +458,82 @@ mod tests { } } + /// The configuration a replica's `hello` receives from `s`. + async fn hello_config(s: &Source) -> metadata::DatabaseConfig { + let replication = Replication { + service: ReplicationLogService::new(s.store.clone(), None, None, false, false, true), + token: None, + }; + replication + .service + .hello(replication.request(HelloRequest { + handshake_version: Some(1), + })) + .await + .unwrap() + .into_inner() + .config + .expect("hello always carries the configuration") + } + + /// `hello` carries the live fence while one is active (`docs/NAMESPACE_FENCE.md` section + /// 6.2), and nothing before the fence and once it is released. The stored configuration + /// never carries it. + #[tokio::test] + async fn hello_carries_replicated_fence() { + let s = Source::new().await; + let unfenced = hello_config(&s).await; + assert_eq!(unfenced.fence, None); + assert!(!unfenced.block_writes); + + let acquired = s.execute(s.acquire(OP, 1, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&acquired), FenceOutcome::Applied); + let gate = s.fence.gate(); + assert_eq!(gate.state(), FenceState::SourceWriteFenced); + let fenced = hello_config(&s).await; + assert_eq!( + fenced.fence, + Some(ReplicatedFence { + state: "SOURCE_WRITE_FENCED".into(), + revision: gate.revision(), + }) + ); + // Apart from the fence, `hello` sends the namespace's logical configuration unchanged: + // the legacy `block_*` mirror lives only in the stored config row (section 13.2). + let stored = s + .store + .with("ns".into(), |ns| { + metadata::DatabaseConfig::from(ns.config().as_ref()) + }) + .await + .unwrap(); + assert_eq!(stored.fence, None); + assert_eq!( + metadata::DatabaseConfig { + fence: None, + ..fenced + }, + stored + ); + + let released = s + .execute(FenceRequest { + namespace: "ns".into(), + operation_id: OP, + command_id: Uuid::from_u128(2), + expected_state: gate.state(), + expected_revision: gate.revision(), + command: FenceCommand::ReleaseSourceWriteFence, + }) + .await + .unwrap(); + assert_eq!(fence_outcome(&released), FenceOutcome::Applied); + assert_eq!(s.fence.gate().state(), FenceState::Released); + let after = hello_config(&s).await; + assert_eq!(after.fence, None); + assert!(!after.block_writes); + } + fn assert_read_fenced_status(status: &tonic::Status) { assert_eq!(status.code(), tonic::Code::FailedPrecondition, "{status:?}"); assert_eq!( diff --git a/libsql-server/src/rpc/proxy.rs b/libsql-server/src/rpc/proxy.rs index 1ab5d1d19a..b189b1143e 100644 --- a/libsql-server/src/rpc/proxy.rs +++ b/libsql-server/src/rpc/proxy.rs @@ -40,7 +40,14 @@ pub mod rpc { impl From for Error { fn from(other: SqldError) -> Self { + // A fence denial is an ordinary SQL error to an older replica, and carries its + // stable code in the additive `stable_code` field for one that maps it + // (`docs/NAMESPACE_FENCE.md` section 6.1). + let stable_code = other + .fence_error() + .and_then(|e| e.outcome().proxy_stable_code()); let code = match other { + _ if stable_code.is_some() => ErrorCode::SqlError, SqldError::LibSqlInvalidQueryParams(_) => ErrorCode::SqlError, SqldError::LibSqlTxTimeout => ErrorCode::TxTimeout, SqldError::LibSqlTxBusy => ErrorCode::TxBusy, @@ -57,6 +64,7 @@ pub mod rpc { message: other.to_string(), code: code as i32, extended_code, + stable_code: stable_code.map(Into::into), } } } @@ -64,6 +72,7 @@ pub mod rpc { impl From for ErrorCode { fn from(other: SqldError) -> Self { match other { + _ if other.fence_error().is_some() => ErrorCode::SqlError, SqldError::LibSqlInvalidQueryParams(_) => ErrorCode::SqlError, SqldError::LibSqlTxTimeout => ErrorCode::TxTimeout, SqldError::LibSqlTxBusy => ErrorCode::TxBusy, @@ -315,6 +324,9 @@ impl ProxyService { Ok(Ok(None)) => self.user_auth_strategy.clone(), Err(e) => match e.as_ref() { 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 @@ -567,6 +579,24 @@ 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. +fn fence_status(e: &crate::error::Error) -> Option { + e.fence_error()?.to_grpc_status() +} + +/// The status for an error looking up the namespace a proxy request names. +fn namespace_status(e: crate::error::Error) -> tonic::Status { + if let crate::error::Error::NamespaceDoesntExist(_) = e { + tonic::Status::failed_precondition(NAMESPACE_DOESNT_EXIST) + } else if let Some(status) = fence_status(&e) { + status + } else { + tonic::Status::internal(e.to_string()) + } +} + #[tonic::async_trait] impl Proxy for ProxyService { type StreamExecStream = Pin> + Send>>; @@ -585,18 +615,13 @@ impl Proxy for ProxyService { (connection_maker, notifier) }) .await - .map_err(|e| { - if let crate::error::Error::NamespaceDoesntExist(_) = e { - tonic::Status::failed_precondition(NAMESPACE_DOESNT_EXIST) - } else { - tonic::Status::internal(e.to_string()) - } - })?; + .map_err(namespace_status)?; - let conn = connection_maker - .create() - .await - .map_err(|e| tonic::Status::unavailable(format!("Unable to create DB: {:?}", e)))?; + let conn = connection_maker.create().await.map_err(|e| { + fence_status(&e).unwrap_or_else(|| { + tonic::Status::unavailable(format!("Unable to create DB: {:?}", e)) + }) + })?; let stream = make_proxy_stream(conn, ctx, req.into_inner()); @@ -617,13 +642,7 @@ impl Proxy for ProxyService { .namespaces .with(ctx.namespace().clone(), |ns| ns.db.connection_maker()) .await - .map_err(|e| { - if let crate::error::Error::NamespaceDoesntExist(_) = e { - tonic::Status::failed_precondition(NAMESPACE_DOESNT_EXIST) - } else { - tonic::Status::internal(e.to_string()) - } - })?; + .map_err(namespace_status)?; let conn = { let lock = self.clients.upgradable_read().await; @@ -645,7 +664,9 @@ impl Proxy for ProxyService { conn } Err(e) => { - return Err(tonic::Status::new(tonic::Code::Internal, e.to_string())) + return Err(fence_status(&e).unwrap_or_else(|| { + tonic::Status::new(tonic::Code::Internal, e.to_string()) + })) } } } @@ -659,7 +680,11 @@ impl Proxy for ProxyService { .execute_program(pgm, ctx, builder, None) .await // TODO: this is no necessarily a permission denied error! - .map_err(|e| tonic::Status::new(tonic::Code::PermissionDenied, e.to_string()))?; + .map_err(|e| { + fence_status(&e).unwrap_or_else(|| { + tonic::Status::new(tonic::Code::PermissionDenied, e.to_string()) + }) + })?; Ok(tonic::Response::new(builder.into_ret())) } @@ -691,13 +716,7 @@ impl Proxy for ProxyService { .namespaces .with(ctx.namespace().clone(), |ns| ns.db.connection_maker()) .await - .map_err(|e| { - if let crate::error::Error::NamespaceDoesntExist(_) = e { - tonic::Status::failed_precondition(NAMESPACE_DOESNT_EXIST) - } else { - tonic::Status::internal(e.to_string()) - } - })?; + .map_err(namespace_status)?; let DescribeRequest { client_id, stmt } = req.into_inner(); let client_id = Uuid::from_str(&client_id).unwrap(); @@ -719,7 +738,11 @@ impl Proxy for ProxyService { lock.insert(client_id, conn.clone()); conn } - Err(e) => return Err(tonic::Status::new(tonic::Code::Internal, e.to_string())), + Err(e) => { + return Err(fence_status(&e).unwrap_or_else(|| { + tonic::Status::new(tonic::Code::Internal, e.to_string()) + })) + } } } }; @@ -755,3 +778,289 @@ impl Proxy for ProxyService { })) } } + +/// Fence denials on the proxy protocol (`docs/NAMESPACE_FENCE.md` section 6.1): the primary +/// fills the additive `stable_code` of a step or program error and answers a namespace it +/// refuses with the typed `FAILED_PRECONDITION` status; the replica side turns both back into +/// the same fence denial, and an older primary's errors keep their untyped mapping. +#[cfg(test)] +mod fence_tests { + use libsql_replication::rpc::proxy::error::ErrorCode; + use libsql_replication::rpc::proxy::exec_resp; + use libsql_replication::rpc::proxy::{ + resp_step, ExecReq, ProgramResp, RespStep, StepError, StreamProgramReq, + }; + use libsql_replication::rpc::replication::NAMESPACE_METADATA_KEY; + use tempfile::tempdir; + use tokio_stream::wrappers::ReceiverStream; + use tokio_stream::StreamExt as _; + use tonic::metadata::BinaryMetadataValue; + + use super::*; + use crate::connection::program::Program; + use crate::error::Error; + use crate::namespace::fence::drain::tests::{fence_outcome, Source, LONG}; + use crate::namespace::fence::outcome::{FenceError, FenceOutcome, GRPC_FENCE_CODE_METADATA}; + use crate::namespace::fence::read::tests::{fenced_source, read_fence}; + use crate::namespace::fence::store as fence_store; + use crate::query_result_builder::QueryBuilderConfig; + + const WRITE: &str = "insert into t values (3)"; + const READ: &str = "select count(*) from t"; + + fn service(s: &NamespaceStore) -> ProxyService { + ProxyService::new(s.clone(), None, false) + } + + fn request(ns: &str, msg: T) -> tonic::Request { + let mut req = tonic::Request::new(msg); + req.metadata_mut().insert_bin( + NAMESPACE_METADATA_KEY, + BinaryMetadataValue::from_bytes(ns.as_bytes()), + ); + Authenticated::FullAccess.upgrade_grpc_request(&mut req); + req + } + + fn program_req(sql: &str) -> rpc::ProgramReq { + rpc::ProgramReq { + client_id: Uuid::new_v4().to_string(), + pgm: Some(Program::seq(&[sql]).into()), + } + } + + #[track_caller] + fn assert_fence_status(status: &tonic::Status, outcome: FenceOutcome) { + assert_eq!(status.code(), tonic::Code::FailedPrecondition, "{status:?}"); + assert_eq!( + status + .metadata() + .get(GRPC_FENCE_CODE_METADATA) + .and_then(|v| v.to_str().ok()), + Some(outcome.as_str()), + "{status:?}" + ); + assert_eq!( + FenceError::outcome_from_grpc_status(status), + Some(outcome), + "{status:?}" + ); + } + + #[track_caller] + fn assert_proxy_error(error: &rpc::Error, outcome: FenceOutcome) { + assert_eq!(error.code, ErrorCode::SqlError as i32, "{error:?}"); + assert_eq!(error.stable_code.as_deref(), Some(outcome.as_str())); + assert!(error.message.starts_with(outcome.as_str()), "{error:?}"); + } + + /// The step errors of a streamed program, run on `ns` through the primary's proxy stream. + async fn stream_program(s: &Source, sql: &str) -> Vec { + let conn = s + .store + .with("ns".into(), |ns| ns.db.connection_maker()) + .await + .unwrap() + .create() + .await + .unwrap(); + let ctx = RequestContext::new( + Authenticated::FullAccess, + "ns".into(), + s.store.meta_store().clone(), + ); + let (snd, rcv) = tokio::sync::mpsc::channel(1); + let stream = make_proxy_stream(conn, ctx, ReceiverStream::new(rcv)); + tokio::pin!(stream); + snd.send(Ok(ExecReq { + request_id: 0, + request: Some(libsql_replication::rpc::proxy::exec_req::Request::Execute( + StreamProgramReq { + pgm: Some(Program::seq(&[sql]).into()), + }, + )), + })) + .await + .unwrap(); + // The request stream stays open until the program is answered: a closed request stream + // ends the proxy stream. + let mut responses = Vec::new(); + while let Some(resp) = stream.next().await { + let resp = resp.unwrap().response.unwrap(); + let last = match &resp { + exec_resp::Response::ProgramResp(p) => p + .steps + .iter() + .any(|s| matches!(s.step, Some(resp_step::Step::Finish(_)))), + _ => true, + }; + responses.push(resp); + if last { + break; + } + } + drop(snd); + responses + } + + fn step_errors(responses: &[exec_resp::Response]) -> Vec<&rpc::Error> { + responses + .iter() + .filter_map(|r| match r { + exec_resp::Response::ProgramResp(p) => Some(p), + _ => None, + }) + .flat_map(|p| &p.steps) + .filter_map(|s| match &s.step { + Some(resp_step::Step::StepError(StepError { error: Some(e) })) => Some(e), + _ => None, + }) + .collect() + } + + /// Section 6.1 on the primary: a write refused at the WAL gate is a step error with the + /// stable code, a read refused by the read fence is a program error with it (streamed) or + /// the typed status (unary), and a namespace whose fence state is unknown is refused with + /// the typed status before any connection is made. + #[tokio::test] + async fn rpc_codes() { + let s = fenced_source().await; + + // Write-fenced: the write step carries MIGRATION_WRITE_FENCED, reads are served. + let streamed = stream_program(&s, WRITE).await; + let errors = step_errors(&streamed); + assert_eq!(errors.len(), 1, "{streamed:?}"); + assert_proxy_error(errors[0], FenceOutcome::MigrationWriteFenced); + assert!(step_errors(&stream_program(&s, READ).await).is_empty()); + + let unary = service(&s.store) + .execute(request("ns", program_req(WRITE))) + .await + .unwrap() + .into_inner(); + match &unary.results[..] { + [QueryResult { + row_result: Some(RowResult::Error(e)), + }] => assert_proxy_error(e, FenceOutcome::MigrationWriteFenced), + other => panic!("{other:?}"), + } + + // Read-fenced: the whole program is refused. + let fenced = s.execute(read_fence(&s, 2, LONG)).await.unwrap(); + assert_eq!(fence_outcome(&fenced), FenceOutcome::Applied); + match &stream_program(&s, READ).await[..] { + [exec_resp::Response::Error(e)] => { + assert_proxy_error(e, FenceOutcome::MigrationReadFenced) + } + other => panic!("{other:?}"), + } + let status = service(&s.store) + .execute(request("ns", program_req(READ))) + .await + .unwrap_err(); + assert_fence_status(&status, FenceOutcome::MigrationReadFenced); + + // A namespace whose fence state cannot be established. + let tmp = tempdir().unwrap(); + let broken = tmp.path().join("dbs").join("broken"); + std::fs::create_dir_all(&broken).unwrap(); + std::fs::write(broken.join(fence_store::MARKER_FILE_NAME), b"garbage").unwrap(); + let store = crate::namespace::open_test_store(tmp.path()).await; + let status = service(&store) + .execute(request("broken", program_req(READ))) + .await + .unwrap_err(); + assert_fence_status(&status, FenceOutcome::FenceStateUnavailable); + assert_eq!( + namespace_status(Error::NamespaceFence( + store + .with("broken".into(), |_| ()) + .await + .unwrap_err() + .fence_error() + .unwrap() + .clone() + )) + .code(), + tonic::Code::FailedPrecondition + ); + // A connection a fence refuses is never reported as `UNAVAILABLE`. + let e = Error::NamespaceFence(FenceError::new(FenceOutcome::MigrationWriteFenced, "x")); + assert_fence_status( + &fence_status(&e).unwrap(), + FenceOutcome::MigrationWriteFenced, + ); + assert!(fence_status(&Error::LibSqlTxBusy).is_none()); + } + + /// Section 6.1 on the replica: a proxied fence denial becomes the fence error the replica + /// answers its client with, whether it arrives as a step error, a program error or a + /// status; an older primary's errors (no `stable_code`) keep their untyped mapping. + #[tokio::test] + async fn replica_maps_proxied_denials() { + let denial = FenceError::new(FenceOutcome::MigrationWriteFenced, "writes are fenced"); + let proxied: rpc::Error = Error::NamespaceFence(denial.clone()).into(); + assert_eq!(proxied.code, ErrorCode::SqlError as i32); + match Error::from_proxy_error(proxied.clone()) { + Error::NamespaceFence(e) => assert_eq!(e, denial), + other => panic!("{other:?}"), + } + // An older primary: same error without the additive field. + let old = rpc::Error { + stable_code: None, + ..proxied.clone() + }; + assert!(matches!( + Error::from_proxy_error(old.clone()), + Error::RpcQueryError(e) if e == old + )); + // A code this server does not know stays untyped too. + let unknown = rpc::Error { + stable_code: Some("SOMETHING_NEW".into()), + ..proxied.clone() + }; + assert!(matches!( + Error::from_proxy_error(unknown), + Error::RpcQueryError(_) + )); + // Other errors are unchanged. + let busy: rpc::Error = Error::LibSqlTxBusy.into(); + assert_eq!(busy.stable_code, None); + assert_eq!(busy.code, ErrorCode::TxBusy as i32); + + // Status at connection time: typed denial, not retried; anything else unchanged. + match Error::from_proxy_status(denial.to_grpc_status().unwrap()) { + Error::NamespaceFence(e) => assert_eq!(e, denial), + other => panic!("{other:?}"), + } + assert!(matches!( + Error::from_proxy_status(tonic::Status::failed_precondition("x")), + Error::RpcQueryExecutionError(_) + )); + + // A step error applied to the replica's result builder. + let resp = ProgramResp { + steps: [ + resp_step::Step::Init(Default::default()), + resp_step::Step::BeginStep(Default::default()), + resp_step::Step::StepError(StepError { + error: Some(proxied), + }), + resp_step::Step::FinishStep(Default::default()), + resp_step::Step::Finish(Default::default()), + ] + .into_iter() + .map(|step| RespStep { step: Some(step) }) + .collect(), + }; + let mut builder = crate::http::user::result_builder::JsonHttpPayloadBuilder::new(); + crate::rpc::streaming_exec::apply_program_resp_to_builder( + &QueryBuilderConfig::default(), + &mut builder, + resp, + |_, _| (), + ) + .unwrap(); + assert_eq!(builder.take_fence_denial(), Some(denial)); + } +} diff --git a/libsql-server/src/rpc/replication/replication_log.rs b/libsql-server/src/rpc/replication/replication_log.rs index 4457c49e7a..7b6c5d51d7 100644 --- a/libsql-server/src/rpc/replication/replication_log.rs +++ b/libsql-server/src/rpc/replication/replication_log.rs @@ -8,6 +8,7 @@ use bytes::Bytes; use chrono::{DateTime, Utc}; use futures::stream::BoxStream; use futures_core::Future; +use libsql_replication::rpc::metadata; pub use libsql_replication::rpc::replication as rpc; use libsql_replication::rpc::replication::log_offset::WalFlavor; use libsql_replication::rpc::replication::replication_log_server::ReplicationLog; @@ -431,7 +432,7 @@ impl ReplicationLog for ReplicationLogService { guard.insert((replica_addr, namespace.clone())); } } - let (logger, config, version, _, _, _) = self + let (logger, config, version, _, _, fence) = self .logger_from_namespace(namespace, "hello", &req, false) .await?; @@ -451,7 +452,13 @@ impl ReplicationLog for ReplicationLogService { generation_id: self.generation_id.to_string(), generation_start_index: 0, current_replication_index: *logger.new_frame_notifier.borrow(), - config: Some(config.as_ref().into()), + config: Some(metadata::DatabaseConfig { + // The live fence, for a replica that applies it (section 6.2 of + // `docs/NAMESPACE_FENCE.md`); older replicas skip it and see the legacy + // `block_*` mirror. The stored configuration never carries it. + fence: fence.gate().replicated(), + ..config.as_ref().into() + }), }; Ok(tonic::Response::new(response)) diff --git a/libsql-server/src/rpc/streaming_exec.rs b/libsql-server/src/rpc/streaming_exec.rs index 482ce2074c..351fd7cb8d 100644 --- a/libsql-server/src/rpc/streaming_exec.rs +++ b/libsql-server/src/rpc/streaming_exec.rs @@ -225,7 +225,7 @@ pub fn apply_program_resp_to_builder( last_insert_rowid, }) => builder.finish_step(affected_row_count, last_insert_rowid)?, Step::StepError(StepError { error: Some(err) }) => { - builder.step_error(crate::error::Error::RpcQueryError(err))? + builder.step_error(crate::error::Error::from_proxy_error(err))? } Step::ColsDescription(ColsDescription { columns }) => { let cols = columns.iter().map(|c| Column { diff --git a/libsql-server/src/schema/error.rs b/libsql-server/src/schema/error.rs index 528251a2b3..b54cc7e87d 100644 --- a/libsql-server/src/schema/error.rs +++ b/libsql-server/src/schema/error.rs @@ -60,7 +60,7 @@ impl IntoResponse for &Error { self.format_err(StatusCode::BAD_REQUEST) } Error::MigrationExecuteError(e) => e.as_ref().into_response(), - Error::NamespaceFence(e) => self.format_err(e.outcome().admin_http_status()), + Error::NamespaceFence(e) => crate::error::fence_error_response(e), _ => self.format_err(StatusCode::INTERNAL_SERVER_ERROR), } } diff --git a/libsql-server/tests/fence/admin.rs b/libsql-server/tests/fence/admin.rs index 2c64935c46..5a5a2aec03 100644 --- a/libsql-server/tests/fence/admin.rs +++ b/libsql-server/tests/fence/admin.rs @@ -26,7 +26,7 @@ fn capabilities() { assert_eq!(body["fence_protocol_version"], 1); assert_eq!(body["enabled"], true); assert_eq!(body["active_fences"], 0); - assert_eq!(body["proxy_stable_code"], false); + assert_eq!(body["proxy_stable_code"], true); let commands: Vec<&str> = body["commands"] .as_array() .unwrap() diff --git a/libsql-server/tests/fence/lifecycle.rs b/libsql-server/tests/fence/lifecycle.rs index 90ac6d9749..2ec5f77e54 100644 --- a/libsql-server/tests/fence/lifecycle.rs +++ b/libsql-server/tests/fence/lifecycle.rs @@ -145,6 +145,7 @@ fn lifecycle_rejected_while_fenced() { message.starts_with(code), "{what} on {ns}: expected {code}, got {body}" ); + assert_eq!(body["code"], code, "{what} on {ns}: {body}"); } // Nothing moved: same state and revision, and no copy was created. let (status, body) = admin.inspect(ns).await?; diff --git a/libsql-server/tests/fence/mod.rs b/libsql-server/tests/fence/mod.rs index c8be3cde94..2e9cbae183 100644 --- a/libsql-server/tests/fence/mod.rs +++ b/libsql-server/tests/fence/mod.rs @@ -4,12 +4,17 @@ mod admin; mod lifecycle; +mod protocol; use std::path::PathBuf; use std::time::Duration; use hyper::StatusCode; -use libsql_server::config::{AdminApiConfig, MetaStoreConfig, RpcServerConfig, UserApiConfig}; +use libsql_server::auth::user_auth_strategies::http_basic::HttpBasic; +use libsql_server::auth::Auth; +use libsql_server::config::{ + AdminApiConfig, MetaStoreConfig, RpcClientConfig, RpcServerConfig, UserApiConfig, +}; use s3s::header::AUTHORIZATION; use serde_json::{json, Value}; use turmoil::{Builder, Sim}; @@ -26,6 +31,8 @@ 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>, } impl Default for Primary { @@ -33,6 +40,7 @@ impl Default for Primary { Self { admin_key: Some(ADMIN_KEY), fence_enabled: true, + user_credential: None, } } } @@ -49,13 +57,20 @@ pub fn make_primary(sim: &mut Sim, path: PathBuf, primary: Primary) { let Primary { admin_key, fence_enabled, + user_credential, } = primary; sim.host("primary", move || { let path = path.clone(); async move { let server = TestServer { path: path.into(), - user_api_config: UserApiConfig::default(), + 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, @@ -80,6 +95,37 @@ pub fn make_primary(sim: &mut Sim, path: PathBuf, primary: Primary) { }); } +/// A replica of `primary` on host `replica0`: user API on 8080, admin API (no auth key) on +/// 9090. It creates a namespace lazily, on first use, by replicating it from the primary. +pub fn make_replica(sim: &mut Sim, path: PathBuf) { + init_tracing(); + sim.host("replica0", move || { + let path = path.clone(); + async move { + let server = TestServer { + path: path.into(), + user_api_config: UserApiConfig::default(), + admin_api_config: Some(AdminApiConfig { + acceptor: TurmoilAcceptor::bind(([0, 0, 0, 0], 9090)).await?, + connector: TurmoilConnector, + disable_metrics: true, + auth_key: None, + }), + rpc_client_config: Some(RpcClientConfig { + remote_url: "http://primary:4567".into(), + connector: TurmoilConnector, + tls_config: None, + }), + disable_namespaces: false, + disable_default_namespace: true, + ..Default::default() + }; + server.start_sim(8080).await?; + Ok(()) + } + }); +} + /// The admin API of `primary`, authenticating with `key` when there is one. pub struct Admin { client: Client, diff --git a/libsql-server/tests/fence/protocol.rs b/libsql-server/tests/fence/protocol.rs new file mode 100644 index 0000000000..7cf970fda3 --- /dev/null +++ b/libsql-server/tests/fence/protocol.rs @@ -0,0 +1,976 @@ +//! Typed fence outcomes on the user-facing protocols: the legacy HTTP API, Hrana over HTTP +//! (`/v1`, `/v2`, `/v3`, cursors), Hrana over WebSocket and `/dump` +//! (`docs/NAMESPACE_FENCE.md` section 6; section 17 rows 3 and 18). + +use futures::SinkExt as _; +use hyper::{Body, Method, Request, StatusCode}; +use serde_json::{json, Value}; +use tempfile::tempdir; +use tokio_stream::StreamExt as _; +use tokio_tungstenite::tungstenite::{self, client::IntoClientRequest}; +use turmoil::net::TcpStream; +use uuid::Uuid; + +use super::{ + acquire_body, command_body, load_and_log_id, make_primary, make_replica, sim, state_of, Admin, + Primary, ADMIN_KEY, +}; +use crate::common::net::TurmoilConnector; + +const WRITE_FENCED: &str = "MIGRATION_WRITE_FENCED"; +const READ_FENCED: &str = "MIGRATION_READ_FENCED"; +const QUARANTINED: &str = "MIGRATION_TARGET_QUARANTINED"; +const UNAVAILABLE: &str = "FENCE_STATE_UNAVAILABLE"; + +fn uuid(n: u128) -> Uuid { + Uuid::from_u128(n) +} + +/// The user API of `primary` (or of another host), for any namespace, with an optional +/// basic-auth credential. +struct User { + client: hyper::Client, + auth: Option, + host: &'static str, +} + +impl User { + fn new() -> Self { + Self::with_auth(None) + } + + fn with_auth(credential: Option<&str>) -> Self { + Self { + client: hyper::Client::builder().build(TurmoilConnector), + auth: credential.map(|c| format!("basic {c}")), + host: "primary", + } + } + + /// The user API of `host` instead. + fn on(host: &'static str) -> Self { + Self { + host, + ..Self::new() + } + } + + async fn request( + &self, + method: Method, + ns: &str, + path: &str, + body: Option, + ) -> anyhow::Result<(StatusCode, String)> { + let mut request = Request::builder() + .method(method) + .uri(format!("http://{ns}.{}:8080{path}", self.host)); + if let Some(auth) = &self.auth { + request = request.header("authorization", auth.as_str()); + } + let request = match body { + Some(body) => request + .header("content-type", "application/json") + .body(Body::from(serde_json::to_vec(&body)?))?, + None => request.body(Body::empty())?, + }; + let response = self.client.request(request).await?; + let status = response.status(); + let body = hyper::body::to_bytes(response.into_body()).await?; + Ok((status, String::from_utf8_lossy(&body).into_owned())) + } + + async fn post(&self, ns: &str, path: &str, body: Value) -> anyhow::Result<(StatusCode, Value)> { + let (status, body) = self.request(Method::POST, ns, path, Some(body)).await?; + let value = serde_json::from_str(&body).unwrap_or(Value::String(body)); + Ok((status, value)) + } + + /// Legacy API: one batch of statements. + async fn legacy(&self, ns: &str, sql: &[&str]) -> anyhow::Result<(StatusCode, Value)> { + self.post(ns, "/", json!({ "statements": sql })).await + } + + /// Hrana 1 over HTTP: one statement. + async fn execute(&self, ns: &str, sql: &str) -> anyhow::Result<(StatusCode, Value)> { + self.post(ns, "/v1/execute", json!({ "stmt": { "sql": sql } })) + .await + } + + /// Hrana 1 over HTTP: one batch, each statement conditional on nothing. + async fn batch(&self, ns: &str, sql: &[&str]) -> anyhow::Result<(StatusCode, Value)> { + self.post(ns, "/v1/batch", json!({ "batch": batch(sql) })) + .await + } + + /// Hrana 2/3 over HTTP: one pipeline. + async fn pipeline( + &self, + ns: &str, + version: u8, + baton: Option<&str>, + requests: Value, + ) -> anyhow::Result<(StatusCode, Value)> { + self.post( + ns, + &format!("/v{version}/pipeline"), + json!({ "baton": baton, "requests": requests }), + ) + .await + } + + /// Hrana 3 over HTTP: a cursor over one batch, as the list of entries it returned. + async fn cursor(&self, ns: &str, sql: &[&str]) -> anyhow::Result<(StatusCode, Vec)> { + let (status, body) = self + .request( + Method::POST, + ns, + "/v3/cursor", + Some(json!({ "baton": null, "batch": batch(sql) })), + ) + .await?; + let entries = body + .lines() + .filter(|line| !line.trim().is_empty()) + .map(|line| serde_json::from_str(line).unwrap_or(Value::String(line.into()))) + .collect(); + Ok((status, entries)) + } + + async fn dump(&self, ns: &str) -> anyhow::Result<(StatusCode, String)> { + self.request(Method::GET, ns, "/dump", None).await + } +} + +/// A Hrana batch of unconditional steps. +fn batch(sql: &[&str]) -> Value { + json!({ "steps": sql.iter().map(|sql| json!({ "stmt": { "sql": sql } })).collect::>() }) +} + +fn execute_req(sql: &str) -> Value { + json!({ "type": "execute", "stmt": { "sql": sql } }) +} + +fn batch_req(sql: &[&str]) -> Value { + json!({ "type": "batch", "batch": batch(sql) }) +} + +/// A user-HTTP refusal by the fence: `423`, and the stable code in the additive `code` field. +#[track_caller] +fn assert_locked(what: &str, (status, body): &(StatusCode, Value), code: &str) { + assert_eq!(*status, StatusCode::LOCKED, "{what}: {body}"); + assert_eq!(body["code"], code, "{what}: {body}"); +} + +/// A Hrana error object carrying `code`. +#[track_caller] +fn assert_hrana_error(what: &str, error: &Value, code: &str) { + assert_eq!(error["code"], code, "{what}: {error}"); + assert!( + error["message"].as_str().unwrap_or_default().contains(code), + "{what}: {error}" + ); +} + +/// `ns` created, loaded with table `t` holding one row, and write-fenced by `op`. Returns the +/// fence's revision. +async fn write_fenced(admin: &Admin, ns: &str, op: Uuid) -> anyhow::Result { + admin.create_namespace(ns).await?; + let log_id = load_and_log_id(admin, ns).await?; + let (status, body) = admin + .command( + ns, + "source/acquire-write-fence", + acquire_body(op, uuid(op.as_u128() + 1), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body).0, "SOURCE_WRITE_FENCED", "{body}"); + Ok(state_of(&body).1) +} + +/// Runs the source command `route` from `state` at `revision` and returns the new revision. +async fn source_command( + admin: &Admin, + ns: &str, + op: Uuid, + command: u128, + route: &str, + (state, revision): (&str, u64), + expect: &str, +) -> anyhow::Result { + let (status, body) = admin + .command( + ns, + route, + command_body( + op, + uuid(command), + state, + revision, + json!({ "drain_policy": { "deadline_ms": 5000 } }), + ), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{route}: {body}"); + assert_eq!(state_of(&body).0, expect, "{route}: {body}"); + Ok(state_of(&body).1) +} + +async fn read_fence(admin: &Admin, ns: &str, op: Uuid, revision: u64) -> anyhow::Result { + source_command( + admin, + ns, + op, + op.as_u128() + 2, + "source/set-read-fence", + ("SOURCE_WRITE_FENCED", revision), + "SOURCE_READ_FENCED", + ) + .await +} + +async fn clear_read_fence(admin: &Admin, ns: &str, op: Uuid, revision: u64) -> anyhow::Result { + let (status, body) = admin + .command( + ns, + "source/clear-read-fence", + command_body( + op, + uuid(op.as_u128() + 3), + "SOURCE_READ_FENCED", + revision, + json!({}), + ), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body).0, "SOURCE_WRITE_FENCED", "{body}"); + Ok(state_of(&body).1) +} + +async fn release(admin: &Admin, ns: &str, op: Uuid, revision: u64) -> anyhow::Result<()> { + let (status, body) = admin + .command( + ns, + "source/release-write-fence", + command_body( + op, + uuid(op.as_u128() + 4), + "SOURCE_WRITE_FENCED", + revision, + json!({}), + ), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body).0, "RELEASED", "{body}"); + Ok(()) +} + +async fn quarantined_target(admin: &Admin, ns: &str, op: Uuid) -> anyhow::Result<()> { + let (status, body) = admin + .command( + ns, + "target/create-quarantined", + command_body(op, uuid(op.as_u128() + 1), "ABSENT", 0, json!({})), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(state_of(&body).0, "TARGET_QUARANTINED", "{body}"); + Ok(()) +} + +/// Every user HTTP entry point answers a fence denial with `423` and the stable code, whether +/// the fence refused one statement (a write) or the whole request (a read, a quarantined +/// target, a fence state the server cannot establish). +#[test] +fn http_codes() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + // A namespace directory whose fence marker cannot be decoded: its fence state cannot be + // established, so the server refuses it (section 13.3). + let broken = tmp.path().join("dbs").join("broken"); + std::fs::create_dir_all(&broken).unwrap(); + std::fs::write(broken.join(".fence"), b"garbage").unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let user = User::new(); + let op = uuid(0x100); + let rev = write_fenced(&admin, "src", op).await?; + + // Write-fenced: writes are refused, reads are served. + assert_locked( + "legacy write", + &user.legacy("src", &["insert into t values (2)"]).await?, + WRITE_FENCED, + ); + assert_locked( + "legacy read then write", + &user + .legacy("src", &["select * from t", "insert into t values (2)"]) + .await?, + WRITE_FENCED, + ); + let (status, body) = user.legacy("src", &["select * from t"]).await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_locked( + "v1 execute write", + &user.execute("src", "insert into t values (2)").await?, + WRITE_FENCED, + ); + let (status, body) = user.execute("src", "select * from t").await?; + assert_eq!(status, StatusCode::OK, "{body}"); + // A Hrana batch reports the refused step in its own error, as for any step error. + let (status, body) = user + .batch("src", &["select * from t", "insert into t values (2)"]) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert!(body["result"]["step_results"][0].is_object(), "{body}"); + assert_hrana_error( + "v1 batch write step", + &body["result"]["step_errors"][1], + WRITE_FENCED, + ); + + // Read-fenced: reads are refused as a whole. + read_fence(&admin, "src", op, rev).await?; + assert_locked( + "legacy read", + &user.legacy("src", &["select * from t"]).await?, + READ_FENCED, + ); + assert_locked( + "v1 execute read", + &user.execute("src", "select * from t").await?, + READ_FENCED, + ); + assert_locked( + "v1 batch read", + &user.batch("src", &["select * from t"]).await?, + READ_FENCED, + ); + + // A quarantined target serves nothing. + quarantined_target(&admin, "tgt", uuid(0x200)).await?; + assert_locked( + "legacy on target", + &user.legacy("tgt", &["select 1"]).await?, + QUARANTINED, + ); + assert_locked( + "v1 execute on target", + &user.execute("tgt", "select 1").await?, + QUARANTINED, + ); + assert_locked( + "v1 batch on target", + &user.batch("tgt", &["select 1"]).await?, + QUARANTINED, + ); + + // A namespace whose fence state is unknown is refused before a connection exists. + for (what, response) in [ + ("legacy", user.legacy("broken", &["select 1"]).await?), + ("v1 execute", user.execute("broken", "select 1").await?), + ( + "v2 pipeline", + user.pipeline("broken", 2, None, json!([execute_req("select 1")])) + .await?, + ), + ] { + assert_locked(what, &response, UNAVAILABLE); + assert_eq!(response.1["detail"], "corrupt_record", "{what}"); + } + Ok(()) + }); + sim.run().unwrap(); +} + +/// Hrana over HTTP (`/v2`, `/v3`, `/v3/cursor`): a fence denial is a Hrana error with the +/// stable code, and the stream stays usable. +#[test] +fn hrana_http_codes() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let user = User::new(); + let op = uuid(0x100); + let rev = write_fenced(&admin, "src", op).await?; + + for version in [2, 3] { + let what = format!("v{version}"); + let (status, body) = user + .pipeline( + "src", + version, + None, + json!([ + execute_req("insert into t values (2)"), + execute_req("select * from t"), + batch_req(&["select * from t", "insert into t values (2)"]), + ]), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{what}: {body}"); + let results = &body["results"]; + assert_eq!(results[0]["type"], "error", "{what}: {body}"); + assert_hrana_error(&what, &results[0]["error"], WRITE_FENCED); + assert_eq!(results[1]["type"], "ok", "{what}: {body}"); + assert_eq!(results[2]["type"], "ok", "{what}: {body}"); + assert_hrana_error( + &what, + &results[2]["response"]["result"]["step_errors"][1], + WRITE_FENCED, + ); + // The stream survives the denial. + let baton = body["baton"].as_str().expect("the stream stays open"); + let (status, body) = user + .pipeline( + "src", + version, + Some(baton), + json!([execute_req("select * from t"), { "type": "close" }]), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{what}: {body}"); + assert_eq!(body["results"][0]["type"], "ok", "{what}: {body}"); + } + let (status, entries) = user + .cursor("src", &["select * from t", "insert into t values (2)"]) + .await?; + assert_eq!(status, StatusCode::OK, "{entries:?}"); + let step_error = entries + .iter() + .find(|e| e["type"] == "step_error") + .unwrap_or_else(|| panic!("no step error: {entries:?}")); + assert_eq!(step_error["step"], 1, "{entries:?}"); + assert_hrana_error("cursor write step", &step_error["error"], WRITE_FENCED); + + read_fence(&admin, "src", op, rev).await?; + for version in [2, 3] { + let what = format!("v{version} read-fenced"); + let (status, body) = user + .pipeline( + "src", + version, + None, + json!([ + execute_req("select * from t"), + batch_req(&["select * from t"]), + { "type": "close" }, + ]), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{what}: {body}"); + let results = &body["results"]; + assert_hrana_error(&what, &results[0]["error"], READ_FENCED); + assert_hrana_error(&what, &results[1]["error"], READ_FENCED); + assert_eq!(results[2]["type"], "ok", "{what}: {body}"); + } + let (status, entries) = user.cursor("src", &["select * from t"]).await?; + assert_eq!(status, StatusCode::OK, "{entries:?}"); + let error = entries + .iter() + .find(|e| e["type"] == "error") + .unwrap_or_else(|| panic!("no error entry: {entries:?}")); + assert_hrana_error("cursor read", &error["error"], READ_FENCED); + + quarantined_target(&admin, "tgt", uuid(0x200)).await?; + let (status, body) = user + .pipeline("tgt", 3, None, json!([execute_req("select 1")])) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_hrana_error("target", &body["results"][0]["error"], QUARANTINED); + Ok(()) + }); + sim.run().unwrap(); +} + +/// Acceptance test (section 17 row 3): in a batch, the steps before a write run, the write +/// step is refused with the fence code, and conditions see the refusal as a failed step. +#[test] +fn batch_denied_mid_batch() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let user = User::new(); + write_fenced(&admin, "src", uuid(0x100)).await?; + + let batch = json!({ + "steps": [ + { "stmt": { "sql": "select count(*) from t" } }, + { "stmt": { "sql": "insert into t values (2)" } }, + { + "condition": { "type": "ok", "step": 1 }, + "stmt": { "sql": "insert into t values (3)" }, + }, + { + "condition": { "type": "not", "cond": { "type": "ok", "step": 1 } }, + "stmt": { "sql": "select count(*) from t" }, + }, + ], + }); + let (status, body) = user + .pipeline( + "src", + 3, + None, + json!([{ "type": "batch", "batch": batch }, { "type": "close" }]), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let result = &body["results"][0]["response"]["result"]; + let count = |step: usize| result["step_results"][step]["rows"][0][0]["value"].clone(); + assert_eq!(count(0), "1", "{body}"); + assert!(result["step_errors"][0].is_null(), "{body}"); + assert!(result["step_results"][1].is_null(), "{body}"); + assert_hrana_error("write step", &result["step_errors"][1], WRITE_FENCED); + // The step conditional on the write did not run; the one conditional on its failure did. + assert!(result["step_results"][2].is_null(), "{body}"); + assert!(result["step_errors"][2].is_null(), "{body}"); + assert_eq!(count(3), "1", "{body}"); + Ok(()) + }); + sim.run().unwrap(); +} + +type Ws = tokio_tungstenite::WebSocketStream; + +/// A Hrana 3 WebSocket session to `ns`, after `hello`. +async fn ws_connect(ns: &str) -> anyhow::Result { + let mut request = format!("ws://{ns}.primary:8080").into_client_request()?; + request + .headers_mut() + .insert("sec-websocket-protocol", "hrana3".parse()?); + let conn = TcpStream::connect("primary:8080").await?; + let (mut ws, _) = tokio_tungstenite::client_async(request, conn).await?; + ws.send(tungstenite::Message::Text( + json!({ "type": "hello", "jwt": null }).to_string(), + )) + .await?; + let hello = ws_next(&mut ws).await?; + assert_eq!(hello["type"], "hello_ok", "{hello}"); + Ok(ws) +} + +async fn ws_next(ws: &mut Ws) -> anyhow::Result { + match ws.try_next().await? { + Some(tungstenite::Message::Text(text)) => Ok(serde_json::from_str(&text)?), + other => anyhow::bail!("unexpected WebSocket message {other:?}"), + } +} + +/// Sends one request and returns the server's answer to it. +async fn ws_request(ws: &mut Ws, request_id: u64, request: Value) -> anyhow::Result { + ws.send(tungstenite::Message::Text( + json!({ "type": "request", "request_id": request_id, "request": request }).to_string(), + )) + .await?; + let response = ws_next(ws).await?; + assert_eq!(response["request_id"], request_id, "{response}"); + Ok(response) +} + +async fn ws_execute(ws: &mut Ws, request_id: u64, sql: &str) -> anyhow::Result { + ws_request( + ws, + request_id, + json!({ "type": "execute", "stream_id": 1, "stmt": { "sql": sql } }), + ) + .await +} + +#[track_caller] +fn assert_ws_ok(what: &str, response: &Value) { + assert_eq!(response["type"], "response_ok", "{what}: {response}"); +} + +#[track_caller] +fn assert_ws_error(what: &str, response: &Value, code: &str) { + assert_eq!(response["type"], "response_error", "{what}: {response}"); + assert_hrana_error(what, &response["error"], code); +} + +/// Hrana over WebSocket: fence denials are request errors with the stable code; the connection +/// and its stream stay usable. +#[test] +fn hrana_ws_codes() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let op = uuid(0x100); + let rev = write_fenced(&admin, "src", op).await?; + + let mut ws = ws_connect("src").await?; + let opened = + ws_request(&mut ws, 1, json!({ "type": "open_stream", "stream_id": 1 })).await?; + assert_ws_ok("open_stream", &opened); + let denied = ws_execute(&mut ws, 2, "insert into t values (2)").await?; + assert_ws_error("write", &denied, WRITE_FENCED); + assert_ws_ok("read", &ws_execute(&mut ws, 3, "select * from t").await?); + let batched = ws_request( + &mut ws, + 4, + json!({ + "type": "batch", + "stream_id": 1, + "batch": batch(&["select * from t", "insert into t values (2)"]), + }), + ) + .await?; + assert_ws_ok("batch", &batched); + assert_hrana_error( + "batch write step", + &batched["response"]["result"]["step_errors"][1], + WRITE_FENCED, + ); + + let rev = read_fence(&admin, "src", op, rev).await?; + let denied = ws_execute(&mut ws, 5, "select * from t").await?; + assert_ws_error("read-fenced read", &denied, READ_FENCED); + let denied = ws_request( + &mut ws, + 6, + json!({ "type": "batch", "stream_id": 1, "batch": batch(&["select * from t"]) }), + ) + .await?; + assert_ws_error("read-fenced batch", &denied, READ_FENCED); + + // Clearing the read fence makes the same stream readable again. + clear_read_fence(&admin, "src", op, rev).await?; + assert_ws_ok( + "read after clear", + &ws_execute(&mut ws, 7, "select * from t").await?, + ); + Ok(()) + }); + sim.run().unwrap(); +} + +/// Acceptance test (section 17 row 3): a WebSocket session whose transaction began before the +/// fence cannot write while the fence holds, nor after it is released; only a transaction +/// begun after the release writes. +#[test] +fn old_ws_session_cannot_write() { + 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 mut ws = ws_connect("src").await?; + let opened = + ws_request(&mut ws, 1, json!({ "type": "open_stream", "stream_id": 1 })).await?; + assert_ws_ok("open_stream", &opened); + assert_ws_ok("begin", &ws_execute(&mut ws, 2, "begin").await?); + assert_ws_ok("read", &ws_execute(&mut ws, 3, "select * from t").await?); + + let op = uuid(0x100); + let (status, body) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(op, uuid(0x101), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let rev = state_of(&body).1; + + let denied = ws_execute(&mut ws, 4, "insert into t values (2)").await?; + assert_ws_error("write while fenced", &denied, WRITE_FENCED); + + release(&admin, "src", op, rev).await?; + // The transaction still began under the earlier write generation. + let denied = ws_execute(&mut ws, 5, "insert into t values (3)").await?; + assert_ws_error("write after release", &denied, WRITE_FENCED); + assert_ws_ok("rollback", &ws_execute(&mut ws, 6, "rollback").await?); + // A fresh transaction on the same session writes. + assert_ws_ok( + "fresh write", + &ws_execute(&mut ws, 7, "insert into t values (4)").await?, + ); + let rows = ws_execute(&mut ws, 8, "select x from t order by x").await?; + assert_ws_ok("rows", &rows); + let values: Vec = rows["response"]["result"]["rows"] + .as_array() + .unwrap() + .iter() + .map(|row| row[0]["value"].clone()) + .collect(); + assert_eq!(values, vec![json!("1"), json!("4")], "{rows}"); + Ok(()) + }); + sim.run().unwrap(); +} + +/// `/dump`: served in full under a write fence; refused with `423` and the code under a read +/// fence and on a quarantined target. (An export cancelled mid-way is an aborted response +/// body: `namespace::fence::stream::tests::dump_response_aborted_on_cancel`.) +#[test] +fn dump_codes() { + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary(&mut sim, tmp.path().to_path_buf(), Primary::default()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let user = User::new(); + let op = uuid(0x100); + let rev = write_fenced(&admin, "src", op).await?; + + let (status, dump) = user.dump("src").await?; + assert_eq!(status, StatusCode::OK, "{dump}"); + assert!(dump.contains("INSERT INTO t"), "{dump}"); + assert!(dump.trim_end().ends_with("COMMIT;"), "{dump}"); + + read_fence(&admin, "src", op, rev).await?; + let (status, body) = user.dump("src").await?; + assert_locked( + "read-fenced dump", + &(status, serde_json::from_str(&body)?), + READ_FENCED, + ); + + quarantined_target(&admin, "tgt", uuid(0x200)).await?; + let (status, body) = user.dump("tgt").await?; + assert_locked( + "target dump", + &(status, serde_json::from_str(&body)?), + QUARANTINED, + ); + Ok(()) + }); + sim.run().unwrap(); +} + +/// Section 17 row 18: authentication failures and missing namespaces keep their own statuses +/// and carry no fence code, so a client can tell them from a fence denial. +#[test] +fn auth_and_not_found_distinct() { + const CREDENTIAL: &str = "dXNlcjpwYXNz"; + let mut sim = sim(); + let tmp = tempdir().unwrap(); + make_primary( + &mut sim, + tmp.path().to_path_buf(), + Primary { + user_credential: Some(CREDENTIAL), + ..Default::default() + }, + ); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let user = User::with_auth(Some(CREDENTIAL)); + admin.create_namespace("src").await?; + for sql in ["create table t (x)", "insert into t values (1)"] { + let (status, body) = user.execute("src", sql).await?; + assert_eq!(status, StatusCode::OK, "{body}"); + } + let (_, body) = admin.inspect("src").await?; + let log_id = body["fence"]["incarnation"]["current_log_id"] + .as_str() + .unwrap() + .to_string(); + let (status, body) = admin + .command( + "src", + "source/acquire-write-fence", + acquire_body(uuid(0x100), uuid(0x101), &log_id), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + + let write = "insert into t values (2)"; + let anonymous = User::new(); + let wrong = User::with_auth(Some("d3Jvbmc6d3Jvbmc=")); + for (what, response, expected) in [ + ( + "no credential", + anonymous.execute("src", write).await?, + StatusCode::UNAUTHORIZED, + ), + ( + "wrong credential", + wrong.execute("src", write).await?, + StatusCode::UNAUTHORIZED, + ), + ( + "no credential, legacy", + anonymous.legacy("src", &[write]).await?, + StatusCode::UNAUTHORIZED, + ), + ( + "no credential, v3", + anonymous + .pipeline("src", 3, None, json!([execute_req(write)])) + .await?, + StatusCode::UNAUTHORIZED, + ), + ( + "missing namespace", + user.execute("nope", write).await?, + StatusCode::NOT_FOUND, + ), + ( + "missing namespace, v3", + user.pipeline("nope", 3, None, json!([execute_req(write)])) + .await?, + StatusCode::NOT_FOUND, + ), + ] { + let (status, body) = &response; + assert_eq!(*status, expected, "{what}: {body}"); + assert!( + body.get("code").map_or(true, |code| !code + .as_str() + .unwrap_or_default() + .starts_with("MIGRATION_") + && code != UNAVAILABLE), + "{what}: {body}" + ); + } + // With the credential, on the existing namespace, it is the fence that answers. + assert_locked( + "fenced write", + &user.execute("src", write).await?, + WRITE_FENCED, + ); + assert_locked( + "fenced legacy write", + &user.legacy("src", &[write]).await?, + WRITE_FENCED, + ); + let (status, body) = user + .pipeline("src", 3, None, json!([execute_req(write)])) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_hrana_error( + "fenced v3 write", + &body["results"][0]["error"], + WRITE_FENCED, + ); + Ok(()) + }); + sim.run().unwrap(); +} + +/// The replica's count of writes it delegated to the primary for `ns`. +async fn delegated_writes(ns: &str) -> anyhow::Result { + let resp = crate::common::http::Client::new() + .get(&format!("http://replica0:9090/v1/namespaces/{ns}/stats")) + .await?; + let body: Value = resp.json().await?; + body["write_requests_delegated"] + .as_u64() + .ok_or_else(|| anyhow::anyhow!("no write_requests_delegated: {body}")) +} + +/// A write sent to a replica is proxied to the primary; when the primary's fence refuses it, +/// the replica answers with the primary's denial: `423` and the stable code on HTTP, the +/// stable code as the Hrana error code (section 6.1). Reads stay local and are served, and +/// writes through the replica work again once the fence is released. +#[test] +fn replica_proxy_preserves_code() { + let mut sim = sim(); + let primary = tempdir().unwrap(); + let replica = tempdir().unwrap(); + make_primary(&mut sim, primary.path().to_path_buf(), Primary::default()); + make_replica(&mut sim, replica.path().to_path_buf()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let user = User::on("replica0"); + let op = uuid(0x100); + let rev = write_fenced(&admin, "src", op).await?; + let write = "insert into t values (2)"; + + assert_locked( + "legacy write", + &user.legacy("src", &[write]).await?, + WRITE_FENCED, + ); + assert_locked( + "v1 execute", + &user.execute("src", write).await?, + WRITE_FENCED, + ); + // A refused step of a `/v1` batch is a step error, as on the primary. + let (status, body) = user.batch("src", &["select 1", write]).await?; + assert_eq!(status, StatusCode::OK, "v1 batch: {body}"); + assert!( + body["result"]["step_errors"][0].is_null(), + "v1 batch: {body}" + ); + assert_hrana_error("v1 batch", &body["result"]["step_errors"][1], WRITE_FENCED); + for version in [2, 3] { + let what = format!("v{version}"); + let (status, body) = user + .pipeline( + "src", + version, + None, + json!([execute_req(write), execute_req("select * from t")]), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{what}: {body}"); + let results = &body["results"]; + assert_eq!(results[0]["type"], "error", "{what}: {body}"); + assert_hrana_error(&what, &results[0]["error"], WRITE_FENCED); + assert_eq!(results[1]["type"], "ok", "{what}: {body}"); + } + let (status, body) = user.legacy("src", &["select count(*) from t"]).await?; + assert_eq!(status, StatusCode::OK, "{body}"); + + release(&admin, "src", op, rev).await?; + let (status, body) = user.execute("src", write).await?; + assert_eq!(status, StatusCode::OK, "after release: {body}"); + Ok(()) + }); + sim.run().unwrap(); +} + +/// A fence denial is final for the request: the replica delegates each refused write to the +/// primary exactly once and answers at once, with no reconnect or retry (the write proxy's +/// only retry, of `UNAVAILABLE`, backs off 500 ms first). +#[test] +fn denial_not_retried() { + let mut sim = sim(); + let primary = tempdir().unwrap(); + let replica = tempdir().unwrap(); + make_primary(&mut sim, primary.path().to_path_buf(), Primary::default()); + make_replica(&mut sim, replica.path().to_path_buf()); + sim.client("client", async { + let admin = Admin::new(Some(ADMIN_KEY)); + let user = User::on("replica0"); + write_fenced(&admin, "src", uuid(0x100)).await?; + // Load the namespace on the replica. + let (status, body) = user.legacy("src", &["select 1"]).await?; + assert_eq!(status, StatusCode::OK, "{body}"); + + let before = delegated_writes("src").await?; + for i in 0..3u64 { + let started = tokio::time::Instant::now(); + assert_locked( + "write", + &user.execute("src", "insert into t values (2)").await?, + WRITE_FENCED, + ); + assert!( + started.elapsed() < std::time::Duration::from_millis(500), + "{:?}", + started.elapsed() + ); + assert_eq!(delegated_writes("src").await?, before + i + 1); + } + Ok(()) + }); + sim.run().unwrap(); +}