openrtc 2.5.0

OpenRTC: a Rust-first P2P runtime for device discovery, signaling, and iroh/QUIC networking.
Documentation
//! Browser reliable object-stream pump for the Draft 14 MoQ Iroh carrier.

#![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;
/// Retains browser reliable object streams and bounded packet pumps. The admitted
/// Rust lifecycle actor remains responsible for eligibility and replacement.
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 {
    /// `duplex` is a Draft 14 MoQ reliable object-stream duplex. The browser
    /// adapter first negotiates both directed tracks and exact Draft 14 support.
    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,
        })
    }

    /// Wait until both directed MoQ tracks have carried the two-object
    /// readiness train. Browser initiators call this before any candidate dial;
    /// responders install their pumps immediately and let the initiator gate.
    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;
        }
    }

    /// Send the authenticated terminal marker until its authenticated ACK or
    /// the bounded local-attempt ceiling. This performs no gateway or
    /// control-plane operation.
    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();
        // Cancelling the reader actively resolves a pending `read()` so the
        // retired carrier cannot retain its task or deliver a late terminal
        // callback after a replacement generation has taken ownership.
        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);
        }
    }
}