Packages

macula

4.1.1
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 stream.rs
Raw

native/macula_quic/src/stream.rs

use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Mutex, RwLock};
use rustler::{Encoder, Env, LocalPid, NifResult, ResourceArc, Term};
use tokio::sync::Notify;
use tokio::task::JoinHandle;
use crate::{atoms, connection::ConnectionResource, message, runtime};
/// Opaque stream handle exposed to Erlang via ResourceArc.
pub struct StreamResource {
send: Mutex<Option<quinn::SendStream>>,
recv: Mutex<Option<quinn::RecvStream>>,
recv_task: Mutex<Option<JoinHandle<()>>>,
pub conn: ResourceArc<ConnectionResource>,
pub owner: RwLock<LocalPid>,
pub active: AtomicBool,
active_notify: Notify,
pub closed: AtomicBool,
}
impl StreamResource {
pub fn new(
send: quinn::SendStream,
recv: quinn::RecvStream,
conn: ResourceArc<ConnectionResource>,
owner: LocalPid,
) -> Self {
Self {
send: Mutex::new(Some(send)),
recv: Mutex::new(Some(recv)),
recv_task: Mutex::new(None),
conn,
owner: RwLock::new(owner),
active: AtomicBool::new(false),
active_notify: Notify::new(),
closed: AtomicBool::new(false),
}
}
/// Start the background read loop. Takes the recv stream from self.
/// Must be called after the ResourceArc is created.
pub fn start_recv_loop(self_arc: ResourceArc<Self>) {
let mut recv_opt = self_arc.recv.lock().unwrap();
let mut recv = match recv_opt.take() {
Some(r) => r,
None => return, // Already started or no recv stream
};
drop(recv_opt);
let stream_arc = self_arc.clone();
let handle = runtime::rt().spawn(async move {
let mut buf = vec![0u8; 65536];
loop {
if stream_arc.closed.load(Ordering::Relaxed) {
break;
}
// Wait for active mode
if !stream_arc.active.load(Ordering::Relaxed) {
stream_arc.active_notify.notified().await;
continue;
}
match recv.read(&mut buf).await {
Ok(Some(n)) => {
let data = buf[..n].to_vec();
let owner = *stream_arc.owner.read().unwrap();
message::send_data(&owner, data, stream_arc.clone());
}
Ok(None) => {
// Peer finished sending
let owner = *stream_arc.owner.read().unwrap();
message::send_event(
&owner,
atoms::peer_send_shutdown(),
stream_arc.clone(),
atoms::none(),
);
break;
}
Err(_e) => {
let owner = *stream_arc.owner.read().unwrap();
message::send_event(
&owner,
atoms::stream_closed(),
stream_arc.clone(),
atoms::none(), // simplified for now
);
break;
}
}
}
});
let mut task = self_arc.recv_task.lock().unwrap();
*task = Some(handle);
}
/// Wake the recv loop when active mode is enabled.
pub fn notify_active(&self) {
self.active_notify.notify_one();
}
}
impl Drop for StreamResource {
fn drop(&mut self) {
self.closed.store(true, Ordering::SeqCst);
if let Some(task) = self.recv_task.lock().unwrap().take() {
task.abort();
}
}
}
/// NIF: send(StreamRef, Data) -> ok | {error, Reason}
#[rustler::nif(schedule = "DirtyCpu")]
fn nif_send<'a>(
env: Env<'a>,
stream: ResourceArc<StreamResource>,
data: rustler::Binary<'a>,
) -> NifResult<Term<'a>> {
if stream.closed.load(Ordering::Relaxed) {
return Ok((atoms::error(), atoms::already_closed()).encode(env));
}
let bytes = data.as_slice().to_vec();
let mut guard = stream.send.lock().unwrap();
let send_stream = match guard.as_mut() {
Some(s) => s,
None => return Ok((atoms::error(), atoms::stream_finished()).encode(env)),
};
// Clone the send stream reference for the async block
// Actually, we need to do the write inside block_on with a mutable ref
let result = runtime::rt().block_on(async {
send_stream
.write_all(&bytes)
.await
.map_err(|e| format!("{}", e))
});
drop(guard); // release lock
match result {
Ok(()) => Ok(atoms::ok().encode(env)),
Err(e) => Ok((atoms::error(), e).encode(env)),
}
}
/// NIF: async_send(StreamRef, Data) -> ok | {error, Reason}
#[rustler::nif]
fn nif_async_send<'a>(
env: Env<'a>,
stream: ResourceArc<StreamResource>,
data: rustler::Binary<'a>,
) -> NifResult<Term<'a>> {
if stream.closed.load(Ordering::Relaxed) {
return Ok((atoms::error(), atoms::already_closed()).encode(env));
}
let bytes = data.as_slice().to_vec();
// For async_send, we block briefly to queue the write (Quinn buffers internally).
// This avoids the MutexGuard-across-await Send issue.
let mut guard = stream.send.lock().unwrap();
if let Some(send_stream) = guard.as_mut() {
let _ = runtime::rt().block_on(send_stream.write_all(&bytes));
}
drop(guard);
Ok(atoms::ok().encode(env))
}
/// NIF: close_stream(StreamRef) -> ok
#[rustler::nif]
fn nif_close_stream<'a>(
env: Env<'a>,
stream: ResourceArc<StreamResource>,
) -> NifResult<Term<'a>> {
stream.closed.store(true, Ordering::SeqCst);
if let Some(task) = stream.recv_task.lock().unwrap().take() {
task.abort();
}
// Finish the send stream gracefully
let mut guard = stream.send.lock().unwrap();
if let Some(mut send_stream) = guard.take() {
let _ = send_stream.finish();
}
Ok(atoms::ok().encode(env))
}
/// NIF: setopt(StreamRef, active, true|false) -> ok
#[rustler::nif]
fn nif_setopt_active<'a>(
env: Env<'a>,
stream: ResourceArc<StreamResource>,
value: bool,
) -> NifResult<Term<'a>> {
stream.active.store(value, Ordering::SeqCst);
if value {
stream.notify_active();
}
Ok(atoms::ok().encode(env))
}
/// NIF: controlling_process(StreamRef, NewPid) -> ok
#[rustler::nif]
fn nif_controlling_process<'a>(
env: Env<'a>,
stream: ResourceArc<StreamResource>,
new_owner: LocalPid,
) -> NifResult<Term<'a>> {
let mut owner = stream.owner.write().unwrap();
*owner = new_owner;
Ok(atoms::ok().encode(env))
}