use std::sync::{Arc, Mutex, OnceLock};
use crate::error::{BootstrapError, BootstrapResult};
use crate::{ClusterHandle, TestCluster};
use super::fixtures::ensure_worker_env;
static SHARED_CLUSTER_HANDLE: OnceLock<Mutex<SharedHandleState>> = OnceLock::new();
enum SharedHandleState {
Uninitialised,
Initialised(&'static ClusterHandle),
Failed(Arc<BootstrapError>),
}
pub fn shared_cluster_handle() -> BootstrapResult<&'static ClusterHandle> {
let mutex = SHARED_CLUSTER_HANDLE.get_or_init(|| Mutex::new(SharedHandleState::Uninitialised));
let mut guard = mutex
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
match &*guard {
SharedHandleState::Initialised(handle) => Ok(*handle),
SharedHandleState::Failed(original_err) => {
let report = color_eyre::eyre::eyre!(
"shared cluster initialisation previously failed: {:?}",
original_err
);
Err(BootstrapError::new(original_err.kind(), report))
}
SharedHandleState::Uninitialised => {
let worker_guard = ensure_worker_env();
match TestCluster::new_split() {
Ok((handle, cluster_guard)) => {
let guarded = cluster_guard.with_worker_guard(worker_guard);
best_effort_register_shutdown_hook_for_handle(&handle);
std::mem::forget(guarded);
let leaked: &'static ClusterHandle = Box::leak(Box::new(handle));
*guard = SharedHandleState::Initialised(leaked);
Ok(leaked)
}
Err(err) => {
let stored = Arc::new(BootstrapError::new(
err.kind(),
color_eyre::eyre::eyre!("bootstrap failed: {:?}", err),
));
*guard = SharedHandleState::Failed(stored);
Err(err)
}
}
}
}
}
fn best_effort_register_shutdown_hook_for_handle(handle: &ClusterHandle) {
if let Err(err) = handle.register_shutdown_on_exit() {
tracing::debug!(
target: crate::observability::LOG_TARGET,
error = %err,
"shutdown hook registration failed; postmaster may be orphaned on exit"
);
}
}
fn best_effort_register_shutdown_hook(cluster: &TestCluster) {
if let Err(err) = cluster.register_shutdown_on_exit() {
tracing::debug!(
target: crate::observability::LOG_TARGET,
error = %err,
"failed to register shutdown hook (may already be registered)"
);
}
}
static SHARED_CLUSTER: OnceLock<Mutex<SharedClusterState>> = OnceLock::new();
enum SharedClusterState {
Uninitialised,
Initialised(SharedClusterPtr),
Failed(Arc<BootstrapError>),
}
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::Uninitialised));
let mut guard = mutex
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
match &*guard {
SharedClusterState::Initialised(ptr) => {
Ok(unsafe { &*ptr.0 })
}
SharedClusterState::Failed(original_err) => {
let report = color_eyre::eyre::eyre!(
"shared cluster initialisation previously failed: {:?}",
original_err
);
Err(BootstrapError::new(original_err.kind(), report))
}
SharedClusterState::Uninitialised => {
let worker_guard = ensure_worker_env();
match TestCluster::new() {
Ok(new_cluster) => {
let guarded_cluster = new_cluster.with_worker_guard(worker_guard);
let leaked: &'static TestCluster = Box::leak(Box::new(guarded_cluster));
best_effort_register_shutdown_hook(leaked);
let ptr = SharedClusterPtr(std::ptr::from_ref::<TestCluster>(leaked));
*guard = SharedClusterState::Initialised(ptr);
Ok(leaked)
}
Err(err) => {
let stored = Arc::new(BootstrapError::new(
err.kind(),
color_eyre::eyre::eyre!("bootstrap failed: {:?}", err),
));
*guard = SharedClusterState::Failed(stored);
Err(err)
}
}
}
}
}