Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 40 additions & 1 deletion libsql-server/src/namespace/configurator/replica.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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();

Expand All @@ -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(
Expand Down Expand Up @@ -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(_) => (),
}

Expand All @@ -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
Expand Down
44 changes: 44 additions & 0 deletions libsql-server/src/namespace/fence/controller.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<CommandKey>,
/// 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<FenceError>,
}

impl GateSnapshot {
Expand All @@ -77,6 +82,7 @@ impl GateSnapshot {
installing: None,
closing_reads: None,
creating_target: None,
primary_denial: None,
}
}

Expand Down Expand Up @@ -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(())
}

Expand Down Expand Up @@ -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<FenceError>) -> 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 {
Expand Down
5 changes: 3 additions & 2 deletions libsql-server/src/namespace/fence/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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;
Expand Down
60 changes: 60 additions & 0 deletions libsql-server/src/namespace/fence/registry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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> {
Expand Down Expand Up @@ -245,4 +268,41 @@ mod tests {
assert!(registry.remove(&"ns".into()).is_some());
assert!(!Arc::ptr_eq(&registry.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());
}
}
Loading
Loading