use crate::ha::HARole;
use crate::ha::{HAError, Result};
use crate::pubsub;
use core::sync::atomic::{AtomicBool, AtomicU64, Ordering};
#[cfg(feature = "log")]
use crate::log::{debug, error};
const HEARTBEAT_TOPIC: u16 = 3;
#[repr(C)]
pub struct HeartbeatPacket {
node_id: u64,
timestamp: u64,
role: u8,
crc32: u32,
}
fn handle_heartbeat_callback(topic_id: u16, data: &[u8]) -> bool {
if topic_id != HEARTBEAT_TOPIC {
return false;
}
#[cfg(feature = "log")]
debug!("Heartbeat callback received data, len: {}", data.len());
if let Some(packet) = HeartbeatPacket::from_bytes(data) {
if !packet.verify_crc() {
#[cfg(feature = "log")]
debug!("Heartbeat CRC check failed");
return true;
}
let node_id = packet.node_id();
let role = packet.role();
let timestamp = packet.timestamp();
#[cfg(feature = "log")]
debug!(
"Received heartbeat, node_id: {}, role: {:?}, timestamp: {}",
node_id, role, timestamp
);
}
true
}
pub struct HeartbeatMonitor {
node_id: u64,
role: HARole,
heartbeat_interval: u64,
failure_detection_time: u64,
last_heartbeat_time: AtomicU64,
master_alive: AtomicBool,
is_initialized: bool,
receiver_running: AtomicBool,
sender_running: AtomicBool,
}
impl HeartbeatPacket {
pub fn new(node_id: u64, role: HARole) -> Self {
let timestamp = crate::platform::get_timestamp_us() / 1000; let role_u8 = role as u8;
let mut packet = Self {
node_id,
timestamp,
role: role_u8,
crc32: 0,
};
packet.update_crc();
packet
}
pub fn update_crc(&mut self) {
self.crc32 = 0;
let mut crc_data = [0u8; 21];
unsafe {
core::ptr::copy_nonoverlapping(
&self.node_id as *const u64 as *const u8,
crc_data.as_mut_ptr(),
8,
);
core::ptr::copy_nonoverlapping(
&self.timestamp as *const u64 as *const u8,
crc_data.as_mut_ptr().add(8),
8,
);
*crc_data.as_mut_ptr().add(16) = self.role;
};
self.crc32 = crate::pubsub::crc32::calculate_crc32(&crc_data);
}
pub fn verify_crc(&self) -> bool {
let mut crc_data = [0u8; 21];
unsafe {
core::ptr::copy_nonoverlapping(
&self.node_id as *const u64 as *const u8,
crc_data.as_mut_ptr(),
8,
);
core::ptr::copy_nonoverlapping(
&self.timestamp as *const u64 as *const u8,
crc_data.as_mut_ptr().add(8),
8,
);
*crc_data.as_mut_ptr().add(16) = self.role;
};
let calculated_crc = crate::pubsub::crc32::calculate_crc32(&crc_data);
calculated_crc == self.crc32
}
pub fn to_bytes(&self) -> alloc::vec::Vec<u8> {
let mut bytes = [0u8; core::mem::size_of::<Self>()];
let (node_id_bytes, rest) = bytes.split_at_mut(8);
node_id_bytes.copy_from_slice(&self.node_id.to_le_bytes());
let (timestamp_bytes, rest) = rest.split_at_mut(8);
timestamp_bytes.copy_from_slice(&self.timestamp.to_le_bytes());
rest[0] = self.role;
let crc32_bytes = &mut rest[4..8];
crc32_bytes.copy_from_slice(&self.crc32.to_le_bytes());
bytes.to_vec()
}
pub fn from_bytes(bytes: &[u8]) -> Option<Self> {
if bytes.len() != core::mem::size_of::<Self>() {
return None;
}
let mut node_id_bytes = [0u8; 8];
node_id_bytes.copy_from_slice(&bytes[0..8]);
let node_id = u64::from_le_bytes(node_id_bytes);
let mut timestamp_bytes = [0u8; 8];
timestamp_bytes.copy_from_slice(&bytes[8..16]);
let timestamp = u64::from_le_bytes(timestamp_bytes);
let role = bytes[16];
let mut crc32_bytes = [0u8; 4];
crc32_bytes.copy_from_slice(&bytes[20..24]);
let crc32 = u32::from_le_bytes(crc32_bytes);
Some(Self {
node_id,
timestamp,
role,
crc32,
})
}
pub fn node_id(&self) -> u64 {
self.node_id
}
pub fn timestamp(&self) -> u64 {
self.timestamp
}
pub fn role(&self) -> HARole {
match self.role {
0 => HARole::Master,
1 => HARole::Slave,
2 => HARole::Auto,
_ => HARole::Auto,
}
}
}
impl HeartbeatMonitor {
pub fn new(heartbeat_interval: u64, failure_detection_time: u64) -> Result<Self> {
Ok(Self {
node_id: 0, role: HARole::Auto,
heartbeat_interval,
failure_detection_time,
last_heartbeat_time: AtomicU64::new(crate::platform::get_timestamp_us() / 1000),
master_alive: AtomicBool::new(true),
is_initialized: false,
receiver_running: AtomicBool::new(false),
sender_running: AtomicBool::new(false),
})
}
pub fn set_node_id(&mut self, node_id: u64) {
self.node_id = node_id;
}
pub fn set_role(&mut self, role: HARole) {
self.role = role;
}
pub fn init(&self) -> Result<()> {
#[cfg(feature = "log")]
debug!(
"Heartbeat monitor initialized, role: {:?}, node_id: {}",
self.role, self.node_id
);
self.init_pubsub()?;
#[cfg(feature = "log")]
debug!("Heartbeat monitor pubsub initialized");
Ok(())
}
pub fn init_master(&self) -> Result<()> {
#[cfg(feature = "log")]
debug!("Initializing master node, starting heartbeat sender");
self.start_heartbeat_sender()?;
#[cfg(feature = "log")]
debug!("Master node heartbeat sender started");
Ok(())
}
pub fn init_slave(&self) -> Result<()> {
self.start_heartbeat_receiver()?;
self.start_heartbeat_check()?;
Ok(())
}
fn init_pubsub(&self) -> Result<()> {
#[cfg(feature = "log")]
debug!("Heartbeat using existing pubsub system");
Ok(())
}
fn start_heartbeat_sender(&self) -> Result<()> {
if self.sender_running.load(Ordering::Relaxed) {
return Ok(());
}
self.sender_running.store(true, Ordering::Relaxed);
Ok(())
}
fn start_heartbeat_receiver(&self) -> Result<()> {
if self.receiver_running.load(Ordering::Relaxed) {
return Ok(());
}
self.receiver_running.store(true, Ordering::Relaxed);
match pubsub::subscribe(HEARTBEAT_TOPIC, handle_heartbeat_callback) {
Ok(_) => {
#[cfg(feature = "log")]
debug!("Successfully subscribed to heartbeat topic");
}
Err(e) => {
#[cfg(feature = "log")]
error!("Failed to subscribe to heartbeat topic: {:?}", e);
}
}
Ok(())
}
fn start_heartbeat_check(&self) -> Result<()> {
Ok(())
}
fn handle_heartbeat(&self, data: &[u8]) {
if let Some(packet) = HeartbeatPacket::from_bytes(data) {
if !packet.verify_crc() {
#[cfg(feature = "log")]
debug!("Heartbeat CRC check failed");
return;
}
let node_id = packet.node_id();
let role = packet.role();
let timestamp = packet.timestamp();
#[cfg(feature = "log")]
debug!(
"Received heartbeat, node_id: {}, role: {:?}, timestamp: {}",
node_id, role, timestamp
);
let now = crate::platform::get_timestamp_us() / 1000; self.last_heartbeat_time.store(now, Ordering::Relaxed);
self.master_alive.store(true, Ordering::Relaxed);
#[cfg(feature = "log")]
debug!("Updated master alive status to true");
}
}
pub fn check_status(&self) -> Result<()> {
if self.role != HARole::Slave {
return Ok(());
}
let now = crate::platform::get_timestamp_us() / 1000; let last_heartbeat = self.last_heartbeat_time.load(Ordering::Relaxed);
if now - last_heartbeat > self.failure_detection_time {
self.master_alive.store(false, Ordering::Relaxed);
return Err(HAError::HeartbeatTimeout);
}
Ok(())
}
fn send_heartbeat(&self) -> Result<()> {
let packet = HeartbeatPacket::new(self.node_id, self.role);
let timestamp = packet.timestamp();
#[cfg(feature = "log")]
debug!(
"Sending heartbeat, node_id: {}, role: {:?}, timestamp: {}",
self.node_id, self.role, timestamp
);
let bytes = packet.to_bytes();
let mut buffer = [0u8; core::mem::size_of::<HeartbeatPacket>()];
let copy_len = core::cmp::min(bytes.len(), buffer.len());
unsafe {
core::ptr::copy_nonoverlapping(bytes.as_ptr(), buffer.as_mut_ptr(), copy_len);
}
match pubsub::publish(HEARTBEAT_TOPIC, &buffer) {
Ok(_) => {
#[cfg(feature = "log")]
debug!("Heartbeat sent successfully");
Ok(())
}
Err(e) => {
#[cfg(feature = "log")]
error!("Failed to send heartbeat: {:?}", e);
Err(HAError::NetworkError)
}
}
}
pub fn shutdown(&self) -> Result<()> {
self.sender_running.store(false, Ordering::Relaxed);
self.receiver_running.store(false, Ordering::Relaxed);
Ok(())
}
pub fn is_master_alive(&self) -> bool {
self.master_alive.load(Ordering::Relaxed)
}
pub fn get_last_heartbeat_time(&self) -> u64 {
self.last_heartbeat_time.load(Ordering::Relaxed)
}
}