#![cfg(all(target_arch = "wasm32", feature = "transport-moq"))]
use std::{
cell::{Cell, RefCell},
rc::Rc,
};
use js_sys::{Function, Promise, Reflect, Uint8Array};
use wasm_bindgen::{JsCast, JsValue};
use wasm_bindgen_futures::{spawn_local, JsFuture};
use crate::{
iroh_carrier::{
segment_packet, CarrierControl, CarrierFrame, CarrierFrameExpectation, CarrierReassembler,
CARRIER_CONTROL_PACKET_ID, CARRIER_HEADER_LEN, MAX_INNER_PACKET_BYTES,
},
packet_carrier_transport::PacketCarrierSession,
};
const MOQ_CARRIER_OBJECT_PAYLOAD_CEILING: usize = CARRIER_HEADER_LEN + MAX_INNER_PACKET_BYTES;
pub struct WasmMoqCarrierSession {
_duplex: JsValue,
reader: JsValue,
writer: JsValue,
closed: Rc<Cell<bool>>,
_session: Rc<PacketCarrierSession>,
expected: CarrierFrameExpectation,
application_key: [u8; 32],
terminal_acknowledged: Rc<Cell<bool>>,
inbound_readiness: Rc<Cell<u8>>,
outbound_readiness_sent: Rc<Cell<bool>>,
}
impl std::fmt::Debug for WasmMoqCarrierSession {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("WasmMoqCarrierSession")
.finish_non_exhaustive()
}
}
impl WasmMoqCarrierSession {
pub fn attach(
duplex: JsValue,
session: PacketCarrierSession,
expected: CarrierFrameExpectation,
application_key: [u8; 32],
on_terminal: Rc<dyn Fn(&'static str)>,
) -> Result<Self, JsValue> {
let readable = Reflect::get(&duplex, &JsValue::from_str("readable"))?;
let writable = Reflect::get(&duplex, &JsValue::from_str("writable"))?;
let reader = method(&readable, "getReader")?.call0(&readable)?;
let writer = method(&writable, "getWriter")?.call0(&writable)?;
let message_ceiling = Reflect::get(&duplex, &JsValue::from_str("maxObjectSize"))
.ok()
.and_then(|value| value.as_f64())
.filter(|value| value.is_finite() && *value >= 0.0)
.map(|value| value.floor() as usize)
.unwrap_or(MOQ_CARRIER_OBJECT_PAYLOAD_CEILING)
.min(MOQ_CARRIER_OBJECT_PAYLOAD_CEILING);
if message_ceiling <= CARRIER_HEADER_LEN {
return Err(JsValue::from_str(
"MoQ object-stream payload ceiling is too small for carrier framing",
));
}
let session = Rc::new(session);
let closed = Rc::new(Cell::new(false));
let inbound_session = session.clone();
let inbound_closed = closed.clone();
let inbound_terminal = on_terminal.clone();
let terminal_notified = Rc::new(Cell::new(false));
let terminal_acknowledged = Rc::new(Cell::new(false));
let inbound_terminal_notified = terminal_notified.clone();
let inbound_terminal_acknowledged = terminal_acknowledged.clone();
let inbound_reader = reader.clone();
let inbound_writer = writer.clone();
let reassembler = Rc::new(RefCell::new(CarrierReassembler::default()));
let inbound_readiness = Rc::new(Cell::new(0_u8));
let outbound_readiness_sent = Rc::new(Cell::new(false));
let inbound_readiness_task = inbound_readiness.clone();
web_sys::console::debug_1(&JsValue::from_str(
"[OpenRTC][MoQ carrier][pump] attached browser packet pumps",
));
spawn_local(async move {
let mut first_inbound = true;
while !inbound_closed.get() {
let Ok(read_result) = call_promise(&inbound_reader, "read", None).await else {
break;
};
if Reflect::get(&read_result, &JsValue::from_str("done"))
.ok()
.and_then(|value| value.as_bool())
.unwrap_or(false)
{
break;
}
let Ok(value) = Reflect::get(&read_result, &JsValue::from_str("value")) else {
break;
};
let bytes = Uint8Array::new(&value).to_vec();
if first_inbound {
first_inbound = false;
web_sys::console::debug_1(&JsValue::from_str(&format!(
"[OpenRTC][MoQ carrier][pump] received first browser object bytes={}",
bytes.len(),
)));
}
let observed = crate::iroh_carrier::observe_carrier_readiness_object(0, &bytes);
if observed != 0 {
let previous = inbound_readiness_task.get();
let updated = previous | observed;
inbound_readiness_task.set(updated);
if updated != previous {
web_sys::console::debug_1(&JsValue::from_str(&format!(
"[OpenRTC][MoQ carrier][pump] received readiness object bytes={} mask={updated:#04b}",
bytes.len(),
)));
}
continue;
}
let frame = CarrierFrame::decode(&bytes, expected);
let Ok(frame) = frame else {
continue;
};
match frame.terminal_control_kind(expected, &application_key) {
Ok(Some(CarrierControl::SessionTokenRevokedAck)) => {
inbound_terminal_acknowledged.set(true);
continue;
}
Ok(Some(control @ CarrierControl::SessionTokenRevoked)) => {
let ack = CarrierFrame::terminal_control(
expected,
&application_key,
CarrierControl::SessionTokenRevokedAck,
)
.encode();
if let Ok(ack) = ack {
let value = Uint8Array::from(ack.as_slice());
let _ =
call_promise(&inbound_writer, "write", Some(value.as_ref())).await;
}
if !inbound_terminal_notified.replace(true) {
inbound_session.close();
if let Some(reason) = control.lifecycle_reason() {
inbound_terminal(reason);
}
}
continue;
}
Err(_) if frame.header.packet_id == CARRIER_CONTROL_PACKET_ID => continue,
_ => {}
}
let now_ms = js_sys::Date::now().max(0.0) as u64;
let Ok(Some(packet)) = reassembler.borrow_mut().push(frame, now_ms) else {
continue;
};
if inbound_session.deliver_inbound(packet).await.is_err() {
break;
}
}
if !inbound_closed.get() && !inbound_terminal_notified.replace(true) {
inbound_session.close();
inbound_terminal("iroh-carrier-ended");
}
});
let outbound_session = session.clone();
let outbound_closed = closed.clone();
let outbound_terminal = on_terminal;
let outbound_terminal_notified = terminal_notified;
let outbound_writer = writer.clone();
let outbound_inbound_readiness = inbound_readiness.clone();
let outbound_readiness_task = outbound_readiness_sent.clone();
spawn_local(async move {
let mut logged_first_readiness_train = false;
while !outbound_closed.get()
&& (!outbound_readiness_task.get()
|| !crate::iroh_carrier::carrier_readiness_is_complete(
outbound_inbound_readiness.get(),
))
{
let mut sent = true;
for (index, object) in crate::iroh_carrier::carrier_readiness_objects()
.into_iter()
.enumerate()
{
let value = Uint8Array::from(object.as_slice());
if let Err(error) =
call_promise(&outbound_writer, "write", Some(value.as_ref())).await
{
web_sys::console::error_2(
&JsValue::from_str(&format!(
"[OpenRTC][MoQ carrier][pump] readiness object send failed index={index} bytes={}",
object.len(),
)),
&error,
);
sent = false;
break;
}
}
if sent {
outbound_readiness_task.set(true);
if !logged_first_readiness_train {
logged_first_readiness_train = true;
web_sys::console::debug_1(&JsValue::from_str(
"[OpenRTC][MoQ carrier][pump] sent readiness object train objects=2",
));
}
}
gloo_timers::future::sleep(std::time::Duration::from_millis(100)).await;
}
if outbound_closed.get() {
return;
}
web_sys::console::debug_1(&JsValue::from_str(
"[OpenRTC][MoQ carrier][pump] bilateral object readiness proven",
));
let mut packet_id = 0_u64;
let mut first_outbound = true;
'packets: while !outbound_closed.get() {
let Ok(packet) = outbound_session.recv_outbound().await else {
break;
};
if first_outbound {
first_outbound = false;
web_sys::console::debug_1(&JsValue::from_str(&format!(
"[OpenRTC][MoQ carrier][pump] sending first browser Iroh packet bytes={}",
packet.len(),
)));
}
packet_id = packet_id.wrapping_add(1);
if packet_id == CARRIER_CONTROL_PACKET_ID {
packet_id = 0;
}
let Ok(frames) = segment_packet(&packet, expected, packet_id, message_ceiling)
else {
break;
};
for frame in frames {
let Ok(encoded) = frame.encode() else {
break 'packets;
};
let value = Uint8Array::from(encoded.as_slice());
if call_promise(&outbound_writer, "write", Some(value.as_ref()))
.await
.is_err()
{
break 'packets;
}
}
}
if !outbound_closed.get() && !outbound_terminal_notified.replace(true) {
outbound_session.close();
outbound_terminal("iroh-carrier-ended");
}
});
Ok(Self {
_duplex: duplex,
reader,
writer,
closed,
_session: session,
expected,
application_key,
terminal_acknowledged,
inbound_readiness,
outbound_readiness_sent,
})
}
pub async fn wait_for_peer_data_bidirectional_readiness(
&self,
timeout: std::time::Duration,
) -> Result<(), JsValue> {
let deadline = js_sys::Date::now() + timeout.as_millis() as f64;
loop {
if self.closed.get() {
return Err(JsValue::from_str(
"MoQ carrier closed before bilateral object readiness",
));
}
if self.outbound_readiness_sent.get()
&& crate::iroh_carrier::carrier_readiness_is_complete(self.inbound_readiness.get())
{
return Ok(());
}
if js_sys::Date::now() >= deadline {
return Err(JsValue::from_str(
"MoQ bilateral object readiness timed out before candidate dial",
));
}
gloo_timers::future::sleep(std::time::Duration::from_millis(10)).await;
}
}
pub async fn send_terminal(&self, reason: &str) -> Result<(), JsValue> {
let control = match reason {
crate::lifecycle_reason::REASON_SESSION_TOKEN_REVOKED => {
CarrierControl::SessionTokenRevoked
}
_ => return Ok(()),
};
let encoded = CarrierFrame::terminal_control(self.expected, &self.application_key, control)
.encode()
.map_err(|error| JsValue::from_str(&error.to_string()))?;
self.terminal_acknowledged.set(false);
for _ in 0..12 {
let value = Uint8Array::from(encoded.as_slice());
call_promise(&self.writer, "write", Some(value.as_ref())).await?;
gloo_timers::future::sleep(std::time::Duration::from_millis(25)).await;
if self.terminal_acknowledged.get() {
break;
}
}
Ok(())
}
}
fn method(target: &JsValue, name: &str) -> Result<Function, JsValue> {
Reflect::get(target, &JsValue::from_str(name))?.dyn_into::<Function>()
}
async fn call_promise(
target: &JsValue,
name: &str,
argument: Option<&JsValue>,
) -> Result<JsValue, JsValue> {
let value = match argument {
Some(argument) => method(target, name)?.call1(target, argument)?,
None => method(target, name)?.call0(target)?,
};
JsFuture::from(value.dyn_into::<Promise>()?).await
}
impl Drop for WasmMoqCarrierSession {
fn drop(&mut self) {
self.closed.set(true);
self._session.close();
if let Ok(cancel) = method(&self.reader, "cancel") {
let _ = cancel.call0(&self.reader);
}
if let Ok(close) = method(&self.writer, "close") {
let _ = close.call0(&self.writer);
}
}
}