use std::sync::Arc;
use crate::distribution::connection::DistConnection;
use crate::distribution::control_link::{
ControlOp, encode_exit_frame, encode_link_frame, encode_unlink_frame,
};
use crate::distribution::remote_link::RemoteLinkError;
use crate::distribution::sender::{ControlOutbound, DistSender};
use crate::etf::encode::EncodeError;
use crate::process::{ExitReason, RemotePid};
use super::SharedState;
pub(super) fn send_link(
shared: &SharedState,
caller_pid: u64,
target: RemotePid,
) -> Result<(), RemoteLinkError> {
if !wire_encodable(caller_pid, target) {
return Err(RemoteLinkError::BadTarget);
}
let Some(dist) = shared.distribution() else {
return Err(RemoteLinkError::NoConnection);
};
let Some(sender) = dist.sender() else {
return Err(RemoteLinkError::NoConnection);
};
let Some(connection) = dist.connections().get_connection(target.node) else {
return Err(RemoteLinkError::NoConnection);
};
let frame = encode_link_frame(
shared.local_node.name,
caller_pid,
target,
&shared.atom_table,
);
enqueue_pinned(sender, connection, frame);
Ok(())
}
pub(super) fn send_unlink(shared: &SharedState, caller_pid: u64, target: RemotePid) {
if !wire_encodable(caller_pid, target) {
return;
}
let Some(dist) = shared.distribution() else {
return;
};
let Some(sender) = dist.sender() else {
return;
};
let Some(connection) = dist.connections().get_connection(target.node) else {
return;
};
let frame = encode_unlink_frame(
shared.local_node.name,
caller_pid,
target,
&shared.atom_table,
);
enqueue_pinned(sender, connection, frame);
}
pub(super) fn send_exit_linked(
shared: &SharedState,
from_pid: u64,
target: RemotePid,
reason: ExitReason,
) {
send_exit(shared, ControlOp::Exit, from_pid, target, reason);
}
pub(super) fn send_exit2(
shared: &SharedState,
caller_pid: u64,
target: RemotePid,
reason: ExitReason,
) {
send_exit(shared, ControlOp::Exit2, caller_pid, target, reason);
}
fn send_exit(
shared: &SharedState,
op: ControlOp,
from_pid: u64,
target: RemotePid,
reason: ExitReason,
) {
if !wire_encodable(from_pid, target) {
return;
}
let Some(dist) = shared.distribution() else {
return;
};
let Some(sender) = dist.sender() else {
return;
};
let Some(connection) = dist.connections().get_connection(target.node) else {
return;
};
let frame = encode_exit_frame(
op,
shared.local_node.name,
from_pid,
target,
reason,
&shared.atom_table,
);
enqueue_pinned(sender, connection, frame);
}
fn wire_encodable(from_pid: u64, target: RemotePid) -> bool {
u32::try_from(from_pid).is_ok()
&& u32::try_from(target.pid_number).is_ok()
&& u32::try_from(target.serial).is_ok()
}
fn enqueue_pinned(
sender: &DistSender,
connection: Arc<DistConnection>,
frame: Result<Vec<u8>, EncodeError>,
) {
match frame {
Ok(frame) => {
let _ = sender.enqueue_control(ControlOutbound {
connection,
frame: Arc::from(frame.into_boxed_slice()),
});
}
Err(_) => {
connection.mark_down_control_overflow();
}
}
}