Merge pull request #61 from logos-blockchain/fix/local-cluster-monitor-cleanup

fix(tf): cluster monitor cleanup
This commit is contained in:
Andrus Salumets 2026-08-04 15:31:45 +07:00 committed by GitHub
commit db880d6115
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
4 changed files with 69 additions and 5 deletions

View File

@ -29,6 +29,7 @@ impl<E: LocalDeployerEnv> LocalClusterOwner<E> {
fn close(&self) {
if !self.closed.swap(true, Ordering::AcqRel) {
self.nodes.stop_all();
E::cleanup_local_cluster(self.nodes.deployment());
}
}
@ -68,6 +69,7 @@ impl<E: LocalDeployerEnv> Clone for LocalCluster<E> {
impl<E: LocalDeployerEnv> LocalCluster<E> {
pub(crate) fn empty(deployment: E::Deployment, keep_tempdir: bool) -> Self {
E::prepare_local_cluster(&deployment);
let nodes = NodeManager::new_with_seed(
deployment.clone(),
NodeClients::default(),

View File

@ -75,6 +75,12 @@ pub trait LocalDeployerEnv: Application + Sized
where
<Self as Application>::NodeConfig: Clone + Send + Sync + 'static,
{
/// Prepares application-specific resources associated with a local cluster.
fn prepare_local_cluster(_deployment: &Self::Deployment) {}
/// Releases application-specific resources associated with a local cluster.
fn cleanup_local_cluster(_deployment: &Self::Deployment) {}
/// Returns named ports that should be reserved in addition to the main
/// network port for each node.
fn local_port_names() -> &'static [&'static str] {

View File

@ -85,6 +85,10 @@ pub struct NodeManagerSeed {
}
impl<E: LocalDeployerEnv> NodeManager<E> {
pub(crate) const fn deployment(&self) -> &E::Deployment {
&self.descriptors
}
pub async fn spawn_initial_nodes(
descriptors: &E::Deployment,
keep_tempdir: bool,

View File

@ -268,7 +268,14 @@ async fn run_retry_attempt<E: LocalDeployerEnv>(
#[cfg(test)]
mod tests {
use std::{collections::HashMap, path::Path};
use std::{
collections::HashMap,
path::Path,
sync::{
Arc,
atomic::{AtomicUsize, Ordering},
},
};
use testing_framework_core::{
scenario::{
@ -278,14 +285,22 @@ mod tests {
topology::DeploymentDescriptor,
};
use super::LocalClusterProvisioner;
use super::{LocalClusterProvisioner, LocalClusterProvisionerError};
use crate::{
BuiltNodeConfig, LaunchSpec, LocalDeployerEnv, NodeConfigEntry, NodeEndpoints,
ProcessSpawnError,
};
#[derive(Clone)]
struct EmptyDeployment;
#[derive(Default)]
struct LifecycleCalls {
prepare: AtomicUsize,
cleanup: AtomicUsize,
}
#[derive(Clone, Default)]
struct EmptyDeployment {
lifecycle: Arc<LifecycleCalls>,
}
impl DeploymentDescriptor for EmptyDeployment {
fn node_count(&self) -> usize {
@ -305,12 +320,23 @@ mod tests {
type NodeConfig = EmptyConfig;
fn external_node_client(source: &ExternalNodeSource) -> Result<Self::NodeClient, DynError> {
if source.endpoint().is_empty() {
return Err("empty test endpoint".into());
}
Ok(source.endpoint().to_owned())
}
}
#[async_trait::async_trait]
impl LocalDeployerEnv for TestEnv {
fn prepare_local_cluster(deployment: &Self::Deployment) {
deployment.lifecycle.prepare.fetch_add(1, Ordering::Relaxed);
}
fn cleanup_local_cluster(deployment: &Self::Deployment) {
deployment.lifecycle.cleanup.fetch_add(1, Ordering::Relaxed);
}
fn build_node_config(
_topology: &Self::Deployment,
_index: usize,
@ -346,7 +372,9 @@ mod tests {
#[tokio::test]
async fn on_demand_managed_unit_exposes_shared_control_and_cleanup() {
let request = ClusterRequest::managed(EmptyDeployment)
let deployment = EmptyDeployment::default();
let lifecycle = Arc::clone(&deployment.lifecycle);
let request = ClusterRequest::managed(deployment)
.with_start_mode(ClusterStartMode::OnDemand)
.with_control(ClusterControlRequest::Full);
let provisioned = LocalClusterProvisioner
@ -356,6 +384,7 @@ mod tests {
let (cluster, mut unit) = provisioned.into_parts();
let cluster = cluster.expect("managed unit should expose its concrete local handle");
assert_eq!(lifecycle.prepare.load(Ordering::Relaxed), 1);
assert_eq!(
unit.control_profile(),
ClusterControlProfile::ManualControlled
@ -366,7 +395,30 @@ mod tests {
unit.take_cleanup()
.expect("managed unit should own cleanup")
.cleanup();
assert_eq!(lifecycle.cleanup.load(Ordering::Relaxed), 1);
assert!(cluster.stop_all().is_err(), "cleanup must lock out clones");
drop(cluster);
assert_eq!(lifecycle.cleanup.load(Ordering::Relaxed), 1);
}
#[tokio::test]
async fn provisioning_failure_cleans_up_prepared_cluster_once() {
let deployment = EmptyDeployment::default();
let lifecycle = Arc::clone(&deployment.lifecycle);
let invalid_source = ExternalNodeSource::new("invalid".into(), String::new());
let request = ClusterRequest::managed(deployment)
.with_start_mode(ClusterStartMode::OnDemand)
.with_external_nodes(vec![invalid_source]);
let error = LocalClusterProvisioner
.provision::<TestEnv>(request, true)
.await
.err()
.expect("invalid external source should fail provisioning");
assert!(matches!(error, LocalClusterProvisionerError::Source { .. }));
assert_eq!(lifecycle.prepare.load(Ordering::Relaxed), 1);
assert_eq!(lifecycle.cleanup.load(Ordering::Relaxed), 1);
}
#[tokio::test]