From 5f0ba5e262dd4ae19aeb1ca3803d0eb111f9b335 Mon Sep 17 00:00:00 2001 From: River Date: Wed, 30 Sep 2026 10:24:22 +0000 Subject: [PATCH 1/2] libsql-server: honour the replicated fence on replica servers A replica server whose primary refuses replication of a namespace with a fence code (a refused `hello`, `log_entries` or `snapshot`, or a stream the primary ended with the typed status) now refuses local reads and streams of its copy with the same code, and asks the reads it had already admitted to stop. The denial is published on the namespace's fence controller on the replica, so every user protocol reports it as it does on the primary; writes keep going to the primary, which refuses them itself. A `hello` the primary answers lifts it. The refusal is no longer retried every second by the handshake loop: the replica's replication loop waits 1 s, doubling to 15 s, until the primary answers again, and counts each refused call in `libsql_server_replica_fence_refusals_total{code}`. Co-authored-by: Tomasz Szymczyszyn --- .../src/namespace/configurator/replica.rs | 15 + .../src/namespace/fence/controller.rs | 44 +++ libsql-server/src/namespace/fence/mod.rs | 5 +- libsql-server/src/namespace/fence/replica.rs | 267 +++++++++++++ .../src/replication/replicator_client.rs | 150 ++++++-- libsql-server/tests/fence/protocol.rs | 361 ++++++++++++++++++ 6 files changed, 815 insertions(+), 27 deletions(-) create mode 100644 libsql-server/src/namespace/fence/replica.rs diff --git a/libsql-server/src/namespace/configurator/replica.rs b/libsql-server/src/namespace/configurator/replica.rs index adea0fd406..986d5880f2 100644 --- a/libsql-server/src/namespace/configurator/replica.rs +++ b/libsql-server/src/namespace/configurator/replica.rs @@ -19,6 +19,7 @@ use crate::database::{Database, ReplicaDatabase}; use crate::namespace::broadcasters::BroadcasterHandle; use crate::namespace::configurator::helpers::{make_stats, run_storage_monitor}; use crate::namespace::fence::controller::FenceController; +use crate::namespace::fence::replica::{refusal_backoff, PrimaryFenceRefusal}; use crate::namespace::meta_store::MetaStoreHandle; use crate::namespace::{Namespace, NamespaceBottomlessDbIdInit, RestoreOption}; use crate::namespace::{NamespaceName, NamespaceStore, ResetCb, ResetOp, ResolveNamespacePathFn}; @@ -77,6 +78,7 @@ impl ConfigureNamespace for ReplicaConfigurator { meta_store_handle.clone(), store.clone(), WalImpl::new_sqlite(&db_path, new_frame_sender).await?, + fence.clone(), ) .await?; let mut replicator = libsql_replication::replicator::Replicator::new_sqlite( @@ -126,6 +128,19 @@ impl ConfigureNamespace for ReplicaConfigurator { loop { match replicator.run().await { err @ Error::Fatal(_) => Err(err)?, + e @ Error::Internal(_) if PrimaryFenceRefusal::of(&e).is_some() => { + // The primary's fence refuses replication of this namespace + // (`docs/NAMESPACE_FENCE.md` section 6.2). The client has published + // the local read denial; retry at a capped, growing interval rather + // than at once, until the primary answers `hello` again. + let refusals = replicator.client_mut().fence_refusals(); + let delay = refusal_backoff(refusals); + tracing::debug!( + "{e}; retrying replication of {namespace} in {delay:?} \ + ({refusals} refusals in a row)" + ); + tokio::time::sleep(delay).await; + } _err @ Error::NamespaceDoesntExist => { // TODO(lucio): there is a bug where a primary will report that a valid // namespace doesn't exist when it does and causes the replicate to diff --git a/libsql-server/src/namespace/fence/controller.rs b/libsql-server/src/namespace/fence/controller.rs index e61421f341..bc358818a4 100644 --- a/libsql-server/src/namespace/fence/controller.rs +++ b/libsql-server/src/namespace/fence/controller.rs @@ -66,6 +66,11 @@ pub struct GateSnapshot { /// observability is refused with `MIGRATION_TARGET_QUARANTINED`, and the namespace is not /// set up. Never persisted; replaced by the record the command's commit publishes. pub creating_target: Option, + /// On a replica server only: the primary's fence denies reads of this namespace (section + /// 6.2), as the replicator last learned it from a refused replication call or from the + /// fence `hello` replicated. Normal reads and streams of the local copy are refused with + /// it. Never persisted and never set on a primary. + pub primary_denial: Option, } impl GateSnapshot { @@ -77,6 +82,7 @@ impl GateSnapshot { installing: None, closing_reads: None, creating_target: None, + primary_denial: None, } } @@ -156,6 +162,11 @@ impl GateSnapshot { )); } } + if let Some(denial) = &self.primary_denial { + if matches!(class, OperationClass::NormalRead | OperationClass::Stream) { + return Err(denial.clone()); + } + } Ok(()) } @@ -502,6 +513,39 @@ impl FenceController { asked } + /// On a replica server: publish what the replicator learned of the primary's fence + /// (section 6.2). `Some` denies normal reads and streams of the local copy with that error + /// and asks every read lease held now to stop, so that work admitted before the replica + /// learned of the fence does not outlive it; `None` admits them again. The denial is + /// published under the lease lock, so a read admitted concurrently is either refused or + /// counted and cancelled. A denial with the code already published leaves the gate as it + /// is. Returns whether the gate changed. + pub fn observe_primary(&self, denial: Option) -> bool { + let leases = self.read_leases.lock(); + let deny = denial.is_some(); + let changed = self.gate.send_if_modified(|gate| { + // A denial with the same code is the same denial, whichever call reported it. + let same = match (&gate.primary_denial, &denial) { + (Some(old), Some(new)) => old.outcome() == new.outcome(), + (None, None) => true, + _ => false, + }; + if same { + return false; + } + gate.primary_denial = denial; + true + }); + if changed && deny { + for entry in leases.live.values() { + if !entry.cancelled.swap(true, Ordering::AcqRel) { + (entry.cancel)(); + } + } + } + changed + } + /// Notified on every read-lease release. Enable the notification before checking /// [`read_lease_counts`](Self::read_lease_counts), so a release in between is not missed. pub(crate) fn read_released(&self) -> &Notify { diff --git a/libsql-server/src/namespace/fence/mod.rs b/libsql-server/src/namespace/fence/mod.rs index 43c384659a..8c89d79f5e 100644 --- a/libsql-server/src/namespace/fence/mod.rs +++ b/libsql-server/src/namespace/fence/mod.rs @@ -11,8 +11,8 @@ //! on them: the per-namespace [`controller`] with its gate and read leases, the positive write //! [`drain`], the source [`read`] fence and its //! [`stream`] leases for dump and replication, quarantined migration [`target`]s with their -//! [`capability`]-scoped [`import`] sessions and seal drain, the [`registry`] that holds the controllers outside the namespace cache, and the test [`hooks`] -//! on their paths. +//! [`capability`]-scoped [`import`] sessions and seal drain, the [`registry`] that holds the controllers outside the namespace cache, the +//! [`replica`]-server view of a primary's fence, and the test [`hooks`] on their paths. // The persistence, controller and protocol layers that consume these types land in the // following commits of this series; until then most of the module is unused by the rest of @@ -29,6 +29,7 @@ pub mod outcome; pub mod read; pub mod record; pub mod registry; +pub mod replica; pub mod state; pub mod store; pub mod stream; diff --git a/libsql-server/src/namespace/fence/replica.rs b/libsql-server/src/namespace/fence/replica.rs new file mode 100644 index 0000000000..bf198f5bce --- /dev/null +++ b/libsql-server/src/namespace/fence/replica.rs @@ -0,0 +1,267 @@ +//! The primary's fence as a replica server sees it (`docs/NAMESPACE_FENCE.md` section 6.2). +//! +//! A replica server holds a copy of a primary's namespace and serves reads of it locally. When +//! the primary's fence denies reads (a source read fence, a quarantined or aborted target, a +//! fence state the primary cannot establish), the primary refuses the replica's replication +//! calls and ends its streams with a typed `FAILED_PRECONDITION` status, and a `hello` it does +//! answer carries the fence it has. This module turns both into the local read denial the +//! replica publishes on its own fence controller +//! ([`FenceController::observe_primary`](super::controller::FenceController::observe_primary)), +//! and paces the replicator's reconnects while the primary keeps refusing. + +use std::time::Duration; + +use libsql_replication::replicator::Error as ReplicatorError; +use libsql_replication::rpc::metadata::ReplicatedFence; + +use super::outcome::{FenceError, OutcomeKind}; +use super::state::{FenceState, OperationClass}; + +/// The first pause after the primary refuses replication with a fence code. +pub const REFUSAL_BACKOFF_INITIAL: Duration = Duration::from_secs(1); +/// The longest pause between two replication attempts the primary's fence refused. It bounds +/// how long a replica keeps denying reads after the primary admits them again. +pub const REFUSAL_BACKOFF_MAX: Duration = Duration::from_secs(15); + +/// A replication call the primary refused, or a stream it ended, because its fence denies +/// replication of the namespace. Carried in [`ReplicatorError::Internal`], which the +/// replicator's handshake loop does not retry by itself, so that the replica's own loop can +/// pace the next attempt. +#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)] +#[error("the primary refused replication: {0}")] +pub struct PrimaryFenceRefusal(pub FenceError); + +impl PrimaryFenceRefusal { + /// The refusal a replication status reports: a data-plane fence denial in the typed form of + /// section 6 (`FAILED_PRECONDITION` with the stable code in `x-libsql-fence-code`). `None` + /// for every other status, which keeps its existing handling. + pub fn from_status(status: &tonic::Status) -> Option { + let denial = FenceError::from_grpc_status(status)?; + (denial.outcome().kind() == OutcomeKind::DataPlane).then_some(Self(denial)) + } + + /// The refusal `error` carries, if it is one. + pub fn of(error: &ReplicatorError) -> Option<&Self> { + match error { + ReplicatorError::Internal(e) => e.downcast_ref::(), + _ => None, + } + } + + /// The local read denial it implies. Replication is `Stream` work, whose column of the + /// permission matrix equals the normal-read column, so every data-plane refusal of it + /// means the primary denies reads. + pub fn local_denial(&self) -> FenceError { + FenceError::new( + self.0.outcome(), + format!( + "the primary denies reads of this namespace: {}", + self.0.message() + ), + ) + } +} + +/// `status` as a replicator error: a [`PrimaryFenceRefusal`] for a fence denial, the +/// replicator's own mapping otherwise. +pub fn replicator_error(status: tonic::Status) -> ReplicatorError { + match PrimaryFenceRefusal::from_status(&status) { + Some(refusal) => ReplicatorError::Internal(Box::new(refusal)), + None => status.into(), + } +} + +/// The local read denial the fence a primary's `hello` replicated implies: `Some` only for a +/// known state whose normal-read column denies. The primary answers `hello` only where it +/// admits streams, so this is `None` in practice; a state this server does not know is not +/// a denial for the same reason. +pub fn denial_from_hello(fence: Option<&ReplicatedFence>) -> Option { + let fence = fence?; + let state = fence.state.parse::().ok()?; + let outcome = state.permits(OperationClass::NormalRead).err()?; + Some(FenceError::new( + outcome, + format!( + "the primary's fence is {} at revision {}", + fence.state, fence.revision + ), + )) +} + +/// The pause before the next replication attempt after `consecutive` refusals in a row +/// (counting the one just received): doubling from [`REFUSAL_BACKOFF_INITIAL`] up to +/// [`REFUSAL_BACKOFF_MAX`]. +pub fn refusal_backoff(consecutive: u32) -> Duration { + let doublings = consecutive.saturating_sub(1).min(16); + REFUSAL_BACKOFF_INITIAL + .saturating_mul(1 << doublings) + .min(REFUSAL_BACKOFF_MAX) +} + +#[cfg(test)] +mod tests { + use tonic::Code; + + use super::super::outcome::FenceOutcome; + use super::*; + + fn status(outcome: FenceOutcome) -> tonic::Status { + FenceError::new(outcome, "no") + .to_grpc_status() + .expect("data-plane outcomes have a gRPC mapping") + } + + #[test] + fn refusal_from_typed_status_only() { + for outcome in [ + FenceOutcome::MigrationReadFenced, + FenceOutcome::MigrationTargetQuarantined, + FenceOutcome::FenceStateUnavailable, + ] { + let refusal = PrimaryFenceRefusal::from_status(&status(outcome)).unwrap(); + assert_eq!(refusal.0.outcome(), outcome); + let denial = refusal.local_denial(); + assert_eq!(denial.outcome(), outcome); + assert!( + denial.message().contains("the primary denies reads"), + "{denial}" + ); + + let error = replicator_error(status(outcome)); + assert_eq!(PrimaryFenceRefusal::of(&error), Some(&refusal)); + } + + // Untyped statuses keep the replicator's own mapping. + for status in [ + tonic::Status::new(Code::FailedPrecondition, "MIGRATION_READ_FENCED: no"), + tonic::Status::new(Code::Unavailable, "down"), + tonic::Status::new( + Code::FailedPrecondition, + libsql_replication::rpc::replication::NAMESPACE_DOESNT_EXIST, + ), + ] { + assert!(PrimaryFenceRefusal::from_status(&status).is_none()); + assert!(PrimaryFenceRefusal::of(&replicator_error(status)).is_none()); + } + assert!(matches!( + replicator_error(tonic::Status::new( + Code::FailedPrecondition, + libsql_replication::rpc::replication::NEED_SNAPSHOT_ERROR_MSG + )), + ReplicatorError::NeedSnapshot + )); + + // A control outcome is never a replication refusal. + let mut control = tonic::Status::new(Code::FailedPrecondition, "x"); + control.metadata_mut().insert( + super::super::outcome::GRPC_FENCE_CODE_METADATA, + tonic::metadata::MetadataValue::from_static("FENCE_PRECONDITION_FAILED"), + ); + assert!(PrimaryFenceRefusal::from_status(&control).is_none()); + } + + #[test] + fn hello_fence_denies_only_read_denying_states() { + let fence = |state: &str| ReplicatedFence { + state: state.into(), + revision: 7, + }; + assert_eq!(denial_from_hello(None), None); + for state in [ + "SOURCE_DRAINING", + "SOURCE_WRITE_FENCED", + "TARGET_WRITE_FENCED", + "SOMETHING_NEWER", + ] { + assert_eq!(denial_from_hello(Some(&fence(state))), None, "{state}"); + } + for (state, outcome) in [ + ("SOURCE_READ_DRAINING", FenceOutcome::MigrationReadFenced), + ("SOURCE_READ_FENCED", FenceOutcome::MigrationReadFenced), + ( + "TARGET_QUARANTINED", + FenceOutcome::MigrationTargetQuarantined, + ), + ("UNKNOWN_UNAVAILABLE", FenceOutcome::FenceStateUnavailable), + ] { + let denial = denial_from_hello(Some(&fence(state))).unwrap(); + assert_eq!(denial.outcome(), outcome, "{state}"); + assert!(denial.message().contains("revision 7"), "{denial}"); + } + } + + #[test] + fn backoff_doubles_to_its_cap() { + let delays: Vec<_> = (1..=7).map(refusal_backoff).collect(); + assert_eq!( + delays, + [1, 2, 4, 8, 15, 15, 15].map(Duration::from_secs).to_vec() + ); + assert_eq!(refusal_backoff(0), REFUSAL_BACKOFF_INITIAL); + assert_eq!(refusal_backoff(u32::MAX), REFUSAL_BACKOFF_MAX); + } + + #[test] + fn observed_denial_refuses_local_reads_and_cancels_leases() { + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::Arc; + + use super::super::controller::{FenceController, LeaseKind}; + use crate::namespace::NamespaceName; + + let fence = FenceController::unfenced(NamespaceName::from_string("ns".into()).unwrap()); + let cancels = Arc::new(AtomicUsize::new(0)); + let lease = fence + .acquire_read_lease(OperationClass::NormalRead, LeaseKind::Sql, { + let cancels = cancels.clone(); + move || { + cancels.fetch_add(1, Ordering::SeqCst); + } + }) + .unwrap(); + let generation = fence.write_generation(); + let refusal = + PrimaryFenceRefusal::from_status(&status(FenceOutcome::MigrationReadFenced)).unwrap(); + + assert!(fence.observe_primary(Some(refusal.local_denial()))); + assert_eq!(cancels.load(Ordering::SeqCst), 1); + assert!(lease.cancelled_by_fence()); + for class in [OperationClass::NormalRead, OperationClass::Stream] { + let err = fence.permits(class).unwrap_err(); + assert_eq!( + err.outcome(), + FenceOutcome::MigrationReadFenced, + "{class:?}" + ); + } + // Writes are the primary's to refuse; maintenance goes on. + for class in [OperationClass::NormalWrite, OperationClass::Maintenance] { + assert!(fence.permits(class).is_ok(), "{class:?}"); + } + assert!(fence + .acquire_read_lease(OperationClass::Stream, LeaseKind::Dump, || ()) + .is_err()); + // The same code again, from another call, is the same denial. + assert!(!fence.observe_primary(Some(refusal.local_denial()))); + assert_eq!(cancels.load(Ordering::SeqCst), 1); + // A different code replaces it. + let quarantined = + PrimaryFenceRefusal::from_status(&status(FenceOutcome::MigrationTargetQuarantined)) + .unwrap(); + assert!(fence.observe_primary(Some(quarantined.local_denial()))); + assert_eq!( + fence + .permits(OperationClass::NormalRead) + .unwrap_err() + .outcome(), + FenceOutcome::MigrationTargetQuarantined + ); + + assert!(fence.observe_primary(None)); + assert!(fence.permits(OperationClass::NormalRead).is_ok()); + assert!(!fence.observe_primary(None)); + // Only reads were affected: the write generation never moved. + assert_eq!(fence.write_generation(), generation); + drop(lease); + } +} diff --git a/libsql-server/src/replication/replicator_client.rs b/libsql-server/src/replication/replicator_client.rs index fb8154824d..541f3b68e6 100644 --- a/libsql-server/src/replication/replicator_client.rs +++ b/libsql-server/src/replication/replicator_client.rs @@ -1,5 +1,7 @@ use std::path::Path; use std::pin::Pin; +use std::sync::atomic::{AtomicU32, Ordering}; +use std::sync::Arc; use bytes::Bytes; use chrono::{DateTime, Utc}; @@ -23,6 +25,9 @@ use crate::connection::config::DatabaseConfig; use crate::metrics::{ REPLICATION_LATENCY, REPLICATION_LATENCY_CACHE_MISS, REPLICATION_LATENCY_OUT_OF_SYNC, }; +use crate::namespace::fence::controller::FenceController; +use crate::namespace::fence::outcome::FenceError; +use crate::namespace::fence::replica::{self, PrimaryFenceRefusal}; use crate::namespace::meta_store::MetaStoreHandle; use crate::namespace::{NamespaceName, NamespaceStore}; use crate::replication::FrameNo; @@ -107,6 +112,12 @@ pub struct Client { store: NamespaceStore, wal_impl: WalImpl, first_sync_since_handshake: bool, + /// The namespace's fence controller on this replica server, on which the primary's fence + /// is published as a local read denial (`docs/NAMESPACE_FENCE.md` section 6.2). + fence: Arc, + /// Replication calls the primary's fence refused since the last `hello` it answered. Shared + /// with active frame streams so a refusal delivered as their terminal status is counted too. + fence_refusals: Arc, } impl Client { @@ -116,6 +127,7 @@ impl Client { meta_store_handle: MetaStoreHandle, store: NamespaceStore, wal_flavor: WalImpl, + fence: Arc, ) -> crate::Result { Ok(Self { namespace, @@ -126,6 +138,88 @@ impl Client { store, wal_impl: wal_flavor, first_sync_since_handshake: true, + fence, + fence_refusals: Arc::new(AtomicU32::new(0)), + }) + } + + /// Replication calls the primary's fence refused in a row, since the last `hello` it + /// answered. The replica's replication loop paces its reconnects by it. + pub(crate) fn fence_refusals(&self) -> u32 { + self.fence_refusals.load(Ordering::Relaxed) + } + + /// Publish what the primary said of its fence as this replica's local read denial, logging + /// when that changes. + fn observe_primary_fence(&self, denial: Option) { + let denies = denial.as_ref().map(|d| d.outcome()); + if self.fence.observe_primary(denial) { + match denies { + Some(code) => tracing::warn!( + namespace = %self.namespace, + "the primary's namespace fence denies reads ({code}): local reads of this \ + replica are refused until the primary admits replication again" + ), + None => tracing::info!( + namespace = %self.namespace, + "the primary admits replication again: local reads are served" + ), + } + } + } + + /// Map a status of a replication call: a fence refusal is published as the local read + /// denial, counted, and returned as a [`PrimaryFenceRefusal`]. + fn status_error(&mut self, status: Status) -> Error { + let error = replica::replicator_error(status); + if let Some(refusal) = PrimaryFenceRefusal::of(&error) { + self.fence_refusals + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |count| { + Some(count.saturating_add(1)) + }) + .ok(); + metrics::increment_counter!( + "libsql_server_replica_fence_refusals_total", + "code" => refusal.0.outcome().as_str(), + ); + tracing::debug!(namespace = %self.namespace, "{refusal}"); + self.observe_primary_fence(Some(refusal.local_denial())); + } + error + } + + /// A stream of the primary's that ends with a fence refusal publishes it as the local read + /// denial. + fn fenced_frames( + &self, + stream: tonic::Streaming, + ) -> impl Stream> + Send + 'static { + let fence = self.fence.clone(); + let fence_refusals = self.fence_refusals.clone(); + let namespace = self.namespace.clone(); + stream.map_err(move |status| { + let error = replica::replicator_error(status); + if let Some(refusal) = PrimaryFenceRefusal::of(&error) { + fence_refusals + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |count| { + Some(count.saturating_add(1)) + }) + .ok(); + metrics::increment_counter!( + "libsql_server_replica_fence_refusals_total", + "code" => refusal.0.outcome().as_str(), + ); + if fence.observe_primary(Some(refusal.local_denial())) { + tracing::warn!( + namespace = %namespace, + "the primary ended replication because its namespace fence denies \ + reads ({}): local reads of this replica are refused until the primary \ + admits replication again", + refusal.0.outcome() + ); + } + } + error }) } @@ -164,9 +258,17 @@ impl ReplicatorClient for Client { self.first_sync_since_handshake = true; tracing::debug!("Attempting to perform handshake with primary."); let req = self.make_request(HelloRequest::new()); - let resp = self.client.hello(req).await?; + let resp = match self.client.hello(req).await { + Ok(resp) => resp, + Err(status) => return Err(self.status_error(status)), + }; let hello = resp.into_inner(); verify_session_token(&hello.session_token).map_err(Error::Client)?; + // The primary answers `hello` only where its fence admits replication. + self.fence_refusals.store(0, Ordering::Relaxed); + self.observe_primary_fence(replica::denial_from_hello( + hello.config.as_ref().and_then(|c| c.fence.as_ref()), + )); self.primary_replication_index = hello.current_replication_index; self.session_token.replace(hello.session_token.clone()); @@ -207,32 +309,30 @@ impl ReplicatorClient for Client { }; let req = self.make_request(offset); - let stream = self - .client - .log_entries(req) - .await? - .into_inner() - .inspect_ok(|f| { - match f.timestamp { - Some(ts_millis) => { - if let Some(commited_at) = DateTime::from_timestamp_millis(ts_millis) { - let lat = Utc::now() - commited_at; - match lat.to_std() { - Ok(lat) => { - // we can record negative values if the clocks are out-of-sync. There is not - // point in recording those values. - REPLICATION_LATENCY.record(lat); - } - Err(_) => { - REPLICATION_LATENCY_OUT_OF_SYNC.increment(1); - } + let stream = match self.client.log_entries(req).await { + Ok(resp) => resp.into_inner(), + Err(status) => return Err(self.status_error(status)), + }; + let stream = self.fenced_frames(stream).inspect_ok(|f| { + match f.timestamp { + Some(ts_millis) => { + if let Some(commited_at) = DateTime::from_timestamp_millis(ts_millis) { + let lat = Utc::now() - commited_at; + match lat.to_std() { + Ok(lat) => { + // we can record negative values if the clocks are out-of-sync. There is not + // point in recording those values. + REPLICATION_LATENCY.record(lat); + } + Err(_) => { + REPLICATION_LATENCY_OUT_OF_SYNC.increment(1); } } } - None => REPLICATION_LATENCY_CACHE_MISS.increment(1), } - }) - .map_err(Into::into); + None => REPLICATION_LATENCY_CACHE_MISS.increment(1), + } + }); Ok(Box::pin(stream)) } @@ -245,11 +345,11 @@ impl ReplicatorClient for Client { let req = self.make_request(offset); match self.client.snapshot(req).await { Ok(resp) => { - let stream = resp.into_inner().map_err(Into::into); + let stream = self.fenced_frames(resp.into_inner()); Ok(Box::pin(stream)) } Err(e) if e.code() == Code::Unavailable => Err(Error::SnapshotPending), - Err(e) => return Err(e.into()), + Err(e) => Err(self.status_error(e)), } } diff --git a/libsql-server/tests/fence/protocol.rs b/libsql-server/tests/fence/protocol.rs index 7cf970fda3..a7b7b85d0f 100644 --- a/libsql-server/tests/fence/protocol.rs +++ b/libsql-server/tests/fence/protocol.rs @@ -974,3 +974,364 @@ fn denial_not_retried() { }); sim.run().unwrap(); } + +/// The primary's internal replication service (the one replica servers use), as a raw client +/// that sees statuses and their metadata. +struct Replication { + client: libsql_replication::rpc::replication::replication_log_client::ReplicationLogClient< + tonic::transport::Channel, + >, + ns: String, + token: Option, +} + +impl Replication { + fn new(ns: &str) -> anyhow::Result { + use tower::ServiceExt as _; + let uri = tonic::transport::Uri::from_static("http://primary:4567"); + let channel = tonic::transport::Channel::builder(uri.clone()).connect_with_connector_lazy( + TurmoilConnector.map_err(|e| -> Box { e.into() }), + ); + Ok(Self { + client: libsql_replication::rpc::replication::replication_log_client::ReplicationLogClient::with_origin(channel, uri), + ns: ns.into(), + token: None, + }) + } + + fn request(&self, msg: T) -> tonic::Request { + use libsql_replication::rpc::replication::{NAMESPACE_METADATA_KEY, SESSION_TOKEN_KEY}; + let mut req = tonic::Request::new(msg); + req.metadata_mut().insert_bin( + NAMESPACE_METADATA_KEY, + tonic::metadata::BinaryMetadataValue::from_bytes(self.ns.as_bytes()), + ); + if let Some(token) = &self.token { + req.metadata_mut().insert( + SESSION_TOKEN_KEY, + tonic::metadata::AsciiMetadataValue::try_from(token.as_ref()).unwrap(), + ); + } + req + } + + /// `hello`; on success keeps the session token and returns the replicated fence. + async fn hello( + &mut self, + ) -> Result, tonic::Status> { + let req = self.request(libsql_replication::rpc::replication::HelloRequest::new()); + let hello = self.client.hello(req).await?.into_inner(); + self.token = Some(hello.session_token.clone()); + Ok(hello.config.and_then(|c| c.fence)) + } + + fn offset(&self) -> tonic::Request { + self.request(libsql_replication::rpc::replication::LogOffset { + next_offset: 0, + wal_flavor: None, + }) + } + + async fn log_entries( + &mut self, + ) -> Result, tonic::Status> { + let req = self.offset(); + Ok(self.client.log_entries(req).await?.into_inner()) + } + + async fn snapshot( + &mut self, + ) -> Result, tonic::Status> { + let req = self.offset(); + Ok(self.client.snapshot(req).await?.into_inner()) + } +} + +/// A status in the typed form of section 6: `FAILED_PRECONDITION`, the stable code in +/// `x-libsql-fence-code`, and the code prefixing the message. +#[track_caller] +fn assert_fence_status(what: &str, status: &tonic::Status, code: &str) { + assert_eq!( + status.code(), + tonic::Code::FailedPrecondition, + "{what}: {status:?}" + ); + assert_eq!( + status + .metadata() + .get("x-libsql-fence-code") + .and_then(|v| v.to_str().ok()), + Some(code), + "{what}: {status:?}" + ); + assert!(status.message().starts_with(code), "{what}: {status:?}"); +} + +/// Replication as a raw peer of the primary sees it (sections 6.2 and 9): `hello` carries the +/// replicated fence while the primary admits replication; under a read fence an open stream +/// ends with the typed status and every call is refused with it; a quarantined target refuses +/// with its own code; after the read fence is cleared `hello` is answered again. +#[test] +fn replication_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); + admin.create_namespace("plain").await?; + load_and_log_id(&admin, "plain").await?; + assert_eq!(Replication::new("plain")?.hello().await?, None, "unfenced"); + + let mut repl = Replication::new("src")?; + let rev = write_fenced(&admin, "src", op).await?; + let fence = repl.hello().await?.expect("hello carries the write fence"); + assert_eq!( + (fence.state.as_str(), fence.revision), + ("SOURCE_WRITE_FENCED", rev) + ); + let mut tail = repl.log_entries().await?; + let frame = tail + .next() + .await + .expect("a frame") + .expect("frames are served"); + assert!(!frame.data.is_empty()); + + let rev = read_fence(&admin, "src", op, rev).await?; + // The open stream ends with the terminal status (frames already buffered first). + let ended = loop { + match tail.next().await { + Some(Ok(_)) => continue, + Some(Err(status)) => break status, + None => panic!("the stream ended without a status"), + } + }; + assert_fence_status("open log_entries", &ended, READ_FENCED); + assert!(tail.next().await.is_none()); + assert_fence_status("hello", &repl.hello().await.unwrap_err(), READ_FENCED); + assert_fence_status( + "log_entries", + &repl.log_entries().await.unwrap_err(), + READ_FENCED, + ); + assert_fence_status("snapshot", &repl.snapshot().await.unwrap_err(), READ_FENCED); + + let rev = clear_read_fence(&admin, "src", op, rev).await?; + let fence = repl.hello().await?.expect("hello carries the write fence"); + assert_eq!( + (fence.state.as_str(), fence.revision), + ("SOURCE_WRITE_FENCED", rev) + ); + release(&admin, "src", op, rev).await?; + assert_eq!(repl.hello().await?, None, "released"); + + quarantined_target(&admin, "dst", uuid(0x200)).await?; + let mut target = Replication::new("dst")?; + assert_fence_status( + "target hello", + &target.hello().await.unwrap_err(), + QUARANTINED, + ); + Ok(()) + }); + sim.run().unwrap(); +} + +/// The first `/v2` result of `sql` on `user`'s host: `Ok(())` or the Hrana error code. +async fn v2_read(user: &User, ns: &str, sql: &str) -> anyhow::Result> { + let (status, body) = user + .pipeline(ns, 2, None, json!([execute_req(sql)])) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let result = &body["results"][0]; + Ok(match result["type"].as_str() { + Some("ok") => Ok(()), + _ => Err(result["error"]["code"] + .as_str() + .unwrap_or_else(|| panic!("no error code: {body}")) + .to_string()), + }) +} + +/// Poll the legacy API on `user`'s host with `sql` until it answers `status`, for at most +/// `within` of simulated time, and return that answer and how long it took. +async fn poll_until( + user: &User, + ns: &str, + sql: &str, + status: StatusCode, + within: std::time::Duration, +) -> anyhow::Result<(Value, std::time::Duration)> { + let started = tokio::time::Instant::now(); + loop { + let (got, body) = user.legacy(ns, &[sql]).await?; + if got == status { + return Ok((body, started.elapsed())); + } + assert!( + started.elapsed() < within, + "still {got} after {:?}: {body}", + started.elapsed() + ); + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + } +} + +/// `src` write-fenced on the primary and loaded on the replica, then read-fenced; returns once +/// the replica denies local reads, with the fence's revision. +async fn replica_read_fenced(admin: &Admin, replica: &User, op: Uuid) -> anyhow::Result { + let rev = write_fenced(admin, "src", op).await?; + let (status, body) = replica.legacy("src", &["select count(*) from t"]).await?; + assert_eq!( + status, + StatusCode::OK, + "replica before the read fence: {body}" + ); + let rev = read_fence(admin, "src", op, rev).await?; + // The replica learns of the fence when its replication stream ends, which the primary's + // read drain waits for; the denial is published as the terminal status arrives. + let (body, took) = poll_until( + replica, + "src", + "select count(*) from t", + StatusCode::LOCKED, + std::time::Duration::from_secs(2), + ) + .await?; + assert_eq!(body["code"], READ_FENCED, "{body}"); + assert!( + took < std::time::Duration::from_millis(500), + "took {took:?}" + ); + Ok(rev) +} + +/// A replica server denies local reads of a namespace whose primary is read-fenced, on every +/// read surface, with the primary's code (section 6.2). +#[test] +fn replica_reads_denied_while_source_read_fenced() { + 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"); + replica_read_fenced(&admin, &user, uuid(0x100)).await?; + + assert_locked( + "legacy read", + &user.legacy("src", &["select * from t"]).await?, + READ_FENCED, + ); + assert_locked( + "v1 execute", + &user.execute("src", "select * from t").await?, + READ_FENCED, + ); + assert_eq!( + v2_read(&user, "src", "select * from t").await?, + Err(READ_FENCED.to_string()) + ); + // `/dump` is not served by a replica server at all ("database is not a primary"). + // Still denied later: the replica does not forget while the primary keeps refusing. + tokio::time::sleep(std::time::Duration::from_secs(30)).await; + assert_locked( + "legacy read later", + &user.legacy("src", &["select * from t"]).await?, + READ_FENCED, + ); + Ok(()) + }); + sim.run().unwrap(); +} + +/// While the primary's fence refuses replication, the replica retries at a growing interval +/// capped at 15 s, not every second or in a tight loop: over 60 s of simulated time it makes +/// a handful of attempts (1 + 2 + 4 + 8 + 15 + 15 + 15 s of pauses), each counted. +#[test] +fn replica_backs_off_on_fence_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"); + replica_read_fenced(&admin, &user, uuid(0x100)).await?; + + let refusals = || { + crate::common::snapshot_metrics() + .get_counter_label( + "libsql_server_replica_fence_refusals_total", + ("code", READ_FENCED), + ) + .unwrap_or(0) + }; + let before = refusals(); + assert!(before >= 1, "the ended stream is counted"); + tokio::time::sleep(std::time::Duration::from_secs(60)).await; + let attempts = refusals() - before; + assert!( + (4..=9).contains(&attempts), + "{attempts} refused attempts in 60 s" + ); + Ok(()) + }); + sim.run().unwrap(); +} + +/// Once the primary clears the read fence the replica answers `hello` again within the +/// back-off cap, serves local reads, and replicates new writes once the source is released. +#[test] +fn replica_resumes_after_clear_read_fence() { + 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 = replica_read_fenced(&admin, &user, op).await?; + // Let the back-off grow to its cap before clearing. + tokio::time::sleep(std::time::Duration::from_secs(40)).await; + + let rev = clear_read_fence(&admin, "src", op, rev).await?; + let (_, took) = poll_until( + &user, + "src", + "select count(*) from t", + StatusCode::OK, + std::time::Duration::from_secs(20), + ) + .await?; + assert!(took <= std::time::Duration::from_secs(16), "took {took:?}"); + assert_eq!(v2_read(&user, "src", "select * from t").await?, Ok(())); + + release(&admin, "src", op, rev).await?; + let (status, body) = User::new() + .legacy("src", &["insert into t values (2)"]) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let started = tokio::time::Instant::now(); + loop { + let (status, body) = user.legacy("src", &["select count(*) from t"]).await?; + assert_eq!(status, StatusCode::OK, "{body}"); + if body[0]["results"]["rows"][0][0] == 2 { + break; + } + assert!( + started.elapsed() < std::time::Duration::from_secs(10), + "not replicated: {body}" + ); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + } + Ok(()) + }); + sim.run().unwrap(); +} From 993abd66e12c3b559cde6af3d71ced4e38ab653f Mon Sep 17 00:00:00 2001 From: River Date: Wed, 30 Sep 2026 10:36:32 +0000 Subject: [PATCH 2/2] libsql-server: do not create a replica namespace the primary's fence refuses A replica server creating a namespace lazily for a name the primary's fence refuses (a quarantined or aborted target, a read-fenced source, a fence state the primary cannot establish) now fails the request at once with the primary's code (Error::NamespaceFence: 423 and the stable code) instead of retrying the handshake, and leaves no local namespace behind: the namespace directory the setup created is removed (a directory that already existed is kept), and NamespaceStore::with forgets the metastore entry handle() added and the controller the attempt created, when neither holds anything durable or is in use by another attempt. A later request, once the primary admits the name, creates it normally. Co-authored-by: Tomasz Szymczyszyn --- .../src/namespace/configurator/replica.rs | 26 +++- libsql-server/src/namespace/fence/registry.rs | 60 +++++++++ libsql-server/src/namespace/meta_store.rs | 69 +++++++++- libsql-server/src/namespace/store.rs | 41 +++++- libsql-server/tests/fence/protocol.rs | 125 ++++++++++++++++++ 5 files changed, 313 insertions(+), 8 deletions(-) diff --git a/libsql-server/src/namespace/configurator/replica.rs b/libsql-server/src/namespace/configurator/replica.rs index 986d5880f2..9846e652f5 100644 --- a/libsql-server/src/namespace/configurator/replica.rs +++ b/libsql-server/src/namespace/configurator/replica.rs @@ -67,6 +67,9 @@ impl ConfigureNamespace for ReplicaConfigurator { Box::pin(async move { tracing::debug!("creating replica namespace"); let db_path = self.base.base_path.join("dbs").join(name.as_str()); + // Whether this server already had a copy of the namespace: a directory this setup + // creates is removed again if the primary's fence refuses the namespace. + let had_copy = db_path.try_exists()?; let channel = self.channel.clone(); let uri = self.uri.clone(); @@ -112,7 +115,28 @@ impl ConfigureNamespace for ReplicaConfigurator { ) .await; } - Err(e) => Err(e)?, + Err(e) => { + if let Some(refusal) = PrimaryFenceRefusal::of(&e) { + // The primary's fence refuses the namespace (a quarantined target, a + // read-fenced source, a fence state it cannot establish; section 13.3): + // fail at once with its code rather than retrying the handshake, and + // leave no local copy behind that this setup created. + let denial = refusal.0.clone(); + drop(replicator); + if !had_copy { + if let Err(e) = tokio::fs::remove_dir_all(&db_path).await { + if e.kind() != std::io::ErrorKind::NotFound { + tracing::warn!( + "failed to remove {} after the primary refused {name}: {e}", + db_path.display() + ); + } + } + } + return Err(crate::Error::NamespaceFence(denial)); + } + Err(e)? + } Ok(_) => (), } diff --git a/libsql-server/src/namespace/fence/registry.rs b/libsql-server/src/namespace/fence/registry.rs index 24dd0ea23e..7dd52425ba 100644 --- a/libsql-server/src/namespace/fence/registry.rs +++ b/libsql-server/src/namespace/fence/registry.rs @@ -66,6 +66,29 @@ impl FenceRegistry { self.controllers.lock().remove(namespace) } + /// Forget `namespace`'s controller if it holds nothing worth keeping: no fence record and + /// no in-memory gate (only what a replica learned of its primary's fence, which the next + /// answered `hello` would replace), and nothing but the registry refers to it. For a name + /// whose setup failed before it was ever served, such as a replica's lazy creation that + /// the primary's fence refused. Returns whether it was forgotten. + pub fn forget_idle(&self, namespace: &NamespaceName) -> bool { + let mut controllers = self.controllers.lock(); + let idle = controllers.get(namespace).is_some_and(|controller| { + let gate = controller.gate(); + // Under the registry lock nobody can take another reference to it. + Arc::strong_count(controller) == 1 + && matches!(gate.fence, StoredFence::None { .. }) + && gate.indeterminate.is_none() + && gate.installing.is_none() + && gate.closing_reads.is_none() + && gate.creating_target.is_none() + }); + if idle { + controllers.remove(namespace); + } + idle + } + /// Refuse a namespace whose fence state is `UNKNOWN_UNAVAILABLE`, or that is being created /// as a quarantined target, before any work is done to serve it. pub fn check_available(&self, namespace: &NamespaceName) -> Result<(), FenceError> { @@ -245,4 +268,41 @@ mod tests { assert!(registry.remove(&"ns".into()).is_some()); assert!(!Arc::ptr_eq(®istry.controller(&"ns".into()), &a)); } + + /// A replica's lazy creation that the primary refused leaves nothing behind in the registry, + /// unless the controller holds fence state or somebody else still refers to it. + #[test] + fn forget_idle_only_unreferenced_plain_controllers() { + let registry = FenceRegistry::default(); + assert!(!registry.forget_idle(&"missing".into())); + + // What a refused replication taught it does not keep it. + let refused = registry.controller(&"refused".into()); + refused.observe_primary(Some(FenceError::new( + FenceOutcome::MigrationTargetQuarantined, + "quarantined on the primary", + ))); + // Still referenced: kept. + assert!(!registry.forget_idle(&"refused".into())); + drop(refused); + assert!(registry.forget_idle(&"refused".into())); + assert!(registry.get(&"refused".into()).is_none()); + // The next use starts from a fresh UNFENCED controller. + assert!(registry + .controller(&"refused".into()) + .permits(OperationClass::NormalRead) + .is_ok()); + + // A controller with fence state is never forgotten. + let registry = FenceRegistry::seeded([( + "lost".into(), + StoredFence::Unavailable { + detail: FenceDetail::CorruptRecord, + reason: "test".into(), + marker: None, + }, + )]); + assert!(!registry.forget_idle(&"lost".into())); + assert!(registry.get(&"lost".into()).is_some()); + } } diff --git a/libsql-server/src/namespace/meta_store.rs b/libsql-server/src/namespace/meta_store.rs index 5f67045640..8665745526 100644 --- a/libsql-server/src/namespace/meta_store.rs +++ b/libsql-server/src/namespace/meta_store.rs @@ -15,7 +15,7 @@ use libsql_sys::wal::{ }; use parking_lot::Mutex; use prost::Message; -use rusqlite::TransactionBehavior; +use rusqlite::{OptionalExtension, TransactionBehavior}; use tokio::sync::oneshot; use tokio::sync::{ mpsc, @@ -1303,6 +1303,41 @@ impl MetaStore { r } + /// Take out the in-memory entry that [`handle`](Self::handle) put in the map for a + /// namespace whose creation then failed, so that [`exists`](Self::exists) and + /// [`lookup`](Self::lookup) do not report a namespace that was never created + /// (`docs/NAMESPACE_FENCE.md` section 13.3, replica lazy creation). Only an entry that no + /// handle is subscribed to any more and that has no stored config row is removed: a config + /// that was persisted, or a creation of the same name still in progress, keeps its entry. + /// Returns whether the entry was removed. + pub async fn forget_unstored(&self, namespace: NamespaceName) -> Result { + let inner = self.inner.clone(); + tokio::task::spawn_blocking(move || -> std::result::Result { + // The connection lock first, as everywhere else that takes both: a config being + // persisted concurrently is either already in its row here, or finds no entry when + // it publishes and inserts its own. + let conn = inner.conn.blocking_lock(); + let stored = conn + .query_row( + "SELECT 1 FROM namespace_configs WHERE namespace = ?1", + [namespace.as_str()], + |_| Ok(()), + ) + .optional()? + .is_some(); + let mut configs = inner.configs.blocking_lock(); + match configs.get(&namespace) { + Some(sender) if !stored && sender.receiver_count() == 0 => { + configs.remove(&namespace); + Ok(true) + } + _ => Ok(false), + } + }) + .await? + .map_err(fence_store_error) + } + // TODO: we need to either make sure that the metastore is restored // before we start accepting connections or we need to contact bottomless // here to check if a namespace exists. Preferably the former. @@ -1739,6 +1774,38 @@ mod fence_tests { } } + /// The in-memory entry a failed creation left is forgotten, so `exists()` and `lookup()` do + /// not report the name; a stored config, or a handle still held, keeps the entry. + #[tokio::test] + async fn forget_unstored_only_unused_unstored_entries() { + let tmp = tempdir().unwrap(); + let meta = open(tmp.path(), true).await; + let ns = NamespaceName::from("lazy"); + + assert!(!meta.forget_unstored(ns.clone()).await.unwrap()); + + let handle = meta.handle(ns.clone()).await.unwrap(); + assert!(meta.exists(&ns).await); + // A handle is still held (a creation in progress): kept. + assert!(!meta.forget_unstored(ns.clone()).await.unwrap()); + assert!(meta.exists(&ns).await); + drop(handle); + assert!(meta.forget_unstored(ns.clone()).await.unwrap()); + assert!(!meta.exists(&ns).await); + assert!(meta.lookup(&ns).await.unwrap().is_none()); + + // A stored config is never forgotten. + let stored = NamespaceName::from("stored"); + meta.handle(stored.clone()) + .await + .unwrap() + .store(DatabaseConfig::default()) + .await + .unwrap(); + assert!(!meta.forget_unstored(stored.clone()).await.unwrap()); + assert!(meta.lookup(&stored).await.unwrap().is_some()); + } + fn request( ns: &'static str, op: Uuid, diff --git a/libsql-server/src/namespace/store.rs b/libsql-server/src/namespace/store.rs index 9a83252b53..aa22de9e26 100644 --- a/libsql-server/src/namespace/store.rs +++ b/libsql-server/src/namespace/store.rs @@ -400,17 +400,46 @@ impl NamespaceStore { // A lookup that cannot create: only the default namespace and lazy creation create a // namespace here, and those refuse a name whose fence state is not established. - let handle = match self.inner.metadata.lookup(&namespace).await? { - Some(handle) => handle, + let (handle, created) = match self.inner.metadata.lookup(&namespace).await? { + Some(handle) => (handle, false), None if namespace == NamespaceName::default() || self.inner.allow_lazy_creation => { - self.inner.metadata.handle(namespace.clone()).await? + (self.inner.metadata.handle(namespace.clone()).await?, true) } None => return Err(Error::NamespaceDoesntExist(namespace.to_string())), }; - f(self + let entry = match self .load_namespace(&namespace, handle, RestoreOption::Latest) - .await?) - .await + .await + { + Ok(entry) => entry, + Err(e) => { + if created && e.fence_error().is_some() { + self.forget_refused_creation(&namespace).await; + } + return Err(e); + } + }; + f(entry).await + } + + /// Undo what a lazy creation that a fence refused left in memory (on a replica server, the + /// primary refused to replicate the name; `docs/NAMESPACE_FENCE.md` section 13.3): the + /// metastore entry [`MetaStore::handle`] added, which would otherwise make the name look + /// like an existing namespace, and the controller the attempt created. Both are kept when + /// they hold anything durable or are in use by another attempt. + async fn forget_refused_creation(&self, namespace: &NamespaceName) { + match self.inner.metadata.forget_unstored(namespace.clone()).await { + Ok(forgotten) => { + let controller = self.inner.fences.forget_idle(namespace); + tracing::debug!( + "refused creation of {namespace}: metastore entry forgotten: {forgotten}, \ + controller forgotten: {controller}" + ); + } + Err(e) => { + tracing::warn!("failed to forget the refused creation of {namespace}: {e}") + } + } } fn resolve_attach_fn(&self) -> ResolveNamespacePathFn { diff --git a/libsql-server/tests/fence/protocol.rs b/libsql-server/tests/fence/protocol.rs index a7b7b85d0f..50ae54e158 100644 --- a/libsql-server/tests/fence/protocol.rs +++ b/libsql-server/tests/fence/protocol.rs @@ -1335,3 +1335,128 @@ fn replica_resumes_after_clear_read_fence() { }); sim.run().unwrap(); } + +/// Walks the quarantined target `ns` of `op` to `TARGET_WRITE_FENCED` (readable): seal the +/// (empty) import, record a successful validation, publish. Returns the revision. +async fn publish_target(admin: &Admin, ns: &str, op: Uuid) -> anyhow::Result { + let (_, body) = admin.inspect(ns).await?; + let (state, mut rev) = state_of(&body); + assert_eq!(state, "TARGET_QUARANTINED", "{body}"); + for (n, route, from, extra, to) in [ + ( + 2, + "target/seal-import", + "TARGET_QUARANTINED", + json!({}), + "TARGET_VALIDATING", + ), + ( + 3, + "target/validation-receipt", + "TARGET_VALIDATING", + json!({ "result": "ok", "summary": "empty" }), + "TARGET_VALIDATING", + ), + ( + 4, + "target/publish-readable", + "TARGET_VALIDATING", + json!({}), + "TARGET_WRITE_FENCED", + ), + ] { + let (status, body) = admin + .command( + ns, + route, + command_body(op, uuid(op.as_u128() + n), from, rev, extra), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{route}: {body}"); + assert_eq!(state_of(&body).0, to, "{route}: {body}"); + rev = state_of(&body).1; + } + Ok(rev) +} + +/// A replica server creating a namespace lazily, for a name the primary's fence refuses to +/// replicate (here a quarantined target), fails the request at once with the primary's code +/// instead of retrying the handshake, and leaves no local copy: no namespace directory it +/// created (a directory that was already there stays) and no metastore entry. Once the target +/// is published the same name is created and served normally (section 13.3). +#[test] +fn replica_lazy_creation_refused_by_fence() { + let mut sim = sim(); + let primary = tempdir().unwrap(); + let replica = tempdir().unwrap(); + // A directory that was on the replica before: a refused creation must not delete it. + let kept = replica.path().join("dbs").join("kept"); + std::fs::create_dir_all(&kept).unwrap(); + std::fs::write(kept.join("sentinel"), b"x").unwrap(); + let dbs = replica.path().join("dbs"); + make_primary(&mut sim, primary.path().to_path_buf(), Primary::default()); + make_replica(&mut sim, replica.path().to_path_buf()); + sim.client("client", async move { + let admin = Admin::new(Some(ADMIN_KEY)); + let user = User::on("replica0"); + let op = uuid(0x100); + quarantined_target(&admin, "tgt", op).await?; + quarantined_target(&admin, "kept", uuid(0x200)).await?; + + for attempt in 0..2 { + let started = tokio::time::Instant::now(); + assert_locked( + &format!("replica read of a quarantined target, attempt {attempt}"), + &user.legacy("tgt", &["select 1"]).await?, + QUARANTINED, + ); + // Refused at the first handshake, not after a second of retries per attempt. + let took = started.elapsed(); + assert!( + took < std::time::Duration::from_millis(500), + "took {took:?}" + ); + assert!( + !dbs.join("tgt").exists(), + "the refused creation left dbs/tgt behind" + ); + } + assert_locked( + "replica read of a quarantined target with a directory", + &user.legacy("kept", &["select 1"]).await?, + QUARANTINED, + ); + assert!( + kept.join("sentinel").exists(), + "a pre-existing directory was removed" + ); + + let rev = publish_target(&admin, "tgt", op).await?; + let (status, body) = user.legacy("tgt", &["select 1"]).await?; + assert_eq!(status, StatusCode::OK, "after publication: {body}"); + assert!( + dbs.join("tgt").exists(), + "published target not created on the replica" + ); + + // Writes through the replica reach the primary once writes are enabled. + let (status, body) = admin + .command( + "tgt", + "target/enable-writes", + command_body( + op, + uuid(op.as_u128() + 5), + "TARGET_WRITE_FENCED", + rev, + json!({}), + ), + ) + .await?; + assert_eq!(status, StatusCode::OK, "{body}"); + let (status, body) = user.legacy("tgt", &["create table t (x)"]).await?; + assert_eq!(status, StatusCode::OK, "write through the replica: {body}"); + Ok(()) + }); + sim.run().unwrap(); +}