diff --git a/libsql-server/src/namespace/configurator/replica.rs b/libsql-server/src/namespace/configurator/replica.rs index adea0fd406..9846e652f5 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}; @@ -66,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(); @@ -77,6 +81,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( @@ -110,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(_) => (), } @@ -126,6 +152,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/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/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/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/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..50ae54e158 100644 --- a/libsql-server/tests/fence/protocol.rs +++ b/libsql-server/tests/fence/protocol.rs @@ -974,3 +974,489 @@ 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(); +} + +/// 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(); +}