use log::{debug, error, info, warn};
use crate::communication::connection::{InMsgQueue, OutMsgQueue};
use crate::communication::messages::{
AlertMessage, ControlMessage, GenericPackageBuilder, SimpleCodifier, Messages,
NoCoordinatedMessage, PackageBuilder,
};
use crate::communication::router::Router;
use crate::coordination::{Stopwatch, Timer};
use crate::membership::client::{MembershipClient, MembershipManager};
use crate::membership::metadata::{CommunicationMetadata, NoCoordinatedMetadata, NodeMetadata, SerializableMetadata};
use crate::nodes::utils::{connect_with_members, SystemCommunicationService};
use crate::nodes::{CommunicationSDK, SerializationSDK, SystemClusterId, SystemNodeId};
use crate::properties::{init_buffer, EncodedMetadata, EncodedMetadataVec, SystemBufferVec, CLUSTER_NODE_COUNT};
pub trait ActuatorWorker {
fn actuate(&mut self, alert: AlertMessage);
}
pub struct ActuatorNode<
Clock1: Stopwatch,
Clock2: Timer,
Worker: ActuatorWorker,
CommsSDK: CommunicationSDK,
Membership: MembershipManager<CLUSTER_NODE_COUNT, NodeMetadata=NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>>,
SerdeSDK: SerializationSDK<
CommsSDK::Address,
NoCoordinatedMessage,
NoCoordinatedMetadata,
CommsSDK::Metadata,
>,
> {
id: SystemNodeId,
cluster: SystemClusterId,
keep_alive: Clock1,
idle_timer: Clock2,
worker: Worker,
metadata: NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>,
membership_manager: Membership,
communication_service: SystemCommunicationService<CommsSDK>,
message_codifier: SerdeSDK::GeneralMessageCodifier,
}
impl<
Clock1: Stopwatch,
Clock2: Timer,
Actuator: ActuatorWorker,
Membership: MembershipManager<CLUSTER_NODE_COUNT, NodeMetadata=NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>>,
CommsSDK: CommunicationSDK,
SerdeSDK: SerializationSDK<
CommsSDK::Address,
NoCoordinatedMessage,
NoCoordinatedMetadata,
CommsSDK::Metadata,
>,
> ActuatorNode<Clock1, Clock2, Actuator, CommsSDK, Membership, SerdeSDK>
where Membership: MembershipClient<Address=CommsSDK::Address>,
CommsSDK::Metadata: SerializableMetadata
{
pub fn new(
id: SystemNodeId,
cluster: SystemClusterId,
keep_alive: Clock1,
idle_timer: Clock2,
worker: Actuator,
metadata: NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>,
membership_manager: Membership,
communication_service: SystemCommunicationService<CommsSDK>,
message_codifier: SerdeSDK::GeneralMessageCodifier,
) -> ActuatorNode<Clock1, Clock2, Actuator, CommsSDK, Membership, SerdeSDK> {
if idle_timer.as_secs() > keep_alive.as_secs() {
error!("The idle-time ({}s) can't be higher than the keep-alive ({}s). Aborting.", idle_timer.as_secs(), keep_alive.as_secs());
panic!()
}
ActuatorNode {
id,
cluster,
idle_timer,
keep_alive,
worker,
metadata,
membership_manager,
communication_service,
message_codifier,
}
}
pub fn connect_with_members(&mut self) {
connect_with_members::<
NoCoordinatedMessage,
CommsSDK,
NoCoordinatedMetadata,
Membership,
SerdeSDK::GeneralMessageCodifier,
CLUSTER_NODE_COUNT,
>(
self.encode_metadata(),
&mut self.membership_manager,
&mut self.communication_service,
&mut self.message_codifier,
);
}
pub fn receive_message(&mut self) -> Option<(SystemNodeId, Messages<NoCoordinatedMessage>)> {
self.communication_service.pop().and_then(|pkg| {
Some((
pkg.header.from,
self.message_codifier.decode(pkg.body.as_slice()).ok()?,
))
})
}
fn encode_metadata(&self) -> EncodedMetadata {
let mut buffer = init_buffer();
let node = if let Ok(bytes) = SerdeSDK::SystemMetadataCodifier::default()
.encode(&self.metadata.node, &mut buffer)
{
EncodedMetadataVec::from_slice(&buffer[..bytes]).unwrap()
} else {
panic!("Couldn't serialize Node #{} system metadata.", self.id);
};
EncodedMetadata { node, coordination: None }
}
fn store_metadata(
&mut self,
previous_id: SystemNodeId,
encoded_metadata: EncodedMetadata,
) -> Option<SystemNodeId> {
let metadata = SerdeSDK::SystemMetadataCodifier::default().decode(encoded_metadata.node.as_slice()).ok()?;
if !metadata.same_cluster(self.cluster) {
info!("Ignoring metadata from node #{} (role: {:?}) because it belongs belongs to a different cluster/s: {}.", metadata.id, metadata.role, metadata.log_clusters());
debug!("Removing router entry from node #{} (temporary entry for node #{}).", previous_id, metadata.id);
self.communication_service.router.unregister(previous_id);
return None;
}
if self.id == metadata.id {
error!("Received metadata from a node with identical ID #{} (previously registerd as #{}). Aborting.", self.id, previous_id);
panic!();
}
info!("Updating metadata for node #{}: {}", metadata.id, metadata);
self.communication_service.update(metadata.id, previous_id, metadata.communication.address());
debug!("Registering node #{} with type {:?} into Membership Manager.", metadata.id, metadata.role);
let member_metadata = NodeMetadata { node: metadata.clone(), coordination: None };
self.membership_manager.register_member(metadata.id, member_metadata);
Some(metadata.id)
}
pub fn send_metadata(&mut self, to: SystemNodeId) -> Option<()> {
let metadata = self.encode_metadata();
let message: Messages<NoCoordinatedMessage> = Messages::System(ControlMessage::Metadata(metadata));
self.send_message(to, message)
}
pub fn send_message(&mut self, to: SystemNodeId, message: Messages<NoCoordinatedMessage>) -> Option<()> {
let mut buffer = init_buffer();
let bytes = self.message_codifier.encode(&message, &mut buffer).ok()?;
let encoded = SystemBufferVec::from_slice(&buffer[..bytes]).ok()?;
let package = GenericPackageBuilder::default()
.to(to)
.from(self.id)
.with_message(encoded)
.build();
self.communication_service.push(package.ok()?);
Some(())
}
pub fn run(&mut self) {
self.connect_with_members();
loop {
if self.keep_alive.is_timeout() {
self.membership_manager.keep_alive();
self.keep_alive.restart();
}
if let Some((id, message)) = self.receive_message() {
debug!("Message received from node #{}: {:?}", id, message);
match message {
Messages::Actuator(alert) => {
info!("Actuating over new alert from node #{id}.");
self.worker.actuate(alert);
}
Messages::System(ControlMessage::Metadata(metadata)) => {
info!("Metadata from node #{id} received.");
if let Some(new_id) = self.store_metadata(id, metadata) {
info!("Replying with ACK to node #{new_id}.");
let message: Messages<NoCoordinatedMessage> = Messages::System(ControlMessage::ACK);
if self.send_message(new_id, message).is_none() {
error!("Couldn't send ACK to node #{new_id}.");
}
}
}
Messages::System(ControlMessage::GetMetadata(metadata)) => {
info!("Metadata request from node #{id} received.");
if let Some(new_id) = self.store_metadata(id, metadata) {
info!("Replying with metadata to node #{new_id}.");
if self.send_metadata(new_id).is_none() {
error!("Couldn't send metadata to node #{new_id}.");
}
}
}
Messages::System(ControlMessage::Exit) => {
info!("Closing Actuator #{}.", self.id);
break;
}
Messages::System(ControlMessage::ACK) => info!("ACK received from node #{id}."),
_ => warn!("Unexpected message received: {message:?}")
}
} else {
self.idle_timer.wait();
}
}
}
}