use std::sync::Arc;
use nodedb::ServerConfig;
use nodedb::bootstrap;
use nodedb::control::cluster::ClusterHandle;
use nodedb::control::state::SharedState;
pub(crate) struct BackgroundLoops {
pub(crate) raft_ready_rx: Option<tokio::sync::watch::Receiver<bool>>,
pub(crate) _lease_renewal: Option<tokio::task::JoinHandle<()>>,
pub(crate) _event_plane: nodedb::event::EventPlane,
}
pub(crate) struct BackgroundLoopsInputs<'a> {
pub(crate) cluster_handle: Option<&'a ClusterHandle>,
pub(crate) wal: Arc<nodedb::wal::WalManager>,
pub(crate) event_consumers: Vec<nodedb::event::bus::EventConsumerRx>,
pub(crate) watermark_store: Arc<nodedb::event::watermark::WatermarkStore>,
pub(crate) trigger_dlq: Arc<std::sync::Mutex<nodedb::event::trigger::TriggerDlq>>,
pub(crate) num_cores: usize,
}
pub(crate) fn spawn(
shared: &Arc<SharedState>,
config: &ServerConfig,
shutdown_rx: tokio::sync::watch::Receiver<bool>,
inputs: BackgroundLoopsInputs<'_>,
) -> anyhow::Result<BackgroundLoops> {
let BackgroundLoopsInputs {
cluster_handle,
wal,
event_consumers,
watermark_store,
trigger_dlq,
num_cores,
} = inputs;
let raft_ready_rx: Option<tokio::sync::watch::Receiver<bool>> =
if let Some(handle) = cluster_handle {
Some(nodedb::control::cluster::start_raft(
handle,
Arc::clone(shared),
&config.server.data_dir,
&config.tuning.cluster_transport,
)?)
} else {
None
};
let _lease_renewal = nodedb::control::lease::LeaseRenewalLoop::spawn(
Arc::clone(shared),
&config.tuning.cluster_transport,
shutdown_rx.clone(),
)
.map(|(join, metrics)| {
shared.loop_metrics_registry.register(metrics);
join
});
bootstrap::background_loops::spawn_response_poller(shared);
let _event_plane = bootstrap::background_loops::spawn_background_loops(
shared,
bootstrap::background_loops::EventPlaneComponents {
wal: Arc::clone(&wal),
event_consumers,
watermark_store,
trigger_dlq,
},
config,
num_cores,
shutdown_rx.clone(),
);
Ok(BackgroundLoops {
raft_ready_rx,
_lease_renewal,
_event_plane,
})
}