use std::sync::{Arc, Mutex, OnceLock};
use super::{
fixtures::ensure_worker_env,
shared_singleton_core::{SharedInitState, get_or_try_init_shared},
};
use crate::{
ClusterHandle,
TestCluster,
error::{BootstrapError, BootstrapResult},
};
static SHARED_CLUSTER_HANDLE: OnceLock<Mutex<SharedHandleState>> = OnceLock::new();
type SharedHandleState = SharedInitState<&'static ClusterHandle, Arc<BootstrapError>>;
pub fn shared_cluster_handle() -> BootstrapResult<&'static ClusterHandle> {
let mutex = SHARED_CLUSTER_HANDLE.get_or_init(|| Mutex::new(SharedInitState::Uninitialized));
let mut guard = mutex
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
get_or_try_init_shared(
&mut guard,
initialize_shared_cluster_handle,
cached_shared_handle_error,
)
}
fn initialize_shared_cluster_handle()
-> Result<&'static ClusterHandle, (Arc<BootstrapError>, BootstrapError)> {
let worker_guard = ensure_worker_env();
match TestCluster::new_split() {
Ok((handle, cluster_guard)) => {
log_shared_cluster_event("shared cluster handle initialized");
let guarded = cluster_guard.with_worker_guard(worker_guard);
leak_shared_handle_after_shutdown_hook(
handle,
guarded,
ClusterHandle::register_shutdown_on_exit,
)
}
Err(err) => {
log_shared_cluster_error(&err, "shared cluster handle initialization failed");
let stored = Arc::new(BootstrapError::new(
err.kind(),
color_eyre::eyre::eyre!("bootstrap failed: {:?}", err),
));
Err((stored, err))
}
}
}
fn cached_shared_handle_error(original_err: &Arc<BootstrapError>) -> BootstrapError {
log_cached_shared_handle_error(original_err);
let report = color_eyre::eyre::eyre!(
"shared cluster initialization previously failed: {:?}",
original_err
);
BootstrapError::new(original_err.kind(), report)
}
fn register_shutdown_hook_or_fail<T, Register>(
target: &T,
register_shutdown_hook: Register,
failure_message: &'static str,
) -> Result<(), (Arc<BootstrapError>, BootstrapError)>
where
Register: FnOnce(&T) -> BootstrapResult<()>,
{
if let Err(err) = register_shutdown_hook(target) {
log_shared_cluster_error(&err, failure_message);
let stored = Arc::new(BootstrapError::new(
err.kind(),
color_eyre::eyre::eyre!("shutdown hook registration failed: {:?}", err),
));
return Err((stored, err));
}
Ok(())
}
fn leak_shared_handle_after_shutdown_hook<Guard, Register>(
handle: ClusterHandle,
guarded: Guard,
register_shutdown_hook: Register,
) -> Result<&'static ClusterHandle, (Arc<BootstrapError>, BootstrapError)>
where
Register: FnOnce(&ClusterHandle) -> BootstrapResult<()>,
{
register_shutdown_hook_or_fail(
&handle,
register_shutdown_hook,
"shared cluster shutdown hook registration failed; dropping guard",
)?;
std::mem::forget(guarded);
log_shared_cluster_event("shared cluster shutdown hook registered; leaking handle");
Ok(Box::leak(Box::new(handle)))
}
fn log_shared_cluster_event(message: &'static str) {
tracing::debug!(
target: crate::observability::LOG_TARGET,
"{message}"
);
}
fn log_shared_cluster_error(err: &BootstrapError, message: &'static str) {
tracing::error!(
target: crate::observability::LOG_TARGET,
error = ?err,
"{message}"
);
}
fn log_cached_shared_handle_error(original_err: &Arc<BootstrapError>) {
tracing::debug!(
target: crate::observability::LOG_TARGET,
error = ?original_err,
"replaying cached shared cluster initialization failure"
);
}
static SHARED_CLUSTER: OnceLock<Mutex<SharedClusterState>> = OnceLock::new();
type SharedClusterState = SharedInitState<SharedClusterPtr, Arc<BootstrapError>>;
#[derive(Clone, Copy)]
struct SharedClusterPtr(*const TestCluster);
unsafe impl Send for SharedClusterPtr {}
unsafe impl Sync for SharedClusterPtr {}
pub fn shared_cluster() -> BootstrapResult<&'static TestCluster> {
let mutex = SHARED_CLUSTER.get_or_init(|| Mutex::new(SharedClusterState::Uninitialized));
let mut guard = mutex
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let ptr = get_or_try_init_shared(
&mut guard,
initialize_shared_cluster,
cached_shared_handle_error,
)?;
Ok(unsafe { &*ptr.0 })
}
fn initialize_shared_cluster() -> Result<SharedClusterPtr, (Arc<BootstrapError>, BootstrapError)> {
let worker_guard = ensure_worker_env();
match TestCluster::new() {
Ok(new_cluster) => {
log_shared_cluster_event("shared cluster initialized");
leak_shared_cluster_after_shutdown_hook(
new_cluster.with_worker_guard(worker_guard),
|cluster| cluster.handle.register_shutdown_on_exit(),
)
.map(SharedClusterPtr)
}
Err(err) => {
log_shared_cluster_error(&err, "shared cluster initialization failed");
let stored = Arc::new(BootstrapError::new(
err.kind(),
color_eyre::eyre::eyre!("bootstrap failed: {:?}", err),
));
Err((stored, err))
}
}
}
fn leak_shared_cluster_after_shutdown_hook<Cluster, Register>(
cluster: Cluster,
register_shutdown_hook: Register,
) -> Result<*const Cluster, (Arc<BootstrapError>, BootstrapError)>
where
Register: FnOnce(&Cluster) -> BootstrapResult<()>,
{
register_shutdown_hook_or_fail(
&cluster,
register_shutdown_hook,
"shared cluster shutdown hook registration failed; dropping cluster",
)?;
let leaked = Box::leak(Box::new(cluster));
log_shared_cluster_event("shared cluster shutdown hook registered; leaking cluster");
Ok(std::ptr::from_ref(leaked))
}
#[cfg(test)]
mod tests {
use std::sync::{
Arc,
atomic::{AtomicBool, Ordering},
};
use super::*;
use crate::{ExecutionPrivileges, test_support::dummy_settings};
struct DropProbe(Arc<AtomicBool>);
impl Drop for DropProbe {
fn drop(&mut self) { self.0.store(true, Ordering::SeqCst); }
}
#[test]
fn registration_failure_drops_guard_instead_of_leaking_shared_handle() {
let was_dropped = Arc::new(AtomicBool::new(false));
let handle = ClusterHandle::from(dummy_settings(ExecutionPrivileges::Unprivileged));
let guarded = DropProbe(Arc::clone(&was_dropped));
let result = leak_shared_handle_after_shutdown_hook(handle, guarded, |_| {
Err(BootstrapError::from(color_eyre::eyre::eyre!(
"injected shutdown hook failure"
)))
});
assert!(
result.is_err(),
"shared-handle initialization should fail when shutdown hook registration fails"
);
assert!(
was_dropped.load(Ordering::SeqCst),
concat!(
"guard must drop on registration failure so postmaster.pid is removed ",
"and no orphaned PostgreSQL process remains"
)
);
}
#[test]
fn registration_failure_drops_guard_instead_of_leaking_shared_cluster() {
let was_dropped = Arc::new(AtomicBool::new(false));
let cluster = DropProbe(Arc::clone(&was_dropped));
let result = leak_shared_cluster_after_shutdown_hook(cluster, |_| {
Err(BootstrapError::from(color_eyre::eyre::eyre!(
"injected shutdown hook failure"
)))
});
assert!(
result.is_err(),
"shared-cluster initialization should fail when shutdown hook registration fails"
);
assert!(
was_dropped.load(Ordering::SeqCst),
"cluster must drop when shutdown hook registration fails"
);
}
}