use std::sync::Arc;
use tracing::info;
use crate::control::cluster::handle::ClusterHandle;
use crate::control::state::SharedState;
pub(super) struct Hooks {
pub(super) quarantine_hook:
Arc<crate::control::cluster::snapshot_hook::RaftSnapshotQuarantineHook>,
pub(super) snapshot_builder: Arc<dyn nodedb_cluster::SnapshotBuilder>,
pub(super) snapshot_applier: Arc<dyn nodedb_cluster::SnapshotApplier>,
pub(super) shuffle_receiver: Arc<dyn nodedb_cluster::ShuffleReceiver>,
pub(super) shuffle_producer: Arc<dyn nodedb_cluster::ShuffleProducer>,
pub(super) shuffle_consumer: Arc<dyn nodedb_cluster::ShuffleConsumer>,
pub(super) shuffle_aggregator: Arc<dyn nodedb_cluster::ShuffleAggregator>,
pub(super) assign_remote_surrogate: Arc<dyn nodedb_cluster::AssignRemoteSurrogate>,
pub(super) calvin_submit: Arc<dyn nodedb_cluster::CalvinSubmit>,
pub(super) calvin_submit_inbox: Arc<dyn nodedb_cluster::CalvinSubmitInbox>,
pub(super) reserve_read: Arc<dyn nodedb_cluster::ReserveRead>,
pub(super) release_reservation: Arc<dyn nodedb_cluster::ReleaseReservation>,
}
pub(super) fn build_hooks(
handle: &ClusterHandle,
shared: &Arc<SharedState>,
data_dir: &std::path::Path,
) -> crate::Result<Hooks> {
let quarantine_hook = Arc::new(
crate::control::cluster::snapshot_hook::RaftSnapshotQuarantineHook {
registry: Arc::clone(&shared.quarantine_registry),
},
);
let snapshot_builder: Arc<dyn nodedb_cluster::SnapshotBuilder> = Arc::new(
crate::control::cluster::snapshot_builder::DataPlaneSnapshotBuilder::new(shared.clone()),
);
let snapshot_applier_concrete =
crate::control::cluster::snapshot_applier::DataPlaneSnapshotApplier::new(shared.clone());
let restored = tokio::task::block_in_place(|| {
tokio::runtime::Handle::current().block_on(
crate::control::cluster::boot_restore::restore_persisted_snapshots(
data_dir,
&snapshot_applier_concrete,
),
)
})?;
if restored > 0 {
info!(
node_id = handle.node_id,
restored, "follower boot-restore re-installed persisted snapshots"
);
}
let snapshot_applier: Arc<dyn nodedb_cluster::SnapshotApplier> =
Arc::new(snapshot_applier_concrete);
let shuffle_receiver: Arc<dyn nodedb_cluster::ShuffleReceiver> = Arc::new(
crate::control::server::shuffle::RegistryShuffleReceiver::new(Arc::clone(
&shared.shuffle_registry,
)),
);
let shuffle_producer: Arc<dyn nodedb_cluster::ShuffleProducer> =
Arc::new(crate::control::server::shuffle::RegistryShuffleProducer::new(shared.clone()));
let shuffle_consumer: Arc<dyn nodedb_cluster::ShuffleConsumer> =
Arc::new(crate::control::server::shuffle::RegistryShuffleConsumer::new(shared.clone()));
let shuffle_aggregator: Arc<dyn nodedb_cluster::ShuffleAggregator> =
Arc::new(crate::control::server::shuffle::RegistryShuffleAggregator::new(shared.clone()));
let assign_remote_surrogate: Arc<dyn nodedb_cluster::AssignRemoteSurrogate> = Arc::new(
crate::control::server::surrogate_exchange::RegistryAssignRemoteSurrogate::new(
shared.clone(),
),
);
let calvin_submit: Arc<dyn nodedb_cluster::CalvinSubmit> =
Arc::new(crate::control::server::calvin_submit::RegistryCalvinSubmit::new(shared.clone()));
let calvin_submit_inbox: Arc<dyn nodedb_cluster::CalvinSubmitInbox> = Arc::new(
crate::control::server::calvin_submit::RegistryCalvinSubmitInbox::new(shared.clone()),
);
let reserve_read: Arc<dyn nodedb_cluster::ReserveRead> =
Arc::new(crate::control::server::reservation::RegistryReserveRead::new(shared.clone()));
let release_reservation: Arc<dyn nodedb_cluster::ReleaseReservation> = Arc::new(
crate::control::server::reservation::RegistryReleaseReservation::new(shared.clone()),
);
Ok(Hooks {
quarantine_hook,
snapshot_builder,
snapshot_applier,
shuffle_receiver,
shuffle_producer,
shuffle_consumer,
shuffle_aggregator,
assign_remote_surrogate,
calvin_submit,
calvin_submit_inbox,
reserve_read,
release_reservation,
})
}