From de2b7200094811fa397e28cfc843597c4dddccf1 Mon Sep 17 00:00:00 2001 From: River Date: Wed, 30 Sep 2026 09:26:10 +0000 Subject: [PATCH 1/3] libsql-replication: add stable error code and replicated fence to protocols Two additive proto3 fields for the namespace fence: - `proxy.Error.stable_code` (tag 4): a stable machine-readable outcome such as `MIGRATION_WRITE_FENCED`, so a replica can return the same typed outcome the primary would have. Absent means "no typed outcome"; older peers skip it. - `metadata.DatabaseConfig.fence` (tag 14, `ReplicatedFence { state, revision }`): the live fence as the primary's replication `hello` sees it. It is filled only by `hello` and only while a fence is active, and is never part of a stored configuration. The primary now fills `DatabaseConfig.fence` in `hello` from the namespace's fence gate. Filling `stable_code` on fence denials and mapping it on the replica side come in a later change; until then no server sets it, and capability discovery keeps reporting `proxy_stable_code: false`. Tests cover wire compatibility in both directions (an absent field encodes exactly as before) and that `hello` carries the fence only while one is active. Co-authored-by: Tomasz Szymczyszyn --- libsql-replication/proto/metadata.proto | 13 ++ libsql-replication/proto/proxy.proto | 5 + libsql-replication/src/generated/metadata.rs | 17 +++ libsql-replication/src/generated/proxy.rs | 6 + libsql-replication/src/rpc.rs | 115 ++++++++++++++++++ libsql-server/src/connection/config.rs | 3 + .../src/namespace/fence/controller.rs | 13 ++ libsql-server/src/namespace/fence/stream.rs | 79 ++++++++++++ libsql-server/src/rpc/proxy.rs | 1 + .../src/rpc/replication/replication_log.rs | 11 +- 10 files changed, 261 insertions(+), 2 deletions(-) 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/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/stream.rs b/libsql-server/src/namespace/fence/stream.rs index 2a043d89fa..bd9ae56097 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; @@ -413,6 +416,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..32bc6aeb96 100644 --- a/libsql-server/src/rpc/proxy.rs +++ b/libsql-server/src/rpc/proxy.rs @@ -57,6 +57,7 @@ pub mod rpc { message: other.to_string(), code: code as i32, extended_code, + stable_code: None, } } } 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)) From c04a426b5b1fd56b2c4469c44790c6df0e650e65 Mon Sep 17 00:00:00 2001 From: River Date: Wed, 30 Sep 2026 09:48:16 +0000 Subject: [PATCH 2/3] libsql-server: typed fence outcomes across HTTP, Hrana and dump A fence denial reaching the user-facing protocols is now a typed answer everywhere instead of a generic or fatal error: - HTTP error bodies for fence errors gain an additive `code` field (and `detail` when there is one), with `423` for data-plane denials, through every wrapper the error arrives in. The legacy `/` API answers a batch with a fenced step as a whole with `423`. - Hrana gains `StmtError::Fence` and `BatchError::Fence`, whose code is the stable fence code, so step denials and whole-request denials (a read under a read fence, a quarantined target) are Hrana errors on `/v1`, `/v2`, `/v3`, cursors and WebSockets, and the stream stays usable. - `/v1/execute` and `/v1/batch` answer whole-request denials with `423` and the code. Integration tests cover each protocol, an old WebSocket transaction, a batch denied mid-way, `/dump`, and that `401`/`404` stay distinct; a lib test shows a dump cancelled by the read drain fails its response body. Co-authored-by: Tomasz Szymczyszyn --- libsql-server/src/error.rs | 85 +- libsql-server/src/hrana/batch.rs | 6 + libsql-server/src/hrana/stmt.rs | 6 + .../src/http/user/hrana_over_http_1.rs | 17 + libsql-server/src/http/user/mod.rs | 5 +- libsql-server/src/http/user/result_builder.rs | 13 + libsql-server/src/namespace/fence/outcome.rs | 21 + libsql-server/src/namespace/fence/stream.rs | 42 + libsql-server/src/schema/error.rs | 2 +- libsql-server/tests/fence/lifecycle.rs | 1 + libsql-server/tests/fence/mod.rs | 15 +- libsql-server/tests/fence/protocol.rs | 853 ++++++++++++++++++ 12 files changed, 1062 insertions(+), 4 deletions(-) create mode 100644 libsql-server/tests/fence/protocol.rs diff --git a/libsql-server/src/error.rs b/libsql-server/src/error.rs index f0cb631769..1408c67432 100644 --- a/libsql-server/src/error.rs +++ b/libsql-server/src/error.rs @@ -151,6 +151,29 @@ 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, + } + } +} + +/// 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 +249,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 +361,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..294c378e01 100644 --- a/libsql-server/src/http/user/mod.rs +++ b/libsql-server/src/http/user/mod.rs @@ -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/outcome.rs b/libsql-server/src/namespace/fence/outcome.rs index cf4df0c9af..debe1ba88f 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 { diff --git a/libsql-server/src/namespace/fence/stream.rs b/libsql-server/src/namespace/fence/stream.rs index bd9ae56097..2711105a2a 100644 --- a/libsql-server/src/namespace/fence/stream.rs +++ b/libsql-server/src/namespace/fence/stream.rs @@ -326,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")] 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/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..32952a9ca9 100644 --- a/libsql-server/tests/fence/mod.rs +++ b/libsql-server/tests/fence/mod.rs @@ -4,11 +4,14 @@ mod admin; mod lifecycle; +mod protocol; use std::path::PathBuf; 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, RpcServerConfig, UserApiConfig}; use s3s::header::AUTHORIZATION; use serde_json::{json, Value}; @@ -26,6 +29,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 +38,7 @@ impl Default for Primary { Self { admin_key: Some(ADMIN_KEY), fence_enabled: true, + user_credential: None, } } } @@ -49,13 +55,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, diff --git a/libsql-server/tests/fence/protocol.rs b/libsql-server/tests/fence/protocol.rs new file mode 100644 index 0000000000..d7255fbd29 --- /dev/null +++ b/libsql-server/tests/fence/protocol.rs @@ -0,0 +1,853 @@ +//! 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, 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`, for any namespace, with an optional basic-auth credential. +struct User { + client: hyper::Client, + auth: Option, +} + +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}")), + } + } + + 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}.primary:8080{path}")); + 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(); +} From daf31f791a83d9901e316ad49074c5944276af26 Mon Sep 17 00:00:00 2001 From: River Date: Wed, 30 Sep 2026 10:06:56 +0000 Subject: [PATCH 3/3] libsql-server: carry fence outcomes through RPC and the replica write proxy The primary's proxy service now fills the additive `stable_code` field (with `code = SQL_ERROR`) for fence denials, on step errors and program errors, streamed and unary. A fence denial before a program runs (the namespace or JWT-key lookup, connection creation, a unary program refused as a whole) is the typed `FAILED_PRECONDITION` status with the stable code in `x-libsql-fence-code`, never `UNAVAILABLE`, which the write proxy retries without bound. On a replica, the write proxy maps a proxied error or status carrying a fence code back to `Error::NamespaceFence`, so the replica answers its client exactly as the primary would: `423` with the `code` field on the HTTP APIs and the stable code as the Hrana error code. Errors from an older primary, which never sets the field, keep their old mapping. Capability discovery now reports `proxy_stable_code: true`. Co-authored-by: Tomasz Szymczyszyn --- libsql-server/src/connection/write_proxy.rs | 8 +- libsql-server/src/error.rs | 26 ++ libsql-server/src/http/user/mod.rs | 2 +- libsql-server/src/namespace/fence/mod.rs | 2 +- libsql-server/src/namespace/fence/outcome.rs | 75 ++++ libsql-server/src/rpc/proxy.rs | 366 +++++++++++++++++-- libsql-server/src/rpc/streaming_exec.rs | 2 +- libsql-server/tests/fence/admin.rs | 2 +- libsql-server/tests/fence/mod.rs | 35 +- libsql-server/tests/fence/protocol.rs | 131 ++++++- 10 files changed, 608 insertions(+), 41 deletions(-) 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 1408c67432..b6903c4a82 100644 --- a/libsql-server/src/error.rs +++ b/libsql-server/src/error.rs @@ -162,6 +162,32 @@ impl 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 diff --git a/libsql-server/src/http/user/mod.rs b/libsql-server/src/http/user/mod.rs index 294c378e01..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] 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 debe1ba88f..723b7b5ab6 100644 --- a/libsql-server/src/namespace/fence/outcome.rs +++ b/libsql-server/src/namespace/fence/outcome.rs @@ -318,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)] @@ -425,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/rpc/proxy.rs b/libsql-server/src/rpc/proxy.rs index 32bc6aeb96..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,7 +64,7 @@ pub mod rpc { message: other.to_string(), code: code as i32, extended_code, - stable_code: None, + stable_code: stable_code.map(Into::into), } } } @@ -65,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, @@ -316,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 @@ -568,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>>; @@ -586,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()); @@ -618,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; @@ -646,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()) + })) } } } @@ -660,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())) } @@ -692,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(); @@ -720,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()) + })) + } } } }; @@ -756,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/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/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/mod.rs b/libsql-server/tests/fence/mod.rs index 32952a9ca9..2e9cbae183 100644 --- a/libsql-server/tests/fence/mod.rs +++ b/libsql-server/tests/fence/mod.rs @@ -12,7 +12,9 @@ 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, RpcServerConfig, UserApiConfig}; +use libsql_server::config::{ + AdminApiConfig, MetaStoreConfig, RpcClientConfig, RpcServerConfig, UserApiConfig, +}; use s3s::header::AUTHORIZATION; use serde_json::{json, Value}; use turmoil::{Builder, Sim}; @@ -93,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 index d7255fbd29..7cf970fda3 100644 --- a/libsql-server/tests/fence/protocol.rs +++ b/libsql-server/tests/fence/protocol.rs @@ -12,8 +12,8 @@ use turmoil::net::TcpStream; use uuid::Uuid; use super::{ - acquire_body, command_body, load_and_log_id, make_primary, sim, state_of, Admin, Primary, - ADMIN_KEY, + acquire_body, command_body, load_and_log_id, make_primary, make_replica, sim, state_of, Admin, + Primary, ADMIN_KEY, }; use crate::common::net::TurmoilConnector; @@ -26,10 +26,12 @@ fn uuid(n: u128) -> Uuid { Uuid::from_u128(n) } -/// The user API of `primary`, for any namespace, with an optional basic-auth credential. +/// 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 { @@ -41,6 +43,15 @@ impl User { 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() } } @@ -53,7 +64,7 @@ impl User { ) -> anyhow::Result<(StatusCode, String)> { let mut request = Request::builder() .method(method) - .uri(format!("http://{ns}.primary:8080{path}")); + .uri(format!("http://{ns}.{}:8080{path}", self.host)); if let Some(auth) = &self.auth { request = request.header("authorization", auth.as_str()); } @@ -851,3 +862,115 @@ fn auth_and_not_found_distinct() { }); 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(); +}