Packages

macula

4.4.10
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_quic src endpoint.rs
Raw

native/macula_quic/src/endpoint.rs

use std::net::{IpAddr, SocketAddr};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, RwLock};
use rustler::{Encoder, Env, LocalPid, NifResult, ResourceArc, Term};
use tokio::task::JoinHandle;
use crate::{atoms, config, message, runtime};
/// Opaque listener handle exposed to Erlang via ResourceArc.
pub struct ListenerResource {
pub endpoint: quinn::Endpoint,
pub local_addr: SocketAddr,
pub owner: RwLock<LocalPid>,
accept_task: Mutex<Option<JoinHandle<()>>>,
pub closed: AtomicBool,
}
impl ListenerResource {
pub fn new(endpoint: quinn::Endpoint, local_addr: SocketAddr, owner: LocalPid) -> Self {
Self {
endpoint,
local_addr,
owner: RwLock::new(owner),
accept_task: Mutex::new(None),
closed: AtomicBool::new(false),
}
}
pub fn set_accept_task(&self, handle: JoinHandle<()>) {
let mut task = self.accept_task.lock().unwrap();
*task = Some(handle);
}
}
impl Drop for ListenerResource {
fn drop(&mut self) {
self.closed.store(true, Ordering::SeqCst);
if let Some(task) = self.accept_task.lock().unwrap().take() {
task.abort();
}
self.endpoint.close(0u32.into(), b"shutdown");
}
}
/// NIF: listen(BindAddr, Port, Opts) -> {ok, ListenerRef} | {error, Reason}
#[rustler::nif(schedule = "DirtyCpu")]
fn nif_listen<'a>(
env: Env<'a>,
bind_addr: String,
port: u32,
certfile: String,
keyfile: String,
alpn: Vec<String>,
idle_timeout_ms: u64,
keep_alive_ms: u64,
bidi_streams: u32,
uni_streams: u32,
) -> NifResult<Term<'a>> {
let caller = env.pid();
let addr: IpAddr = bind_addr
.parse()
.map_err(|e| rustler::Error::Term(Box::new(format!("invalid bind_addr: {}", e))))?;
let server_config = config::build_server_config(
&certfile,
&keyfile,
&alpn,
idle_timeout_ms,
keep_alive_ms,
bidi_streams,
uni_streams,
)
.map_err(|e| rustler::Error::Term(Box::new(e)))?;
let socket = config::create_bound_socket(addr, port as u16)
.map_err(|e| rustler::Error::Term(Box::new(e)))?;
// Quinn's TokioRuntime requires a tokio context (Handle::current()).
// We enter the runtime context here since NIFs run on BEAM scheduler threads.
let _guard = runtime::rt().enter();
let endpoint = quinn::Endpoint::new(
quinn::EndpointConfig::default(),
Some(server_config),
socket,
Arc::new(quinn::TokioRuntime),
)
.map_err(|e| rustler::Error::Term(Box::new(format!("endpoint create: {}", e))))?;
let local_addr = endpoint
.local_addr()
.map_err(|e| rustler::Error::Term(Box::new(format!("local_addr: {}", e))))?;
let resource = ResourceArc::new(ListenerResource::new(endpoint, local_addr, caller));
Ok((atoms::ok(), resource).encode(env))
}
/// NIF: close_listener(ListenerRef) -> ok
#[rustler::nif]
fn nif_close_listener<'a>(
env: Env<'a>,
listener: ResourceArc<ListenerResource>,
) -> NifResult<Term<'a>> {
listener.closed.store(true, Ordering::SeqCst);
if let Some(task) = listener.accept_task.lock().unwrap().take() {
task.abort();
}
listener.endpoint.close(0u32.into(), b"shutdown");
Ok(atoms::ok().encode(env))
}
/// NIF: async_accept(ListenerRef) -> ok
/// Starts the accept loop. Each new connection delivers {quic, new_conn, ConnRef, Info}.
#[rustler::nif]
fn nif_async_accept<'a>(
env: Env<'a>,
listener: ResourceArc<ListenerResource>,
) -> NifResult<Term<'a>> {
let endpoint = listener.endpoint.clone();
let listener_arc = listener.clone();
let handle = runtime::rt().spawn(async move {
loop {
if listener_arc.closed.load(Ordering::Relaxed) {
break;
}
match endpoint.accept().await {
Some(incoming) => {
let listener_ref = listener_arc.clone();
// Spawn per-connection task for handshake
tokio::spawn(async move {
if listener_ref.closed.load(Ordering::Relaxed) {
return;
}
let remote_addr = incoming.remote_address().to_string();
match incoming.await {
Ok(connection) => {
let conn_resource = ResourceArc::new(
crate::connection::ConnectionResource::new(
connection,
*listener_ref.owner.read().unwrap(),
),
);
let owner = *listener_ref.owner.read().unwrap();
message::send_new_conn(&owner, conn_resource, remote_addr);
}
Err(e) => {
eprintln!("[macula_quic] accept handshake failed: {}", e);
}
}
});
}
None => break, // Endpoint closed
}
}
});
listener.set_accept_task(handle);
Ok(atoms::ok().encode(env))
}