Packages

macula

4.4.5
7.0.0 6.0.0 5.2.2 5.2.1 5.2.0 5.1.0 5.0.0 4.8.0 4.7.1 4.7.0 4.6.0 4.5.0 4.4.10 4.4.9 4.4.8 4.4.7 4.4.6 4.4.5 4.4.4 4.4.3 4.4.2 4.4.1 4.4.0 4.3.1 4.3.0 4.2.9 4.2.8 4.2.7 4.2.6 4.2.5 4.2.4 4.2.3 4.2.2 4.2.1 4.2.0 4.1.1 4.1.0 4.0.0 3.16.0 3.15.3 3.15.2 3.15.1 3.14.0 3.13.0 3.12.1 3.12.0 3.11.1 3.11.0 3.10.3 3.10.2 3.10.1 3.9.0 3.8.0 3.7.0 3.5.0 3.4.0 3.3.0 3.2.0 3.1.0 3.0.0 2.1.1 2.1.0 2.0.0 1.5.2 1.5.1 1.4.30 1.4.29 1.4.28 1.4.27 1.4.26 1.4.25 1.4.24 1.4.23 1.4.22 1.4.21 1.4.20 1.4.19 1.4.18 1.4.17 1.4.16 1.4.15 1.4.14 1.4.13 1.4.11 1.4.10 1.4.9 1.4.8 1.4.7 1.4.6 1.4.5 1.4.4 1.4.3 1.4.2 1.4.1 1.4.0 1.3.1 1.3.0 1.2.0 1.1.0 1.0.10 1.0.9 1.0.8 1.0.7 1.0.6 1.0.5 1.0.4 1.0.3 1.0.2 1.0.1 1.0.0 0.48.6 0.48.5 0.48.4 0.48.3 0.48.2 0.48.1 0.48.0 0.47.1 0.47.0 0.46.3 0.46.1 0.46.0 0.45.3 0.45.2 0.45.1 0.45.0 0.44.2 0.44.1 0.44.0 0.43.3 0.43.2 0.43.1 0.43.0 0.42.9 0.42.8 0.42.7 0.42.6 0.42.5 0.42.4 0.42.3 0.42.2 0.42.1 0.42.0 0.41.1 0.41.0 0.40.1 0.40.0 0.39.9 0.39.8 0.39.7 0.39.6 0.39.5 0.39.4 0.39.3 0.39.2 0.39.1 0.39.0 0.38.8 0.38.7 0.38.6 0.38.5 0.38.4 0.38.3 0.38.2 0.38.1 0.38.0 0.37.7 0.37.6 0.37.5 0.37.4 0.37.3 0.37.2 0.37.1 0.37.0 0.36.6 0.36.5 0.36.4 0.36.3 0.36.2 0.36.1 0.36.0 0.35.4 0.35.3 0.35.2 0.35.1 0.35.0 0.34.1 0.34.0 0.33.1 0.33.0 0.32.5 0.32.4 0.32.3 0.32.2 0.32.1 0.32.0 0.31.9 0.31.8 0.31.7 0.31.6 0.31.5 0.31.4 0.31.3 0.31.2 0.31.1 0.31.0 0.30.10 0.30.9 0.30.8 0.30.7 0.30.6 0.30.5 0.30.4 0.30.3 0.30.2 0.30.1 0.30.0 0.29.0 0.28.3 0.28.2 0.28.1 0.28.0 0.27.1 0.27.0 0.26.1 0.26.0 0.25.6 0.25.5 0.25.4 0.25.3 0.25.2 0.25.1 0.25.0 0.24.6 0.24.5 0.24.4 0.24.3 0.24.2 0.24.1 0.24.0 0.23.3 0.23.2 0.23.1 0.23.0 0.22.12 0.22.11 0.22.10 0.22.9 0.22.8 0.22.7 0.22.6 0.22.5 0.22.4 0.22.3 0.22.2 0.22.1 0.22.0 0.21.7 0.21.6 0.21.5 0.21.4 0.21.2 0.21.1 0.21.0 0.20.25 0.20.24 0.20.23 0.20.22 0.20.21 0.20.20 0.20.19 0.20.18 0.20.17 0.20.16 0.20.15 0.20.14 0.20.13 0.20.12 0.20.11 0.20.10 0.20.9 0.20.8 0.20.7 0.20.6 0.20.5 0.20.3 0.20.2 0.20.1 0.20.0 0.19.2 0.19.1 0.19.0 0.18.1 0.18.0 0.17.4 0.17.3 0.17.2 0.17.1 0.17.0 0.16.6 0.16.5 0.16.4 0.16.3 0.16.2 0.16.1 0.16.0 0.15.1 0.15.0 0.14.3 0.14.2 0.14.1 0.14.0 0.12.6 0.12.5 0.12.3 0.11.3 0.10.2 0.10.1 0.10.0 0.9.2 0.9.1 0.9.0 0.8.25 0.8.24 0.8.23 0.8.22 0.8.21 0.8.20 0.8.19 0.8.18 0.8.17 0.8.16 0.8.15 0.8.14 0.8.13 0.8.12 0.8.11 0.8.10 0.8.9 0.8.8 0.8.7 0.8.6 0.8.5 0.8.4 0.8.3 0.8.2 0.8.1 0.8.0 0.7.30 0.7.29 0.7.28 0.7.27 0.7.26 0.7.25 0.7.24 0.7.23 0.7.22 0.7.21 0.7.20 0.7.19 0.7.18 0.7.17 0.7.16 0.7.15 0.7.14 0.7.13 0.7.12 0.7.11 0.7.10 0.7.9 0.7.8 0.7.7 0.7.6 0.7.5 0.7.4 0.7.3 0.7.2 0.7.1 0.7.0 0.6.7 0.6.6 0.6.5 0.6.4 0.6.3 0.6.2 0.6.1 0.6.0 0.5.0 0.4.4 0.4.3 0.4.2 0.4.1 0.4.0 0.3.4 0.3.3 0.3.2 0.3.1

