Packages

Elixir NIF wrapper for the Reticulum cryptographic mesh networking stack.

Current section

Files

Jump to
ex_reticulum native ex_reticulum_nif src transport.rs
Raw

native/ex_reticulum_nif/src/transport.rs

use rustler::types::LocalPid;
use rustler::types::tuple::make_tuple;
use rustler::{Atom, Binary, Encoder, Env, OwnedBinary, OwnedEnv, ResourceArc, Term};
use reticulum::destination::{DestinationDesc, DestinationName};
use reticulum::hash::AddressHash;
use reticulum::iface::tcp_client::TcpClient;
use reticulum::iface::tcp_server::TcpServer;
use reticulum::iface::udp::UdpInterface;
use reticulum::transport::{Transport, TransportConfig};
use crate::atoms;
use crate::resources::{DestinationResource, IdentityKind, IdentityResource, LinkResource, TransportResource};
pub fn runtime() -> &'static tokio::runtime::Runtime {
static RUNTIME: std::sync::OnceLock<tokio::runtime::Runtime> = std::sync::OnceLock::new();
RUNTIME.get_or_init(|| {
tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.expect("Failed to create Tokio runtime")
})
}
#[rustler::nif(schedule = "DirtyIo")]
fn transport_new(
env: Env,
name: String,
identity: ResourceArc<IdentityResource>,
broadcast: bool,
retransmit: bool,
reroute_eager: bool,
restart_outlinks: bool,
announce_forever: bool,
pid: LocalPid,
) -> Term {
let id_lock = match identity.0.lock() {
Ok(l) => l,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
let private_id = match &*id_lock {
IdentityKind::Private(p) => p.clone(),
IdentityKind::Public(_) => return (atoms::error(), atoms::invalid_argument()).encode(env),
};
drop(id_lock);
let mut config = TransportConfig::new(&name, &private_id, broadcast);
config.set_retransmit(retransmit);
config.set_reroute_eager(reroute_eager);
config.set_restart_outlinks(restart_outlinks);
config.set_announce_forever(announce_forever);
let rt = runtime();
// Transport::new internally calls tokio::spawn, so it must run within the runtime context
let (transport, announce_rx, in_link_rx, out_link_rx, data_rx) = rt.block_on(async {
let transport = Transport::new(config);
let announce_rx = transport.recv_announces().await;
let in_link_rx = transport.in_link_events();
let out_link_rx = transport.out_link_events();
let data_rx = transport.received_data_events();
(transport, announce_rx, in_link_rx, out_link_rx, data_rx)
});
spawn_announce_forwarder(rt, pid.clone(), announce_rx);
spawn_link_event_forwarder(rt, pid.clone(), in_link_rx, true);
spawn_link_event_forwarder(rt, pid.clone(), out_link_rx, false);
spawn_data_forwarder(rt, pid, data_rx);
let resource = ResourceArc::new(TransportResource {
transport: std::sync::Mutex::new(transport),
});
(atoms::ok(), resource).encode(env)
}
fn spawn_announce_forwarder(
rt: &tokio::runtime::Runtime,
pid: LocalPid,
mut rx: tokio::sync::broadcast::Receiver<reticulum::transport::AnnounceEvent>,
) {
rt.spawn(async move {
loop {
let event = match rx.recv().await {
Ok(event) => event,
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
};
let dest = event.destination.lock().await;
let desc = dest.desc;
let addr_bytes: Vec<u8> = desc.address_hash.as_slice().to_vec();
let name_hash_bytes: Vec<u8> = desc.name.as_name_hash_slice().to_vec();
let identity = desc.identity;
let id_addr: Vec<u8> = identity.address_hash.as_slice().to_vec();
let pub_key: Vec<u8> = identity.public_key_bytes().to_vec();
let ver_key: Vec<u8> = identity.verifying_key_bytes().to_vec();
let app_data_vec: Vec<u8> = event.app_data.as_slice().to_vec();
let app_data_len = event.app_data.len();
drop(dest);
let mut msg_env = OwnedEnv::new();
let result = msg_env.send_and_clear(&pid, |env| {
let addr = match make_binary(env, &addr_bytes) {
Some(b) => b.encode(env),
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
let name_hash_bin = match make_binary(env, &name_hash_bytes) {
Some(b) => b.encode(env),
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
let id_addr_bin = match make_binary(env, &id_addr) {
Some(b) => b.encode(env),
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
let pub_key_bin = match make_binary(env, &pub_key) {
Some(b) => b.encode(env),
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
let ver_key_bin = match make_binary(env, &ver_key) {
Some(b) => b.encode(env),
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
let app_data = if app_data_len == 0 {
atoms::nil().encode(env)
} else {
match make_binary(env, &app_data_vec[..app_data_len]) {
Some(b) => b.encode(env),
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
}
};
make_tuple(env, &[
atoms::ex_reticulum_event().encode(env),
atoms::announce().encode(env),
addr,
name_hash_bin,
pub_key_bin,
ver_key_bin,
id_addr_bin,
app_data,
])
});
if result.is_err() {
break;
}
}
});
}
fn spawn_link_event_forwarder(
rt: &tokio::runtime::Runtime,
pid: LocalPid,
mut rx: tokio::sync::broadcast::Receiver<reticulum::destination::link::LinkEventData>,
is_inbound: bool,
) {
rt.spawn(async move {
loop {
let event = match rx.recv().await {
Ok(event) => event,
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
};
let link_id: Vec<u8> = event.id.as_slice().to_vec();
let addr: Vec<u8> = event.address_hash.as_slice().to_vec();
let (event_atom, payload): (Atom, Option<Vec<u8>>) = match &event.event {
reticulum::destination::link::LinkEvent::Activated => (atoms::activated(), None),
reticulum::destination::link::LinkEvent::Data(p) => {
(atoms::data(), Some(p.as_slice().to_vec()))
}
reticulum::destination::link::LinkEvent::Proof(h) => {
(atoms::proof(), Some(h.as_slice().to_vec()))
}
reticulum::destination::link::LinkEvent::Closed => (atoms::closed(), None),
};
let direction = if is_inbound {
atoms::r#in()
} else {
atoms::out()
};
let mut msg_env = OwnedEnv::new();
let result = msg_env.send_and_clear(&pid, |env| {
let link_id_bin = match make_binary(env, &link_id) {
Some(b) => b,
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
let addr_bin = match make_binary(env, &addr) {
Some(b) => b,
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
let payload_term = match &payload {
Some(data) => match make_binary(env, data) {
Some(b) => b.encode(env),
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
},
None => atoms::nil().encode(env),
};
(
atoms::ex_reticulum_event(),
atoms::link_event(),
direction,
link_id_bin,
addr_bin,
event_atom,
payload_term,
)
.encode(env)
});
if result.is_err() {
break;
}
}
});
}
fn spawn_data_forwarder(
rt: &tokio::runtime::Runtime,
pid: LocalPid,
mut rx: tokio::sync::broadcast::Receiver<reticulum::transport::ReceivedData>,
) {
rt.spawn(async move {
loop {
let event = match rx.recv().await {
Ok(event) => event,
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
};
let dest: Vec<u8> = event.destination.as_slice().to_vec();
let data: Vec<u8> = event.data.as_slice().to_vec();
let data_len = event.data.len();
let mut msg_env = OwnedEnv::new();
let result = msg_env.send_and_clear(&pid, |env| {
let dest_bin = match make_binary(env, &dest) {
Some(b) => b,
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
let data_bin = match make_binary(env, &data[..data_len]) {
Some(b) => b,
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
(atoms::ex_reticulum_event(), atoms::data(), dest_bin, data_bin).encode(env)
});
if result.is_err() {
break;
}
}
});
}
fn make_binary<'a>(env: Env<'a>, data: &[u8]) -> Option<Binary<'a>> {
let mut bin = OwnedBinary::new(data.len())?;
bin.as_mut_slice().copy_from_slice(data);
Some(Binary::from_owned(bin, env))
}
#[rustler::nif(schedule = "DirtyIo")]
fn transport_add_interface<'a>(
env: Env<'a>,
transport: ResourceArc<TransportResource>,
iface_type: Atom,
addr: String,
extra: Term<'a>,
) -> Term<'a> {
let tp = match transport.transport.lock() {
Ok(t) => t,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
let iface_manager = tp.iface_manager();
let rt = runtime();
if iface_type == atoms::tcp_client() {
rt.block_on(async {
let mut mgr = iface_manager.lock().await;
mgr.spawn(TcpClient::new(&addr), TcpClient::spawn);
});
(atoms::ok(), atoms::added()).encode(env)
} else if iface_type == atoms::tcp_server() {
let mgr_clone = iface_manager.clone();
rt.block_on(async {
let mut mgr = iface_manager.lock().await;
mgr.spawn(TcpServer::new(&addr, mgr_clone), TcpServer::spawn);
});
(atoms::ok(), atoms::added()).encode(env)
} else if iface_type == atoms::udp() {
let forward: Option<String> = if extra.is_atom() {
None
} else {
extra.decode().ok()
};
rt.block_on(async {
let mut mgr = iface_manager.lock().await;
mgr.spawn(
UdpInterface::new(addr, forward),
UdpInterface::spawn,
);
});
(atoms::ok(), atoms::added()).encode(env)
} else {
(atoms::error(), atoms::invalid_argument()).encode(env)
}
}
#[rustler::nif(schedule = "DirtyIo")]
fn transport_add_destination(
env: Env,
transport: ResourceArc<TransportResource>,
identity: ResourceArc<IdentityResource>,
app_name: String,
aspects: String,
) -> Term {
let id_lock = match identity.0.lock() {
Ok(l) => l,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
let private_id = match &*id_lock {
IdentityKind::Private(p) => p.clone(),
IdentityKind::Public(_) => return (atoms::error(), atoms::invalid_argument()).encode(env),
};
drop(id_lock);
let name = DestinationName::new(&app_name, &aspects);
let mut tp = match transport.transport.lock() {
Ok(t) => t,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
let dest = runtime().block_on(tp.add_destination(private_id, name));
let address_hash = runtime().block_on(async { dest.lock().await.desc.address_hash });
let addr_bin = match make_binary(env, address_hash.as_slice()) {
Some(b) => b,
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
let resource = ResourceArc::new(DestinationResource {
destination: dest,
address_hash,
});
(atoms::ok(), resource, addr_bin).encode(env)
}
#[rustler::nif(schedule = "DirtyIo")]
fn transport_announce<'a>(
env: Env<'a>,
transport: ResourceArc<TransportResource>,
destination: ResourceArc<DestinationResource>,
app_data: Term<'a>,
) -> Term<'a> {
let tp = match transport.transport.lock() {
Ok(t) => t,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
let app_data_opt: Option<Vec<u8>> = if app_data.is_atom() {
None
} else {
app_data.decode::<Binary>().ok().map(|b| b.as_slice().to_vec())
};
runtime().block_on(tp.send_announce(
&destination.destination,
app_data_opt.as_deref(),
));
(atoms::ok(), atoms::announced()).encode(env)
}
#[rustler::nif(schedule = "DirtyIo")]
fn transport_has_destination<'a>(
env: Env<'a>,
transport: ResourceArc<TransportResource>,
address_hash: Binary<'a>,
) -> Term<'a> {
if address_hash.as_slice().len() != 16 {
return (atoms::error(), atoms::invalid_argument()).encode(env);
}
let mut hash_bytes = [0u8; 16];
hash_bytes.copy_from_slice(address_hash.as_slice());
let hash = AddressHash::new(hash_bytes);
let tp = match transport.transport.lock() {
Ok(t) => t,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
let result = runtime().block_on(tp.has_destination(&hash));
(atoms::ok(), result).encode(env)
}
#[rustler::nif(schedule = "DirtyIo")]
fn transport_request_path<'a>(
env: Env<'a>,
transport: ResourceArc<TransportResource>,
address_hash: Binary<'a>,
) -> Term<'a> {
if address_hash.as_slice().len() != 16 {
return (atoms::error(), atoms::invalid_argument()).encode(env);
}
let mut hash_bytes = [0u8; 16];
hash_bytes.copy_from_slice(address_hash.as_slice());
let hash = AddressHash::new(hash_bytes);
let tp = match transport.transport.lock() {
Ok(t) => t,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
runtime().block_on(tp.request_path(&hash, None, None));
(atoms::ok(), atoms::sent()).encode(env)
}
#[rustler::nif(schedule = "DirtyIo")]
fn transport_link<'a>(
env: Env<'a>,
transport: ResourceArc<TransportResource>,
address_hash: Binary<'a>,
identity_ref: ResourceArc<IdentityResource>,
name_hash: Binary<'a>,
) -> Term<'a> {
if address_hash.as_slice().len() != 16 {
return (atoms::error(), atoms::invalid_argument()).encode(env);
}
let id_lock = match identity_ref.0.lock() {
Ok(l) => l,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
let identity = match &*id_lock {
IdentityKind::Public(p) => *p,
IdentityKind::Private(p) => *p.as_identity(),
};
drop(id_lock);
let mut hash_bytes = [0u8; 16];
hash_bytes.copy_from_slice(address_hash.as_slice());
let addr = AddressHash::new(hash_bytes);
let name = DestinationName::new_from_hash_slice(name_hash.as_slice());
let desc = DestinationDesc {
identity,
address_hash: addr,
name,
};
let tp = match transport.transport.lock() {
Ok(t) => t,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
let link = runtime().block_on(tp.link(desc));
let link_id = runtime().block_on(async { *link.lock().await.id() });
let id_bin = match make_binary(env, link_id.as_slice()) {
Some(b) => b,
None => return (atoms::error(), atoms::out_of_memory()).encode(env),
};
let resource = ResourceArc::new(LinkResource {
link,
id: link_id,
});
(atoms::ok(), resource, id_bin).encode(env)
}
#[rustler::nif(schedule = "DirtyIo")]
fn transport_link_send<'a>(
env: Env<'a>,
transport: ResourceArc<TransportResource>,
link: ResourceArc<LinkResource>,
data: Binary<'a>,
) -> Term<'a> {
let rt = runtime();
let packet = rt.block_on(async {
link.link.lock().await.data_packet(data.as_slice())
});
match packet {
Ok(packet) => {
let tp = match transport.transport.lock() {
Ok(t) => t,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
rt.block_on(tp.send_packet(packet));
(atoms::ok(), atoms::sent()).encode(env)
}
Err(reticulum::error::RnsError::LinkClosed) => {
(atoms::error(), atoms::link_closed()).encode(env)
}
Err(_) => (atoms::error(), atoms::packet_error()).encode(env),
}
}
#[rustler::nif(schedule = "DirtyIo")]
fn transport_link_close(
env: Env,
transport: ResourceArc<TransportResource>,
link: ResourceArc<LinkResource>,
) -> Term {
let tp = match transport.transport.lock() {
Ok(t) => t,
Err(_) => return (atoms::error(), atoms::lock_error()).encode(env),
};
match runtime().block_on(tp.link_close(link.id)) {
Ok(()) => (atoms::ok(), atoms::closed()).encode(env),
Err(_) => (atoms::error(), atoms::link_closed()).encode(env),
}
}