use heapless::Vec;
use log::{debug, error, info, warn};
use crate::communication::connection::{InMsgQueue, OutMsgQueue};
use crate::communication::messages::{ControlMessage, GenericPackageBuilder, SimpleCodifier, Messages, NoCoordinatedMessage, PackageBuilder, EventMessage};
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, CLUSTERS_PER_TRANSLATOR, EncodedMetadata, SystemBufferVec, CLUSTER_NODE_COUNT};
const TRANSLATOR_NODE_COUNT: usize = CLUSTER_NODE_COUNT * CLUSTERS_PER_TRANSLATOR;
pub trait ServerModule {
fn inbound(&mut self) -> Option<(SystemNodeId, Messages<NoCoordinatedMessage>)>;
fn outbound(&mut self, message: EventMessage);
}
pub struct TranslatorNode<
Clock1: Stopwatch,
Clock2: Timer,
CommsSDK: CommunicationSDK,
SerdeSDK: SerializationSDK<
CommsSDK::Address,
NoCoordinatedMessage,
NoCoordinatedMetadata,
CommsSDK::Metadata,
>,
Membership: MembershipManager<TRANSLATOR_NODE_COUNT, NodeMetadata=NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>>,
ServerHandler: ServerModule,
> {
id: SystemNodeId,
clusters: Vec<SystemClusterId, CLUSTERS_PER_TRANSLATOR>,
keep_alive: Clock1,
idle_timer: Clock2,
metadata: EncodedMetadata,
node_codifier: SerdeSDK::SystemMetadataCodifier,
message_codifier: SerdeSDK::GeneralMessageCodifier,
membership_manager: Membership,
communication_service: SystemCommunicationService<CommsSDK>,
server_handler: ServerHandler,
}
impl<
Clock1: Stopwatch,
Clock2: Timer,
CommsSDK: CommunicationSDK,
SerdeSDK: SerializationSDK<
CommsSDK::Address,
NoCoordinatedMessage,
NoCoordinatedMetadata,
CommsSDK::Metadata,
>,
Membership: MembershipManager<TRANSLATOR_NODE_COUNT, NodeMetadata=NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>>,
ServerHandler: ServerModule,
> TranslatorNode<Clock1, Clock2, CommsSDK, SerdeSDK, Membership, ServerHandler>
where Membership: MembershipClient<Address=CommsSDK::Address>,
CommsSDK::Metadata: SerializableMetadata,
{
#[allow(clippy::too_many_arguments)]
pub fn new(
id: SystemNodeId,
clusters: Vec<SystemClusterId, CLUSTERS_PER_TRANSLATOR>,
keep_alive: Clock1,
idle_timer: Clock2,
metadata: EncodedMetadata,
node_codifier: SerdeSDK::SystemMetadataCodifier,
message_codifier: SerdeSDK::GeneralMessageCodifier,
membership_manager: Membership,
communication_service: SystemCommunicationService<CommsSDK>,
server_handler: ServerHandler,
) -> TranslatorNode<Clock1, Clock2, CommsSDK, SerdeSDK, Membership, ServerHandler> {
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!()
}
TranslatorNode {
id,
clusters,
keep_alive,
idle_timer,
metadata,
node_codifier,
message_codifier,
membership_manager,
communication_service,
server_handler,
}
}
fn store_metadata(
&mut self,
previous_id: SystemNodeId,
encoded_metadata: EncodedMetadata,
) -> Option<SystemNodeId> {
let metadata = self.node_codifier.decode(encoded_metadata.node.as_slice()).ok()?;
if !self.clusters.iter().any(|cluster| metadata.same_cluster(*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)
}
fn send_message(&mut self, peer_id: SystemNodeId, message: Messages<NoCoordinatedMessage>) {
let mut buffer = init_buffer();
if let Ok(codified) = self.message_codifier.encode(&message, &mut buffer) {
if let Ok(vec) = SystemBufferVec::from_slice(&buffer[..codified]) {
let package = GenericPackageBuilder::default()
.to(peer_id)
.from(self.id)
.with_message(vec)
.build();
self.communication_service.push(package.unwrap());
}
}
}
fn check_for_server_messages(&mut self) -> bool {
return if let Some((id, message)) = self.server_handler.inbound() {
info!("About to send a message to node #{}.", id);
debug!("Message: {:?}", message);
self.send_message(id, message);
true
} else {
false
}
}
fn check_for_nodes_messages(&mut self) -> bool {
if let Some(package) = self.communication_service.pop() {
if let Some(message) = self.message_codifier.decode(&package.body).ok() {
let sender_id = package.header.from;
match message {
Messages::System(ControlMessage::Metadata(encoded_metadata)) => {
info!("Metadata from node #{sender_id} received.");
if let Some(new_id) = self.store_metadata(sender_id, encoded_metadata) {
info!("Replying with ACK to node #{new_id}.");
let response = Messages::System(ControlMessage::ACK);
self.send_message(new_id, response);
}
}
Messages::System(ControlMessage::GetMetadata(encoded_metadata)) => {
info!("Metadata request from node #{sender_id} received.");
if let Some(new_id) = self.store_metadata(sender_id, encoded_metadata) {
info!("Replying with metadata to node #{new_id}.");
let response = Messages::System(ControlMessage::Metadata(self.metadata.clone()));
self.send_message(new_id, response);
}
}
Messages::Reporter(message) => {
info!("Received event from node #{sender_id}.");
self.server_handler.outbound(message);
}
Messages::System(ControlMessage::ACK) => info!("ACK received from node #{sender_id}."),
_ => warn!("Unexpected message received: {message:?}"),
}
} else {
warn!("Couldn't decode package: {package:?}");
}
true
} else {
false
}
}
pub fn connect_with_members(&mut self) {
connect_with_members::<
NoCoordinatedMessage,
CommsSDK,
NoCoordinatedMetadata,
Membership,
SerdeSDK::GeneralMessageCodifier,
TRANSLATOR_NODE_COUNT,
>(
self.metadata.clone(),
&mut self.membership_manager,
&mut self.communication_service,
&mut self.message_codifier,
);
}
pub fn run(&mut self) {
self.connect_with_members();
info!("Starting to process messages.");
loop {
if self.keep_alive.is_timeout() {
self.membership_manager.keep_alive();
self.keep_alive.restart();
}
let received1 = self.check_for_nodes_messages();
let received2 = self.check_for_server_messages();
if !received1 && !received2 {
self.idle_timer.wait();
}
}
}
}