Packages

macula

1.4.18
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 connection.rs
Raw

native/macula_quic/src/connection.rs

use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Mutex, RwLock};
use rustler::{Encoder, Env, LocalPid, NifResult, ResourceArc, Term};
use tokio::task::JoinHandle;
use crate::{atoms, config, message, runtime, stream};
/// Opaque connection handle exposed to Erlang via ResourceArc.
pub struct ConnectionResource {
pub connection: quinn::Connection,
pub owner: RwLock<LocalPid>,
stream_accept_task: Mutex<Option<JoinHandle<()>>>,
pub closed: AtomicBool,
}
impl ConnectionResource {
pub fn new(connection: quinn::Connection, owner: LocalPid) -> Self {
Self {
connection,
owner: RwLock::new(owner),
stream_accept_task: Mutex::new(None),
closed: AtomicBool::new(false),
}
}
pub fn set_stream_accept_task(&self, handle: JoinHandle<()>) {
let mut task = self.stream_accept_task.lock().unwrap();
*task = Some(handle);
}
}
impl Drop for ConnectionResource {
fn drop(&mut self) {
self.closed.store(true, Ordering::SeqCst);
if let Some(task) = self.stream_accept_task.lock().unwrap().take() {
task.abort();
}
self.connection.close(0u32.into(), b"closed");
}
}
/// NIF: connect(Host, Port, Opts) -> {ok, ConnRef} | {error, Reason}
///
/// Blocks the dirty scheduler until handshake completes (up to timeout).
#[rustler::nif(schedule = "DirtyCpu")]
fn nif_connect<'a>(
env: Env<'a>,
host: String,
port: u32,
alpn: Vec<String>,
verify: bool,
idle_timeout_ms: u64,
keep_alive_ms: u64,
timeout_ms: u64,
) -> NifResult<Term<'a>> {
let caller = env.pid();
let client_config = config::build_client_config(&alpn, verify, idle_timeout_ms, keep_alive_ms)
.map_err(|e| rustler::Error::Term(Box::new(e)))?;
let result = runtime::rt().block_on(async {
// Resolve hostname
let addr_str = format!("{}:{}", host, port);
let addrs: Vec<std::net::SocketAddr> = tokio::net::lookup_host(&addr_str)
.await
.map_err(|e| format!("resolve {}: {}", addr_str, e))?
.collect();
let remote_addr = addrs
.first()
.ok_or_else(|| format!("no addresses for {}", addr_str))?;
// Create client endpoint — match address family to remote
let local_bind: std::net::SocketAddr = if remote_addr.is_ipv6() {
"[::]:0".parse().unwrap()
} else {
"0.0.0.0:0".parse().unwrap()
};
let mut endpoint = quinn::Endpoint::client(local_bind)
.map_err(|e| format!("client endpoint: {}", e))?;
endpoint.set_default_client_config(client_config);
// Connect with timeout
let connecting = endpoint
.connect(*remote_addr, &host)
.map_err(|e| format!("connect: {}", e))?;
let connection = tokio::time::timeout(
std::time::Duration::from_millis(timeout_ms),
connecting,
)
.await
.map_err(|_| "connection_timeout".to_string())?
.map_err(|e| format!("handshake: {}", e))?;
Ok::<quinn::Connection, String>(connection)
});
match result {
Ok(connection) => {
let resource = ResourceArc::new(ConnectionResource::new(connection, caller));
Ok((atoms::ok(), resource).encode(env))
}
Err(e) => Ok((atoms::error(), e).encode(env)),
}
}
/// NIF: open_stream(ConnRef) -> {ok, StreamRef} | {error, Reason}
#[rustler::nif(schedule = "DirtyCpu")]
fn nif_open_stream<'a>(
env: Env<'a>,
conn: ResourceArc<ConnectionResource>,
) -> NifResult<Term<'a>> {
if conn.closed.load(Ordering::Relaxed) {
return Ok((atoms::error(), atoms::already_closed()).encode(env));
}
let caller = env.pid();
let connection = conn.connection.clone();
let result = runtime::rt().block_on(async {
let (send, recv) = connection
.open_bi()
.await
.map_err(|e| format!("open_bi: {}", e))?;
Ok::<(quinn::SendStream, quinn::RecvStream), String>((send, recv))
});
match result {
Ok((send, recv)) => {
let resource = ResourceArc::new(stream::StreamResource::new(
send, recv, conn.clone(), caller,
));
stream::StreamResource::start_recv_loop(resource.clone());
Ok((atoms::ok(), resource).encode(env))
}
Err(e) => Ok((atoms::error(), e).encode(env)),
}
}
/// NIF: close_connection(ConnRef) -> ok
#[rustler::nif]
fn nif_close_connection<'a>(
env: Env<'a>,
conn: ResourceArc<ConnectionResource>,
) -> NifResult<Term<'a>> {
conn.closed.store(true, Ordering::SeqCst);
if let Some(task) = conn.stream_accept_task.lock().unwrap().take() {
task.abort();
}
conn.connection.close(0u32.into(), b"closed");
Ok(atoms::ok().encode(env))
}
/// NIF: async_accept_stream(ConnRef) -> ok
/// Starts stream accept loop. Delivers {quic, new_stream, StreamRef, Props}.
#[rustler::nif]
fn nif_async_accept_stream<'a>(
env: Env<'a>,
conn: ResourceArc<ConnectionResource>,
) -> NifResult<Term<'a>> {
let connection = conn.connection.clone();
let conn_arc = conn.clone();
let handle = runtime::rt().spawn(async move {
loop {
if conn_arc.closed.load(Ordering::Relaxed) {
break;
}
match connection.accept_bi().await {
Ok((send, recv)) => {
let owner = *conn_arc.owner.read().unwrap();
let stream_resource = ResourceArc::new(stream::StreamResource::new(
send,
recv,
conn_arc.clone(),
owner,
));
stream::StreamResource::start_recv_loop(stream_resource.clone());
message::send_new_stream(&owner, stream_resource, conn_arc.clone(), 0);
}
Err(_) => break, // Connection closed
}
}
});
conn.set_stream_accept_task(handle);
Ok(atoms::ok().encode(env))
}
/// NIF: controlling_process_conn(ConnRef, NewPid) -> ok
#[rustler::nif]
fn nif_controlling_process_conn<'a>(
env: Env<'a>,
conn: ResourceArc<ConnectionResource>,
new_owner: LocalPid,
) -> NifResult<Term<'a>> {
let mut owner = conn.owner.write().unwrap();
*owner = new_owner;
Ok(atoms::ok().encode(env))
}
/// NIF: peername(ConnRef) -> {ok, {IP, Port}} | {error, Reason}
#[rustler::nif]
fn nif_peername<'a>(
env: Env<'a>,
conn: ResourceArc<ConnectionResource>,
) -> NifResult<Term<'a>> {
let addr = conn.connection.remote_address();
let ip = addr.ip().to_string();
let port = addr.port() as u32;
Ok((atoms::ok(), (ip, port)).encode(env))
}