use std::sync::{Arc, Weak};
use crate::distribution::control::encode_pg_update_frame;
use crate::distribution::pg::{PgPropagation, PgUpdate};
use crate::distribution::sender::DistOutbound;
use super::SharedState;
pub(super) struct SchedulerPgPropagation {
pub(super) shared: Weak<SharedState>,
}
impl PgPropagation for SchedulerPgPropagation {
fn broadcast(&self, update: PgUpdate) {
let Some(shared) = self.shared.upgrade() else {
return;
};
let Some(sender) = &shared.dist_sender else {
return;
};
let local_node = shared.local_node.name;
let Ok(frame) = encode_pg_update_frame(update, local_node, &shared.atom_table) else {
return;
};
let frame: Arc<[u8]> = Arc::from(frame.into_boxed_slice());
for node in shared.distribution_connections.connected_nodes() {
sender.enqueue(DistOutbound::ToNode {
node,
frame: Arc::clone(&frame),
});
}
}
}