Packages

macula

5.2.0
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::{Binary, 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).
///
/// `verify_pubkey` is a 32-byte Ed25519 pubkey to pin against the
/// leaf cert's SubjectPublicKeyInfo. An empty binary disables
/// pinning and falls back to `verify` semantics (system-CA or skip).
///
/// `verify_pubkey` is `Binary<'a>` rather than `Vec<u8>` because
/// rustler's `Vec<u8>` decoder requires a list term and rejects
/// Erlang binaries (which is how every caller passes pubkeys).
/// See cert.rs:nif_generate_self_signed_cert for the same pattern.
#[rustler::nif(schedule = "DirtyIo")]
fn nif_connect<'a>(
env: Env<'a>,
host: String,
port: u32,
alpn: Vec<String>,
verify: bool,
verify_pubkey: Binary<'a>,
idle_timeout_ms: u64,
keep_alive_ms: u64,
timeout_ms: u64,
) -> NifResult<Term<'a>> {
let caller = env.pid();
let pinned = if verify_pubkey.is_empty() {
None
} else {
Some(verify_pubkey.as_slice().to_vec())
};
let client_config =
config::build_client_config(&alpn, verify, pinned, idle_timeout_ms, keep_alive_ms)
.map_err(|e| rustler::Error::Term(Box::new(e)))?;
let result: Result<quinn::Connection, String> = runtime::rt().block_on(async {
// One deadline covers the WHOLE operation — DNS resolution,
// endpoint acquisition and the CONNECT/HELLO handshake — so a
// stall in any stage (a hung resolver, a black-holed handshake)
// always returns within `timeout_ms` instead of parking the
// scheduler forever. The old code only wrapped the handshake,
// leaving `lookup_host` unbounded.
let fut = async {
// Strip square brackets if the caller passed `[ipv6]` form
// (used by the pubkey-pin path where the host string is a
// synthetic `[ipv6]` derived from the target pubkey). The
// bare IP works for both DNS resolution and SNI.
let host_str: &str = host.trim_start_matches('[').trim_end_matches(']');
// Two-arg lookup_host avoids the bracket+colon parsing the
// single-string form requires for IPv6.
let addrs: Vec<std::net::SocketAddr> =
tokio::net::lookup_host((host_str, port as u16))
.await
.map_err(|e| format!("resolve {}:{}: {}", host_str, port, e))?
.collect();
let remote_addr = *addrs
.first()
.ok_or_else(|| format!("no addresses for {}:{}", host_str, port))?;
// Shared client endpoint per address family (see
// `runtime::client_endpoint`) — reused across all dials so we
// don't leak a socket + driver task per connection.
let endpoint = runtime::client_endpoint(remote_addr.is_ipv6())?;
// Per-dial client config (verify / ALPN / pubkey-pin) via
// `connect_with`; SNI = bare host string (rustls ServerName
// accepts a literal IP address as a valid name).
let connection = endpoint
.connect_with(client_config, remote_addr, host_str)
.map_err(|e| format!("connect: {}", e))?
.await
.map_err(|e| format!("handshake: {}", e))?;
Ok::<quinn::Connection, String>(connection)
};
match tokio::time::timeout(std::time::Duration::from_millis(timeout_ms), fut).await {
Ok(inner) => inner,
Err(_) => Err("connection_timeout".to_string()),
}
});
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}
///
/// Network-IO bound (`open_bi` awaits stream flow-control credit), so it
/// runs on a dirty-IO scheduler — not dirty-CPU. Dirty-CPU schedulers are
/// scarce (one per core) and a blocking wait there starves everything;
/// dirty-IO is the correct class and there are far more of them.
#[rustler::nif(schedule = "DirtyIo")]
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))
}
/// NIF: max_datagram_size(ConnRef) -> {ok, Bytes} | {error, already_closed}
///
/// Returns the current path MTU on this connection as tracked by
/// Quinn's path-state machine. Reflects DPLPMTUD probing (RFC 8899)
/// once the connection has been up long enough; before that, returns
/// Quinn's initial-MTU default (typically 1200 for IPv6).
///
/// Misnamed for historical reasons — semantics is path MTU in bytes,
/// not max QUIC datagram payload size. Phase 4.2.
#[rustler::nif]
fn nif_max_datagram_size<'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 stats = conn.connection.stats();
let mtu = stats.path.current_mtu as u64;
Ok((atoms::ok(), mtu).encode(env))
}