Packages

macula

4.4.3
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 = "DirtyCpu")]
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 = runtime::rt().block_on(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))?;
// 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. SNI = bare host string (rustls
// ServerName accepts a literal IP address as a valid name).
let connecting = endpoint
.connect(*remote_addr, host_str)
.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))
}
/// 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))
}