use std::sync::Arc;
use std::thread::JoinHandle;
use tracing::{info, warn};
use crate::data::eventfd::{EventFd, EventFdNotifier};
use crate::data::executor::core_loop::CoreLoop;
use super::boot_replay::replay_wal_and_rebuild_indexes;
use super::boot_restore::load_boot_checkpoints;
use super::boot_seed::seed_catalog_state;
use super::event_loop::run_event_loop;
use super::params::SpawnCoreParams;
pub fn spawn_core(
params: SpawnCoreParams<'_>,
) -> std::io::Result<(JoinHandle<()>, EventFdNotifier)> {
let SpawnCoreParams {
core_id,
request_rx,
response_tx,
data_dir,
wal_records,
tombstones,
num_cores,
compaction_config,
system_metrics,
event_producer,
governor,
quiesce,
hlc,
array_catalog,
quarantine_registry,
maintenance_budget,
doc_config_seed,
vector_index_param_seed,
columnar_schema_seed,
replay_done,
} = params;
let data_dir = data_dir.to_path_buf();
let efd = EventFd::new().map_err(std::io::Error::other)?;
let notifier = efd.notifier();
let handle = std::thread::Builder::new()
.name(format!("data-core-{core_id}"))
.spawn(move || {
match nodedb_mem::arena::pin_thread_arena(core_id as u32) {
Ok(arena) => info!(core_id, arena, "pinned to jemalloc arena"),
Err(e) => warn!(core_id, error = %e, "failed to pin jemalloc arena, continuing with default"),
}
let mut core = CoreLoop::open_with_array_catalog(
core_id,
request_rx,
response_tx,
&data_dir,
hlc,
array_catalog,
)
.expect("failed to open CoreLoop engines");
wire_core_dependencies(
&mut core,
WiredDependencies {
governor,
maintenance_budget,
system_metrics,
event_producer,
quiesce,
quarantine_registry,
},
);
let checkpoint_interval = compaction_config.checkpoint_interval;
core.set_compaction_config(
compaction_config.interval,
compaction_config.tombstone_threshold,
);
core.set_query_tuning(compaction_config.query);
core.set_graph_tuning(compaction_config.graph);
core.set_timeseries_tuning(compaction_config.timeseries);
load_boot_checkpoints(&mut core)
.expect("boot checkpoint load failed: corrupt or unreadable checkpoint");
seed_catalog_state(
&mut core,
&doc_config_seed,
&vector_index_param_seed,
&columnar_schema_seed,
);
replay_wal_and_rebuild_indexes(
&mut core,
&wal_records,
num_cores,
&tombstones,
&vector_index_param_seed,
);
let _ = replay_done.send(());
info!(core_id, "data plane core started (eventfd-driven)");
run_event_loop(&mut core, core_id, &efd, checkpoint_interval);
})?;
Ok((handle, notifier))
}
struct WiredDependencies {
governor: Arc<nodedb_mem::MemoryGovernor>,
maintenance_budget: Arc<crate::control::maintenance::MaintenanceBudgetTracker>,
system_metrics: Option<Arc<crate::control::metrics::SystemMetrics>>,
event_producer: Option<crate::event::bus::EventProducer>,
quiesce: Option<Arc<crate::bridge::quiesce::CollectionQuiesce>>,
quarantine_registry: Arc<crate::storage::quarantine::QuarantineRegistry>,
}
fn wire_core_dependencies(core: &mut CoreLoop, deps: WiredDependencies) {
let WiredDependencies {
governor,
maintenance_budget,
system_metrics,
event_producer,
quiesce,
quarantine_registry,
} = deps;
core.set_governor(governor);
core.set_maintenance_budget(maintenance_budget);
if let Some(m) = system_metrics {
core.set_metrics(m);
}
if let Some(ep) = event_producer {
core.set_event_producer(ep);
}
if let Some(q) = quiesce {
core.set_quiesce(q);
}
core.set_quarantine_registry(quarantine_registry);
}