use std::collections::HashMap;
use std::future::Future;
use std::panic::AssertUnwindSafe;
use std::pin::Pin;
use std::sync::{Arc, Mutex};
use futures_util::stream::{FuturesUnordered, StreamExt};
use futures_util::FutureExt;
use tokio::sync::mpsc::{self, UnboundedReceiver, UnboundedSender};
use tokio::sync::oneshot;
use crate::engine::Departure;
use crate::interfaces::IfacContext;
use crate::interfaces::{
ConnectionView, InterfaceId, InterfaceKind, InterfaceOriginKind, InterfaceSnapshot, Membership,
ReportsStatus, StatusView,
};
use crate::manifold::driver::{
tokio_grant_lane, AddInterfaceCommand, HostCommand, TokioInterfaceSeam,
};
use crate::manifold::interface_seam::{frame_cap_for, Interface};
use crate::node_introspection::{InterfaceIfacSnapshot, InterfaceInventoryEntry};
use super::super::ManuallyAttached;
use super::PrnsNodeHandle;
const HOST_LANE_DEPTH: usize = 256;
fn lane_depth_for(_slot_cap: usize) -> usize {
HOST_LANE_DEPTH
}
#[derive(Clone)]
struct RuntimeIfac {
context: IfacContext,
network_name: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct InterfaceAttachmentMetadata {
pub name: Option<String>,
pub origin: InterfaceOriginKind,
}
#[derive(Clone, Copy)]
struct InterfacePlacement {
membership: Membership,
origin: InterfaceOriginKind,
}
impl RuntimeIfac {
fn snapshot(&self) -> InterfaceIfacSnapshot {
InterfaceIfacSnapshot {
signature: self.context.ifac_signature(),
size: self.context.ifac_size(),
network_name: self.network_name.clone(),
}
}
}
impl PrnsNodeHandle {
pub fn add_interface<I>(&self, interface: I) -> AttachedInterface
where
I: Interface + ReportsStatus + Send + 'static,
{
self.add_interface_with_metadata(
interface,
InterfaceAttachmentMetadata {
name: None,
origin: InterfaceOriginKind::Configured,
},
)
}
pub fn add_interface_with_ifac<I>(&self, interface: I, ifac: IfacContext) -> AttachedInterface
where
I: Interface + ReportsStatus + Send + 'static,
{
self.add_interface_with_ifac_name(interface, ifac, None)
}
pub fn add_interface_with_ifac_name<I>(
&self,
interface: I,
ifac: IfacContext,
network_name: Option<String>,
) -> AttachedInterface
where
I: Interface + ReportsStatus + Send + 'static,
{
self.add_interface_with_metadata_and_ifac_name(
interface,
InterfaceAttachmentMetadata {
name: None,
origin: InterfaceOriginKind::Configured,
},
ifac,
network_name,
)
}
pub fn add_interface_with_metadata<I>(
&self,
interface: I,
metadata: InterfaceAttachmentMetadata,
) -> AttachedInterface
where
I: Interface + ReportsStatus + Send + 'static,
{
self.add_interface_access(interface, metadata, None)
}
pub fn add_interface_with_metadata_and_ifac_name<I>(
&self,
interface: I,
metadata: InterfaceAttachmentMetadata,
ifac: IfacContext,
network_name: Option<String>,
) -> AttachedInterface
where
I: Interface + ReportsStatus + Send + 'static,
{
self.add_interface_access(
interface,
metadata,
Some(RuntimeIfac {
context: ifac,
network_name,
}),
)
}
fn add_interface_access<I>(
&self,
interface: I,
metadata: InterfaceAttachmentMetadata,
ifac: Option<RuntimeIfac>,
) -> AttachedInterface
where
I: Interface + ReportsStatus + Send + 'static,
{
let InterfaceAttachmentMetadata { name, origin } = metadata;
let placement = InterfacePlacement {
membership: Membership::Independent,
origin,
};
let descriptor = interface.descriptor();
let view = interface.status_view();
let connection = interface.connection_view();
let attached = attach_interface(
&self.commands,
&self.iface_build,
&self.notify_tx,
interface,
InterfaceWiring {
descriptor,
placement,
connection,
ifac: ifac.as_ref().map(|access| access.context.clone()),
},
);
register_status(
&self.interfaces,
attached.id(),
view.map(|view| RegisteredInterface {
view,
placement,
mode: descriptor.mode,
gravity: descriptor.gravity,
ifac: ifac.as_ref().map(RuntimeIfac::snapshot),
name,
}),
);
attached
}
#[must_use]
pub fn interface_inventory(&self) -> std::vec::Vec<InterfaceInventoryEntry> {
let Ok(map) = self.interfaces.lock() else {
return std::vec::Vec::new();
};
map.values()
.flat_map(|registered| {
let placement = registered.placement;
let ifac = registered.ifac.clone();
let name = registered.name.clone();
(registered.view)().into_iter().map(move |vitals| {
let counts = self.store.counts(vitals.id);
InterfaceInventoryEntry {
name: name.clone(),
origin: placement.origin,
snapshot: InterfaceSnapshot {
id: vitals.id,
mode: registered.mode,
gravity: registered.gravity,
connection: vitals.connection,
failure_reason: vitals.failure_reason,
rx_bytes: vitals.rx_bytes,
tx_bytes: vitals.tx_bytes,
transfer_rates: vitals.transfer_rates,
destinations: counts.destinations,
links: counts.links,
transported_links: counts.transported_links,
membership: placement.membership,
},
ifac: ifac.clone(),
}
})
})
.collect()
}
#[must_use]
pub fn set_interface_name(&self, id: InterfaceId, name: impl Into<String>) -> bool {
let Ok(mut interfaces) = self.interfaces.lock() else {
return false;
};
let Some(interface) = interfaces.get_mut(&id) else {
return false;
};
interface.name = Some(name.into());
true
}
#[must_use]
pub fn interfaces(&self) -> std::vec::Vec<InterfaceSnapshot> {
self.interface_inventory()
.into_iter()
.map(|entry| entry.snapshot)
.collect()
}
pub fn supervise<S>(&self, supervisor: S) -> AttachedSupervisor
where
S: InterfaceSupervisor + ReportsStatus + Send + 'static,
{
self.supervise_access(supervisor, None)
}
pub fn supervise_with_ifac<S>(&self, supervisor: S, ifac: IfacContext) -> AttachedSupervisor
where
S: InterfaceSupervisor + ReportsStatus + Send + 'static,
{
self.supervise_with_ifac_name(supervisor, ifac, None)
}
pub fn supervise_with_ifac_name<S>(
&self,
supervisor: S,
ifac: IfacContext,
network_name: Option<String>,
) -> AttachedSupervisor
where
S: InterfaceSupervisor + ReportsStatus + Send + 'static,
{
self.supervise_access(
supervisor,
Some(RuntimeIfac {
context: ifac,
network_name,
}),
)
}
fn supervise_access<S>(&self, supervisor: S, ifac: Option<RuntimeIfac>) -> AttachedSupervisor
where
S: InterfaceSupervisor + ReportsStatus + Send + 'static,
{
let id = InterfaceId::from_channel_tag(S::KIND, supervisor.channel_tag());
let placement = InterfacePlacement {
membership: Membership::Independent,
origin: InterfaceOriginKind::Configured,
};
let policy = supervisor.policy();
let view = supervisor.status_view();
let ifac_status = ifac.as_ref().map(RuntimeIfac::snapshot);
let fleet = Fleet {
supervisor_id: id,
commands: self.commands.clone(),
iface_build: self.iface_build.clone(),
notify_tx: self.notify_tx.clone(),
interfaces: self.interfaces.clone(),
ifac,
entropy: self.entropy,
};
let build: Box<dyn FnOnce() -> Pin<Box<dyn Future<Output = ()>>> + Send> =
Box::new(move || Box::pin(supervisor.run(fleet)));
let _ = self.iface_build.send(DriverMsg::Add {
id,
supervisor: None,
build,
});
register_status(
&self.interfaces,
id,
view.map(|view| RegisteredInterface {
view,
placement,
mode: policy.mode,
gravity: policy.gravity,
ifac: ifac_status,
name: None,
}),
);
AttachedSupervisor {
id,
iface_build: self.iface_build.clone(),
}
}
pub fn remove_interface(&self, id: InterfaceId) {
let _ = self.commands.send(HostCommand::RemoveInterface {
id,
departure: Departure::MayReturn,
});
let _ = self.iface_build.send(DriverMsg::Stop { id });
}
pub fn forget_interface(&self, id: InterfaceId) {
let _ = self.commands.send(HostCommand::RemoveInterface {
id,
departure: Departure::Forgotten,
});
let _ = self.iface_build.send(DriverMsg::Stop { id });
}
pub fn attach<A: Attachable>(&self, attachable: A) -> A::Attached {
attachable.attach_to(self)
}
pub fn attach_with_ifac<A: Attachable>(&self, attachable: A, ifac: IfacContext) -> A::Attached {
attachable.attach_to_with_ifac(self, ifac, None)
}
pub fn attach_with_ifac_name<A: Attachable>(
&self,
attachable: A,
ifac: IfacContext,
network_name: Option<String>,
) -> A::Attached {
attachable.attach_to_with_ifac(self, ifac, network_name)
}
}
pub trait Attachable {
type Attached;
fn attach_to(self, handle: &PrnsNodeHandle) -> Self::Attached;
fn attach_to_with_ifac(
self,
handle: &PrnsNodeHandle,
ifac: IfacContext,
network_name: Option<String>,
) -> Self::Attached;
}
pub trait AttachIntent {
fn attach(self, handle: &PrnsNodeHandle);
}
impl AttachIntent for ManuallyAttached {
fn attach(self, _handle: &PrnsNodeHandle) {}
}
impl<F: FnOnce(&PrnsNodeHandle)> AttachIntent for F {
fn attach(self, handle: &PrnsNodeHandle) {
self(handle)
}
}
pub struct AttachedInterface {
id: InterfaceId,
commands: UnboundedSender<HostCommand>,
iface_build: UnboundedSender<DriverMsg>,
}
impl AttachedInterface {
#[must_use]
pub fn id(&self) -> InterfaceId {
self.id
}
pub fn teardown(self) {
let _ = self.commands.send(HostCommand::RemoveInterface {
id: self.id,
departure: Departure::MayReturn,
});
let _ = self.iface_build.send(DriverMsg::Stop { id: self.id });
}
}
pub struct AttachedSupervisor {
id: InterfaceId,
iface_build: UnboundedSender<DriverMsg>,
}
impl AttachedSupervisor {
#[must_use]
pub fn id(&self) -> InterfaceId {
self.id
}
pub fn teardown(self) {
let _ = self.iface_build.send(DriverMsg::Stop { id: self.id });
}
}
struct InterfaceWiring {
descriptor: crate::interfaces::InterfaceDescriptor,
placement: InterfacePlacement,
connection: Option<ConnectionView>,
ifac: Option<IfacContext>,
}
fn attach_interface<I>(
commands: &UnboundedSender<HostCommand>,
iface_build: &UnboundedSender<DriverMsg>,
notify_tx: &UnboundedSender<InterfaceId>,
interface: I,
wiring: InterfaceWiring,
) -> AttachedInterface
where
I: Interface + Send + 'static,
{
let InterfaceWiring {
descriptor,
placement,
connection,
ifac,
} = wiring;
let id = descriptor.id;
let supervisor = match placement.membership {
Membership::Independent => None,
Membership::FleetMember { supervisor_id } => Some(supervisor_id),
};
let logical_interface = supervisor.unwrap_or(id);
let slot_cap = frame_cap_for(&descriptor);
let depth = lane_depth_for(slot_cap);
let (in_producer, in_consumer) = tokio_grant_lane(slot_cap, depth);
let (out_producer, out_consumer) = tokio_grant_lane(slot_cap, depth);
let seam = TokioInterfaceSeam::new(id, in_producer, notify_tx.clone(), out_consumer)
.with_origin(placement.origin)
.with_commands(commands.clone());
let build: Box<dyn FnOnce() -> Pin<Box<dyn Future<Output = ()>>> + Send> =
Box::new(move || Box::pin(interface.run(seam)));
let _ = commands.send(HostCommand::AddInterface(AddInterfaceCommand {
descriptor,
logical_interface,
inbound: in_consumer,
egress: out_producer,
connection,
ifac,
}));
let _ = iface_build.send(DriverMsg::Add {
id,
supervisor,
build,
});
AttachedInterface {
id,
commands: commands.clone(),
iface_build: iface_build.clone(),
}
}
pub struct Fleet {
supervisor_id: InterfaceId,
commands: UnboundedSender<HostCommand>,
iface_build: UnboundedSender<DriverMsg>,
notify_tx: UnboundedSender<InterfaceId>,
interfaces: Arc<Mutex<HashMap<InterfaceId, RegisteredInterface>>>,
ifac: Option<RuntimeIfac>,
entropy: crate::manifold::driver::TokioEntropy,
}
impl Fleet {
pub fn fill_entropy(&self, bytes: &mut [u8]) {
self.entropy.fill(bytes);
}
pub fn add<I>(&self, interface: I) -> AttachedInterface
where
I: Interface + ReportsStatus + Send + 'static,
{
let view = interface.status_view();
let connection = interface.connection_view();
let descriptor = interface.descriptor();
let placement = InterfacePlacement {
membership: Membership::FleetMember {
supervisor_id: self.supervisor_id,
},
origin: InterfaceOriginKind::Configured,
};
let attached = attach_interface(
&self.commands,
&self.iface_build,
&self.notify_tx,
interface,
InterfaceWiring {
descriptor,
placement,
connection,
ifac: self.ifac.as_ref().map(|access| access.context.clone()),
},
);
register_status(
&self.interfaces,
attached.id(),
view.map(|view| RegisteredInterface {
view,
placement,
mode: descriptor.mode,
gravity: descriptor.gravity,
ifac: self.ifac.as_ref().map(RuntimeIfac::snapshot),
name: None,
}),
);
attached
}
#[must_use]
pub fn detached(supervisor_id: InterfaceId) -> (Self, DetachedFleet) {
let (commands, commands_rx) = mpsc::unbounded_channel();
let (iface_build, iface_build_rx) = mpsc::unbounded_channel();
let (notify_tx, notify_rx) = mpsc::unbounded_channel();
let fleet = Fleet {
supervisor_id,
commands,
iface_build,
notify_tx,
interfaces: Arc::new(Mutex::new(HashMap::new())),
ifac: None,
entropy: crate::manifold::driver::TokioEntropy,
};
let tail = DetachedFleet {
_commands: commands_rx,
_iface_build: iface_build_rx,
_notify: notify_rx,
};
(fleet, tail)
}
}
pub struct DetachedFleet {
_commands: UnboundedReceiver<HostCommand>,
_iface_build: UnboundedReceiver<DriverMsg>,
_notify: UnboundedReceiver<InterfaceId>,
}
#[allow(async_fn_in_trait)]
pub trait InterfaceSupervisor {
const KIND: InterfaceKind;
fn channel_tag(&self) -> &[u8];
fn policy(&self) -> crate::interfaces::EffectiveInterfacePolicy;
async fn run(self, fleet: Fleet);
}
pub(super) enum DriverMsg {
Add {
id: InterfaceId,
supervisor: Option<InterfaceId>,
build: Box<dyn FnOnce() -> Pin<Box<dyn Future<Output = ()>>> + Send>,
},
Stop {
id: InterfaceId,
},
}
pub(super) async fn drive_interfaces(
initial: std::vec::Vec<Pin<Box<dyn Future<Output = ()>>>>,
mut messages: UnboundedReceiver<DriverMsg>,
commands: UnboundedSender<HostCommand>,
interfaces: Arc<Mutex<HashMap<InterfaceId, RegisteredInterface>>>,
) {
let mut futures: FuturesUnordered<Pin<Box<dyn Future<Output = Option<InterfaceId>>>>> = initial
.into_iter()
.map(
|run| -> Pin<Box<dyn Future<Output = Option<InterfaceId>>>> {
Box::pin(async move {
run.await;
None
})
},
)
.collect();
let mut stops: HashMap<InterfaceId, oneshot::Sender<()>> = HashMap::new();
let mut supervisor_of: HashMap<InterfaceId, InterfaceId> = HashMap::new();
let mut open = true;
loop {
if !open && futures.is_empty() {
return;
}
tokio::select! {
message = messages.recv(), if open => match message {
Some(DriverMsg::Add { id, supervisor, build }) => {
if let Some(supervisor_id) = supervisor {
let _ = supervisor_of.insert(id, supervisor_id);
}
let (stop_tx, stop_rx) = oneshot::channel();
let run = std::panic::catch_unwind(AssertUnwindSafe(build));
let guarded: Pin<Box<dyn Future<Output = Option<InterfaceId>>>> = match run {
Ok(run) => Box::pin(async move {
tokio::select! {
_ = AssertUnwindSafe(run).catch_unwind() => {}
_ = stop_rx => {}
}
Some(id)
}),
Err(_) => Box::pin(async move { Some(id) }),
};
futures.push(guarded);
stops.insert(id, stop_tx);
}
Some(DriverMsg::Stop { id }) => {
let stopped = stop_interface(&mut stops, id);
supervisor_of.remove(&id);
forget_status(&interfaces, id);
stop_supervised_members(
&mut stops,
&mut supervisor_of,
&interfaces,
&commands,
id,
);
if stopped {
drain_stopped_interface(
&mut futures,
&mut stops,
&mut supervisor_of,
&interfaces,
&commands,
id,
)
.await;
}
}
None => open = false,
},
done = futures.next(), if !futures.is_empty() => {
if let Some(Some(id)) = done {
complete_interface(
&mut stops,
&mut supervisor_of,
&interfaces,
&commands,
id,
);
}
}
}
}
}
pub(super) struct RegisteredInterface {
view: StatusView,
placement: InterfacePlacement,
mode: crate::interfaces::InterfaceMode,
gravity: crate::interfaces::InterfaceGravity,
ifac: Option<InterfaceIfacSnapshot>,
name: Option<String>,
}
fn register_status(
interfaces: &Arc<Mutex<HashMap<InterfaceId, RegisteredInterface>>>,
id: InterfaceId,
registered: Option<RegisteredInterface>,
) {
if let (Some(registered), Ok(mut map)) = (registered, interfaces.lock()) {
map.insert(id, registered);
}
}
fn forget_status(
interfaces: &Arc<Mutex<HashMap<InterfaceId, RegisteredInterface>>>,
id: InterfaceId,
) {
if let Ok(mut map) = interfaces.lock() {
map.remove(&id);
}
}
fn stop_interface(stops: &mut HashMap<InterfaceId, oneshot::Sender<()>>, id: InterfaceId) -> bool {
if let Some(stop) = stops.remove(&id) {
let _ = stop.send(());
true
} else {
false
}
}
async fn drain_stopped_interface(
futures: &mut FuturesUnordered<Pin<Box<dyn Future<Output = Option<InterfaceId>>>>>,
stops: &mut HashMap<InterfaceId, oneshot::Sender<()>>,
supervisor_of: &mut HashMap<InterfaceId, InterfaceId>,
interfaces: &Arc<Mutex<HashMap<InterfaceId, RegisteredInterface>>>,
commands: &UnboundedSender<HostCommand>,
stopped_id: InterfaceId,
) {
while let Some(done) = futures.next().await {
let Some(id) = done else {
continue;
};
complete_interface(stops, supervisor_of, interfaces, commands, id);
if id == stopped_id {
return;
}
}
}
fn complete_interface(
stops: &mut HashMap<InterfaceId, oneshot::Sender<()>>,
supervisor_of: &mut HashMap<InterfaceId, InterfaceId>,
interfaces: &Arc<Mutex<HashMap<InterfaceId, RegisteredInterface>>>,
commands: &UnboundedSender<HostCommand>,
id: InterfaceId,
) {
if stops.remove(&id).is_some() {
supervisor_of.remove(&id);
forget_status(interfaces, id);
let _ = commands.send(HostCommand::RemoveInterface {
id,
departure: Departure::MayReturn,
});
stop_supervised_members(stops, supervisor_of, interfaces, commands, id);
}
}
fn stop_supervised_members(
stops: &mut HashMap<InterfaceId, oneshot::Sender<()>>,
supervisor_of: &mut HashMap<InterfaceId, InterfaceId>,
interfaces: &Arc<Mutex<HashMap<InterfaceId, RegisteredInterface>>>,
commands: &UnboundedSender<HostCommand>,
supervisor_id: InterfaceId,
) {
let members: std::vec::Vec<InterfaceId> = supervisor_of
.iter()
.filter(|(_, supervisor)| **supervisor == supervisor_id)
.map(|(member, _)| *member)
.collect();
for member in members {
let _ = stop_interface(stops, member);
supervisor_of.remove(&member);
forget_status(interfaces, member);
let _ = commands.send(HostCommand::RemoveInterface {
id: member,
departure: Departure::MayReturn,
});
}
}
#[cfg(test)]
mod tests;