Macula HTTP/3 Mesh SDK — connect, subscribe, publish, call, advertise

Current section

Files

Jump to
macula native macula_tun_nif src tun.rs
Raw

native/macula_tun_nif/src/tun.rs

//! TUN device lifecycle + packet I/O.
//!
//! Resource model: `TunResource` is the BEAM-side opaque handle. While the
//! resource is alive, the underlying TUN device exists. When BEAM GC drops
//! the resource, the Drop impl signals the reader thread (if any) and
//! closes the device.
//!
//! Reader thread: `nif_start_reader/2` spawns a Rust thread that does
//! blocking reads on the TUN fd and sends each packet as
//! `{macula_net_packet, Handle, Payload}` to the registered Pid via
//! `OwnedEnv::send_and_clear`. The thread exits when the stop flag is
//! set or when reads start returning errors (device closed).
use rustler::{Binary, Encoder, Env, LocalPid, NewBinary, OwnedEnv, ResourceArc, Term};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::thread;
use std::time::Duration;
use tun_rs::{DeviceBuilder, SyncDevice};
mod atoms {
rustler::atoms! {
ok,
error,
macula_net_packet,
already_started,
}
}
/// BEAM-visible opaque handle.
pub struct TunResource {
inner: Mutex<Option<TunInner>>,
}
struct TunInner {
device: Arc<SyncDevice>,
name: String,
reader: Option<ReaderState>,
}
struct ReaderState {
stop: Arc<AtomicBool>,
handle: Option<thread::JoinHandle<()>>,
}
impl Drop for TunResource {
fn drop(&mut self) {
if let Ok(mut guard) = self.inner.lock() {
if let Some(inner) = guard.take() {
stop_reader(inner);
}
}
}
}
fn stop_reader(mut inner: TunInner) {
if let Some(reader) = inner.reader.take() {
reader.stop.store(true, Ordering::Release);
if let Some(h) = reader.handle {
let _ = h.join();
}
}
drop(inner.device);
}
// =============================================================================
// NIF: open
// =============================================================================
/// Open a TUN device.
///
/// Args:
/// if_name :: binary — desired interface name (e.g. "macula0")
/// mtu :: integer — IPv6 MTU (>= 1280)
///
/// Returns: `{ok, ResourceHandle}` | `{error, Reason}`.
///
/// Requires `CAP_NET_ADMIN`.
#[rustler::nif]
pub fn nif_open<'a>(
env: Env<'a>,
if_name: Binary<'a>,
mtu: u32,
) -> Term<'a> {
let name = match std::str::from_utf8(if_name.as_slice()) {
Ok(s) => s.to_string(),
Err(_) => return (atoms::error(), "invalid_utf8_ifname".to_string()).encode(env),
};
let device = DeviceBuilder::new()
.name(&name)
.mtu(mtu as u16)
.build_sync();
let device = match device {
Ok(d) => d,
Err(e) => return (atoms::error(), e.to_string()).encode(env),
};
let actual_name = device.name().unwrap_or_else(|_| name.clone());
let resource = ResourceArc::new(TunResource {
inner: Mutex::new(Some(TunInner {
device: Arc::new(device),
name: actual_name,
reader: None,
})),
});
(atoms::ok(), resource).encode(env)
}
// =============================================================================
// NIF: name
// =============================================================================
#[rustler::nif]
pub fn nif_name<'a>(env: Env<'a>, handle: ResourceArc<TunResource>) -> Term<'a> {
let guard = match handle.inner.lock() {
Ok(g) => g,
Err(_) => return (atoms::error(), "lock_poisoned".to_string()).encode(env),
};
match guard.as_ref() {
Some(inner) => (atoms::ok(), inner.name.clone()).encode(env),
None => (atoms::error(), "closed".to_string()).encode(env),
}
}
// =============================================================================
// NIF: close
// =============================================================================
/// Close the TUN device early. Otherwise it's closed when BEAM GCs the
/// resource. Idempotent.
#[rustler::nif]
pub fn nif_close<'a>(env: Env<'a>, handle: ResourceArc<TunResource>) -> Term<'a> {
let mut guard = match handle.inner.lock() {
Ok(g) => g,
Err(_) => return (atoms::error(), "lock_poisoned".to_string()).encode(env),
};
if let Some(inner) = guard.take() {
stop_reader(inner);
}
atoms::ok().encode(env)
}
// =============================================================================
// NIF: write
// =============================================================================
/// Write a raw IPv6 packet to the TUN device.
///
/// `packet` MUST be a complete IPv6 packet (header + payload). tun-rs
/// writes it verbatim to the kernel TUN fd.
#[rustler::nif]
pub fn nif_write<'a>(
env: Env<'a>,
handle: ResourceArc<TunResource>,
packet: Binary<'a>,
) -> Term<'a> {
let device = {
let guard = match handle.inner.lock() {
Ok(g) => g,
Err(_) => return (atoms::error(), "lock_poisoned".to_string()).encode(env),
};
match guard.as_ref() {
Some(inner) => inner.device.clone(),
None => return (atoms::error(), "closed".to_string()).encode(env),
}
};
match device.send(packet.as_slice()) {
Ok(_n) => atoms::ok().encode(env),
Err(e) => (atoms::error(), e.to_string()).encode(env),
}
}
// =============================================================================
// NIF: start_reader
// =============================================================================
/// Spawn a Rust thread that reads packets from the TUN device and forwards
/// each as `{macula_net_packet, Handle, Payload}` to `pid`.
#[rustler::nif]
pub fn nif_start_reader<'a>(
env: Env<'a>,
handle: ResourceArc<TunResource>,
pid: LocalPid,
) -> Term<'a> {
let mut guard = match handle.inner.lock() {
Ok(g) => g,
Err(_) => return (atoms::error(), "lock_poisoned".to_string()).encode(env),
};
let inner = match guard.as_mut() {
Some(i) => i,
None => return (atoms::error(), "closed".to_string()).encode(env),
};
if inner.reader.is_some() {
return (atoms::error(), atoms::already_started()).encode(env);
}
let stop = Arc::new(AtomicBool::new(false));
let stop_clone = stop.clone();
let device = inner.device.clone();
let handle_clone = handle.clone();
let thread_handle = thread::Builder::new()
.name(format!("macula-tun-reader-{}", inner.name))
.spawn(move || reader_loop(device, handle_clone, pid, stop_clone))
.ok();
inner.reader = Some(ReaderState {
stop,
handle: thread_handle,
});
atoms::ok().encode(env)
}
fn reader_loop(
device: Arc<SyncDevice>,
handle: ResourceArc<TunResource>,
pid: LocalPid,
stop: Arc<AtomicBool>,
) {
let mut buf = vec![0u8; 65535];
while !stop.load(Ordering::Acquire) {
match device.recv(&mut buf) {
Ok(n) if n > 0 => {
let packet_bytes = buf[..n].to_vec();
let mut owned = OwnedEnv::new();
let send_result = owned.send_and_clear(&pid, |env| {
let mut bin = NewBinary::new(env, packet_bytes.len());
bin.as_mut_slice().copy_from_slice(&packet_bytes);
let bin_term: Binary = bin.into();
(
atoms::macula_net_packet(),
handle.clone(),
bin_term,
)
.encode(env)
});
if send_result.is_err() {
break;
}
}
Ok(_) => {}
Err(e) => {
use std::io::ErrorKind;
match e.kind() {
ErrorKind::WouldBlock | ErrorKind::Interrupted => {
thread::sleep(Duration::from_millis(1));
}
_ => break,
}
}
}
}
}