use crate::config::DbConfig;
use crate::ha::heartbeat::HeartbeatMonitor;
use crate::ha::replication::ReplicationManager;
use crate::ha::role::RoleManager;
use crate::ha::sync_handler::SyncHandler;
use crate::ha::sync_receiver::SyncReceiver;
use crate::ha::{HAError, HARole, ReplicationMode, Result, SyncState};
use crate::pubsub;
use crate::pubsub::{init as pubsub_init, PubSubConfig, UdpMode};
use crate::transaction::LogItem;
#[cfg(feature = "log")]
use crate::log::{debug, error, info};
pub struct HAManager {
config: &'static DbConfig,
role_manager: RoleManager,
replication_manager: ReplicationManager,
heartbeat_monitor: HeartbeatMonitor,
sync_handler: Option<SyncHandler>,
sync_receiver: Option<SyncReceiver>,
is_initialized: bool,
}
impl HAManager {
pub fn new(config: &'static DbConfig) -> Result<Self> {
let ha_config = config.ha_config.as_ref().ok_or(HAError::InvalidParameter)?;
let role_manager = RoleManager::new(ha_config.ha_role)?;
let replication_manager = ReplicationManager::new(ha_config.replication_mode)?;
let mut heartbeat_monitor = HeartbeatMonitor::new(
ha_config.heartbeat_interval_ms,
ha_config.failure_detection_ms,
)?;
heartbeat_monitor.set_node_id(ha_config.node_id as u64);
Ok(Self {
config,
role_manager,
replication_manager,
heartbeat_monitor,
sync_handler: None,
sync_receiver: None,
is_initialized: false,
})
}
pub fn init(&mut self) -> Result<()> {
#[cfg(feature = "log")]
debug!("Initializing HA manager");
#[cfg(feature = "log")]
debug!("Calling init_pubsub()");
match self.init_pubsub() {
Ok(_) => {
#[cfg(feature = "log")]
debug!("HA manager pubsub initialized successfully");
}
Err(e) => {
#[cfg(feature = "log")]
error!("HA manager pubsub initialization failed: {:?}", e);
return Err(e);
}
}
#[cfg(feature = "log")]
debug!("Calling role_manager.init()");
match self.role_manager.init() {
Ok(_) => {
#[cfg(feature = "log")]
debug!("HA manager role manager initialized successfully");
}
Err(e) => {
#[cfg(feature = "log")]
error!("HA manager role manager initialization failed: {:?}", e);
return Err(e);
}
}
#[cfg(feature = "log")]
debug!("Calling replication_manager.init()");
match self.replication_manager.init() {
Ok(_) => {
#[cfg(feature = "log")]
debug!("HA manager replication manager initialized successfully");
}
Err(e) => {
#[cfg(feature = "log")]
error!(
"HA manager replication manager initialization failed: {:?}",
e
);
return Err(e);
}
}
#[cfg(feature = "log")]
debug!("Calling heartbeat_monitor.init()");
match self.heartbeat_monitor.init() {
Ok(_) => {
#[cfg(feature = "log")]
debug!("HA manager heartbeat monitor initialized successfully");
}
Err(e) => {
#[cfg(feature = "log")]
error!(
"HA manager heartbeat monitor initialization failed: {:?}",
e
);
return Err(e);
}
}
let role = self.role_manager.get_role();
#[cfg(feature = "log")]
debug!("Current role: {:?}", role);
match role {
HARole::Master => {
#[cfg(feature = "log")]
debug!("Initializing as master node, calling init_master()");
match self.init_master() {
Ok(_) => {
#[cfg(feature = "log")]
debug!("Master node initialized successfully");
}
Err(e) => {
#[cfg(feature = "log")]
error!("Master node initialization failed: {:?}", e);
return Err(e);
}
}
}
HARole::Slave => {
#[cfg(feature = "log")]
debug!("Initializing as slave node, calling init_slave()");
match self.init_slave() {
Ok(_) => {
#[cfg(feature = "log")]
debug!("Slave node initialized successfully");
}
Err(e) => {
#[cfg(feature = "log")]
error!("Slave node initialization failed: {:?}", e);
return Err(e);
}
}
}
HARole::Auto => {
#[cfg(feature = "log")]
debug!("Initializing in auto mode, calling init_auto()");
match self.init_auto() {
Ok(_) => {
#[cfg(feature = "log")]
debug!("Auto mode initialized successfully");
}
Err(e) => {
#[cfg(feature = "log")]
error!("Auto mode initialization failed: {:?}", e);
return Err(e);
}
}
}
}
#[cfg(feature = "log")]
debug!("HA manager initialized successfully");
Ok(())
}
fn init_master(&mut self) -> Result<()> {
self.replication_manager.init_master()?;
self.heartbeat_monitor.init_master()?;
let mut sync_handler = SyncHandler::new();
sync_handler.init()?;
self.sync_handler = Some(sync_handler);
#[cfg(feature = "log")]
info!("Master sync handler initialized");
Ok(())
}
fn init_slave(&mut self) -> Result<()> {
self.replication_manager.init_slave()?;
self.heartbeat_monitor.init_slave()?;
let ha_config = self
.config
.ha_config
.as_ref()
.ok_or(HAError::InvalidParameter)?;
let mut sync_receiver = SyncReceiver::new(ha_config.node_id as u8);
sync_receiver.init()?;
self.sync_receiver = Some(sync_receiver);
#[cfg(feature = "log")]
info!("Slave sync receiver initialized");
self.connect_to_master()?;
Ok(())
}
fn init_auto(&mut self) -> Result<()> {
#[cfg(feature = "log")]
debug!("init_auto: Starting auto mode initialization");
let cluster_available = self.detect_existing_cluster()?;
if cluster_available {
#[cfg(feature = "log")]
debug!("init_auto: Existing cluster detected, initializing as slave");
self.role_manager.set_role(HARole::Slave)?;
self.init_slave()?;
} else {
#[cfg(feature = "log")]
debug!("init_auto: No existing cluster detected, initializing as master");
self.role_manager.set_role(HARole::Master)?;
self.init_master()?;
}
Ok(())
}
fn detect_existing_cluster(&self) -> Result<bool> {
#[cfg(feature = "log")]
debug!("detect_existing_cluster: Checking for existing cluster");
let probe_data = [0u8; 1]; let probe_topic = 99;
#[cfg(feature = "log")]
debug!("detect_existing_cluster: Sending cluster probe");
let _ = pubsub::publish(probe_topic, &probe_data);
#[cfg(feature = "std")]
std::thread::sleep(std::time::Duration::from_millis(100));
#[cfg(feature = "log")]
debug!("detect_existing_cluster: No existing cluster detected (test mode)");
Ok(false)
}
fn connect_to_master(&self) -> Result<()> {
#[cfg(feature = "log")]
debug!("connect_to_master: Starting connection to master");
let ha_config = self
.config
.ha_config
.as_ref()
.ok_or(HAError::InvalidParameter)?;
if ha_config.master_address.is_none() || ha_config.master_port.is_none() {
#[cfg(feature = "log")]
debug!("connect_to_master: Master address or port not set, skipping connection");
return Ok(());
}
let master_address = ha_config
.master_address
.expect("master_address must be set");
let master_port = ha_config.master_port.expect("master_port must be set");
#[cfg(feature = "log")]
debug!(
"connect_to_master: Connecting to master at {}:{}",
master_address, master_port
);
#[cfg(feature = "log")]
debug!("connect_to_master: Establishing connection to master");
#[cfg(feature = "log")]
debug!("connect_to_master: Requesting full sync from master");
self.replication_manager.request_full_sync()?;
#[cfg(feature = "log")]
debug!("connect_to_master: Starting to receive WAL logs");
#[cfg(feature = "log")]
debug!("connect_to_master: Connection to master completed successfully");
Ok(())
}
pub fn replicate_wal(&mut self, log_item: &LogItem) -> Result<()> {
if self.role_manager.get_role() != HARole::Master {
return Ok(());
}
self.replication_manager.replicate_wal(log_item)?;
Ok(())
}
pub fn check_status(&self) -> Result<()> {
match self.heartbeat_monitor.check_status() {
Err(HAError::HeartbeatTimeout) => {
self.handle_heartbeat_timeout()?;
}
Err(e) => {
return Err(e);
}
Ok(_) => {
}
}
self.replication_manager.check_status()?;
Ok(())
}
fn handle_heartbeat_timeout(&self) -> Result<()> {
if self.role_manager.get_role() != HARole::Slave {
return Ok(());
}
Ok(())
}
pub fn shutdown(&mut self) -> Result<()> {
if let Some(ref mut handler) = self.sync_handler {
handler.shutdown()?;
}
if let Some(ref mut receiver) = self.sync_receiver {
receiver.shutdown()?;
}
self.heartbeat_monitor.shutdown()?;
self.replication_manager.shutdown()?;
self.role_manager.shutdown()?;
if let Err(e) = crate::pubsub::shutdown() {
match e {
crate::pubsub::PubSubError::InitFailed => return Err(HAError::InitFailed),
crate::pubsub::PubSubError::NetworkError => return Err(HAError::NetworkError),
crate::pubsub::PubSubError::InvalidParameter => {
return Err(HAError::InvalidParameter)
}
_ => return Err(HAError::ReplicationError),
}
}
Ok(())
}
pub fn get_role(&self) -> HARole {
self.role_manager.get_role()
}
pub fn get_replication_mode(&self) -> ReplicationMode {
self.replication_manager.get_replication_mode()
}
pub fn get_sync_state(&self) -> SyncState {
if let Some(ref handler) = self.sync_handler {
handler.get_state().into()
} else if let Some(ref receiver) = self.sync_receiver {
receiver.get_state()
} else {
SyncState::Idle
}
}
pub fn get_replication_manager(&self) -> &ReplicationManager {
&self.replication_manager
}
pub fn get_replication_manager_mut(&mut self) -> &mut ReplicationManager {
&mut self.replication_manager
}
pub fn promote_to_master(&mut self) -> Result<()> {
if self.role_manager.get_role() == HARole::Master {
return Ok(());
}
self.role_manager.set_role(HARole::Master)?;
self.heartbeat_monitor.set_role(HARole::Master);
self.init_master()?;
Ok(())
}
pub fn demote_to_slave(&mut self) -> Result<()> {
if self.role_manager.get_role() == HARole::Slave {
return Ok(());
}
self.role_manager.set_role(HARole::Slave)?;
self.heartbeat_monitor.set_role(HARole::Slave);
let _ = self.init_slave();
Ok(())
}
fn init_pubsub(&self) -> Result<()> {
let ha_config = self
.config
.ha_config
.as_ref()
.ok_or(HAError::InvalidParameter)?;
let pubsub_config = PubSubConfig {
udp_mode: UdpMode::Broadcast,
multicast_addr: None,
port: ha_config.replication_port, max_topics: 16, max_subscribers_per_topic: 16,
buffer_size: 8192,
enable_nack: true,
retransmit_timeout: core::time::Duration::from_millis(100),
max_retransmits: 3,
heartbeat_interval: core::time::Duration::from_secs(10),
frame_pool_size: 256,
};
match pubsub_init(pubsub_config) {
Ok(_) => {}
Err(e) => {
match e {
crate::pubsub::PubSubError::InitFailed => return Err(HAError::InitFailed),
crate::pubsub::PubSubError::NetworkError => return Err(HAError::NetworkError),
crate::pubsub::PubSubError::InvalidParameter => {
return Err(HAError::InvalidParameter)
}
_ => return Err(HAError::InitFailed), }
}
}
Ok(())
}
}